Files
Migu2.0/MiGu.Server/Fleet/CdmTaskSync.cs
T
wei.wu 125372b157 移除历史状态归一化及相关遗留逻辑
- 移除 LegacyStatusNormalizationMigrator,仅保留 AreaLayoutModeDefaultMigrator
- 删除 WmsDispatchStatus 枚举、常量和全局 using
- 统一所有 DbContext 注入为 MiGuDbContext,移除 PlatformDbContext
- 服务实体操作统一用 IEditableRepository<T>,移除本地实现
- 移除 WmsService 的 MigrateLegacyAsync 方法
- 精简接口模型,移除部分字段和请求体
- 前端保存容器时 status 字段取自已有数据,移除写死默认值
- 更新文档,去除旧说明,统一数据库初始化流程
2026-08-03 09:41:03 +08:00

206 lines
7.1 KiB
C#

using System.Text.Json;
using Microsoft.EntityFrameworkCore;
using Microsoft.Extensions.Options;
using MiGu.Server.Auth;
using MiGu.Server.Launcher;
using MiGu.Server.Persistence;
namespace MiGu.Server.Fleet;
/// <summary>投影 /projection/deliveries 返回的单行(camelCase)。</summary>
public sealed class CdmTaskDto
{
public string id { get; set; } = "";
public string? taskId { get; set; }
public int missionId { get; set; }
public string missionName { get; set; } = "";
public string missionTypeName { get; set; } = "";
public int srcSiteId { get; set; }
public string srcLabel { get; set; } = "";
public int dstSiteId { get; set; }
public string dstLabel { get; set; } = "";
public string status { get; set; } = "";
public string statusCode { get; set; } = "";
public int? carId { get; set; }
public string? carName { get; set; }
public int priority { get; set; }
public string? createTime { get; set; }
public string? startTime { get; set; }
public string? finishTime { get; set; }
public string? stuckReason { get; set; }
public bool overdue { get; set; }
}
/// <summary>
/// 单例:从 SimpleLite 投影拉取 CDM 任务并 upsert 到 platform.db(cdm_tasks),永不删除=保留历史。
/// 同时维护「SimpleLite 是否在线 / 最近同步时间」,供任务页离线降级展示。
/// </summary>
public sealed class CdmTaskSyncer
{
private readonly IServiceScopeFactory _scopeFactory;
private readonly IHttpClientFactory _httpFactory;
private readonly SimpleLiteOptions _sl;
private readonly InternalTokenStore _token;
private readonly ILogger<CdmTaskSyncer> _log;
private readonly SemaphoreSlim _gate = new(1, 1);
private static readonly JsonSerializerOptions JsonOpt = new() { PropertyNameCaseInsensitive = true };
public volatile bool Online;
public DateTimeOffset? LastSyncAt { get; private set; }
public int LastCount { get; private set; }
public CdmTaskSyncer(
IServiceScopeFactory scopeFactory,
IHttpClientFactory httpFactory,
IOptions<SimpleLiteOptions> sl,
InternalTokenStore token,
ILogger<CdmTaskSyncer> log)
{
_scopeFactory = scopeFactory;
_httpFactory = httpFactory;
_sl = sl.Value;
_token = token;
_log = log;
}
/// <summary>拉取 + 落库一次。并发调用时若已有同步在进行则直接跳过(返回当前在线状态)。</summary>
public async Task<bool> SyncOnceAsync(CancellationToken ct)
{
if (!await _gate.WaitAsync(0, ct)) return Online;
try
{
var dtos = await FetchAsync(ct);
if (dtos == null)
{
Online = false;
return false;
}
try
{
using var scope = _scopeFactory.CreateScope();
var db = scope.ServiceProvider.GetRequiredService<MiGuDbContext>();
await UpsertAsync(db, dtos, ct);
}
catch (Exception ex)
{
// 拉取成功即视为在线;落库失败只记日志,不影响在线判定
_log.LogDebug(ex, "cdm upsert failed");
}
Online = true;
LastSyncAt = DateTimeOffset.UtcNow;
LastCount = dtos.Count;
return true;
}
finally
{
_gate.Release();
}
}
private async Task<List<CdmTaskDto>?> FetchAsync(CancellationToken ct)
{
var port = _sl.ProjectionPort > 0 ? _sl.ProjectionPort : 8222;
try
{
using var client = _httpFactory.CreateClient();
client.Timeout = TimeSpan.FromSeconds(10);
using var req = new HttpRequestMessage(
HttpMethod.Get,
$"http://127.0.0.1:{port}/projection/deliveries?includeFinished=true&includeAborted=true");
var token = _token.Token;
if (!string.IsNullOrEmpty(token))
req.Headers.TryAddWithoutValidation("X-Platform-Internal-Token", token);
using var resp = await client.SendAsync(req, ct);
if (!resp.IsSuccessStatusCode) return null;
var text = await resp.Content.ReadAsStringAsync(ct);
return JsonSerializer.Deserialize<List<CdmTaskDto>>(text, JsonOpt) ?? new List<CdmTaskDto>();
}
catch (Exception ex)
{
_log.LogDebug(ex, "cdm fetch failed");
return null;
}
}
private static async Task UpsertAsync(MiGuDbContext db, IReadOnlyList<CdmTaskDto> dtos, CancellationToken ct)
{
var valid = dtos.Where(d => !string.IsNullOrWhiteSpace(d.id)).ToList();
if (valid.Count == 0) return;
var now = DateTimeOffset.UtcNow;
var ids = valid.Select(d => d.id).ToList();
var existing = await db.CdmTasks.Where(t => ids.Contains(t.Id)).ToDictionaryAsync(t => t.Id, ct);
foreach (var d in valid)
{
if (existing.TryGetValue(d.id, out var rec))
{
Map(d, rec);
rec.LastSeenAt = now;
}
else
{
var created = new CdmTaskRecord { Id = d.id, FirstSeenAt = now, LastSeenAt = now };
Map(d, created);
db.CdmTasks.Add(created);
}
}
await db.SaveChangesAsync(ct);
}
private static void Map(CdmTaskDto d, CdmTaskRecord rec)
{
rec.TaskId = string.IsNullOrWhiteSpace(d.taskId) ? null : d.taskId;
rec.MissionId = d.missionId;
rec.MissionName = d.missionName ?? "";
rec.MissionTypeName = d.missionTypeName ?? "";
rec.SrcSiteId = d.srcSiteId;
rec.SrcLabel = d.srcLabel ?? "";
rec.DstSiteId = d.dstSiteId;
rec.DstLabel = d.dstLabel ?? "";
rec.Status = d.status ?? "";
rec.StatusCode = d.statusCode ?? "";
rec.CarId = d.carId;
rec.CarName = d.carName;
rec.Priority = d.priority;
rec.CreateTime = d.createTime;
rec.StartTime = d.startTime;
rec.FinishTime = d.finishTime;
rec.StuckReason = d.stuckReason;
rec.Overdue = d.overdue;
}
}
/// <summary>后台循环:定时把 CDM 任务同步进 platform.db,保证无人打开页面时也能捕获终态历史。</summary>
public sealed class CdmTaskSyncService : BackgroundService
{
private readonly CdmTaskSyncer _syncer;
private readonly ILogger<CdmTaskSyncService> _log;
public CdmTaskSyncService(CdmTaskSyncer syncer, ILogger<CdmTaskSyncService> log)
{
_syncer = syncer;
_log = log;
}
protected override async Task ExecuteAsync(CancellationToken stoppingToken)
{
try { await Task.Delay(TimeSpan.FromSeconds(3), stoppingToken); }
catch { return; }
while (!stoppingToken.IsCancellationRequested)
{
try { await _syncer.SyncOnceAsync(stoppingToken); }
catch (Exception ex) { _log.LogDebug(ex, "cdm sync loop error"); }
try { await Task.Delay(TimeSpan.FromSeconds(8), stoppingToken); }
catch { break; }
}
}
}