Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
85 changes: 69 additions & 16 deletions docs/mongodb.md
Original file line number Diff line number Diff line change
Expand Up @@ -174,23 +174,66 @@ documents.

## Operation semantics

Every row applies to both the synchronous and the asynchronous variant: they run the same driver
operation, on the same session, with the same success criterion and the same exceptions.

| Method | Behavior | Exception |
|---|---|---|
| `FindAsync(key)` | filter on `_id` | — (`null` if absent) |
| `FindAsync(key, include)` | **not supported** | `NotSupportedException` |
| `InsertAsync` | `insertOne` | `MongoWriteException` on duplicate key |
| `InsertManyAsync` | ordered `insertMany`, in a transaction | `MongoBulkWriteException` |
| `UpdateAsync` | replaces the entire document | `DBConcurrencyException` if it does not exist |
| `UpdateManyAsync` | ordered `bulkWrite`, in a transaction | `DBConcurrencyException` if any do not exist |
| `UpsertAsync` | replaces or inserts | `DBConcurrencyException` if not applied |
| `RemoveAsync` | `deleteOne` | `DBConcurrencyException` if it does not exist |

Two points may surprise users coming from relational providers:

- **`UpdateAsync` replaces the entire document**, not only the modified fields: there is no change
| `Find(key)` | filter on `_id` | — (`null` if absent) |
| `Find(key, include)` | as `Find(key)`: embedded navigations are already loaded, see below | `NotSupportedException` for references to other collections |
| `Insert` | `insertOne` | `MongoWriteException` on duplicate key |
| `InsertMany` | ordered `insertMany`, in a transaction | `MongoBulkWriteException` on duplicate key |
| `Update` | replaces the entire document | `DBConcurrencyException` if it does not exist |
| `UpdateMany` | ordered `bulkWrite` of replacements, in a transaction | `DBConcurrencyException` if any do not exist |
| `Upsert` | replaces or inserts | `DBConcurrencyException` if not applied |
| `Remove(entity)` / `Remove(key)` | `deleteOne` | `DBConcurrencyException` if it does not exist |

Some points may surprise users coming from relational providers:

- **`Update` replaces the entire document**, not only the modified fields: there is no change
tracking. For a partial update, use `Collection.UpdateOneAsync` directly with the current session.
- **An update that changes nothing is successful.** The criterion is "the document exists", not
"the document was rewritten".
- **An update or upsert that changes nothing is successful.** The criterion is "the document
exists", not "the document was rewritten".
- **`InsertMany` and `UpdateMany` with an empty sequence do nothing**: no round-trip and no
transaction, so they succeed on a standalone server too. A `null` element in the sequence throws
`ArgumentException`.
- **Duplicate keys** surface as the driver's own `MongoWriteException` / `MongoBulkWriteException`,
unchanged, as the ADO.NET provider does with the database exceptions. In a transaction, nothing
of the failed commit is persisted.
- **Unacknowledged writes** (write concern `w: 0`) are not checked: the server returns no counts,
so a missing document cannot be told apart from a successful write, and no
`DBConcurrencyException` is raised.
- **Cancellation.** Outside a unit of work the token reaches the driver. Inside a unit of work an
operation is only queued: its token is checked when it is queued (an operation already cancelled
is not queued and throws `OperationCanceledException`), while the commit is governed by the token
passed to `SaveAsync`.

### Include

The provider behaves as EF Core does with owned types. **Embedded navigations are always loaded**
with the document, so including them is accepted and has no effect. This keeps provider-agnostic
code, such as generated code that calls `Include`, working unchanged on MongoDB.

```csharp
// all valid, and equivalent to FindAsync(id): Items and their children are in the document
await cartRepository.FindAsync(id, include => include.Include(cart => cart.Items));
await cartRepository.FindAsync(id, include => include
.Include(cart => cart.Items, items => items.Include(item => item.Discount)));
await cartRepository.FindAsync(id, include => include.Include("Items.Discount"));
```

