using System.Collections.Concurrent;
using System.Diagnostics;
using System.Xml.Linq;
using Dpz.Core.AspNetCore;
using Dpz.Core.EnumLibrary;
using Dpz.Core.EnumLibrary.Pipeline;
using Dpz.Core.Public.Entity.Pipeline;
using Dpz.Core.Public.ViewModel;
using Dpz.Core.Service.Mediator.Features.Pipeline.Commands;
using Dpz.Core.Web.Jobs.Hangfire.Build;
using Dpz.Core.Web.Jobs.Hubs;
using Hangfire;
using JetBrains.Annotations;
using Mediator;
using Microsoft.AspNetCore.SignalR;
using Renci.SshNet;
using ZiggyCreatures.Caching.Fusion;

namespace Dpz.Core.Web.Jobs.Hangfire;

[UsedImplicitly]
public class BuildActivator(
    ILogger<BuildActivator> logger,
    IFusionCache fusionCache,
    IHubContext<PipelineHub> hubContext,
    SshConfigurationService sshConfigurationService,
    PipelineVersionProbe versionProbe,
    IMediator mediator,
    IConfiguration configuration
) : JobActivator
{
    public const string BuildLockKey = "Dpz.Core.Web.Jobs.Hangfire.BuildActivator.BuildLock";

    private PipelineStepCollector? _steps;

    public async Task RunBuildAsync(
        BuildOption option,
        string runId,
        VmUserInfo? triggeredBy,
        CancellationToken cancellationToken
    )
    {
        var sshConfiguration = sshConfigurationService.GetSshConfiguration();
        var validateResult = sshConfigurationService.Validate(sshConfiguration);
        if (!validateResult.Success)
        {
            logger.LogWarning(
                "ssh 配置校验失败,已退出本次发布:{Message}",
                validateResult.Message
            );
            await hubContext
                .Clients.Group(runId)
                .SendAsync(
                    "BuildFail",
                    new { runId, message = validateResult.Message },
                    cancellationToken
                );
            return;
        }

        var cache = await fusionCache.TryGetAsync<string>(BuildLockKey, token: cancellationToken);
        if (cache.HasValue && !string.IsNullOrWhiteSpace(cache.Value))
        {
            logger.LogWarning("正在运行实例:{RunId},已退出本次发布", cache.Value);
            return;
        }

        await fusionCache.SetAsync(BuildLockKey, runId, token: cancellationToken);

        logger.LogInformation("正在运行实例:{RunId}", runId);
        var historyRequest = CreateHistoryRequest(option, runId, triggeredBy);
        _steps = new PipelineStepCollector(sshConfiguration!.Build.Host);
        try
        {
            await SendStageChangedAsync(runId, PipelineStage.PullCode, cancellationToken);
            using var client = await ConnectBuildServerAndPullCodeAsync(
                runId,
                sshConfiguration.Build,
                cancellationToken
            );

            var stopwatch = new Stopwatch();
            stopwatch.Start();

            var projectVersions = new Dictionary<BuildProject, string>();
            if (!string.IsNullOrWhiteSpace(option.Tag))
            {
                foreach (var value in Enum.GetValues<BuildProject>())
                {
                    if (value == BuildProject.None)
                    {
                        continue;
                    }

                    if ((option.Project & value) == value)
                    {
                        projectVersions.Add(value, option.Tag);
                    }
                }
            }
            else
            {
                projectVersions = await ReadProjectVersionsAsync(
                    option.Project,
                    sshConfiguration.Build,
                    cancellationToken
                );
            }

            stopwatch.Stop();
            logger.LogInformation(
                "读取版本号成功,耗时:{ReadVersion}ms",
                stopwatch.ElapsedMilliseconds
            );
            historyRequest.Versions = projectVersions
                .Select(x => new PipelineVersion { Project = x.Key, Version = x.Value })
                .ToList();

            await WaitForSameVersionConfirmationAsync(runId, projectVersions, cancellationToken);

            // 开发环境跳过真实的构建、发布
            if (configuration.IsDevelopment)
            {
                await DevelopmentSettingsAsync(runId, historyRequest, cancellationToken);
                return;
            }

            stopwatch.Restart();

            await SendStageChangedAsync(runId, PipelineStage.BuildImages, cancellationToken);

            await BuildPushImagesAsync(
                runId,
                client,
                option.Project,
                sshConfiguration.Build,
                projectVersions,
                cancellationToken
            );

            stopwatch.Stop();
            logger.LogInformation(
                "镜像构建并推送成功,耗时:{BuildPushTime}ms",
                stopwatch.ElapsedMilliseconds
            );

            stopwatch.Restart();
            await SendStageChangedAsync(runId, PipelineStage.Publish, cancellationToken);
            await PublishAsync(
                option.Project,
                sshConfiguration.Servers,
                runId,
                projectVersions,
                cancellationToken
            );
            stopwatch.Stop();
            logger.LogInformation(
                "发布、部署成功,耗时:{DeploymentTime}ms",
                stopwatch.ElapsedMilliseconds
            );
            historyRequest.Status = PipelineStatus.Success;
            await hubContext
                .Clients.Group(runId)
                .SendAsync("BuildCompleted", new { runId }, cancellationToken);
        }
        catch (PipelineVersionNotConfirmedException e)
        {
            historyRequest.Status = PipelineStatus.Abort;
            historyRequest.Error = e.Message;
            await hubContext
                .Clients.Group(runId)
                .SendAsync("BuildFail", new { runId, message = e.Message }, cancellationToken);
        }
        catch (OperationCanceledException)
        {
            historyRequest.Status = PipelineStatus.Abort;
            historyRequest.Error = "CI/CD被中止";
            await hubContext
                .Clients.Group(runId)
                .SendAsync(
                    "BuildFail",
                    new { runId, message = historyRequest.Error },
                    cancellationToken
                );
        }
        catch (Exception e) when (e is not OperationCanceledException)
        {
            await hubContext
                .Clients.Group(runId)
                .SendCoreAsync("BuildFail", [new { runId }], cancellationToken);
            logger.LogError(e, "build fail");
            historyRequest.Status = PipelineStatus.Failure;
            historyRequest.Error = e.Message;
        }
        finally
        {
            // 补充结束时间和总耗时
            var endTime = DateTime.Now;
            historyRequest.EndTime = endTime;
            historyRequest.DurationMs = Math.Max(
                0,
                (long)(endTime - historyRequest.StartTime).TotalMilliseconds
            );
            historyRequest.Steps =
                _steps?.FinalizeSteps(historyRequest.Status, historyRequest.Error) ?? [];
            await SaveHistoryAsync(historyRequest);
            await fusionCache.RemoveAsync(BuildLockKey, token: cancellationToken);
            await fusionCache.RemoveAsync(
                PipelineVersionProbe.WaitCacheKey(runId),
                token: cancellationToken
            );
            await fusionCache.RemoveAsync(
                PipelineVersionProbe.DecisionCacheKey(runId),
                token: cancellationToken
            );
        }
    }

    /// <summary>
    /// 开发环境跳过构建、推送、发布、部署
    /// </summary>
    private async Task DevelopmentSettingsAsync(
        string runId,
        CreatePipelineHistoryRequest historyRequest,
        CancellationToken cancellationToken
    )
    {
        await SendStageChangedAsync(runId, PipelineStage.BuildImages, cancellationToken);
        await Task.Delay(3000, cancellationToken);
        await SendStageChangedAsync(runId, PipelineStage.Publish, cancellationToken);
        await Task.Delay(1200, cancellationToken);

        var random = new Random();
        var value = random.Next(0, 3);
        if (value % 2 != 0)
        {
            historyRequest.Status = PipelineStatus.Failure;
            historyRequest.Error = "CI/CD 随机到错误";
        }
        historyRequest.Status = PipelineStatus.Success;

        await hubContext
            .Clients.Group(runId)
            .SendAsync("BuildCompleted", new { runId }, cancellationToken);
    }

    /// <summary>
    /// 即将发布版本与线上一致时,等待用户二次确认后再继续。
    /// </summary>
    private async Task WaitForSameVersionConfirmationAsync(
        string runId,
        Dictionary<BuildProject, string> projectVersions,
        CancellationToken cancellationToken
    )
    {
        var matches = await versionProbe.FindSameVersionsAsync(projectVersions, cancellationToken);
        if (matches.Count == 0)
        {
            return;
        }

        var waitKey = PipelineVersionProbe.WaitCacheKey(runId);
        var decisionKey = PipelineVersionProbe.DecisionCacheKey(runId);
        var waitDuration = TimeSpan.FromMinutes(3);
        await fusionCache.SetAsync(
            waitKey,
            matches,
            options => options.SetDuration(waitDuration),
            token: cancellationToken
        );
        await fusionCache.RemoveAsync(decisionKey, token: cancellationToken);

        var timeoutSeconds = (int)PipelineVersionProbe.ConfirmTimeout.TotalSeconds;
        await hubContext
            .Clients.Group(runId)
            .SendAsync(
                "VersionConfirmationRequired",
                new
                {
                    runId,
                    timeoutSeconds,
                    matches,
                },
                cancellationToken
            );
        logger.LogInformation(
            "线上版本与即将发布版本一致,等待确认:{RunId},匹配 {Count} 个项目",
            runId,
            matches.Count
        );

        var deadline = DateTime.Now.Add(PipelineVersionProbe.ConfirmTimeout);
        while (DateTime.Now < deadline)
        {
            cancellationToken.ThrowIfCancellationRequested();
            var decision = await fusionCache.TryGetAsync<string>(
                decisionKey,
                token: cancellationToken
            );
            if (decision.HasValue && !string.IsNullOrWhiteSpace(decision.Value))
            {
                await fusionCache.RemoveAsync(waitKey, token: cancellationToken);
                if (
                    string.Equals(
                        decision.Value,
                        PipelineVersionProbe.DecisionConfirm,
                        StringComparison.OrdinalIgnoreCase
                    )
                )
                {
                    logger.LogInformation("用户已确认同版本发布:{RunId}", runId);
                    return;
                }

                throw new PipelineVersionNotConfirmedException("同版本发布未确认");
            }

            await Task.Delay(400, cancellationToken);
        }

        logger.LogWarning("同版本发布确认超时:{RunId}", runId);
        await fusionCache.RemoveAsync(waitKey, token: cancellationToken);
        throw new PipelineVersionNotConfirmedException("同版本发布未确认");
    }

    /// <summary>
    /// 获取项目对应的构建命令
    /// </summary>
    private static List<(BuildProject Project, SshCommand Command)> GetBuildProjectImageCommands(
        BuildProject project,
        SshClient client,
        SshBuildServer server,
        Dictionary<BuildProject, string> projectVersions
    )
    {
        var cdWorkspace = $"cd {server.Workspace.TrimEnd('/')}/src && \\\n";

        return BuiltInCommands
            .LazyBuildCommands.Value.Where(x => (project & x.Key) == x.Key)
            .Select(x =>
            {
                var commandText =
                    $"{cdWorkspace}{x.Value.Replace("<Version>", projectVersions[x.Key])}";
                commandText = WrapBuildCommandWithSudoPassword(commandText, server.Password);
                return (x.Key, client.CreateCommand(commandText));
            })
            .ToList();
    }

    /// <summary>
    /// SSH 命令没有 TTY,sudo 无法交互输入密码。用构建机 SSH 密码走 sudo -S。
    /// </summary>
    private static string WrapBuildCommandWithSudoPassword(string command, string password)
    {
        if (string.IsNullOrEmpty(password))
        {
            return command;
        }

        return $"printf '%s\\n' {QuoteForSingleQuotedShell(password)} | sudo -S -p '' bash -lc "
            + QuoteForSingleQuotedShell(command);
    }

    private static string QuoteForSingleQuotedShell(string value)
    {
        return $"'{value.Replace("'", "'\\''", StringComparison.Ordinal)}'";
    }

    /// <summary>
    /// 在远程SSH 发布、部署应用
    /// </summary>
    private async Task PublishAsync(
        BuildProject project,
        Dictionary<string, SshDeploymentServer> servers,
        string runId,
        Dictionary<BuildProject, string> projectVersions,
        CancellationToken cancellationToken
    )
    {
        var publishTasks = servers.Select(PublishServerAsync).ToArray();
        await Task.WhenAll(publishTasks);

        async Task PublishServerAsync(KeyValuePair<string, SshDeploymentServer> server)
        {
            using var client = new SshClient(
                host: server.Value.Host,
                port: server.Value.Port,
                username: server.Value.Username,
                password: server.Value.Password
            );

            await client.ConnectAsync(cancellationToken);

            foreach (
                var (cmdProject, commandText) in BuiltInCommands.LazyPublishCommands.Value.Where(
                    x => (server.Value.Project & project & x.Key) == x.Key
                )
            )
            {
                cancellationToken.ThrowIfCancellationRequested();
                if (!projectVersions.TryGetValue(cmdProject, out var version))
                {
                    throw new InvalidOperationException(
                        $"当前发布、部署项目{cmdProject}的版本号不存在"
                    );
                }

                var commandTextAs = commandText.Replace("<Version>", version);
                using var command = client.CreateCommand(commandTextAs);
                await ExecuteWithOutputAsync(
                    command,
                    runId,
                    server.Key,
                    cmdProject,
                    cancellationToken
                );

                if (command.ExitStatus is not 0)
                {
                    throw new InvalidOperationException(
                        $"Publish command failed on {server.Key} for {cmdProject}. "
                            + $"Exit status: {command.ExitStatus}. Error: {command.Error}"
                    );
                }

                logger.LogInformation(
                    "发布服务器 {Server} 的项目 {Project} 成功",
                    server.Key,
                    cmdProject
                );
            }
        }
    }

    /// <summary>
    /// 执行远程命令,并在执行期间实时推送终端输出
    /// </summary>
    private async Task ExecuteWithOutputAsync(
        SshCommand command,
        string runId,
        string serverName,
        BuildProject project,
        CancellationToken cancellationToken
    )
    {
        var outputTask = PublishTerminalOutputAsync(
            command.OutputStream,
            runId,
            serverName,
            project,
            cancellationToken
        );
        var errorTask = PublishTerminalOutputAsync(
            command.ExtendedOutputStream,
            runId,
            serverName,
            project,
            cancellationToken
        );

        try
        {
            await command.ExecuteAsync(cancellationToken);
        }
        finally
        {
            await Task.WhenAll(outputTask, errorTask);
        }
    }

    /// <summary>
    /// 逐行读取 SSH 终端输出,并推送到 SignalR 指定分组
    /// </summary>
    private async Task PublishTerminalOutputAsync(
        Stream stream,
        string runId,
        string serverName,
        BuildProject project,
        CancellationToken cancellationToken
    )
    {
        using var reader = new StreamReader(stream);
        while (await reader.ReadLineAsync(cancellationToken) is { } line)
        {
            if (string.IsNullOrWhiteSpace(line))
            {
                continue;
            }
            _steps?.AppendLine(serverName, project, line);
            await PublishMessageAsync(
                new BuildMessage(runId, serverName, project, line),
                cancellationToken
            );
        }
    }

    private async Task PublishMessageAsync(
        BuildMessage message,
        CancellationToken cancellationToken
    )
    {
        await hubContext
            .Clients.Group(message.RunId)
            .SendAsync("ReceiveBuildMessage", message, cancellationToken);
    }

    /// <summary>
    /// 推送构建阶段变更事件
    /// </summary>
    private async Task SendStageChangedAsync(
        string runId,
        PipelineStage stage,
        CancellationToken cancellationToken
    )
    {
        await hubContext
            .Clients.Group(runId)
            .SendAsync("BuildStageChanged", new { runId, stage = (int)stage }, cancellationToken);
    }

    /// <summary>
    /// 连接构建服务器并拉取代码。拉取过程的输出不属于任何单一项目,
    /// 以 (BuildProject)0 标记,前端将其归入「总览」面板
    /// </summary>
    private async Task<SshClient> ConnectBuildServerAndPullCodeAsync(
        string runId,
        SshBuildServer server,
        CancellationToken cancellationToken
    )
    {
        var stopwatch = new Stopwatch();
        stopwatch.Start();
        var client = new SshClient(
            host: server.Host,
            port: server.Port,
            username: server.Username,
            password: server.Password
        );

        await client.ConnectAsync(cancellationToken);
        client.KeepAliveInterval = TimeSpan.FromSeconds(30);
        stopwatch.Stop();

        logger.LogInformation(
            "连接编译服务器:{Host}成功,耗时:{ConnectedTime}ms",
            server.Host,
            stopwatch.ElapsedMilliseconds
        );

        stopwatch.Restart();
        using var pullCommand = client.CreateCommand($"cd {server.Workspace} && git pull");
        await ExecuteWithOutputAsync(
            pullCommand,
            runId,
            server.Host,
            BuildProject.None,
            cancellationToken
        );
        if (pullCommand.ExitStatus is not 0)
        {
            throw new InvalidOperationException(
                $"git pull fail."
                    + $"Exit status: {pullCommand.ExitStatus}. Error: {pullCommand.Error}"
            );
        }

        stopwatch.Stop();
        logger.LogInformation(
            "编译服务器拉取代码成功,耗时:{PullCodeTime}ms",
            stopwatch.ElapsedMilliseconds
        );
        return client;
    }

    /// <summary>
    /// 构建并推送镜像
    /// </summary>
    private async Task BuildPushImagesAsync(
        string runId,
        SshClient client,
        BuildProject project,
        SshBuildServer server,
        Dictionary<BuildProject, string> projectVersions,
        CancellationToken cancellationToken
    )
    {
        var commands = GetBuildProjectImageCommands(project, client, server, projectVersions);
        foreach (var (cmdProject, buildCmd) in commands)
        {
            using var cmdAs = buildCmd;
            try
            {
                await ExecuteWithOutputAsync(
                    cmdAs,
                    runId,
                    server.Host,
                    cmdProject,
                    cancellationToken
                );
                if (cmdAs.ExitStatus is not 0)
                {
                    throw new InvalidOperationException(
                        $"Build or push command failed for {project}. "
                            + $"Exit status: {cmdAs.ExitStatus}. Error: {cmdAs.Error}"
                    );
                }
            }
            catch (Exception e)
            {
                logger.LogError(e, "execute build images fail");
                throw;
            }
        }
    }

    private async Task<Dictionary<BuildProject, string>> ReadProjectVersionsAsync(
        BuildProject project,
        SshBuildServer server,
        CancellationToken cancellationToken
    )
    {
        var projectWorkspace = new Dictionary<BuildProject, string>
        {
            {
                BuildProject.DpzCore,
                $"{server.Workspace.TrimEnd('/')}/src/Dpz.Core.Web/Dpz.Core.Web.csproj"
            },
            {
                BuildProject.DpzWebApi,
                $"{server.Workspace.TrimEnd('/')}/src/Dpz.Core.WebApi/Dpz.Core.WebApi.csproj"
            },
            {
                BuildProject.DpzCoreAuth,
                $"{server.Workspace.TrimEnd('/')}/src/Dpz.Core.Auth/Dpz.Core.Auth.csproj"
            },
        };

        var versions = new ConcurrentDictionary<BuildProject, string>();
        await Parallel.ForEachAsync(
            projectWorkspace.Where(x => (project & x.Key) == x.Key),
            cancellationToken,
            async (x, ct) =>
            {
                using var client = new SftpClient(
                    server.Host,
                    server.Port,
                    server.Username,
                    server.Password
                );
                await client.ConnectAsync(ct);
                using var stream = new MemoryStream();
                await client.DownloadFileAsync(x.Value, stream, ct);
                stream.Position = 0;
                var xmlRoot = await XElement.LoadAsync(stream, LoadOptions.None, ct);
                var version = xmlRoot.Element("PropertyGroup")?.Element("Version")?.Value;
                if (string.IsNullOrWhiteSpace(version))
                {
                    throw new InvalidOperationException($"未读到{x.Key}项目的版本号");
                }
                versions.TryAdd(x.Key, version);
            }
        );
        return versions.ToDictionary();
    }

    /// <summary>
    /// 创建本次流水线入参
    /// </summary>
    private static CreatePipelineHistoryRequest CreateHistoryRequest(
        BuildOption option,
        string runId,
        VmUserInfo? triggeredBy
    )
    {
        var startTime = DateTime.Now;
        var tag = string.IsNullOrWhiteSpace(option.Tag) ? null : option.Tag.Trim();
        return new CreatePipelineHistoryRequest
        {
            RunId = runId,
            StartTime = startTime,
            EndTime = startTime,
            DurationMs = 0,
            Status = PipelineStatus.Failure,
            Project = option.Project,
            VersionMode = tag is null ? PipelineVersionMode.Auto : PipelineVersionMode.Manual,
            Tag = tag,
            TriggeredBy = triggeredBy,
        };
    }

    /// <summary>
    /// 保存本次流水线记录
    /// </summary>
    private async Task SaveHistoryAsync(CreatePipelineHistoryRequest request)
    {
        try
        {
            var result = await mediator.Send(request, CancellationToken.None);
            if (!result.Success)
            {
                logger.LogWarning("CI/CD 历史记录保存失败:{Message}", result.Message);
            }
        }
        catch (Exception e)
        {
            logger.LogError(e, "CI/CD 历史记录保存异常");
        }
    }

    /// <summary>
    /// 按 (服务器, 项目) 收集终端输出行,构建历史记录的步骤
    /// </summary>
    private sealed class PipelineStepCollector(string buildHost)
    {
        /// <summary>
        /// 每个步骤保存的最大输出行数,与前端日志面板上限一致
        /// </summary>
        private const int MaxLinesPerStep = 2000;

        private readonly ConcurrentDictionary<
            (string Server, BuildProject Project),
            PipelineStep
        > _steps = new();

        public void AppendLine(string server, BuildProject project, string line)
        {
            var step = _steps.GetOrAdd(
                (server, project),
                key => new PipelineStep
                {
                    Stage = ResolveStage(key.Server, key.Project),
                    Server = key.Server,
                    Project = key.Project,
                    Status = PipelineStatus.Success,
                    StartTime = DateTime.Now,
                }
            );
            lock (step)
            {
                step.Lines.Add(line);
                if (step.Lines.Count > MaxLinesPerStep)
                {
                    step.Lines.RemoveAt(0);
                }
            }
        }

        /// <summary>
        /// 步骤收尾:补充结束时间和耗时,非成功时统一标记
        /// </summary>
        public List<PipelineStep> FinalizeSteps(PipelineStatus status, string? error)
        {
            var endTime = DateTime.Now;
            var steps = _steps.Values.ToList();
            foreach (var step in steps)
            {
                var startTime = step.StartTime ?? endTime;
                step.EndTime = endTime;
                step.DurationMs = Math.Max(0, (long)(endTime - startTime).TotalMilliseconds);
                if (status != PipelineStatus.Success)
                {
                    step.Status = status;
                    step.Error = error;
                }
            }
            return steps.OrderBy(x => x.StartTime).ToList();
        }

        private PipelineStage ResolveStage(string server, BuildProject project)
        {
            if (server != buildHost)
            {
                return PipelineStage.Publish;
            }
            return project == BuildProject.None
                ? PipelineStage.PullCode
                : PipelineStage.BuildImages;
        }
    }
}
⚠⚠    以下内容为AI分析的结果,请根据实际情况进行判断。

