Files
Migu2.0/MiGu.Server/Wms/WmsTransportTaskService.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

340 lines
14 KiB
C#

using System.Text.Json;
using Microsoft.EntityFrameworkCore;
using MiGu.Server.Persistence;
namespace MiGu.Server.Wms;
public sealed class WmsTransportTaskService
{
private readonly MiGuDbContext _db;
private readonly WmsTransportPlanner _planner;
private readonly WmsService _wms;
private readonly IWmsDispatchAdapter _dispatch;
private readonly JsonSerializerOptions _json = new(JsonSerializerDefaults.Web);
public WmsTransportTaskService(
MiGuDbContext db,
WmsTransportPlanner planner,
WmsService wms,
IWmsDispatchAdapter dispatch)
{
_db = db;
_planner = planner;
_wms = wms;
_dispatch = dispatch;
}
public async Task<List<WmsTransportTask>> ListAsync(string? status = null, string? triggerType = null)
{
var query = _db.WmsTransportTasks.AsNoTracking().AsQueryable();
if (!string.IsNullOrWhiteSpace(status) &&
Enum.TryParse<WmsTransportTaskStatus>(status.Trim(), true, out var st) &&
WmsTransportTaskStatuses.All.Contains(st))
query = query.Where(x => x.Status == st);
if (!string.IsNullOrWhiteSpace(triggerType))
query = query.Where(x => x.BusinessType == triggerType);
return await query.OrderByDescending(x => x.CreatedAt).Take(500).ToListAsync();
}
public Task<TransportCandidatePreview> PreviewCandidatesAsync(WmsTransportRequest request) =>
_planner.PreviewAsync(request);
public async Task<WmsTransportTask> GenerateAsync(WmsTransportRequest request, string actor)
{
var picked = await _planner.PickBestAsync(request)
?? throw new InvalidOperationException("未找到可用的搬运候选");
WmsTransportRule? rule = null;
if (picked.RuleId != Guid.Empty)
rule = await _db.WmsTransportRules.AsNoTracking().FirstOrDefaultAsync(x => x.Id == picked.RuleId);
var options = TransportSelectorParser.ParseTaskOptions(rule?.TaskOptionsJson);
await using var tx = await _db.Database.BeginTransactionAsync();
try
{
await EnsureNoActiveReservationAsync(picked.ContainerId, picked.TargetStorageId);
await EnsureStorageAvailableAsync(picked.SourceStorageId, picked.TargetStorageId);
var snapshot = BuildSnapshot(picked, request);
var status = options.RequireManualConfirm || !options.AutoReserve
? WmsTransportTaskStatuses.Pending
: WmsTransportTaskStatuses.Reserved;
var task = new WmsTransportTask
{
BusinessType = request.TriggerType,
RuleId = picked.RuleId == Guid.Empty ? null : picked.RuleId,
SourceStorageId = picked.SourceStorageId,
TargetStorageId = picked.TargetStorageId,
ContainerId = picked.ContainerId,
MaterialId = picked.MaterialId,
Quantity = picked.Quantity ?? request.Quantity,
Status = status,
Reason = request.Reason.TrimOr(""),
SnapshotJson = JsonSerializer.Serialize(snapshot, _json),
TaskPriority = request.Priority
};
StampCreate(task, actor);
_db.WmsTransportTasks.Add(task);
await _db.SaveChangesAsync();
if (status == WmsTransportTaskStatuses.Reserved)
await CreateReservationAsync(task, options, actor);
await AppendHistoryAsync(task.Id, "", task.Status.ToString(), actor, request.Reason, "", task.SnapshotJson);
await tx.CommitAsync();
return task;
}
catch
{
await tx.RollbackAsync();
throw;
}
}
public async Task<WmsTransportTask> ReserveAsync(Guid taskId, string actor, string reason = "")
{
var task = await FindTaskAsync(taskId);
if (task.Status != WmsTransportTaskStatuses.Pending)
throw new InvalidOperationException("只有 Pending 状态的任务可以预占");
var rule = task.RuleId.HasValue
? await _db.WmsTransportRules.AsNoTracking().FirstOrDefaultAsync(x => x.Id == task.RuleId.Value)
: null;
var options = TransportSelectorParser.ParseTaskOptions(rule?.TaskOptionsJson);
await using var tx = await _db.Database.BeginTransactionAsync();
try
{
await EnsureNoActiveReservationAsync(task.ContainerId, task.TargetStorageId);
await EnsureStorageAvailableAsync(task.SourceStorageId, task.TargetStorageId);
await CreateReservationAsync(task, options, actor);
await ChangeStatusAsync(task, WmsTransportTaskStatuses.Reserved, actor, reason);
await tx.CommitAsync();
return task;
}
catch
{
await tx.RollbackAsync();
throw;
}
}
public async Task<WmsTransportTask> CancelAsync(Guid taskId, string actor, string reason = "")
{
var task = await FindTaskAsync(taskId);
if (task.Status is WmsTransportTaskStatuses.Completed or WmsTransportTaskStatuses.Cancelled)
throw new InvalidOperationException("任务已结束,不能取消");
await ReleaseReservationsAsync(task.Id);
await ChangeStatusAsync(task, WmsTransportTaskStatuses.Cancelled, actor, reason);
return task;
}
public async Task<WmsTransportTask> CompleteAsync(Guid taskId, string actor, string reason = "")
{
var task = await FindTaskAsync(taskId);
if (task.Status is WmsTransportTaskStatuses.Completed or WmsTransportTaskStatuses.Cancelled)
throw new InvalidOperationException("任务已结束");
var source = await _db.Storages.AsNoTracking().FirstOrDefaultAsync(x => x.Id == task.SourceStorageId)
?? throw new InvalidOperationException("源库位不存在");
var target = await _db.Storages.AsNoTracking().FirstOrDefaultAsync(x => x.Id == task.TargetStorageId)
?? throw new InvalidOperationException("目标库位不存在");
if (string.IsNullOrWhiteSpace(target.SiteId))
throw new InvalidOperationException("目标库位缺少 SiteId,无法完成搬运");
await using var tx = await _db.Database.BeginTransactionAsync();
try
{
await _wms.BindOrTransferLocation(new ContainerLocationRequest(
null, null, task.ContainerId, ContainerLocationTypes.Storage.ToString(),
task.TargetStorageId.ToString("D"), ContainerLocationStatuses.Active.ToString(),
DateTimeOffset.UtcNow, "TransportTask", reason, false, "", "{}"), actor);
await ReleaseReservationsAsync(task.Id);
await ChangeStatusAsync(task, WmsTransportTaskStatuses.Completed, actor, reason);
await tx.CommitAsync();
return task;
}
catch
{
await tx.RollbackAsync();
throw;
}
}
public async Task<WmsTransportTask> DispatchAsync(Guid taskId, string actor)
{
// 原子抢占:仅一条 Reserved 且未在下发中的任务能进入 Dispatching,避免双发
var now = DateTimeOffset.UtcNow;
var claimed = await _db.WmsTransportTasks
.Where(x => x.Id == taskId
&& x.Status == WmsTransportTaskStatuses.Reserved
&& x.DispatchStatus != "Dispatching"
&& x.DispatchStatus != "Dispatched")
.ExecuteUpdateAsync(s => s
.SetProperty(x => x.DispatchStatus, "Dispatching")
.SetProperty(x => x.UpdatedAt, now)
.SetProperty(x => x.UpdatedBy, actor));
if (claimed == 0)
{
var existing = await FindTaskAsync(taskId);
if (existing.Status != WmsTransportTaskStatuses.Reserved)
throw new InvalidOperationException("只有 Reserved 状态的任务可以下发");
throw new InvalidOperationException("任务正在下发中,请勿重复操作");
}
var task = await FindTaskAsync(taskId);
var source = await _db.Storages.AsNoTracking().FirstAsync(x => x.Id == task.SourceStorageId);
var target = await _db.Storages.AsNoTracking().FirstAsync(x => x.Id == task.TargetStorageId);
var container = await _db.Containers.AsNoTracking().FirstAsync(x => x.Id == task.ContainerId);
var result = await _dispatch.DispatchAsync(new DispatchTransportCommand(
task.Id, source.SiteId, target.SiteId, container.Code, task.TaskPriority,
task.BusinessType, task.Reason));
if (result.Success)
{
task.DispatchMissionId = result.MissionId ?? "";
task.DeliveryId = result.DeliveryId ?? "";
task.DispatchStatus = "Dispatched";
await ChangeStatusAsync(task, WmsTransportTaskStatuses.Dispatched, actor, "调度下发");
}
else
{
task.ErrorMessage = result.ErrorMessage ?? "下发失败";
task.DispatchStatus = "Failed";
await ChangeStatusAsync(task, WmsTransportTaskStatuses.Failed, actor, result.ErrorMessage ?? "下发失败");
}
return task;
}
private async Task ExpireStaleReservationsAsync()
{
var now = DateTimeOffset.UtcNow;
var rows = await _db.WmsTransportReservations
.Where(x => x.Status == WmsReservationStatuses.Active
&& x.ExpiresAt != null
&& x.ExpiresAt < now)
.ToListAsync();
if (rows.Count == 0) return;
foreach (var row in rows)
row.Status = WmsReservationStatuses.Expired;
await _db.SaveChangesAsync();
}
private async Task EnsureNoActiveReservationAsync(Guid containerId, Guid targetStorageId)
{
await ExpireStaleReservationsAsync();
var now = DateTimeOffset.UtcNow;
if (await _db.WmsTransportReservations.AnyAsync(x =>
x.Status == WmsReservationStatuses.Active
&& (x.ExpiresAt == null || x.ExpiresAt > now)
&& x.ContainerId == containerId))
throw new InvalidOperationException("容器已被其他任务预占");
if (await _db.WmsTransportReservations.AnyAsync(x =>
x.Status == WmsReservationStatuses.Active
&& (x.ExpiresAt == null || x.ExpiresAt > now)
&& x.TargetStorageId == targetStorageId))
throw new InvalidOperationException("目标库位已被其他任务预占");
}
private async Task EnsureStorageAvailableAsync(Guid sourceStorageId, Guid targetStorageId)
{
var targetOccupied = await _db.ContainerLocations.AnyAsync(x =>
x.LocationType == ContainerLocationTypes.Storage && x.LocationId == targetStorageId.ToString("D"));
if (targetOccupied)
throw new InvalidOperationException("目标库位已被占用");
var sourceHasContainer = await _db.ContainerLocations.AnyAsync(x =>
x.LocationType == ContainerLocationTypes.Storage &&
x.LocationId == sourceStorageId.ToString("D"));
if (!sourceHasContainer)
throw new InvalidOperationException("源库位没有容器");
}
private async Task CreateReservationAsync(WmsTransportTask task, TransportTaskOptions options, string actor)
{
var reservation = new WmsTransportReservation
{
TaskId = task.Id,
ContainerId = task.ContainerId,
SourceStorageId = task.SourceStorageId,
TargetStorageId = task.TargetStorageId,
Status = WmsReservationStatuses.Active,
ExpiresAt = options.ReservationTtlSeconds > 0
? DateTimeOffset.UtcNow.AddSeconds(options.ReservationTtlSeconds)
: null
};
StampCreate(reservation, actor);
_db.WmsTransportReservations.Add(reservation);
await _db.SaveChangesAsync();
}
private async Task ReleaseReservationsAsync(Guid taskId)
{
var rows = await _db.WmsTransportReservations.Where(x => x.TaskId == taskId && x.Status == WmsReservationStatuses.Active).ToListAsync();
foreach (var row in rows)
row.Status = WmsReservationStatuses.Released;
await _db.SaveChangesAsync();
}
private async Task ChangeStatusAsync(WmsTransportTask task, WmsTransportTaskStatus toStatus, string actor, string reason)
{
var from = task.Status;
task.Status = toStatus;
task.UpdatedBy = actor;
await _db.SaveChangesAsync();
await AppendHistoryAsync(task.Id, from.ToString(), toStatus.ToString(), actor, reason, task.ErrorMessage, task.SnapshotJson);
}
private async Task AppendHistoryAsync(Guid taskId, string from, string to, string actor, string reason, string error, string snapshot)
{
_db.WmsTransportTaskHistories.Add(new WmsTransportTaskHistory
{
TaskId = taskId,
FromStatus = from,
ToStatus = to,
Operator = actor,
OperatedAt = DateTimeOffset.UtcNow,
Reason = reason.TrimOr(""),
ErrorMessage = error.TrimOr(""),
SnapshotJson = snapshot
});
await _db.SaveChangesAsync();
}
private async Task<WmsTransportTask> FindTaskAsync(Guid id) =>
await _db.WmsTransportTasks.FirstOrDefaultAsync(x => x.Id == id)
?? throw new InvalidOperationException("搬运任务不存在");
private static object BuildSnapshot(TransportCandidatePair picked, WmsTransportRequest request) => new
{
request.TriggerType,
picked.RuleId,
picked.RuleCode,
picked.SourceStorageId,
picked.SourceStorageCode,
picked.TargetStorageId,
picked.TargetStorageCode,
picked.ContainerId,
picked.ContainerCode,
picked.MaterialId,
picked.MaterialCode,
Quantity = picked.Quantity ?? request.Quantity,
GeneratedAt = DateTimeOffset.UtcNow
};
private static void StampCreate(AggregateRoot entity, string actor)
{
entity.CreatedBy = actor;
entity.UpdatedBy = actor;
}
}