using System.Security.Cryptography;
using System.Text;
using System.Text.Json;
using System.Text.Json.Nodes;
using System.Text.RegularExpressions;
using Admin.NET.Plugin.AiDOP.DataPlatform.Executors;
using Admin.NET.Plugin.AiDOP.Entity.DataPlatform;
using Microsoft.Extensions.Logging;
namespace Admin.NET.Plugin.AiDOP.DataPlatform.Inbound;
/// 入站接收主流水线。步骤顺序按任务书 §6.3 定死,勿换序。
public sealed class MdpInboundReceiveService : ITransient
{
private static readonly Regex TableNameRe = new(@"^[A-Za-z0-9_]+$", RegexOptions.Compiled);
private static readonly Regex IsoOffsetRe = new(@"[Zz]|[+\-]\d{2}:\d{2}", RegexOptions.Compiled);
private const int MaxRowBytes = 64 * 1024;
private static readonly TimeSpan InProgressWindow = TimeSpan.FromMinutes(10);
private readonly ISqlSugarClient _db;
private readonly MdpInboundAuthService _auth;
private readonly MdpInboundFieldMapper _mapper;
private readonly MdpStagingWriter _writer;
private readonly MdpInboundModuleTrigger _trigger;
private readonly MdpInboundSnapshotService _snapshots;
private readonly MdmMirrorUpsertService _mirror;
private readonly ILogger _logger;
public MdpInboundReceiveService(
ISqlSugarClient db,
MdpInboundAuthService auth,
MdpInboundFieldMapper mapper,
MdpStagingWriter writer,
MdpInboundModuleTrigger trigger,
MdpInboundSnapshotService snapshots,
MdmMirrorUpsertService mirror,
ILogger logger)
{
_db = db;
_auth = auth;
_mapper = mapper;
_writer = writer;
_trigger = trigger;
_snapshots = snapshots;
_mirror = mirror;
_logger = logger;
}
public async Task ReceiveAsync(MdpInboundReceiveArgs args, CancellationToken ct)
{
var sw = System.Diagnostics.Stopwatch.StartNew();
var entityCode = (args.EntityCode ?? string.Empty).Trim().ToUpperInvariant();
var hash = ToSha256Hex(args.RawBody);
MdpEntity entity;
MdpInboundGrant grant;
try
{
entity = await LoadInboundEntityAsync(entityCode, ct)
?? throw new MdpInboundBatchException(403, "entity not authorized");
var denied = await _auth.CheckAsync(args.AccessKey, args.TenantId, entityCode, args.ClientIp, ct);
if (denied != null)
return Fail(denied.Value.Status, denied.Value.Message);
grant = await _auth.GetGrantAsync(args.AccessKey, entityCode, ct)
?? throw new MdpInboundBatchException(403, "entity not authorized");
}
catch (MdpInboundBatchException ex)
{
return Fail(ex.HttpStatus, ex.Message);
}
MdpInboundParsedBatch batch;
try
{
batch = ParseBatch(args.RawBody, grant, entity);
}
catch (MdpInboundBatchException ex)
{
return Fail(ex.HttpStatus, ex.Message);
}
catch (JsonException)
{
return Fail(400, "invalid json");
}
var source = await _db.Queryable()
.Where(s => s.SourceCode == grant.SourceCode && s.Status == 1)
.FirstAsync(ct);
if (source == null)
return Fail(500, "inbound source not configured");
var sourceTable = string.IsNullOrWhiteSpace(entity.SourceTableName)
? entity.EntityCode
: entity.SourceTableName!;
try
{
await _snapshots.EnsureAcceptAsync(
batch.SnapshotId, args.AccessKey, args.TenantId, entity.EntityCode, batch.Seq, ct);
}
catch (MdpInboundBatchException ex)
{
return Fail(ex.HttpStatus, ex.Message);
}
MdpInboundRequest request;
try
{
request = await BeginRequestAsync(args, entity, hash, batch.SnapshotId, ct);
}
catch (MdpInboundBatchException ex)
{
return Fail(ex.HttpStatus, ex.Message);
}
if (request.Status == "COMMITTED" &&
string.Equals(request.RequestHash, hash, StringComparison.OrdinalIgnoreCase) &&
!string.IsNullOrWhiteSpace(request.ReceiptJson))
{
return Receipt(request.ReceiptJson, replay: true);
}
await TryArchiveEnvelopeAsync(args, request, hash, ct);
try
{
var maps = await _mapper.LoadAsync(args.TenantId, entity.EntityCode, args.AccessKey, ct);
var accepted = new List();
var rejected = new List();
var staleRejected = 0;
var equalTsOverride = 0;
foreach (var el in batch.Rows)
{
var prepared = PrepareRow(el, entity, maps, grant, args.TenantId, rejected, out var batchDeny);
if (batchDeny != null)
{
await MarkFailedAsync(request.Id, batchDeny.Message, sw.ElapsedMilliseconds, ct);
return Fail(batchDeny.HttpStatus, batchDeny.Message);
}
if (prepared == null)
continue;
if (entity.StaleGuardEnabled == 1)
{
var decision = await GuardStaleAsync(
entity, args.TenantId, sourceTable, prepared, grant.SourceCode, ct);
if (decision.RejectStale)
{
staleRejected++;
rejected.Add(new MdpInboundRejectedRow { BizKey = prepared.BizKey, Reason = "stale" });
continue;
}
if (decision.EqualTsOverride)
equalTsOverride++;
prepared = decision.Row ?? prepared;
}
accepted.Add(prepared);
}
var ctx = new MdpPullContext { TenantId = args.TenantId, BatchId = request.SyncBatchId };
var written = new List();
var tran = await _db.Ado.UseTranAsync(async () =>
{
foreach (var row in accepted)
{
await TransferIfNeededAsync(
entity, args.TenantId, sourceTable, row.BizKey, grant.SourceCode, ct);
var n = await _writer.UpsertAsync(
source, entity, sourceTable, row.Dict, row.RawJson, row.BizKey, ctx);
if (n == 0)
{
rejected.Add(new MdpInboundRejectedRow
{
BizKey = row.BizKey,
Reason = "tenant_unresolved"
});
}
else
{
written.Add(row);
}
}
var data = new MdpInboundReceiptData
{
EntityCode = entity.EntityCode,
Accepted = written.Count,
Rejected = rejected.Count,
RejectedRows = rejected,
IdempotentReplay = false,
SyncBatchId = request.SyncBatchId,
RequestId = request.Id,
TransformEnqueued = false,
StaleRejected = staleRejected,
EqualTsOverride = equalTsOverride
};
var resultCode = rejected.Count > 0 ? 2 : 0;
var body = Envelope(resultCode, resultCode == 2 ? "partial" : "ok", data);
var receiptJson = JsonSerializer.Serialize(body, MdpInboundJson.Options);
var now = DateTime.Now;
await _db.Updateable()
.SetColumns(r => new MdpInboundRequest
{
Status = "COMMITTED",
HttpStatus = 202,
ResultCode = resultCode,
Accepted = written.Count,
Rejected = rejected.Count,
StaleRejected = staleRejected,
EqualTsOverride = equalTsOverride,
ReceiptJson = receiptJson,
ErrorMsg = null,
ElapsedMs = (int)sw.ElapsedMilliseconds,
UpdateTime = now
})
.Where(r => r.Id == request.Id)
.ExecuteCommandAsync(ct);
if (!string.IsNullOrWhiteSpace(batch.SnapshotId) && batch.Seq is > 0)
await _snapshots.AdvanceAsync(batch.SnapshotId, batch.Seq.Value, written.Count, ct);
});
if (!tran.IsSuccess)
{
var err = tran.ErrorException?.Message ?? "transaction failed";
await MarkFailedAsync(request.Id, Truncate(err, 1000), sw.ElapsedMilliseconds, ct);
_logger.LogError(tran.ErrorException, "inbound write transaction failed requestId={Id}", request.Id);
return Fail(500, "inbound write failed");
}
var committed = await _db.Queryable()
.Where(r => r.Id == request.Id)
.FirstAsync(ct);
var receiptJson = committed?.ReceiptJson;
try
{
await _mirror.MirrorCommittedAsync(
entity.EntityCode, args.TenantId, grant.FactoryId, grant.SourceCode, written, ct);
}
catch (Exception ex)
{
_logger.LogError(ex, "inbound mirror after commit failed requestId={Id}", request.Id);
}
var enqueued = false;
try
{
enqueued = await _trigger.EnqueueForEntityAsync(
entity.EntityCode, args.TenantId, grant.FactoryId, ct);
}
catch (Exception ex)
{
_logger.LogError(ex, "inbound trigger after commit failed requestId={Id}", request.Id);
}
if (enqueued)
receiptJson = await PatchTransformEnqueuedAsync(request.Id, receiptJson, ct);
return Receipt(receiptJson, replay: false);
}
catch (Exception ex)
{
await MarkFailedAsync(request.Id, Truncate(ex.Message, 1000), sw.ElapsedMilliseconds, ct);
_logger.LogError(ex, "inbound receive failed requestId={Id}", request.Id);
return Fail(500, "inbound receive failed");
}
}
public async Task SchemaAsync(
string entityCode, string accessKey, long tenantId, string clientIp, CancellationToken ct)
{
var code = (entityCode ?? string.Empty).Trim().ToUpperInvariant();
var denied = await _auth.CheckAsync(accessKey, tenantId, code, clientIp, ct);
if (denied != null)
return Fail(denied.Value.Status, denied.Value.Message);
var entity = await LoadInboundEntityAsync(code, ct);
if (entity == null)
return Fail(403, "entity not authorized");
var maps = await _mapper.LoadAsync(tenantId, entity.EntityCode, accessKey, ct);
return Ok(_mapper.BuildSchema(entity, maps));
}
public async Task ReceiptAsync(
string syncBatchId, string accessKey, CancellationToken ct)
{
var row = await _db.Queryable()
.Where(r => r.SyncBatchId == syncBatchId && r.AccessKey == accessKey)
.FirstAsync(ct);
if (row == null)
return Fail(404, "receipt not found");
object receipt = null;
if (!string.IsNullOrWhiteSpace(row.ReceiptJson))
{
try
{
receipt = JsonSerializer.Deserialize