using System;
using System.Collections.Generic;
using System.Linq;
using System.Linq.Expressions;
using System.Runtime.CompilerServices;
using System.Threading;
using System.Threading.Tasks;
using Dpz.Core.Entity.Base;
using Microsoft.Extensions.Configuration;
using Microsoft.Extensions.DependencyInjection;
using MongoDB.Driver;

namespace Dpz.Core.MongodbAccess;

public sealed class Repository<T> : IRepository<T>
    where T : IBaseEntity
{
    private readonly MongodbAccess<T> _access;

    public Repository(string? connectionString)
    {
        _access = new MongodbAccess<T>(connectionString);
    }

    [ActivatorUtilitiesConstructor]
    public Repository(IConfiguration configuration)
    {
        var connectionString = configuration.GetConnectionString("mongodb");
        _access = new MongodbAccess<T>(connectionString);
    }

    public Repository(IConfiguration configuration, string collectionName)
    {
        var connectionString = configuration.GetConnectionString("mongodb");
        _access = string.IsNullOrEmpty(collectionName)
            ? new MongodbAccess<T>(connectionString)
            : new MongodbAccess<T>(connectionString, collectionName);
    }

    public Repository(string? connectionString, string collectionName)
    {
        _access = string.IsNullOrEmpty(collectionName)
            ? new MongodbAccess<T>(connectionString)
            : new MongodbAccess<T>(connectionString, collectionName);
    }

    public IQueryable<T> MongodbQueryable => _access.MongoQueryable;

    public IMongoCollection<T> Collection => _access.Collection;

    public IQueryable<T> SearchFor(Expression<Func<T, bool>> predicate)
    {
        return _access.MongoQueryable.Where(predicate);
    }

    public IFindFluent<T, T> SearchFor(FilterDefinition<T> filter)
    {
        return Collection.Find(filter);
    }

    public IFindFluent<T, T> SearchFor(FilterDefinition<T> filter, FindOptions options)
    {
        return Collection.Find(filter, options);
    }

    public async IAsyncEnumerable<T> SearchForAsync(
        FilterDefinition<T> filter,
        [EnumeratorCancellation] CancellationToken cancellationToken = default
    )
    {
        var result = await Collection.FindAsync(filter, cancellationToken: cancellationToken);
        while (await result.MoveNextAsync(cancellationToken))
        {
            foreach (var item in result.Current)
            {
                yield return item;
            }
        }
    }

    public async IAsyncEnumerable<T> SearchForAsync(
        FilterDefinition<T> filter,
        FindOptions<T> options,
        [EnumeratorCancellation] CancellationToken cancellationToken = default
    )
    {
        var result = await Collection.FindAsync(filter, options, cancellationToken);
        while (await result.MoveNextAsync(cancellationToken))
        {
            foreach (var item in result.Current)
            {
                yield return item;
            }
        }
    }

    public async Task<T?> FindAsync(object id, CancellationToken cancellationToken = default)
    {
        var filter = MongodbExtensions.GetIdPropertyFilter<T>(id);
        return await (
            await Collection.FindAsync(filter, cancellationToken: cancellationToken)
        ).SingleOrDefaultAsync(cancellationToken: cancellationToken);
    }

    public async Task InsertAsync(T entity, CancellationToken cancellationToken = default)
    {
        await Collection.InsertOneAsync(entity, cancellationToken: cancellationToken);
    }

    public async Task InsertAsync(
        IReadOnlyCollection<T> source,
        CancellationToken cancellationToken = default
    )
    {
        if (source.Count > 0)
        {
            await Collection.InsertManyAsync(source, cancellationToken: cancellationToken);
        }
    }

    public async Task<DeleteResult> DeleteAsync(
        Expression<Func<T, bool>> filter,
        CancellationToken cancellationToken = default
    )
    {
        var result = await Collection.DeleteManyAsync(filter, cancellationToken);
        return result;
    }

    public async Task<DeleteResult> DeleteAsync(
        FilterDefinition<T> filter,
        CancellationToken cancellationToken = default
    )
    {
        var result = await Collection.DeleteManyAsync(filter, cancellationToken);
        return result;
    }

    public async Task<DeleteResult> DeleteAsync(
        object id,
        CancellationToken cancellationToken = default
    )
    {
        var filter = MongodbExtensions.GetIdPropertyFilter<T>(id);
        return await Collection.DeleteOneAsync(filter, cancellationToken);
    }

    public async Task<UpdateResult> UpdateAsync(
        Expression<Func<T, bool>> predicate,
        UpdateDefinition<T> update,
        CancellationToken cancellationToken = default
    )
    {
        var result = await Collection.UpdateManyAsync(
            predicate,
            update,
            cancellationToken: cancellationToken
        );
        return result;
    }

    public async Task<ReplaceOneResult> UpdateAsync(
        T entity,
        CancellationToken cancellationToken = default
    )
    {
        var filter = MongodbExtensions.GetIdPropertyFilter(entity);

        var result = await Collection.ReplaceOneAsync(
            filter,
            entity,
            cancellationToken: cancellationToken
        );
        return result;
    }

    public async Task<BulkWriteResult<T>> UpdateAsync(
        IEnumerable<T> entities,
        CancellationToken cancellationToken = default
    )
    {
        var writes = entities
            .Select(x => new ReplaceOneModel<T>(MongodbExtensions.GetIdPropertyFilter(x), x))
            .ToList();

        return await Collection.BulkWriteAsync(writes, cancellationToken: cancellationToken);
    }

    public IMongoDatabase Database => _access.Database;
}
⚠⚠    以下内容为AI分析的结果,请根据实际情况进行判断。

下面按模块逐项解释这段 Repository 代码的功能与设计要点(T 受限于 IBaseEntity):