这段代码定义了一个名为 BuildActivator 的类,继承自 Hangfire 的 JobActivator,用于执行整个 CI/CD 发布流水线。它由 ASP.NET Core 依赖注入构造,集成了 SignalR、SSH/SFTP、FusionCache、Mediator、Hangfire 等组件。


1. 类的主要职责

BuildActivator 是核心的“构建执行器”,其入口方法 RunBuildAsync 负责按顺序完成:

  • SSH 配置校验
  • 全局发布锁控制
  • 拉取代码
  • 读取各项目版本号
  • 处理“版本与线上一致”时的用户确认
  • 构建并推送 Docker 镜像
  • SSH 远程发布部署
  • 保存流水线历史记录
  • 通过 SignalR 实时向客户端推送进度、日志

2. 构造函数与依赖

构造函数采用主构造函数语法,注入以下依赖:

依赖用途
ILogger<BuildActivator>日志记录
IFusionCache分布式/进程内缓存,用于锁和版本确认状态
IHubContext<PipelineHub>向 SignalR 分组发送消息
SshConfigurationService获取和校验 SSH 配置
PipelineVersionProbe探测线上版本,判断是否需要二次确认
IMediator发送保存历史记录的命令
IConfiguration获取环境配置(如是否开发环境)

3. 全局构建锁