Every navigation in the path is validated **before** the query, so the outcome does not depend
on whether the document exists. A request the provider cannot satisfy is never ignored:

| Request | Exception |
|---|---|
| navigation to an entity stored in **another collection** (a type with `[Collection]` / `[Table]`) | `NotSupportedException` |
| member that is **not persisted** (for example `[BsonIgnore]`) | `NotSupportedException` |
| member that does not exist, a scalar value, or an expression that is not a member access (a filtered include such as `x => x.Items.Where(...)`) | `InvalidOperationException` |

References between collections are only keys, and the provider does not resolve them, as the
EF Core providers for document databases also do. Load the referenced aggregate from its own
repository, or query its collection from a specialized repository, passing `Session`.

## Repository

Expand Down Expand Up @@ -221,6 +264,16 @@ The base class exposes `Collection` (`IMongoCollection<TEntity>`), `Database`, a
> Passing `Session` to custom queries is not optional: an operation executed without a session runs
> on an implicit session, and therefore **outside** the transaction of the current unit of work.

The repository is registered and consumed through the CAEP abstractions, like the repositories of
the other providers: `AddData` registers the MongoDB `IDataContext` the constructor needs.

```csharp
builder.Services.AddScoped<IProductRepository, ProductRepository>();

// a plain repository, without specialized queries
builder.Services.AddScoped<IRepository<Category, Guid>, MongoDBRepository<Category, Guid>>();
```

### Mapped repository

When the document must differ from the domain, for example for denormalization, technical fields,
Expand Down Expand Up @@ -305,8 +358,8 @@ MongoDbContainer container = new MongoDbBuilder()

Not yet supported, in order of impact:

- **`Include`** — both for embedded navigations (where it would be a no-op) and inter-aggregate
references. Throws `NotSupportedException`.
- **`Include` of references between collections** and **filtered includes** — see
[Include](#include). Embedded navigations are supported.
- **Optimistic concurrency** — no version token.
- **Multitenancy and soft delete** — available in the EF Core provider, not here.
- **Change tracking** — has no equivalent in the aggregate-based model.
Expand Down
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
<Project Sdk="Microsoft.NET.Sdk">

<PropertyGroup>
<TargetFramework>net7.0</TargetFramework>
<TargetFrameworks>net7.0;net8.0;net9.0;net10.0</TargetFrameworks>
<LangVersion>latest</LangVersion>
<ImplicitUsings>enable</ImplicitUsings>
<Nullable>enable</Nullable>
Expand Down
108 changes: 67 additions & 41 deletions src/Data.MongoDB/DataContext.cs
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@
using CodeArchitects.Platform.Data.MongoDB.Filters;
using CodeArchitects.Platform.Data.MongoDB.Model;
using CodeArchitects.Platform.Data.MongoDB.Model.Implementation;
using CodeArchitects.Platform.Data.MongoDB.Navigation;
using CodeArchitects.Platform.Data.Navigation;
using MongoDB.Driver;
using System.Data;
Expand Down Expand Up @@ -61,14 +62,28 @@ public DataContext(
where TEntity : class
where TKey : IEquatable<TKey>
{
throw IncludeNotSupported<TEntity>();
ValidateInclude(includeAction);

return Find<TEntity, TKey>(key);
}

public Task<TEntity?> FindAsync<TEntity, TKey>(TKey key, IncludeAction<TEntity> includeAction, CancellationToken cancellationToken = default)
where TEntity : class
where TKey : IEquatable<TKey>
{
throw IncludeNotSupported<TEntity>();
ValidateInclude(includeAction);

return FindAsync<TEntity, TKey>(key, cancellationToken);
}

private void ValidateInclude<TEntity>(IncludeAction<TEntity> includeAction)
where TEntity : class
{
if (includeAction is null)
throw new ArgumentNullException(nameof(includeAction));

_ = EnsureEntity<TEntity>();
includeAction(new EmbeddedIncluder<TEntity>(_model));
}

#endregion
Expand All @@ -93,14 +108,14 @@ public void InsertMany<TEntity, TKey>(IEnumerable<TEntity> entities)
where TEntity : class
where TKey : IEquatable<TKey>
{
_stateManager.Execute(InsertManyExecution<TEntity, TKey>(entities), requiresTransaction: true);
ExecuteMany(entities, InsertManyExecution<TEntity, TKey>);
}

public Task InsertManyAsync<TEntity, TKey>(IEnumerable<TEntity> entities, CancellationToken cancellationToken = default)
where TEntity : class
where TKey : IEquatable<TKey>
{
return _stateManager.ExecuteAsync(InsertManyExecution<TEntity, TKey>(entities), requiresTransaction: true, cancellationToken);
return ExecuteManyAsync(entities, InsertManyExecution<TEntity, TKey>, cancellationToken);
}

private Execution InsertExecution<TEntity, TKey>(TEntity entity)
Expand All @@ -118,16 +133,11 @@ private Execution InsertExecution<TEntity, TKey>(TEntity entity)
(session, cancellationToken) => collection.InsertOneAsync(session, entity, cancellationToken: cancellationToken));
}

