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;
}
}
}
这段代码定义了一个名为 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 构建命令构造
GetBuildProjectImageCommands 从 BuiltInCommands.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:
- 定义
DpzCore、DpzWebApi、DpzCoreAuth对应的.csproj路径 - 使用
Parallel.ForEachAsync通过 SFTP 下载文件 - 用
XElement读取<PropertyGroup><Version> - 若缺失版本号则抛出异常
5.10 创建历史请求
CreateHistoryRequest 返回一个 CreatePipelineHistoryRequest 对象,包含 RunId、起止时间、状态、项目、版本模式(手动/自动)和触发者。
5.11 保存历史
SaveHistoryAsync 使用 Mediator 发送命令,保存失败仅记录日志,不中断主流程。
6. 内部类 PipelineStepCollector
用于把所有日志行汇总成结构化“步骤”数据:
- 按键
(Server, Project)存储PipelineStep AppendLine将日志追加到对应步骤,并限制最多2000行(防止内存溢出)FinalizeSteps在任务结束后补充EndTime、DurationMs,若整体失败则同步设置每个步骤的状态与错误ResolveStage根据服务器和项目识别步骤阶段:- 非构建主机 → 发布阶段
- 构建主机且项目为
None→ 拉代码阶段 - 构建主机且指定项目 → 构建镜像阶段
7. 设计亮点总结
- 异步贯穿:使用
async/await和CancellationToken支持取消和超时。 - 实时交互:通过 SignalR 推送阶段变化、日志输出和需要用户确认的事件。
- 强健锁机制:利用缓存实现分布式锁,配合
finally保证无论成功失败都会释放。 - SSH 流式处理:分别读取 stdout/stderr,避免 Unix 管道阻塞。
- 开发友好:开发环境提供模拟执行,方便不依赖真实服务器即可测试流程。
- 可观测性:日志和构建步骤被记录,并能保存到数据库供后续查看。
整体上,这段代码是一个完整的、生产可用的多服务器 Docker 化应用发布流水线实现。