补齐运输规划与任务服务在状态/优先级字段上的读写与调度衔接。 Co-authored-by: Cursor <cursoragent@cursor.com>
339 lines
14 KiB
C#
339 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 PlatformDbContext _db;
|
|
private readonly WmsTransportPlanner _planner;
|
|
private readonly WmsService _wms;
|
|
private readonly IWmsDispatchAdapter _dispatch;
|
|
private readonly JsonSerializerOptions _json = new(JsonSerializerDefaults.Web);
|
|
|
|
public WmsTransportTaskService(
|
|
PlatformDbContext 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))
|
|
query = query.Where(x => x.Status == status);
|
|
if (!string.IsNullOrWhiteSpace(triggerType))
|
|
query = query.Where(x => x.BusinessType == triggerType);
|
|
var rows = await query.ToListAsync();
|
|
return rows.OrderByDescending(x => x.CreatedAt).Take(500).ToList();
|
|
}
|
|
|
|
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, 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,
|
|
task.TargetStorageId.ToString("D"), ContainerLocationStatuses.Active,
|
|
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, string toStatus, string actor, string reason)
|
|
{
|
|
var from = task.Status;
|
|
task.Status = toStatus;
|
|
task.UpdatedBy = actor;
|
|
await _db.SaveChangesAsync();
|
|
await AppendHistoryAsync(task.Id, from, toStatus, 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(EntityBase entity, string actor)
|
|
{
|
|
entity.CreatedBy = actor;
|
|
entity.UpdatedBy = actor;
|
|
}
|
|
}
|