private Execution InsertManyExecution<TEntity, TKey>(IEnumerable<TEntity> entities)
private Execution InsertManyExecution<TEntity, TKey>(TEntity[] documents)
where TEntity : class
where TKey : IEquatable<TKey>
{
if (entities is null)
throw new ArgumentNullException(nameof(entities));

_ = EnsureEntity<TEntity>();
IMongoCollection<TEntity> collection = _collections.GetCollection<TEntity>();
TEntity[] documents = entities as TEntity[] ?? entities.ToArray();

InsertManyOptions options = new() { IsOrdered = true };

Expand Down Expand Up @@ -158,14 +168,14 @@ public void UpdateMany<TEntity, TKey>(IEnumerable<TEntity> entities)
where TEntity : class
where TKey : IEquatable<TKey>
{
_stateManager.Execute(UpdateManyExecution<TEntity, TKey>(entities), requiresTransaction: true);
ExecuteMany(entities, UpdateManyExecution<TEntity, TKey>);
}

public Task UpdateManyAsync<TEntity, TKey>(IEnumerable<TEntity> entities, CancellationToken cancellationToken = default)
where TEntity : class
where TKey : IEquatable<TKey>
{
return _stateManager.ExecuteAsync(UpdateManyExecution<TEntity, TKey>(entities), requiresTransaction: true, cancellationToken);
return ExecuteManyAsync(entities, UpdateManyExecution<TEntity, TKey>, cancellationToken);
}

private Execution UpdateExecution<TEntity, TKey>(TEntity entity)
Expand All @@ -185,40 +195,25 @@ private Execution UpdateExecution<TEntity, TKey>(TEntity entity)
await collection.ReplaceOneAsync(session, filter, entity, cancellationToken: cancellationToken), entityModel, entity));
}

private Execution UpdateManyExecution<TEntity, TKey>(IEnumerable<TEntity> entities)
private Execution UpdateManyExecution<TEntity, TKey>(TEntity[] documents)
where TEntity : class
where TKey : IEquatable<TKey>
{
if (entities is null)
throw new ArgumentNullException(nameof(entities));

IEntityModel entityModel = EnsureEntity<TEntity>();
IMongoCollection<TEntity> collection = _collections.GetCollection<TEntity>();

// A single BulkWrite instead of N ReplaceOne: one round-trip, and an aggregate
// MatchedCount to verify that every document existed.
ReplaceOneModel<TEntity>[] requests = entities
ReplaceOneModel<TEntity>[] requests = documents
.Select(entity => new ReplaceOneModel<TEntity>(_filters.ByEntity<TEntity, TKey>(entityModel, entity), entity))
.ToArray();

BulkWriteOptions options = new() { IsOrdered = true };

return new Execution(
session =>
{
if (requests.Length == 0)
return;

EnsureAllUpdated(collection.BulkWrite(session, requests, options), requests.Length, entityModel);
},
async (session, cancellationToken) =>
{
if (requests.Length == 0)
return;

EnsureAllUpdated(
await collection.BulkWriteAsync(session, requests, options, cancellationToken), requests.Length, entityModel);
});
session => EnsureAllUpdated(collection.BulkWrite(session, requests, options), requests.Length, entityModel),
async (session, cancellationToken) => EnsureAllUpdated(
await collection.BulkWriteAsync(session, requests, options, cancellationToken), requests.Length, entityModel));
}

