| | | 1 | | using ArturRios.Data.MongoDb.Repositories; |
| | | 2 | | using ArturRios.Output; |
| | | 3 | | using MongoDB.Driver; |
| | | 4 | | |
| | | 5 | | namespace ArturRios.Data.MongoDb.Transactions; |
| | | 6 | | |
| | | 7 | | /// <summary> |
| | | 8 | | /// MongoDB implementation of the unit of work. Opens a client session, sets it as the context's |
| | | 9 | | /// ambient session so repository operations enlist, and commits/aborts the transaction. |
| | | 10 | | /// </summary> |
| | | 11 | | /// <param name="client">The Mongo client.</param> |
| | | 12 | | /// <param name="context">The Mongo context whose ambient session is managed.</param> |
| | 16 | 13 | | public class MongoUnitOfWork(IMongoClient client, MongoContext context) : IMongoUnitOfWork, IAsyncMongoUnitOfWork |
| | | 14 | | { |
| | | 15 | | /// <inheritdoc /> |
| | | 16 | | public async Task<ProcessOutput> ExecuteInTransactionAsync(Func<Task> work, CancellationToken ct = default) |
| | | 17 | | { |
| | 10 | 18 | | using var session = await client.StartSessionAsync(cancellationToken: ct).ConfigureAwait(false); |
| | 10 | 19 | | var previousSession = context.Session; |
| | 10 | 20 | | context.Session = session; |
| | 10 | 21 | | session.StartTransaction(); |
| | | 22 | | try |
| | | 23 | | { |
| | 10 | 24 | | await work().ConfigureAwait(false); |
| | 6 | 25 | | await session.CommitTransactionAsync(ct).ConfigureAwait(false); |
| | 6 | 26 | | return ProcessOutput.New; |
| | | 27 | | } |
| | 4 | 28 | | catch (Exception ex) |
| | | 29 | | { |
| | 4 | 30 | | await AbortQuietlyAsync(session).ConfigureAwait(false); |
| | | 31 | | |
| | 4 | 32 | | if (ex is OperationCanceledException) |
| | | 33 | | { |
| | 0 | 34 | | throw; |
| | | 35 | | } |
| | | 36 | | |
| | 4 | 37 | | return ProcessOutput.New.WithError(MongoErrors.Describe(ex)); |
| | | 38 | | } |
| | | 39 | | finally |
| | | 40 | | { |
| | 10 | 41 | | context.Session = previousSession; |
| | | 42 | | } |
| | 10 | 43 | | } |
| | | 44 | | |
| | | 45 | | /// <inheritdoc /> |
| | | 46 | | public async Task<DataOutput<TResult>> ExecuteInTransactionAsync<TResult>(Func<Task<TResult>> work, |
| | | 47 | | CancellationToken ct = default) |
| | | 48 | | { |
| | 2 | 49 | | using var session = await client.StartSessionAsync(cancellationToken: ct).ConfigureAwait(false); |
| | 2 | 50 | | var previousSession = context.Session; |
| | 2 | 51 | | context.Session = session; |
| | 2 | 52 | | session.StartTransaction(); |
| | | 53 | | try |
| | | 54 | | { |
| | 2 | 55 | | var result = await work().ConfigureAwait(false); |
| | 2 | 56 | | await session.CommitTransactionAsync(ct).ConfigureAwait(false); |
| | 2 | 57 | | return DataOutput<TResult>.New.WithData(result); |
| | | 58 | | } |
| | 0 | 59 | | catch (Exception ex) |
| | | 60 | | { |
| | 0 | 61 | | await AbortQuietlyAsync(session).ConfigureAwait(false); |
| | | 62 | | |
| | 0 | 63 | | if (ex is OperationCanceledException) |
| | | 64 | | { |
| | 0 | 65 | | throw; |
| | | 66 | | } |
| | | 67 | | |
| | 0 | 68 | | return DataOutput<TResult>.New.WithError(MongoErrors.Describe(ex)); |
| | | 69 | | } |
| | | 70 | | finally |
| | | 71 | | { |
| | 2 | 72 | | context.Session = previousSession; |
| | | 73 | | } |
| | 2 | 74 | | } |
| | | 75 | | |
| | | 76 | | /// <inheritdoc /> |
| | | 77 | | public ProcessOutput ExecuteInTransaction(Action work) |
| | | 78 | | { |
| | 0 | 79 | | using var session = client.StartSession(); |
| | 0 | 80 | | var previousSession = context.Session; |
| | 0 | 81 | | context.Session = session; |
| | 0 | 82 | | session.StartTransaction(); |
| | | 83 | | try |
| | | 84 | | { |
| | 0 | 85 | | work(); |
| | 0 | 86 | | session.CommitTransaction(); |
| | 0 | 87 | | return ProcessOutput.New; |
| | | 88 | | } |
| | 0 | 89 | | catch (OperationCanceledException) |
| | | 90 | | { |
| | 0 | 91 | | AbortQuietly(session); |
| | | 92 | | |
| | 0 | 93 | | throw; |
| | | 94 | | } |
| | 0 | 95 | | catch (Exception ex) |
| | | 96 | | { |
| | 0 | 97 | | AbortQuietly(session); |
| | | 98 | | |
| | 0 | 99 | | return ProcessOutput.New.WithError(MongoErrors.Describe(ex)); |
| | | 100 | | } |
| | | 101 | | finally |
| | | 102 | | { |
| | 0 | 103 | | context.Session = previousSession; |
| | 0 | 104 | | } |
| | 0 | 105 | | } |
| | | 106 | | |
| | | 107 | | /// <inheritdoc /> |
| | | 108 | | public DataOutput<TResult> ExecuteInTransaction<TResult>(Func<TResult> work) |
| | | 109 | | { |
| | 0 | 110 | | using var session = client.StartSession(); |
| | 0 | 111 | | var previousSession = context.Session; |
| | 0 | 112 | | context.Session = session; |
| | 0 | 113 | | session.StartTransaction(); |
| | | 114 | | try |
| | | 115 | | { |
| | 0 | 116 | | var result = work(); |
| | 0 | 117 | | session.CommitTransaction(); |
| | 0 | 118 | | return DataOutput<TResult>.New.WithData(result); |
| | | 119 | | } |
| | 0 | 120 | | catch (OperationCanceledException) |
| | | 121 | | { |
| | 0 | 122 | | AbortQuietly(session); |
| | | 123 | | |
| | 0 | 124 | | throw; |
| | | 125 | | } |
| | 0 | 126 | | catch (Exception ex) |
| | | 127 | | { |
| | 0 | 128 | | AbortQuietly(session); |
| | | 129 | | |
| | 0 | 130 | | return DataOutput<TResult>.New.WithError(MongoErrors.Describe(ex)); |
| | | 131 | | } |
| | | 132 | | finally |
| | | 133 | | { |
| | 0 | 134 | | context.Session = previousSession; |
| | 0 | 135 | | } |
| | 0 | 136 | | } |
| | | 137 | | |
| | | 138 | | // Abort must never mask the failure that triggered it: the server aborts the transaction itself |
| | | 139 | | // on a write conflict, so an explicit abort can fail on a transaction that is already gone. It |
| | | 140 | | // also runs untied to the caller's token, which may already be canceled. |
| | | 141 | | private static async Task AbortQuietlyAsync(IClientSessionHandle session) |
| | | 142 | | { |
| | | 143 | | try |
| | | 144 | | { |
| | 4 | 145 | | await session.AbortTransactionAsync(CancellationToken.None).ConfigureAwait(false); |
| | 4 | 146 | | } |
| | 0 | 147 | | catch |
| | | 148 | | { |
| | | 149 | | // Already aborted, or the session is gone. Disposing the session completes the cleanup. |
| | 0 | 150 | | } |
| | 4 | 151 | | } |
| | | 152 | | |
| | | 153 | | private static void AbortQuietly(IClientSessionHandle session) |
| | | 154 | | { |
| | | 155 | | try |
| | | 156 | | { |
| | 0 | 157 | | session.AbortTransaction(); |
| | 0 | 158 | | } |
| | 0 | 159 | | catch |
| | | 160 | | { |
| | | 161 | | // Already aborted, or the session is gone. Disposing the session completes the cleanup. |
| | 0 | 162 | | } |
| | 0 | 163 | | } |
| | | 164 | | } |