public const string BuildLockKey = "Dpz.Core.Web.Jobs.Hangfire.BuildActivator.BuildLock";

RunBuildAsync 开始时:

  • 尝试从缓存读取 BuildLockKey
  • 如果已有值,则说明已有构建在运行,直接返回
  • 否则将当前 runId 写入缓存,获得锁

finally 中移除锁,从而保证同一时间只有一个 CI/CD 任务在执行。


4. 构建主流程 RunBuildAsync

4.1 校验 SSH 配置

var sshConfiguration = sshConfigurationService.GetSshConfiguration();
var validateResult = sshConfigurationService.Validate(sshConfiguration);

如果配置不合法,直接推送 BuildFail 消息到 SignalR 分组。

4.2 检查全局锁

如果有其他实例正在执行,记录日志并退出。

4.3 初始化历史记录和步骤收集器

var historyRequest = CreateHistoryRequest(option, runId, triggeredBy);
_steps = new PipelineStepCollector(sshConfiguration!.Build.Host);

historyRequest 记录本次发布的基本信息;_steps 用于按“服务器 + 项目”收集日志行。

4.4 拉取代码

  • 发送阶段事件 PipelineStage.PullCode
  • 连接构建服务器,在其中执行 git pull
  • git pull 退出码非 0,抛出异常