#endregion
Expand Down Expand Up @@ -357,12 +352,43 @@ private IEntityModel EnsureEntity<TEntity>()
return entityModel;
}

private static NotSupportedException IncludeNotSupported<TEntity>()
private void ExecuteMany<TEntity>(IEnumerable<TEntity> entities, Func<TEntity[], Execution> execution)
where TEntity : class
{
return new NotSupportedException(
$"The MongoDB provider does not support Find/FindAsync with Include yet (entity '{typeof(TEntity).Name}'). " +
"Intra-aggregate associations are already embedded in the document and need no Include; " +
"for inter-aggregate references, query the target collection explicitly.");
TEntity[] documents = ToBatch(entities);
if (documents.Length == 0)
return;

_stateManager.Execute(execution(documents), requiresTransaction: true);
}

private Task ExecuteManyAsync<TEntity>(
IEnumerable<TEntity> entities,
Func<TEntity[], Execution> execution,
CancellationToken cancellationToken)
where TEntity : class
{
TEntity[] documents = ToBatch(entities);
if (documents.Length == 0)
return Task.CompletedTask;

return _stateManager.ExecuteAsync(execution(documents), requiresTransaction: true, cancellationToken);
}

private TEntity[] ToBatch<TEntity>(IEnumerable<TEntity> entities)
where TEntity : class
{
if (entities is null)
throw new ArgumentNullException(nameof(entities));

_ = EnsureEntity<TEntity>();
TEntity[] documents = [.. entities];

int index = Array.FindIndex(documents, document => document is null);
if (index >= 0)
throw new ArgumentException($"The element at index {index} is null.", nameof(entities));

return documents;
}

/// <summary>
Expand All @@ -373,7 +399,7 @@ private static void EnsureUpdated<TEntity, TKey>(ReplaceOneResult result, IEntit
where TEntity : class
where TKey : IEquatable<TKey>
{
if (result.IsAcknowledged && result.MatchedCount > 0)
if (!result.IsAcknowledged || result.MatchedCount > 0)
return;

throw new DBConcurrencyException(
Expand All @@ -383,7 +409,7 @@ private static void EnsureUpdated<TEntity, TKey>(ReplaceOneResult result, IEntit

private static void EnsureAllUpdated(BulkWriteResult result, int expected, IEntityModel entityModel)
{
if (result.IsAcknowledged && result.MatchedCount == expected)
if (!result.IsAcknowledged || result.MatchedCount == expected)
return;

throw new DBConcurrencyException(
Expand All @@ -396,7 +422,7 @@ private static void EnsureUpserted<TEntity, TKey>(ReplaceOneResult result, IEnti
where TEntity : class
where TKey : IEquatable<TKey>
{
if (result.IsAcknowledged && (result.MatchedCount > 0 || result.UpsertedId is not null))
if (!result.IsAcknowledged || result.MatchedCount > 0 || result.UpsertedId is not null)
return;

throw new DBConcurrencyException(
Expand All @@ -406,7 +432,7 @@ private static void EnsureUpserted<TEntity, TKey>(ReplaceOneResult result, IEnti

private static void EnsureRemoved(DeleteResult result, IEntityModel entityModel, object? key)
{
if (result.IsAcknowledged && result.DeletedCount > 0)
if (!result.IsAcknowledged || result.DeletedCount > 0)
return;

throw new DBConcurrencyException(
Expand Down
Loading
Loading