Files
Migu2.0/MiGu.Server/Ota/OtaJobRunner.cs
T
wei.wu 1f72488a8d 新增车辆表及OTA权限与逻辑优化
新增车辆任务与报警表,完善实体与DbContext配置。细化OTA权限校验,增强回传会话IP安全。优化OTA上传与设置面板,调度器支持重启恢复。报警采集逻辑支持历史分段。
2026-07-27 15:12:14 +08:00

463 lines
18 KiB
C#
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
using System.Collections.Concurrent;
namespace MiGu.Server.Ota;
public sealed class OtaJobRunner
{
private readonly OtaStore _store;
private readonly WatchDogClient _wd;
private readonly OtaVehicleSource _vehicles;
private readonly ILogger<OtaJobRunner> _log;
private readonly ConcurrentDictionary<string, CancellationTokenSource> _running = new();
public OtaJobRunner(OtaStore store, WatchDogClient wd, OtaVehicleSource vehicles, ILogger<OtaJobRunner> log)
{
_store = store;
_wd = wd;
_vehicles = vehicles;
_log = log;
RecoverInterruptedJobs();
}
/// <summary>进程重启后把落盘中非终态 job 标为 failed,避免 UI 永久显示 running。</summary>
private void RecoverInterruptedJobs()
{
try
{
foreach (var job in _store.ListJobs(500))
{
if (job.Status is not ("running" or "pending")) continue;
job.Status = "failed";
job.Message = string.IsNullOrWhiteSpace(job.Message)
? "进程重启,任务中断"
: job.Message;
job.FinishedAt = DateTimeOffset.UtcNow;
foreach (var step in job.Steps.Where(s => s.Status is "pending" or "running"))
{
step.Status = "failed";
step.Error ??= "进程重启,任务中断";
}
_store.SaveJob(job);
_log.LogWarning("OTA job {Id} marked failed after process restart", job.Id);
}
}
catch (Exception ex)
{
_log.LogWarning(ex, "OTA interrupted job recovery failed");
}
}
public OtaJob EnqueueSync(CreateSyncJobRequest req, string? user)
{
var target = _store.GetTarget() ?? throw new InvalidOperationException("未设置目标版本,请先在版本库激活");
var components = (req.Components is { Count: > 0 } ? req.Components : OtaPathMap.ComponentKeys.ToList())
.Where(c => target.Components.ContainsKey(c))
.ToList();
if (components.Count == 0) throw new InvalidOperationException("目标版本中无选定组件");
if (req.CarIds.Count == 0) throw new InvalidOperationException("请选择车辆");
var settings = _store.GetSettings();
var job = new OtaJob
{
Id = _store.NextJobId(),
Kind = "sync",
Status = "pending",
CreatedAt = DateTimeOffset.UtcNow,
CreatedBy = user,
PackageId = target.PackageId,
CarIds = req.CarIds.ToList(),
Components = components,
RequireLatencyCheck = req.RequireLatencyCheck ?? settings.LatencyEnabled
};
foreach (var carId in job.CarIds)
foreach (var comp in components)
job.Steps.Add(new OtaJobStep { CarId = carId, Component = comp, Status = "pending" });
job.TotalSteps = job.Steps.Count;
_store.SaveJob(job);
_ = Task.Run(() => RunAsync(job.Id));
return job;
}
public OtaJob EnqueueCustomFile(CreateCustomFileJobRequest req, IReadOnlyList<OtaCustomFileItem> files, string? user)
{
if (req.CarIds.Count == 0) throw new InvalidOperationException("请选择车辆");
if (files.Count == 0) throw new InvalidOperationException("请至少添加一个本地文件");
if (string.IsNullOrWhiteSpace(req.RemotePath)) throw new InvalidOperationException("请填写小车内目标路径");
var settings = _store.GetSettings();
var ops = req.RestartOps is { Count: > 0 }
? req.RestartOps.Distinct().ToList()
: new List<int> { req.RestartOp };
if (ops.Count == 0) ops.Add(-1);
var job = new OtaJob
{
Id = _store.NextJobId(),
Kind = "customFile",
Status = "pending",
CreatedAt = DateTimeOffset.UtcNow,
CreatedBy = user,
CarIds = req.CarIds.ToList(),
CustomFileName = files[0].FileName,
CustomRemotePath = req.RemotePath,
CustomRestartOp = ops[0],
CustomLocalPath = files[0].LocalPath,
CustomFiles = files.ToList(),
CustomRestartOps = ops,
RequireLatencyCheck = req.RequireLatencyCheck ?? settings.LatencyEnabled,
Components = new List<string> { "custom" }
};
foreach (var carId in job.CarIds)
job.Steps.Add(new OtaJobStep { CarId = carId, Component = $"custom×{files.Count}", Status = "pending" });
job.TotalSteps = job.Steps.Count;
_store.SaveJob(job);
_ = Task.Run(() => RunAsync(job.Id));
return job;
}
public OtaJob EnqueueConfigPush(CreateConfigPushRequest req, string? user)
{
if (req.CarIds.Count == 0) throw new InvalidOperationException("请选择车辆");
var app = req.App.Trim().ToLowerInvariant();
if (app is not ("medulla" or "detour" or "clumsy"))
throw new InvalidOperationException("app 须为 medulla|detour|clumsy");
var settings = _store.GetSettings();
var job = new OtaJob
{
Id = _store.NextJobId(),
Kind = "configPush",
Status = "pending",
CreatedAt = DateTimeOffset.UtcNow,
CreatedBy = user,
CarIds = req.CarIds.ToList(),
ConfigApp = app,
ConfigJson = req.Json,
RequireLatencyCheck = req.RequireLatencyCheck ?? settings.LatencyEnabled,
Components = new List<string> { $"config:{app}" }
};
foreach (var carId in job.CarIds)
job.Steps.Add(new OtaJobStep { CarId = carId, Component = $"config:{app}", Status = "pending" });
job.TotalSteps = job.Steps.Count;
_store.SaveJob(job);
_ = Task.Run(() => RunAsync(job.Id));
return job;
}
public bool Cancel(string jobId)
{
var job = _store.GetJob(jobId);
if (job == null) return false;
if (job.Status is "succeeded" or "failed" or "partial" or "cancelled") return false;
// 正在运行:只发取消信号,由 RunAsync 统一收尾,避免与运行线程并发写同一 job 文件。
if (_running.TryGetValue(jobId, out var cts))
{
cts.Cancel();
return true;
}
// 尚未开始(Task.Run 排队中):直接落盘取消;RunAsync 启动时有 status 守卫会跳过执行。
foreach (var step in job.Steps.Where(s => s.Status == "pending"))
{
step.Status = "skipped";
step.Error = "已取消";
}
job.Status = "cancelled";
job.FinishedAt = DateTimeOffset.UtcNow;
job.DoneSteps = job.Steps.Count(s => s.Status is "succeeded" or "failed" or "skipped");
_store.SaveJob(job);
return true;
}
public OtaJob? RetryFailed(string jobId, string? user)
{
var old = _store.GetJob(jobId);
if (old == null) return null;
var failedCars = old.Steps.Where(s => s.Status == "failed").Select(s => s.CarId).Distinct().ToList();
if (failedCars.Count == 0) throw new InvalidOperationException("没有失败项可重试");
return old.Kind switch
{
"sync" => EnqueueSync(new CreateSyncJobRequest
{
CarIds = failedCars,
Components = old.Components,
RequireLatencyCheck = old.RequireLatencyCheck
}, user),
"customFile" when ResolveCustomFiles(old).Count > 0 =>
EnqueueCustomFile(new CreateCustomFileJobRequest
{
CarIds = failedCars,
RemotePath = old.CustomRemotePath ?? "",
RestartOps = old.CustomRestartOps is { Count: > 0 }
? old.CustomRestartOps
: new List<int> { old.CustomRestartOp ?? -1 },
RequireLatencyCheck = old.RequireLatencyCheck
}, ResolveCustomFiles(old), user),
"configPush" => EnqueueConfigPush(new CreateConfigPushRequest
{
CarIds = failedCars,
App = old.ConfigApp ?? "",
Json = old.ConfigJson ?? "",
RequireLatencyCheck = old.RequireLatencyCheck
}, user),
_ => throw new InvalidOperationException("无法重试该任务类型或文件已丢失")
};
}
private async Task RunAsync(string jobId)
{
var cts = new CancellationTokenSource();
if (!_running.TryAdd(jobId, cts)) return;
// 单一内存实例贯穿全程;所有「改状态 + 落盘」都在 jobLock 下串行,杜绝并发车批次的丢更新。
var jobLock = new object();
try
{
var job = _store.GetJob(jobId);
if (job == null) return;
// 守卫:排队期间被取消 / 已终结的任务不再执行。
if (job.Status is "cancelled" or "succeeded" or "failed" or "partial") return;
var settings = _store.GetSettings();
var cars = await _vehicles.ListCarsAsync(cts.Token);
var byId = cars.ToDictionary(c => c.Id, StringComparer.OrdinalIgnoreCase);
job.Status = job.RequireLatencyCheck ? "probing" : "running";
_store.SaveJob(job);
// 解析 IP
foreach (var step in job.Steps)
{
if (byId.TryGetValue(step.CarId, out var car))
step.Ip = car.Ip;
}
if (job.RequireLatencyCheck)
{
var threshold = settings.RttThresholdMs;
var overMode = settings.OverThreshold;
foreach (var carId in job.CarIds.Distinct())
{
var ip = job.Steps.FirstOrDefault(s => s.CarId == carId)?.Ip;
if (string.IsNullOrEmpty(ip))
{
SkipCar(job, carId, "无 IP");
continue;
}
var rtt = await _wd.MeasureRttMsAsync(ip, cts.Token);
if (rtt == null || rtt > threshold)
{
if (string.Equals(overMode, "skip", StringComparison.OrdinalIgnoreCase) || rtt == null)
SkipCar(job, carId, rtt == null ? "延迟探测失败" : $"RTT {rtt}ms > {threshold}ms");
}
}
_store.SaveJob(job);
}
job.Status = "running";
_store.SaveJob(job);
var maxCar = Math.Max(1, settings.MaxCar);
var carGroups = job.CarIds.Distinct()
.Where(id => job.Steps.Any(s => s.CarId == id && s.Status == "pending"))
.Chunk(maxCar);
foreach (var batch in carGroups)
{
if (cts.IsCancellationRequested) break;
var tasks = batch.Select(carId => RunCarAsync(job, jobLock, carId, settings.BandwidthKbps, cts.Token));
await Task.WhenAll(tasks);
}
lock (jobLock)
{
if (cts.IsCancellationRequested)
{
foreach (var s in job.Steps.Where(s => s.Status is "pending" or "running"))
{
s.Status = "skipped";
s.Error ??= "已取消";
}
job.Status = "cancelled";
job.FinishedAt = DateTimeOffset.UtcNow;
Recalc(job);
}
else
{
Finalize(job);
}
_store.SaveJob(job);
}
}
catch (Exception ex)
{
_log.LogError(ex, "OTA job {Id} crashed", jobId);
var job = _store.GetJob(jobId);
if (job != null)
{
job.Status = "failed";
job.Message = ex.Message;
job.FinishedAt = DateTimeOffset.UtcNow;
_store.SaveJob(job);
}
}
finally
{
_running.TryRemove(jobId, out _);
cts.Dispose();
}
}
private static void SkipCar(OtaJob job, string carId, string reason)
{
foreach (var step in job.Steps.Where(s => s.CarId == carId && s.Status == "pending"))
{
step.Status = "skipped";
step.Error = reason;
}
}
private static List<OtaCustomFileItem> ResolveCustomFiles(OtaJob job)
{
if (job.CustomFiles is { Count: > 0 })
return job.CustomFiles.Where(f => File.Exists(f.LocalPath)).ToList();
if (!string.IsNullOrEmpty(job.CustomLocalPath) && File.Exists(job.CustomLocalPath))
{
return new List<OtaCustomFileItem>
{
new()
{
LocalPath = job.CustomLocalPath,
FileName = job.CustomFileName ?? Path.GetFileName(job.CustomLocalPath)
}
};
}
return new();
}
private async Task UploadCustomFilesForCarAsync(OtaJob job, string ip, int bandwidth, CancellationToken ct)
{
var files = ResolveCustomFiles(job);
if (files.Count == 0) throw new InvalidOperationException("自定义文件已丢失");
var remote = job.CustomRemotePath ?? "";
var ops = job.CustomRestartOps is { Count: > 0 }
? job.CustomRestartOps
: new List<int> { job.CustomRestartOp ?? -1 };
// 非末文件一律 -1;末文件按所选重启项依次再传(对齐 CarOTA.App)
for (var i = 0; i < files.Count; i++)
{
var f = files[i];
var isLast = i == files.Count - 1;
if (!isLast)
{
await _wd.UploadCustomFileAsync(ip, f.LocalPath, f.FileName, remote, -1, bandwidth, ct);
continue;
}
if (ops.Count == 1 && ops[0] == -1)
{
await _wd.UploadCustomFileAsync(ip, f.LocalPath, f.FileName, remote, -1, bandwidth, ct);
}
else
{
foreach (var op in ops.Where(o => o != -1).DefaultIfEmpty(-1))
await _wd.UploadCustomFileAsync(ip, f.LocalPath, f.FileName, remote, op, bandwidth, ct);
}
}
}
private async Task RunCarAsync(OtaJob job, object jobLock, string carId, int bandwidth, CancellationToken ct)
{
// 只处理本车 step;同批其他车任务并行修改各自 step,共享同一 job 实例,写盘统一在 jobLock 下串行。
var steps = job.Steps.Where(s => s.CarId == carId && s.Status == "pending").ToList();
if (steps.Count == 0) return;
var ip = steps[0].Ip;
if (string.IsNullOrEmpty(ip))
{
lock (jobLock)
{
foreach (var s in steps) { s.Status = "failed"; s.Error = "无 IP"; }
Recalc(job);
_store.SaveJob(job);
}
return;
}
foreach (var step in steps)
{
if (ct.IsCancellationRequested)
{
lock (jobLock)
{
step.Status = "skipped";
step.Error = "已取消";
Recalc(job);
_store.SaveJob(job);
}
continue;
}
lock (jobLock)
{
step.Status = "running";
Recalc(job);
_store.SaveJob(job);
}
try
{
switch (job.Kind)
{
case "sync":
{
var target = _store.GetTarget() ?? throw new InvalidOperationException("目标版本丢失");
if (!target.Components.TryGetValue(step.Component, out var art))
throw new InvalidOperationException($"组件 {step.Component} 不在目标中");
await _wd.UploadComponentAsync(ip, step.Component, art.Path, art.FileName, bandwidth, ct);
break;
}
case "customFile":
await UploadCustomFilesForCarAsync(job, ip, bandwidth, ct);
break;
case "configPush":
await _wd.PutJsonAsync(ip, job.ConfigApp!, job.ConfigJson!, ct);
break;
}
lock (jobLock)
{
step.Status = "succeeded";
step.Error = null;
Recalc(job);
_store.SaveJob(job);
}
}
catch (Exception ex)
{
lock (jobLock)
{
step.Status = ct.IsCancellationRequested ? "skipped" : "failed";
step.Error = ct.IsCancellationRequested ? "已取消" : ex.Message;
Recalc(job);
_store.SaveJob(job);
}
_log.LogWarning(ex, "OTA step fail {Job} {Car} {Comp}", job.Id, carId, step.Component);
}
}
}
private static void Recalc(OtaJob job)
{
job.DoneSteps = job.Steps.Count(s => s.Status is "succeeded" or "failed" or "skipped");
}
private static void Finalize(OtaJob job)
{
Recalc(job);
job.FinishedAt = DateTimeOffset.UtcNow;
var anyFail = job.Steps.Any(s => s.Status == "failed");
var anyOk = job.Steps.Any(s => s.Status == "succeeded");
var anyPending = job.Steps.Any(s => s.Status is "pending" or "running");
if (job.Status == "cancelled") return;
if (anyPending) job.Status = "partial";
else if (anyFail && anyOk) job.Status = "partial";
else if (anyFail) job.Status = "failed";
else if (anyOk) job.Status = "succeeded";
else job.Status = "cancelled";
}
}