A unit of work for MongoDB on .NET, built on the official driver. A vault holds a database's collections: reads return the
driver's own types with your query filters applied, and writes are queued until SaveAsync, which sends them as one
client bulk write, in a transaction. Soft delete, multi-tenancy, a concurrency token, interceptors, change tracking and
migrations come with it, and it logs, traces and measures what it does.
1.0 is in beta. Feedback is welcome in the issues. The changelog maps 0.5's API to 1.0's.
- MongoDB 8.0 or later, as a replica set: a save is one client bulk write, in a transaction.
- .NET 10 or .NET 11, and MongoDB.Driver 3.12.
dotnet add package MongoFlow --prereleaseA vault declares its collections as properties. A keyed collection is looked up by key: by default the member the driver
maps to _id.
public sealed class Order
{
public int Id { get; set; }
public string Customer { get; set; } = "";
public decimal Total { get; set; }
public bool IsDeleted { get; set; }
}
public sealed class ShopVault : MongoVault
{
public IVaultCollection<Order, int> Orders { get; init; } = null!;
}Register it with the database it uses, and resolve it from DI; it's scoped, like a request.
services.AddSingleton<IMongoClient>(new MongoClient("mongodb://localhost:27017/?replicaSet=rs0"));
services.AddMongoVault<ShopVault>(vault => vault
.UseDatabase("shop")
.UseSoftDelete((Order x) => x.IsDeleted));Reads start with one await, which resolves the query filters, and return the driver's types. Writes are queued, and written when the vault is saved.
public sealed class OrderService(ShopVault vault)
{
public async Task<List<Order>> LargeAsync(CancellationToken cancellationToken)
{
var orders = await vault.Orders.QueryAsync(cancellationToken);
return await orders.Where(x => x.Total > 100).ToListAsync(cancellationToken);
}
// Both writes go in one bulk write, in a transaction: the new order is stored and the old one deleted, or neither.
public async Task ReplaceAsync(int oldOrderId, Order newOrder, CancellationToken cancellationToken)
{
vault.Orders.Add(newOrder);
vault.Orders.DeleteByKey(oldOrderId); // soft delete turns this into an update that sets IsDeleted
await vault.SaveAsync(cancellationToken);
}
public async Task DiscountAsync(int orderId, decimal amount, CancellationToken cancellationToken)
{
vault.Orders.UpdateByKey(orderId, Builders<Order>.Update.Inc(x => x.Total, -amount));
await vault.SaveAsync(cancellationToken);
}
}Each vault is configured where it's registered, by the vault itself (IConfigurableVault<TSelf>), or by configurations
that apply to every vault (AddDefaultVaultConfiguration), which a vault can skip. Collections are selected by their
property, so a wrong key type or a collection the vault doesn't declare is a compile error. Everything is applied and
validated once, when the vault is first resolved or migrated.
services.AddMongoVault<IPolicyVault, PolicyVault>(vault => vault
.UseDatabase("policies")
.Collection(x => x.Policies, policies => policies
.Name("insurance_policies")
.Key(p => p.PolicyNumber)));Configuration: registration, keys and composite keys, defaults, and the order settings apply in.
QueryAsync, FindAsync, AggregateAsync and GetByKeyAsync read; Add, AddRange, Replace, Update,
UpdateByKey, UpdateMany, Delete, DeleteByKey and DeleteMany queue writes. A write by key or filter also carries
the query filters, so it can't reach a document a read couldn't see. A query joined with another collection's query
joins only what that query's filters show. MongoCollection is the driver's collection, with nothing applied.
Reads and writes: joins, what a save sends, its result, and what may run in parallel.
With vault.UseChangeTracking(), the documents reads return are tracked, and SaveAsync writes what changed in them:
one update by key, setting the changed fields and nothing else.
var policy = await vault.Policies.GetByKeyAsync("P-1001");
policy!.Status = PolicyStatus.Cancelled;
await vault.SaveAsync(); // { $set: { Status: "Cancelled" } }A save runs in a transaction of its own, or joins the scope's open one. IVaultTransactionManager.BeginAsync opens one
that every vault saved in the scope joins, so saves of several vaults commit together, or roll back together.
await using var transaction = await transactions.BeginAsync();
policies.Claims.Add(claim);
await policies.SaveAsync();
customers.Customers.UpdateByKey(customerId, Builders<Customer>.Update.Inc(x => x.OpenClaims, 1));
await customers.SaveAsync();
await transaction.CommitAsync();Transactions: rollback, and what a failed save does to an open transaction.
Query filters apply to every read and write by key or filter: static, built per query from the request's services, or
asynchronous. Features bundle filters and interceptors behind a FeatureKey, so they can be switched off for one read
or one collection. Soft delete, multi-tenancy and a concurrency token are built in.
vault.UseSoftDelete((ISoftDeletable x) => x.IsDeleted)
.UseMultiTenancy((ITenantOwned x) => x.TenantId, services => services.GetRequiredService<ITenant>().Id)
.UseConcurrencyToken((IVersioned x) => x.Version);
var everything = await vault.Orders.Without(MultiTenancyFeature.Key).QueryAsync();A VaultInterceptor sees a save's writes before they're sent and their results after, and can replace, remove or add
writes, or guard one with a condition the server checks. Its hooks run before the write, after it, after the commit, and
when the save fails.
Interceptors: hooks, write conditions, and an audit trail and an outbox written in the same transaction.
Migrations are classes with a version, applied once, in order, each in a transaction unless it opts out, and recorded. They can be reverted down to a version.
await app.Services.GetRequiredService<IVaultMigrator>().MigrateAllAsync();MongoFlow logs through the ILoggerFactory in DI, and warns about keys and feature fields no index covers. It traces
saves and migrations, and measures saves and transactions, through System.Diagnostics, for OpenTelemetry.
services.AddOpenTelemetry()
.WithTracing(tracing => tracing.AddSource(MongoFlowTelemetry.ActivitySourceName, MongoTelemetry.ActivitySourceName))
.WithMetrics(metrics => metrics.AddMeter(MongoFlowTelemetry.MeterName));MongoFlow.Identity stores Identity's users and roles in a vault, where your features and interceptors apply to them too.
The samples are a small insurance platform
that runs: dotnet run --project samples/MongoFlow.Samples starts MongoDB in Docker and walks through every feature
above, as requests would.
Issues and pull requests are welcome. The tests need Docker for their MongoDB; .claude/rules/ describes how tests and
benchmarks are written.
MIT; see LICENSE.md.