4.5 获取项目版本号

  • 如果指定了 option.Tag,则所有项目都使用该 Tag 作为版本
  • 否则通过 SFTP 读取每个 .csproj 文件中的 <Version> 节点

版本号会写入 historyRequest.Versions

4.6 同版本二次确认

await WaitForSameVersionConfirmationAsync(runId, projectVersions, cancellationToken);

versionProbe.FindSameVersionsAsync 会检查线上版本是否与即将发布的版本一致。如果存在一致项目:

  • 在缓存中设置等待键和确认键
  • 向 SignalR 发送 VersionConfirmationRequired 事件,让用户决定“确认发布”或“取消”
  • 等待最多 3 分钟(实际使用 PipelineVersionProbe.ConfirmTimeout
  • 如果用户确认则继续;否则抛出 PipelineVersionNotConfirmedException

4.7 开发环境模拟

if (configuration.IsDevelopment)
{
    await DevelopmentSettingsAsync(runId, historyRequest, cancellationToken);
    return;
}

在开发环境下跳过真实构建/发布,只模拟阶段切换和随机失败,用于本地测试。

4.8 构建并推送镜像

  • 发送 PipelineStage.BuildImages
  • 根据 option.Project 选出需要构建的项目
  • 生成每个项目的构建命令
  • 通过 SSH 执行命令,并将终端输出实时推送给前端

4.9 发布部署

  • 发送 PipelineStage.Publish
  • 遍历所有服务器(sshConfiguration.Servers
  • 对每台服务器,根据其配置的 Project 与本次发布的 project 取交集,执行对应的部署命令
  • 同样实时推送输出

4.10 完成与错误处理

成功时:

historyRequest.Status = PipelineStatus.Success;
await hubContext.Clients.Group(runId).SendAsync("BuildCompleted", new { runId }, cancellationToken);

异常处理逻辑:

  • PipelineVersionNotConfirmedException → 状态 Abort,推送失败消息
  • OperationCanceledException → 状态 Abort,错误信息 “CI/CD被中止”
  • 其他 Exception → 状态 Failure,记录错误日志,推送失败
  • finally 中统一:
    • 计算结束时间和耗时
    • 调用 _steps.FinalizeSteps 生成步骤列表
    • 调用 SaveHistoryAsync 保存历史
    • 移除缓存中的锁、等待键、确认键

5. 关键辅助方法详解

5.1 开发环境模拟 DevelopmentSettingsAsync

模拟阶段变化,随机失败/成功,用于调试和演示。

5.2 同版本确认等待

WaitForSameVersionConfirmationAsync 具体流程已在上文描述,它使用了 FusionCache 作为决策中转站。

5.3 构建命令构造

GetBuildProjectImageCommandsBuiltInCommands.LazyBuildCommands.Value 获取预定义命令模板,将 <Version> 替换为实际版本号,并拼接 cd workspace/src。若服务器密码存在,则用 sudo 方式执行:

return $"printf '%s\\n' {QuoteForSingleQuotedShell(password)} | sudo -S -p '' bash -lc "
       + QuoteForSingleQuotedShell(command);

5.4 单引号转义

QuoteForSingleQuotedShell 使密码/命令可安全嵌入单引号包裹的 shell 字符串中。

5.5 发布服务器任务

PublishAsync 内定义局部异步函数 PublishServerAsync,所有服务器并行发布(Task.WhenAll)。每台服务器执行所有匹配项目的发布命令,并通过 ExecuteWithOutputAsync 实时读取输出判断退出状态。

5.6 终端输出实时推送

ExecuteWithOutputAsync
    └─ 并行读取流:OutputStream / ExtendedOutputStream
        └─ PublishTerminalOutputAsync
            ├─ 忽略空行
            ├─ _steps?.AppendLine(...) 收集日志
            └─ PublishMessageAsync -> SignalR

5.7 SSH 连接与拉代码

ConnectBuildServerAndPullCodeAsync 创建 SshClient,连接后设置 KeepAlive,执行 git pull,并统计耗时。

5.8 构建镜像

BuildPushImagesAsync 遍历所有命令执行,失败时抛出 InvalidOperationException,包含退出码和错误信息。

5.9 读取项目版本号

ReadProjectVersionsAsync

  • 定义 DpzCoreDpzWebApiDpzCoreAuth 对应的 .csproj 路径
  • 使用 Parallel.ForEachAsync 通过 SFTP 下载文件
  • XElement 读取 <PropertyGroup><Version>
  • 若缺失版本号则抛出异常

5.10 创建历史请求

CreateHistoryRequest 返回一个 CreatePipelineHistoryRequest 对象,包含 RunId、起止时间、状态、项目、版本模式(手动/自动)和触发者。

5.11 保存历史

SaveHistoryAsync 使用 Mediator 发送命令,保存失败仅记录日志,不中断主流程。


6. 内部类 PipelineStepCollector

用于把所有日志行汇总成结构化“步骤”数据:

  • 按键 (Server, Project) 存储 PipelineStep
  • AppendLine 将日志追加到对应步骤,并限制最多 2000 行(防止内存溢出)
  • FinalizeSteps 在任务结束后补充 EndTimeDurationMs,若整体失败则同步设置每个步骤的状态与错误
  • ResolveStage 根据服务器和项目识别步骤阶段:
    • 非构建主机 → 发布阶段
    • 构建主机且项目为 None → 拉代码阶段
    • 构建主机且指定项目 → 构建镜像阶段

7. 设计亮点总结

  • 异步贯穿:使用 async/awaitCancellationToken 支持取消和超时。
  • 实时交互:通过 SignalR 推送阶段变化、日志输出和需要用户确认的事件。
  • 强健锁机制:利用缓存实现分布式锁,配合 finally 保证无论成功失败都会释放。
  • SSH 流式处理:分别读取 stdout/stderr,避免 Unix 管道阻塞。
  • 开发友好:开发环境提供模拟执行,方便不依赖真实服务器即可测试流程。
  • 可观测性:日志和构建步骤被记录,并能保存到数据库供后续查看。

整体上,这段代码是一个完整的、生产可用的多服务器 Docker 化应用发布流水线实现。

评论加载中...