总体说明

  • Repository 是一个针对 MongoDB 的泛型仓储实现,封装了对 MongodbAccess(源码未显示,推测是对 MongoDB Driver 的封装)的访问,提供常用的增删改查操作接口。
  • 依赖项包括 MongoDB.Driver、Microsoft.Extensions.Configuration/DependencyInjection 等。
  • 通过多种构造函数支持直接传入连接字符串、从 IConfiguration 中读取连接字符串,或同时指定 collection 名称。
  • 带有 [ActivatorUtilitiesConstructor] 的构造函数方便在 ASP.NET Core DI 容器中被选为注入构造函数。

构造函数与连接/集合选择

  • Repository(string? connectionString):使用传入的 connectionString 创建 MongodbAccess
  • Repository(IConfiguration configuration):从 configuration.GetConnectionString("mongodb") 获取连接字符串并创建 MongodbAccess(带有 ActivatorUtilitiesConstructor,供 DI 使用)。
  • Repository(IConfiguration configuration, string collectionName):从 IConfiguration 取连接字符串;如果 collectionName 非空则将其传入 MongodbAccess,否则使用默认集合名字。
  • Repository(string? connectionString, string collectionName):同上但直接用传入的 connectionString 和 collectionName。
  • MongodbAccess 负责实际的 MongoDB 连接/集合初始化(本代码未给出)。

公开属性

  • MongodbQueryable:返回 _access.MongoQueryable,允许通过 LINQ 对集合进行查询(基于 MongoDB LINQ 支持)。
  • Collection:返回 IMongoCollection,直接暴露底层集合以便运行原生 MongoDB.Driver 操作。
  • Database:返回 IMongoDatabase(来自 _access),便于执行数据库级操作。

查询方法

  • SearchFor(Expression<Func<T, bool>> predicate):基于 LINQ 表达式返回 IQueryable(延迟查询,交给驱动执行)。
  • SearchFor(FilterDefinition filter):返回 Find 的 IFindFluent<T,T>,可继续链式构造查询(例如 Sort/Limit)。
  • SearchFor(FilterDefinition filter, FindOptions options):带 FindOptions 的 Find,返回 IFindFluent。
  • SearchForAsync(FilterDefinition filter, CancellationToken):异步遍历 FindAsync 结果并通过 IAsyncEnumerable yield 返回所有文档(每批 MoveNextAsync 后 yield 当前批次元素)。
  • SearchForAsync(FilterDefinition filter, FindOptions options, CancellationToken):带 FindOptions 的异步枚举。

单条查找

  • FindAsync(object id, CancellationToken):使用 MongodbExtensions.GetIdPropertyFilter(id) 构建基于实体 id 的过滤器,然后异步查找并返回 SingleOrDefault(找到 1 条返回实体,0 条返回 null;若多条则抛异常)。

插入

  • InsertAsync(T entity):InsertOneAsync 插入单个文档。
  • InsertAsync(IReadOnlyCollection source):如果 source.Count>0 则 InsertManyAsync 批量插入(注意:如果 source 为 null,会抛 NRE;只有在 Count>0 才调用 InsertMany)。

删除

  • DeleteAsync(Expression<Func<T, bool>> filter):DeleteManyAsync,接受表达式筛选器。
  • DeleteAsync(FilterDefinition filter):DeleteManyAsync,接受 FilterDefinition。
  • DeleteAsync(object id):用 id 构建过滤器并 DeleteOneAsync 删除单个文档。

更新 / 替换

  • UpdateAsync(Expression<Func<T, bool>> predicate, UpdateDefinition update):UpdateManyAsync,基于 UpdateDefinition 做部分更新($set 等)。
  • UpdateAsync(T entity):通过 MongodbExtensions.GetIdPropertyFilter(entity) 找到文档并用 ReplaceOneAsync 完整替换该文档(整文档替换)。
  • UpdateAsync(IEnumerable entities):对传入实体集合创建 ReplaceOneModel 列表并 BulkWriteAsync 批量替换(按每个实体的 id 作为筛选条件)。

返回类型(重要)

  • DeleteResult / UpdateResult / ReplaceOneResult / BulkWriteResult:来自 MongoDB.Driver,包含操作影响的结果(已删除/已更新/匹配计数等)。
  • IFindFluent<T,T>:可继续链式操作然后 ToList/First 等。
  • IAsyncEnumerable:支持异步流式消费。

依赖的辅助项

  • MongodbExtensions.GetIdPropertyFilter(...):用于基于实体或 id 构建 id 过滤器(具体实现未给出,但关键用于按照主键查找/删除/替换)。
  • IConfiguration.GetConnectionString("mongodb"):约定从配置中读取名为 "mongodb" 的连接字符串。

注意事项与改进建议(可选)

  • 插入的 IList/集合重载没有对 null 做保护(如果可能传入 null,可先检查以免 NRE)。
  • FindAsync 返回 SingleOrDefaultAsync,若存在多条匹配会抛出异常,调用方应保证 id 唯一或使用 FirstOrDefault 视需求。
  • 当前代码未处理异常或重试策略,生产环境可能需要增加错误处理与日志、索引检查、超时配置等。
  • 对 UpdateAsync(IEnumerable) 的 BulkWrite 没有使用 upsert 标记(若需要可配置 ReplaceOneModel 的 Upsert = true)。
  • 并发/线程安全:MongoDB .NET Driver 的 IMongoCollection/客户端是线程安全的,通常可以复用 Repository 实例。

总结

  • 这是一个围绕 MongoDB 的通用仓储类,封装了常用的 CRUD、批量操作以及基于 LINQ 的查询入口,支持通过配置或直接传参初始化连接与集合,方便在应用中以统一接口操作 MongoDB 文档实体。
评论加载中...