| | | 1 | | using System.Linq.Expressions; |
| | | 2 | | using System.Runtime.CompilerServices; |
| | | 3 | | using ArturRios.Data.MongoDb.Exceptions; |
| | | 4 | | using ArturRios.Data.MongoDb.Interfaces; |
| | | 5 | | using ArturRios.Output; |
| | | 6 | | using Microsoft.Extensions.Logging; |
| | | 7 | | using MongoDB.Bson; |
| | | 8 | | using MongoDB.Bson.Serialization; |
| | | 9 | | using MongoDB.Driver; |
| | | 10 | | |
| | | 11 | | namespace ArturRios.Data.MongoDb.Repositories; |
| | | 12 | | |
| | | 13 | | /// <summary> |
| | | 14 | | /// MongoDB implementation of the document repository contracts. Runs against the |
| | | 15 | | /// <see cref="MongoContext" /> collection and enlists in its ambient session so operations |
| | | 16 | | /// participate in a unit-of-work transaction. Failures are returned as <see cref="DataOutput{T}" />. |
| | | 17 | | /// </summary> |
| | | 18 | | /// <typeparam name="T">The document type.</typeparam> |
| | | 19 | | /// <param name="context">The Mongo context.</param> |
| | | 20 | | /// <param name="logger"> |
| | | 21 | | /// Optional logger. Envelopes never carry driver text, so a failure is otherwise |
| | | 22 | | /// undiagnosable: supply a logger and the full exception, plus the document type and the |
| | | 23 | | /// repository method that failed, is written at <see cref="LogLevel.Error" />. Document |
| | | 24 | | /// contents and key values are never logged. Resolved from DI when logging is registered. |
| | | 25 | | /// </param> |
| | 38 | 26 | | public class MongoDocumentRepository<T>(MongoContext context, ILogger<MongoDocumentRepository<T>>? logger = null) |
| | | 27 | | : IDocumentRepository<T>, IAsyncDocumentRepository<T> where T : Document |
| | | 28 | | { |
| | | 29 | | /// <summary>Message returned when an operation fails with no finer classification.</summary> |
| | | 30 | | protected const string OperationFailedMessage = MongoErrors.GenericMessage; |
| | | 31 | | |
| | | 32 | | /// <summary>Message returned when the failure is transient and the operation may be retried.</summary> |
| | | 33 | | protected const string TransientMessage = MongoErrors.TransientMessage; |
| | | 34 | | |
| | | 35 | | /// <summary>Message returned on an optimistic-concurrency conflict.</summary> |
| | | 36 | | protected const string ConcurrencyMessage = MongoErrors.ConcurrencyMessage; |
| | | 37 | | |
| | | 38 | | /// <summary>Message returned when a write violates a unique index.</summary> |
| | | 39 | | protected const string UniqueViolationMessage = MongoErrors.UniqueViolationMessage; |
| | | 40 | | |
| | | 41 | | // Cached: the serialized BSON element name for VersionedDocument.Version on T |
| | | 42 | | // (respects any element-name convention the consumer registered). |
| | 4 | 43 | | private static readonly string VersionElementName = |
| | 4 | 44 | | BsonClassMap.LookupClassMap(typeof(T)).AllMemberMaps |
| | 8 | 45 | | .FirstOrDefault(m => m.MemberName == nameof(VersionedDocument.Version))?.ElementName |
| | 4 | 46 | | ?? nameof(VersionedDocument.Version); |
| | | 47 | | |
| | | 48 | | /// <summary>The collection for <typeparamref name="T" />.</summary> |
| | 96 | 49 | | protected IMongoCollection<T> Collection => context.GetCollection<T>(); |
| | | 50 | | |
| | 94 | 51 | | private IClientSessionHandle? Session => context.Session; |
| | | 52 | | |
| | | 53 | | /// <inheritdoc /> |
| | | 54 | | public Task<DataOutput<IEnumerable<T>>> GetAllAsync(CancellationToken ct = default) => |
| | 12 | 55 | | GuardedAsync<IEnumerable<T>>(async () => await FindFluent(FilterDefinition<T>.Empty).ToListAsync(ct).ConfigureAw |
| | | 56 | | |
| | | 57 | | /// <inheritdoc /> |
| | | 58 | | public Task<DataOutput<T?>> GetByIdAsync(string id, CancellationToken ct = default) => |
| | 12 | 59 | | GuardedAsync<T?>(async () => await FindFluent(IdFilter(id)).FirstOrDefaultAsync(ct).ConfigureAwait(false)); |
| | | 60 | | |
| | | 61 | | /// <inheritdoc /> |
| | | 62 | | public Task<DataOutput<IEnumerable<T>>> FindAsync(Expression<Func<T, bool>> predicate, |
| | | 63 | | CancellationToken ct = default) => |
| | 4 | 64 | | GuardedAsync<IEnumerable<T>>(async () => await FindFluent(Builders<T>.Filter.Where(predicate)).ToListAsync(ct).C |
| | | 65 | | |
| | | 66 | | /// <inheritdoc /> |
| | | 67 | | public Task<DataOutput<string>> CreateAsync(T document, CancellationToken ct = default) => |
| | 12 | 68 | | GuardedAsync(async () => |
| | 12 | 69 | | { |
| | 12 | 70 | | EnsureId(document); |
| | 12 | 71 | | await InsertOneAsync(document, ct).ConfigureAwait(false); |
| | 12 | 72 | | return document.Id; |
| | 24 | 73 | | }); |
| | | 74 | | |
| | | 75 | | /// <inheritdoc /> |
| | | 76 | | public Task<DataOutput<IEnumerable<string>>> CreateRangeAsync(IEnumerable<T> documents, |
| | | 77 | | CancellationToken ct = default) => |
| | 4 | 78 | | GuardedAsync<IEnumerable<string>>(async () => |
| | 4 | 79 | | { |
| | 4 | 80 | | var list = documents.ToList(); |
| | 24 | 81 | | foreach (var d in list) |
| | 4 | 82 | | { |
| | 8 | 83 | | EnsureId(d); |
| | 4 | 84 | | } |
| | 4 | 85 | | |
| | 4 | 86 | | await InsertManyAsync(list, ct).ConfigureAwait(false); |
| | 12 | 87 | | return list.Select(d => d.Id).ToList(); |
| | 8 | 88 | | }); |
| | | 89 | | |
| | | 90 | | /// <inheritdoc /> |
| | | 91 | | public Task<DataOutput<T>> UpdateAsync(T document, CancellationToken ct = default) => |
| | 2 | 92 | | GuardedAsync(async () => |
| | 2 | 93 | | { |
| | 2 | 94 | | await ReplaceAsync(document, ct).ConfigureAwait(false); |
| | 2 | 95 | | return document; |
| | 4 | 96 | | }); |
| | | 97 | | |
| | | 98 | | /// <inheritdoc /> |
| | | 99 | | public Task<DataOutput<IEnumerable<T>>> |
| | | 100 | | UpdateRangeAsync(IEnumerable<T> documents, CancellationToken ct = default) => |
| | 0 | 101 | | GuardedAsync<IEnumerable<T>>(async () => |
| | 0 | 102 | | { |
| | 0 | 103 | | var list = documents.ToList(); |
| | 0 | 104 | | foreach (var d in list) |
| | 0 | 105 | | { |
| | 0 | 106 | | await ReplaceAsync(d, ct).ConfigureAwait(false); |
| | 0 | 107 | | } |
| | 0 | 108 | | |
| | 0 | 109 | | return list; |
| | 0 | 110 | | }); |
| | | 111 | | |
| | | 112 | | /// <inheritdoc /> |
| | | 113 | | public Task<DataOutput<string>> DeleteAsync(T document, CancellationToken ct = default) => |
| | 2 | 114 | | GuardedAsync(async () => |
| | 2 | 115 | | { |
| | 2 | 116 | | await DeleteManyAsync(IdFilter(document.Id), ct).ConfigureAwait(false); |
| | 2 | 117 | | return document.Id; |
| | 4 | 118 | | }); |
| | | 119 | | |
| | | 120 | | /// <inheritdoc /> |
| | | 121 | | public Task<DataOutput<IEnumerable<string>>> DeleteRangeAsync(IEnumerable<string> ids, |
| | | 122 | | CancellationToken ct = default) => |
| | 2 | 123 | | GuardedAsync<IEnumerable<string>>(async () => |
| | 2 | 124 | | { |
| | 2 | 125 | | var idList = ids.ToList(); |
| | 2 | 126 | | await DeleteManyAsync(Builders<T>.Filter.In(d => d.Id, idList), ct).ConfigureAwait(false); |
| | 2 | 127 | | return idList; |
| | 4 | 128 | | }); |
| | | 129 | | |
| | | 130 | | /// <inheritdoc /> |
| | 2 | 131 | | public IQueryable<T> Query() => Collection.AsQueryable(); |
| | | 132 | | |
| | | 133 | | /// <inheritdoc /> |
| | | 134 | | public DataOutput<IEnumerable<T>> GetAll() => |
| | 16 | 135 | | Guarded(IEnumerable<T> () => FindFluent(FilterDefinition<T>.Empty).ToList()); |
| | | 136 | | |
| | | 137 | | /// <inheritdoc /> |
| | | 138 | | public DataOutput<T?> GetById(string id) => |
| | 24 | 139 | | Guarded<T?>(() => FindFluent(IdFilter(id)).FirstOrDefault()); |
| | | 140 | | |
| | | 141 | | /// <inheritdoc /> |
| | | 142 | | public DataOutput<IEnumerable<T>> Find(Expression<Func<T, bool>> predicate) => |
| | 4 | 143 | | Guarded(IEnumerable<T> () => FindFluent(Builders<T>.Filter.Where(predicate)).ToList()); |
| | | 144 | | |
| | | 145 | | /// <inheritdoc /> |
| | 18 | 146 | | public DataOutput<string> Create(T document) => Guarded(() => |
| | 18 | 147 | | { |
| | 18 | 148 | | EnsureId(document); |
| | 18 | 149 | | InsertOne(document); |
| | 14 | 150 | | return document.Id; |
| | 18 | 151 | | }); |
| | | 152 | | |
| | | 153 | | /// <inheritdoc /> |
| | 4 | 154 | | public DataOutput<IEnumerable<string>> CreateRange(IEnumerable<T> documents) => Guarded(IEnumerable<string> () => |
| | 4 | 155 | | { |
| | 4 | 156 | | var list = documents.ToList(); |
| | 24 | 157 | | foreach (var d in list) |
| | 4 | 158 | | { |
| | 8 | 159 | | EnsureId(d); |
| | 4 | 160 | | } |
| | 4 | 161 | | |
| | 4 | 162 | | InsertMany(list); |
| | 12 | 163 | | return list.Select(d => d.Id).ToList(); |
| | 4 | 164 | | }); |
| | | 165 | | |
| | | 166 | | /// <inheritdoc /> |
| | 10 | 167 | | public DataOutput<T> Update(T document) => Guarded(() => |
| | 10 | 168 | | { |
| | 10 | 169 | | Replace(document); |
| | 6 | 170 | | return document; |
| | 10 | 171 | | }); |
| | | 172 | | |
| | | 173 | | /// <inheritdoc /> |
| | 0 | 174 | | public DataOutput<IEnumerable<T>> UpdateRange(IEnumerable<T> documents) => Guarded(IEnumerable<T> () => |
| | 0 | 175 | | { |
| | 0 | 176 | | var list = documents.ToList(); |
| | 0 | 177 | | foreach (var d in list) |
| | 0 | 178 | | { |
| | 0 | 179 | | Replace(d); |
| | 0 | 180 | | } |
| | 0 | 181 | | |
| | 0 | 182 | | return list; |
| | 0 | 183 | | }); |
| | | 184 | | |
| | | 185 | | /// <inheritdoc /> |
| | 2 | 186 | | public DataOutput<string> Delete(T document) => Guarded(() => |
| | 2 | 187 | | { |
| | 2 | 188 | | DeleteMany(IdFilter(document.Id)); |
| | 2 | 189 | | return document.Id; |
| | 2 | 190 | | }); |
| | | 191 | | |
| | | 192 | | /// <inheritdoc /> |
| | 2 | 193 | | public DataOutput<IEnumerable<string>> DeleteRange(IEnumerable<string> ids) => Guarded(IEnumerable<string> () => |
| | 2 | 194 | | { |
| | 2 | 195 | | var idList = ids.ToList(); |
| | 2 | 196 | | DeleteMany(Builders<T>.Filter.In(d => d.Id, idList)); |
| | 2 | 197 | | return idList; |
| | 2 | 198 | | }); |
| | | 199 | | |
| | | 200 | | // --- session-aware driver helpers (sync) --- |
| | 34 | 201 | | private static FilterDefinition<T> IdFilter(string id) => Builders<T>.Filter.Eq(d => d.Id, id); |
| | | 202 | | |
| | | 203 | | private static void EnsureId(T document) |
| | | 204 | | { |
| | 46 | 205 | | if (string.IsNullOrEmpty(document.Id)) |
| | | 206 | | { |
| | 38 | 207 | | document.Id = ObjectId.GenerateNewId().ToString(); |
| | | 208 | | } |
| | 46 | 209 | | } |
| | | 210 | | |
| | | 211 | | private IFindFluent<T, T> FindFluent(FilterDefinition<T> filter) => |
| | 36 | 212 | | Session is { } s ? Collection.Find(s, filter) : Collection.Find(filter); |
| | | 213 | | |
| | | 214 | | private void InsertOne(T document) |
| | | 215 | | { |
| | 18 | 216 | | if (Session is { } s) |
| | | 217 | | { |
| | 0 | 218 | | Collection.InsertOne(s, document); |
| | | 219 | | } |
| | | 220 | | else |
| | | 221 | | { |
| | 18 | 222 | | Collection.InsertOne(document); |
| | | 223 | | } |
| | 14 | 224 | | } |
| | | 225 | | |
| | | 226 | | private void InsertMany(IEnumerable<T> documents) |
| | | 227 | | { |
| | 4 | 228 | | if (Session is { } s) |
| | | 229 | | { |
| | 0 | 230 | | Collection.InsertMany(s, documents); |
| | | 231 | | } |
| | | 232 | | else |
| | | 233 | | { |
| | 4 | 234 | | Collection.InsertMany(documents); |
| | | 235 | | } |
| | 4 | 236 | | } |
| | | 237 | | |
| | | 238 | | private void DeleteMany(FilterDefinition<T> filter) |
| | | 239 | | { |
| | 4 | 240 | | if (Session is { } s) |
| | | 241 | | { |
| | 0 | 242 | | Collection.DeleteMany(s, filter); |
| | | 243 | | } |
| | | 244 | | else |
| | | 245 | | { |
| | 4 | 246 | | Collection.DeleteMany(filter); |
| | | 247 | | } |
| | 4 | 248 | | } |
| | | 249 | | |
| | | 250 | | // Replace with optimistic-concurrency handling for VersionedDocument. |
| | | 251 | | private void Replace(T document) |
| | | 252 | | { |
| | 10 | 253 | | if (document is VersionedDocument versioned) |
| | | 254 | | { |
| | 8 | 255 | | var expected = versioned.Version; |
| | 8 | 256 | | versioned.Version = expected + 1; |
| | 8 | 257 | | var filter = Builders<T>.Filter.And(IdFilter(document.Id), |
| | 8 | 258 | | Builders<T>.Filter.Eq(VersionElementName, expected)); |
| | | 259 | | |
| | | 260 | | ReplaceOneResult result; |
| | | 261 | | try |
| | | 262 | | { |
| | 8 | 263 | | result = ReplaceOne(filter, document); |
| | 8 | 264 | | } |
| | 0 | 265 | | catch |
| | | 266 | | { |
| | | 267 | | // The write never landed, so the in-memory bump must not survive: keeping it would |
| | | 268 | | // make every retry filter on a version the server never stored. |
| | 0 | 269 | | versioned.Version = expected; |
| | 0 | 270 | | throw; |
| | | 271 | | } |
| | | 272 | | |
| | 8 | 273 | | if (result.MatchedCount == 0) |
| | | 274 | | { |
| | 4 | 275 | | versioned.Version = expected; // roll back the in-memory bump on a failed (stale) update |
| | 4 | 276 | | throw new MongoConcurrencyException(); |
| | | 277 | | } |
| | | 278 | | |
| | 4 | 279 | | return; |
| | | 280 | | } |
| | | 281 | | |
| | 2 | 282 | | ReplaceOne(IdFilter(document.Id), document); |
| | 2 | 283 | | } |
| | | 284 | | |
| | | 285 | | private ReplaceOneResult ReplaceOne(FilterDefinition<T> filter, T document) => |
| | 10 | 286 | | Session is { } s ? Collection.ReplaceOne(s, filter, document) : Collection.ReplaceOne(filter, document); |
| | | 287 | | |
| | | 288 | | /// <summary>Runs a synchronous operation, converting failures to envelope errors.</summary> |
| | | 289 | | /// <param name="operation">The operation to run.</param> |
| | | 290 | | /// <param name="operationName">The calling repository method, used as log context.</param> |
| | | 291 | | protected DataOutput<TResult> Guarded<TResult>(Func<TResult> operation, |
| | | 292 | | [CallerMemberName] string operationName = "") |
| | | 293 | | { |
| | | 294 | | try |
| | | 295 | | { |
| | 58 | 296 | | return DataOutput<TResult>.New.WithData(operation()); |
| | | 297 | | } |
| | 0 | 298 | | catch (OperationCanceledException) |
| | | 299 | | { |
| | 0 | 300 | | throw; |
| | | 301 | | } |
| | 8 | 302 | | catch (Exception ex) |
| | | 303 | | { |
| | 8 | 304 | | return Fail<TResult>(ex, operationName); |
| | | 305 | | } |
| | 58 | 306 | | } |
| | | 307 | | |
| | | 308 | | // --- session-aware driver helpers (async) --- |
| | | 309 | | private Task InsertOneAsync(T document, CancellationToken ct) => |
| | 12 | 310 | | Session is { } s |
| | 12 | 311 | | ? Collection.InsertOneAsync(s, document, null, ct) |
| | 12 | 312 | | : Collection.InsertOneAsync(document, null, ct); |
| | | 313 | | |
| | | 314 | | private Task InsertManyAsync(IEnumerable<T> documents, CancellationToken ct) => |
| | 4 | 315 | | Session is { } s |
| | 4 | 316 | | ? Collection.InsertManyAsync(s, documents, null, ct) |
| | 4 | 317 | | : Collection.InsertManyAsync(documents, null, ct); |
| | | 318 | | |
| | | 319 | | private Task DeleteManyAsync(FilterDefinition<T> filter, CancellationToken ct) => |
| | 4 | 320 | | Session is { } s ? Collection.DeleteManyAsync(s, filter, null, ct) : Collection.DeleteManyAsync(filter, ct); |
| | | 321 | | |
| | | 322 | | private async Task ReplaceAsync(T document, CancellationToken ct) |
| | | 323 | | { |
| | 2 | 324 | | if (document is VersionedDocument versioned) |
| | | 325 | | { |
| | 0 | 326 | | var expected = versioned.Version; |
| | 0 | 327 | | versioned.Version = expected + 1; |
| | 0 | 328 | | var filter = Builders<T>.Filter.And(IdFilter(document.Id), |
| | 0 | 329 | | Builders<T>.Filter.Eq(VersionElementName, expected)); |
| | | 330 | | ReplaceOneResult result; |
| | | 331 | | try |
| | | 332 | | { |
| | 0 | 333 | | result = Session is { } s |
| | 0 | 334 | | ? await Collection.ReplaceOneAsync(s, filter, document, cancellationToken: ct).ConfigureAwait(false) |
| | 0 | 335 | | : await Collection.ReplaceOneAsync(filter, document, cancellationToken: ct).ConfigureAwait(false); |
| | 0 | 336 | | } |
| | 0 | 337 | | catch |
| | | 338 | | { |
| | | 339 | | // The write never landed, so the in-memory bump must not survive: keeping it would |
| | | 340 | | // make every retry filter on a version the server never stored. |
| | 0 | 341 | | versioned.Version = expected; |
| | 0 | 342 | | throw; |
| | | 343 | | } |
| | | 344 | | |
| | 0 | 345 | | if (result.MatchedCount == 0) |
| | | 346 | | { |
| | 0 | 347 | | versioned.Version = expected; // roll back the in-memory bump on a failed (stale) update |
| | 0 | 348 | | throw new MongoConcurrencyException(); |
| | | 349 | | } |
| | | 350 | | |
| | 0 | 351 | | return; |
| | | 352 | | } |
| | | 353 | | |
| | 2 | 354 | | var idFilter = IdFilter(document.Id); |
| | 2 | 355 | | if (Session is { } session) |
| | | 356 | | { |
| | 0 | 357 | | await Collection.ReplaceOneAsync(session, idFilter, document, cancellationToken: ct).ConfigureAwait(false); |
| | | 358 | | } |
| | | 359 | | else |
| | | 360 | | { |
| | 2 | 361 | | await Collection.ReplaceOneAsync(idFilter, document, cancellationToken: ct).ConfigureAwait(false); |
| | | 362 | | } |
| | 2 | 363 | | } |
| | | 364 | | |
| | | 365 | | /// <summary>Runs an asynchronous operation, converting failures to envelope errors.</summary> |
| | | 366 | | /// <param name="operation">The operation to run.</param> |
| | | 367 | | /// <param name="operationName">The calling repository method, used as log context.</param> |
| | | 368 | | protected async Task<DataOutput<TResult>> GuardedAsync<TResult>(Func<Task<TResult>> operation, |
| | | 369 | | [CallerMemberName] string operationName = "") |
| | | 370 | | { |
| | | 371 | | try |
| | | 372 | | { |
| | 36 | 373 | | return DataOutput<TResult>.New.WithData(await operation().ConfigureAwait(false)); |
| | | 374 | | } |
| | 0 | 375 | | catch (OperationCanceledException) |
| | | 376 | | { |
| | 0 | 377 | | throw; |
| | | 378 | | } |
| | 0 | 379 | | catch (Exception ex) |
| | | 380 | | { |
| | 0 | 381 | | return Fail<TResult>(ex, operationName); |
| | | 382 | | } |
| | 36 | 383 | | } |
| | | 384 | | |
| | | 385 | | /// <summary> |
| | | 386 | | /// Logs the failure when a logger is configured, and maps it to an error envelope. |
| | | 387 | | /// Driver text names indexes, collections, key values and cluster endpoints, so it goes to |
| | | 388 | | /// the log and never to the caller. |
| | | 389 | | /// </summary> |
| | | 390 | | /// <param name="ex">The exception caught by a guard.</param> |
| | | 391 | | /// <param name="operationName">The repository method that failed, used as log context.</param> |
| | | 392 | | protected DataOutput<TResult> Fail<TResult>(Exception ex, string operationName = "") |
| | | 393 | | { |
| | | 394 | | // A concurrency conflict is an expected outcome of optimistic locking, not an operational |
| | | 395 | | // fault: logging it at Error would fill the log with routine contention. |
| | 8 | 396 | | var level = ex is MongoConcurrencyException ? LogLevel.Debug : LogLevel.Error; |
| | | 397 | | |
| | 8 | 398 | | logger?.Log(level, ex, "Mongo operation failed. Document: {Document}, operation: {Operation}", |
| | 8 | 399 | | typeof(T).Name, operationName); |
| | | 400 | | |
| | 8 | 401 | | return DataOutput<TResult>.New.WithError(MongoErrors.Describe(ex)); |
| | | 402 | | } |
| | | 403 | | } |