using System.Linq.Expressions; using Elsa.Extensions; using Elsa.MongoDb.Extensions; using JetBrains.Annotations; using MongoDB.Driver; using MongoDB.Driver.Linq; namespace Elsa.MongoDb.Common; /// /// A generic repository class around MongoDb for accessing documents. /// /// The type of the document. [PublicAPI] public class MongoDbStore where TDocument : class { private readonly IMongoCollection _collection; /// public MongoDbStore(IMongoCollection collection) { _collection = collection; } /// /// Returns a queryable collection of documents. /// public IMongoCollection GetCollection() => _collection; /// /// Saves the document. /// /// The document to save. /// The cancellation token. public async Task AddAsync(TDocument document, CancellationToken cancellationToken = default) { await _collection.InsertOneAsync(document, new InsertOneOptions(), cancellationToken); return document; } /// /// Saves a list of documents. /// /// The documents to save. /// The cancellation token. public async Task AddManyAsync(IEnumerable documents, CancellationToken cancellationToken = default) { await _collection.InsertManyAsync(documents, new InsertManyOptions(), cancellationToken); } /// /// Saves the document. /// /// The document to save. /// The cancellation token. public async Task SaveAsync(TDocument document, CancellationToken cancellationToken = default) { return await _collection.FindOneAndReplaceAsync(document.BuildIdFilter(), document, new FindOneAndReplaceOptions{ ReturnDocument = ReturnDocument.After, IsUpsert = true }, cancellationToken); } /// /// Saves the document. /// /// The document to save. /// The selector to use. /// The cancellation token. public async Task SaveAsync(TDocument document, Expression> selector, CancellationToken cancellationToken = default) { return await _collection.FindOneAndReplaceAsync(document.BuildExpression(selector), document, new FindOneAndReplaceOptions{ ReturnDocument = ReturnDocument.After, IsUpsert = true }, cancellationToken); } /// /// Saves the specified documents. /// /// The documents to save. /// The cancellation token. public async Task SaveManyAsync(IEnumerable documents, CancellationToken cancellationToken = default) { var writes = new List>(); foreach (var document in documents) { var replacement = new ReplaceOneModel(document.BuildIdFilter(), document) { IsUpsert = true }; writes.Add(replacement); } await _collection.BulkWriteAsync(writes, cancellationToken: cancellationToken); } /// /// Saves the specified documents. /// /// The documents to save. /// The primary key to use. /// The cancellation token. public async Task SaveManyAsync(IEnumerable documents, string primaryKey = "Id", CancellationToken cancellationToken = default) { var writes = new List>(); foreach (var document in documents) { var replacement = new ReplaceOneModel(document.BuildFilter(primaryKey), document) { IsUpsert = true }; writes.Add(replacement); } await _collection.BulkWriteAsync(writes, cancellationToken: cancellationToken); } /// /// Finds the document matching the specified predicate /// /// The predicate to use. /// The cancellation token. /// The document if found, otherwise null. public async Task FindAsync(Expression> predicate, CancellationToken cancellationToken = default) => await _collection.AsQueryable().Where(predicate).FirstOrDefaultAsync(cancellationToken); /// /// Finds a single document using a query /// /// The query to use /// The cancellation token /// The document if found, otherwise null public async Task FindAsync(Func, IMongoQueryable> query, CancellationToken cancellationToken = default) => await query(_collection.AsQueryable()).FirstOrDefaultAsync(cancellationToken); /// /// Finds a list of documents matching the specified predicate /// public async Task> FindManyAsync(Expression> predicate, CancellationToken cancellationToken = default) => await _collection.AsQueryable().Where(predicate).ToListAsync(cancellationToken); /// /// Queries the database using a query and a selector. /// public async Task> FindManyAsync(Func, IMongoQueryable> query, Expression> selector, CancellationToken cancellationToken = default) => await query(_collection.AsQueryable()).Select(selector).ToListAsync(cancellationToken); /// /// Finds a list of documents using a query /// public async Task> FindManyAsync(Func, IMongoQueryable> query, CancellationToken cancellationToken = default) => await query(_collection.AsQueryable()).ToListAsync(cancellationToken); /// /// Queries the database using a query and a selector. /// public async Task> FindMany(Func, IMongoQueryable> query, Expression> selector, CancellationToken cancellationToken = default) => await query(_collection.AsQueryable()).Select(selector).ToListAsync(cancellationToken); /// /// Counts documents in the collection using a filter. /// public async Task CountAsync(Func, IMongoQueryable> query, CancellationToken cancellationToken = default) => await query(_collection.AsQueryable()).LongCountAsync(cancellationToken); /// /// Counts documents in the collection using a filter and distinct by a key selector. /// public async Task CountAsync(Func, IMongoQueryable> query, Expression> propertySelector, CancellationToken cancellationToken = default) => await query((IMongoQueryable)_collection.AsQueryable().DistinctBy(propertySelector)).LongCountAsync(cancellationToken); /// /// Checks if any documents exist. /// public async Task AnyAsync(Expression> predicate, CancellationToken cancellationToken = default) => await _collection.AsQueryable().Where(predicate).AnyAsync(cancellationToken); /// /// Deletes documents using a predicate. /// /// The number of documents deleted. public async Task DeleteWhereAsync(Expression> predicate, CancellationToken cancellationToken = default) { var documentsToDelete = await _collection.AsQueryable().Where(predicate).ToListAsync(cancellationToken); var count = documentsToDelete.LongCount(); await _collection.DeleteManyAsync(predicate, cancellationToken); return count; } /// /// Deletes documents using a query. /// /// The number of documents deleted. public async Task DeleteWhereAsync(Func, IMongoQueryable> query, Expression> keySelector, CancellationToken cancellationToken = default) { var key = keySelector.GetPropertyName(); return await DeleteWhereAsync(query, key, cancellationToken); } /// /// Deletes documents using a query. /// /// The number of documents deleted. public async Task DeleteWhereAsync(Func, IMongoQueryable> query, string key = "Id", CancellationToken cancellationToken = default) { var documentsToDelete = await query(_collection.AsQueryable()).ToListAsync(cancellationToken); var count = documentsToDelete.LongCount(); var filter = documentsToDelete.BuildIdFilterForList(key); await _collection.DeleteManyAsync(filter, cancellationToken); return count; } }