using Dpz.Core.Backup;
using Dpz.Core.Entity.Base;
using Dpz.Core.Infrastructure;
using Dpz.Core.Public.Entity;
using Dpz.Core.Public.ViewModel;
using Dpz.Core.Service.RepositoryService;
using Hangfire;
using JetBrains.Annotations;
using MongoDB.Bson.Serialization;
using MongoDB.Driver;

namespace Dpz.Core.Web.Jobs.Hangfire;

[UsedImplicitly]
public class BackupActivator(
    IBackupService backupService,
    IBackupRecordService backupRecordService,
    IConfiguration configuration,
    //ISafeFileService safeFileService,
    ILogger<BackupActivator> logger
) : JobActivator
{
    [ProlongExpirationTime]
    public async Task BackupAsync()
    {
        var filter = new Dictionary<string, string>
        {
            {
                nameof(MessageOutboxRecord),
                SerializerFilterDocument(
                    new ExpressionFilterDefinition<MessageOutboxRecord>(x =>
                        x.Status != OutboxMessageStatus.Consumed
                    )
                )
            },
        };

        try
        {
            await BackupMainAsync("mongodb", filter);
            await BackupMainAsync("AgileConfig");
        }
        catch (Exception e)
        {
            logger.LogError(e, "backup fail");
        }
    }

    private static string SerializerFilterDocument<T>(FilterDefinition<T> filter)
        where T : IBaseEntity
    {
        var serializerRegistry = BsonSerializer.SerializerRegistry;
        var documentSerializer = serializerRegistry.GetSerializer<T>();

        var serializerArgs = new RenderArgs<T>(documentSerializer, serializerRegistry);
        var result = filter.Render(serializerArgs);

        return result.ToString();
    }

    private async Task BackupMainAsync(
        string connectionStringName,
        Dictionary<string, string>? filter = null
    )
    {
        logger.LogInformation("Start backing up database: {Database}", connectionStringName);
        var connectionString = configuration.GetConnectionString(connectionStringName);
        backupService.UploadAsync = UploadAsync;
        var backupRecord = await backupService.BackupAsync(connectionString, filter);
        await backupRecordService.AddRecordAsync(backupRecord);
        logger.LogInformation("Database: {Database} backup complete", connectionStringName);
    }

    private async Task<VmBackupRecord?> UploadAsync(BackupResult result)
    {
        if (string.IsNullOrEmpty(result.BackupPath))
        {
            throw new BusinessException("backup path is null");
        }

        if (!File.Exists(result.BackupPath))
        {
            throw new BusinessException("backup file not exist");
        }
        var backupRootPath = configuration.GetValue<string>("BackupRootPath");
        if (string.IsNullOrEmpty(backupRootPath) || !Directory.Exists(backupRootPath))
        {
            throw new BusinessException("backup root path is null");
        }

        var date = DateTime.Now;
        var backupPath = new List<string> { "db", date.Year.ToString(), date.Month.ToString() };
        var fileName = Path.GetFileName(result.BackupPath);
#if DEBUG
        backupPath.Insert(1, "Test");
#endif
        // var cloudFile = new CloudFile
        // {
        //     PathToFile = pathToFile,
        //     Stream = new FileStream(
        //         result.BackupPath,
        //         FileMode.Open,
        //         FileAccess.Read,
        //         FileShare.Read,
        //         1 << 12
        //     ),
        // };
        // var address = await safeFileService.UploadFileForFtpAsync(cloudFile);
        // if (address == null)
        // {
        //     throw new BusinessException("upload backup file fail");
        // }

        var backupPathDir = Path.Combine(backupRootPath, Path.Combine(backupPath.ToArray()));
        var destination = Path.Combine(backupPathDir, fileName);
        try
        {
            if (!Directory.Exists(backupPathDir))
            {
                Directory.CreateDirectory(backupPathDir);
            }
            destination = Path.Combine(backupPathDir, fileName);
            File.Move(result.BackupPath, destination);
            await using var writer = File.CreateText(destination + "__password.txt");
            await writer.WriteAsync(result.ZipPassword);
            await writer.FlushAsync();
            writer.Close();
        }
        catch (Exception e)
        {
            logger.LogError(
                e,
                "move backup file fail, source: {Source}, destination: {Destination}",
                result.BackupPath,
                destination
            );
            throw;
        }

        return new VmBackupRecord
        {
            Filename = destination,
            BackupTime = DateTime.Now,
            Database = result.Database,
            BackupPassword = result.ZipPassword,
        };
    }
}
⚠⚠    以下内容为AI分析的结果,请根据实际情况进行判断。

代码解释

这是一个用于数据库备份的 Hangfire 后台任务激活器类,主要功能是定期备份数据库并将备份文件保存到指定位置。

类结构

主要依赖注入

  • IBackupService: 备份服务
  • IBackupRecordService: 备份记录服务
  • IConfiguration: 配置服务
  • ILogger: 日志记录器

核心方法

1. BackupAsync() - 主备份方法

[ProlongExpirationTime]
public async Task BackupAsync()
  • 使用 [ProlongExpirationTime] 特性延长 Hangfire 任务过期时间
  • 定义过滤器字典,针对 MessageOutboxRecord 表只备份状态不为 Consumed 的记录
  • 依次备份两个数据库:
    • "mongodb" - 应用过滤器
    • "AgileConfig" - 不应用过滤器
  • 异常统一捕获并记录日志

2. SerializerFilterDocument<T>() - 过滤器序列化

private static string SerializerFilterDocument<T>(FilterDefinition<T> filter)
  • 将 MongoDB 的 FilterDefinition 对象序列化为字符串
  • 用于将过滤条件持久化或传递

3. BackupMainAsync() - 执行备份逻辑

private async Task BackupMainAsync(string connectionStringName, Dictionary<string, string>? filter = null)
  • 从配置中获取连接字符串
  • 设置备份服务的上传回调方法
  • 执行备份并保存备份记录
  • 记录开始和完成日志

4. UploadAsync() - 处理备份文件

private async Task<VmBackupRecord?> UploadAsync(BackupResult result)

主要逻辑:

  1. 验证备份文件

    • 检查备份路径是否为空
    • 检查备份文件是否存在
    • 验证备份根目录配置
  2. 构建目标路径

    {BackupRootPath}/db/{Year}/{Month}/{FileName}
    
    • DEBUG 模式下会插入 "Test" 子目录
  3. 文件操作

    • 创建目标目录(如不存在)
    • 移动备份文件到目标位置
    • 创建密码文件({FileName}__password.txt)保存压缩密码
  4. 返回备份记录

    • 返回包含文件名、备份时间、数据库名、密码的 VmBackupRecord 对象

注释掉的代码

代码中包含被注释的 FTP 上传逻辑,说明可能之前使用云存储,现在改为本地文件系统存储。

关键特性

  • ✅ 支持条件过滤备份(可选择性备份某些表的部分数据)
  • ✅ 自动按年月组织备份文件
  • ✅ 备份文件带密码保护,密码单独存储
  • ✅ 完整的异常处理和日志记录
  • ✅ 使用 Hangfire 进行后台任务调度

使用场景

适用于需要定期自动备份 MongoDB 和 AgileConfig 配置数据库的应用系统。

评论加载中...