| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955 |
- 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;
- /// <summary>入站接收主流水线。步骤顺序按任务书 §6.3 定死,勿换序。</summary>
- 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<MdpInboundReceiveService> _logger;
- public MdpInboundReceiveService(
- ISqlSugarClient db,
- MdpInboundAuthService auth,
- MdpInboundFieldMapper mapper,
- MdpStagingWriter writer,
- MdpInboundModuleTrigger trigger,
- MdpInboundSnapshotService snapshots,
- MdmMirrorUpsertService mirror,
- ILogger<MdpInboundReceiveService> logger)
- {
- _db = db;
- _auth = auth;
- _mapper = mapper;
- _writer = writer;
- _trigger = trigger;
- _snapshots = snapshots;
- _mirror = mirror;
- _logger = logger;
- }
- public async Task<MdpInboundOutcome> 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<MdpSource>()
- .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<MdpInboundPreparedRow>();
- var rejected = new List<MdpInboundRejectedRow>();
- 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<MdpInboundPreparedRow>();
- 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<MdpInboundRequest>()
- .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<MdpInboundRequest>()
- .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<MdpInboundOutcome> 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<MdpInboundOutcome> ReceiptAsync(
- string syncBatchId, string accessKey, CancellationToken ct)
- {
- var row = await _db.Queryable<MdpInboundRequest>()
- .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<object>(row.ReceiptJson);
- }
- catch (JsonException)
- {
- receipt = row.ReceiptJson;
- }
- }
- return Ok(new
- {
- status = row.Status,
- accepted = row.Accepted,
- rejected = row.Rejected,
- staleRejected = row.StaleRejected,
- receiptJson = receipt
- });
- }
- /// <summary>NDJSON 分块:每行须为合法接收报文,并带根级 idempotencyKey。</summary>
- public async Task<MdpInboundOutcome> BulkAsync(MdpInboundReceiveArgs args, CancellationToken ct)
- {
- var text = Encoding.UTF8.GetString(args.RawBody ?? []);
- var lines = text.Replace("\r\n", "\n").Split('\n');
- var nonempty = lines.Count(l => !string.IsNullOrWhiteSpace(l));
- if (nonempty == 0)
- return Fail(400, "empty bulk");
- if (nonempty > 100)
- return Fail(413, "bulk line limit exceeded");
- var chunks = new List<object>();
- var lineNo = 0;
- foreach (var rawLine in lines)
- {
- var line = rawLine.Trim();
- if (string.IsNullOrEmpty(line))
- continue;
- lineNo++;
- string idem;
- try
- {
- using var doc = JsonDocument.Parse(line);
- if (doc.RootElement.ValueKind != JsonValueKind.Object
- || !TryGetPropertyIgnoreCase(doc.RootElement, "idempotencyKey", out var keyEl)
- || keyEl.ValueKind != JsonValueKind.String
- || string.IsNullOrWhiteSpace(keyEl.GetString()))
- {
- chunks.Add(new { line = lineNo, httpStatus = 400, message = "missing line idempotencyKey" });
- continue;
- }
- idem = keyEl.GetString();
- }
- catch (JsonException)
- {
- chunks.Add(new { line = lineNo, httpStatus = 400, message = "invalid json" });
- continue;
- }
- var outcome = await ReceiveAsync(new MdpInboundReceiveArgs
- {
- EntityCode = args.EntityCode,
- RawBody = Encoding.UTF8.GetBytes(line),
- AccessKey = args.AccessKey,
- TenantId = args.TenantId,
- IdempotencyKey = idem,
- ClientIp = args.ClientIp,
- Path = args.Path,
- Headers = args.Headers
- }, ct);
- chunks.Add(new { line = lineNo, httpStatus = outcome.HttpStatus, body = outcome.Body });
- }
- return new MdpInboundOutcome
- {
- HttpStatus = 202,
- Body = Envelope(0, "ok", new { chunks })
- };
- }
- /// <summary>日终对账:扫当日请求 + 按 sync_batch_id 取 stg 键,字典序拼接后 sha256。首期不加新索引。</summary>
- public async Task<MdpInboundOutcome> DigestAsync(
- string entityCode, string accessKey, long tenantId, string clientIp, string dateRaw, 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);
- if (!DateTime.TryParse(dateRaw, out var parsed))
- return Fail(400, "date required (yyyy-MM-dd)");
- var from = parsed.Date;
- var to = from.AddDays(1);
- var entity = await LoadInboundEntityAsync(code, ct);
- if (entity == null)
- return Fail(403, "entity not authorized");
- var reqs = await _db.Queryable<MdpInboundRequest>()
- .Where(r => r.AccessKey == accessKey && r.CreateTime >= from && r.CreateTime < to)
- .Where("UPPER(entity_code) = @code", new SugarParameter("@code", code))
- .ToListAsync(ct);
- var accepted = reqs.Where(r => r.Status == "COMMITTED").Sum(r => r.Accepted);
- var rejected = reqs.Sum(r => r.Rejected);
- var batchIds = reqs
- .Where(r => r.Status == "COMMITTED" && !string.IsNullOrWhiteSpace(r.SyncBatchId))
- .Select(r => r.SyncBatchId)
- .Distinct()
- .ToList();
- var keys = new List<string>();
- if (batchIds.Count > 0
- && !string.IsNullOrWhiteSpace(entity.TargetTableName)
- && TableNameRe.IsMatch(entity.TargetTableName))
- {
- keys = await _db.Ado.SqlQueryAsync<string>(
- $"""
- SELECT source_biz_key FROM `{entity.TargetTableName}`
- WHERE tenant_id=@t AND sync_batch_id IN ({string.Join(",", batchIds.Select((_, i) => "@b" + i))})
- """,
- new[] { new SugarParameter("@t", tenantId) }
- .Concat(batchIds.Select((b, i) => new SugarParameter("@b" + i, b)))
- .ToArray());
- }
- var sorted = keys
- .Where(k => !string.IsNullOrWhiteSpace(k))
- .Distinct(StringComparer.Ordinal)
- .OrderBy(k => k, StringComparer.Ordinal)
- .ToList();
- var digest = ToSha256Hex(Encoding.UTF8.GetBytes(string.Join("\n", sorted)));
- return Ok(new
- {
- entityCode = code,
- date = from.ToString("yyyy-MM-dd"),
- requestCount = reqs.Count,
- accepted,
- rejected,
- committedBatches = batchIds.Count,
- bizKeyCount = sorted.Count,
- bizKeyDigest = digest,
- sort = "ordinal-asc join \\n then sha256-hex"
- });
- }
- private async Task<MdpEntity> LoadInboundEntityAsync(string entityCode, CancellationToken ct)
- {
- return await _db.Queryable<MdpEntity>()
- .Where(e => e.InboundEnabled == 1)
- .Where("UPPER(entity_code) = @code", new SugarParameter("@code", entityCode))
- .FirstAsync(ct);
- }
- private static MdpInboundParsedBatch ParseBatch(byte[] rawBody, MdpInboundGrant grant, MdpEntity entity)
- {
- using var doc = JsonDocument.Parse(Encoding.UTF8.GetString(rawBody));
- var root = doc.RootElement;
- JsonElement listEl;
- string snapshotId = null;
- int? parsedSeq = null;
- if (root.ValueKind == JsonValueKind.Array)
- {
- listEl = root;
- }
- else if (root.ValueKind == JsonValueKind.Object
- && TryGetPropertyIgnoreCase(root, "data", out var data)
- && data.ValueKind == JsonValueKind.Object
- && TryGetPropertyIgnoreCase(data, "list", out var list)
- && list.ValueKind == JsonValueKind.Array)
- {
- listEl = list;
- if (TryGetPropertyIgnoreCase(data, "snapshotId", out var snap) && snap.ValueKind == JsonValueKind.String)
- snapshotId = snap.GetString();
- if (TryGetPropertyIgnoreCase(data, "seq", out var seqEl)
- && seqEl.ValueKind == JsonValueKind.Number
- && seqEl.TryGetInt32(out var seqVal))
- parsedSeq = seqVal;
- }
- else
- {
- throw new MdpInboundBatchException(400, "invalid payload: expect data.list or root array");
- }
- var rows = new List<JsonElement>();
- foreach (var item in listEl.EnumerateArray())
- {
- if (item.ValueKind != JsonValueKind.Object)
- throw new MdpInboundBatchException(400, "invalid payload: row must be object");
- AssertNoCaseDuplicateKeys(item);
- var raw = item.GetRawText();
- if (Encoding.UTF8.GetByteCount(raw) > MaxRowBytes)
- throw new MdpInboundBatchException(413, "row too large");
- rows.Add(item.Clone());
- }
- var limit = Math.Min(
- grant.BatchRowLimit > 0 ? grant.BatchRowLimit : 500,
- entity.InboundMaxRows > 0 ? entity.InboundMaxRows : 500);
- if (rows.Count > limit)
- throw new MdpInboundBatchException(413, "row limit exceeded");
- return new MdpInboundParsedBatch { Rows = rows, SnapshotId = snapshotId, Seq = parsedSeq };
- }
- private static void AssertNoCaseDuplicateKeys(JsonElement obj)
- {
- var seen = new HashSet<string>(StringComparer.OrdinalIgnoreCase);
- foreach (var prop in obj.EnumerateObject())
- {
- if (!seen.Add(prop.Name))
- throw new MdpInboundBatchException(400, "duplicate key");
- }
- }
- private async Task<MdpInboundRequest> BeginRequestAsync(
- MdpInboundReceiveArgs args, MdpEntity entity, string hash, string snapshotId, CancellationToken ct)
- {
- var now = DateTime.Now;
- var draft = new MdpInboundRequest
- {
- TenantId = args.TenantId,
- AccessKey = args.AccessKey,
- EntityCode = entity.EntityCode,
- IdempotencyKey = args.IdempotencyKey,
- RequestHash = hash,
- SyncBatchId = "PENDING",
- SnapshotId = snapshotId,
- Status = "RECEIVED",
- ClientIp = args.ClientIp,
- CreateTime = now,
- UpdateTime = now
- };
- try
- {
- var id = await _db.Insertable(draft).ExecuteReturnBigIdentityAsync();
- draft.Id = id;
- draft.SyncBatchId = $"INB_{now:yyyyMMdd}_{id:D6}";
- await _db.Updateable<MdpInboundRequest>()
- .SetColumns(r => new MdpInboundRequest { SyncBatchId = draft.SyncBatchId, UpdateTime = now })
- .Where(r => r.Id == id)
- .ExecuteCommandAsync(ct);
- return draft;
- }
- catch (Exception ex) when (IsDuplicateKey(ex))
- {
- var existing = await _db.Queryable<MdpInboundRequest>()
- .Where(r => r.TenantId == args.TenantId
- && r.AccessKey == args.AccessKey
- && r.IdempotencyKey == args.IdempotencyKey)
- .Where("UPPER(entity_code) = @code", new SugarParameter("@code", entity.EntityCode.ToUpperInvariant()))
- .FirstAsync(ct)
- ?? throw new MdpInboundBatchException(409, "in-progress");
- if (string.Equals(existing.Status, "COMMITTED", StringComparison.OrdinalIgnoreCase))
- {
- if (!string.Equals(existing.RequestHash, hash, StringComparison.OrdinalIgnoreCase))
- throw new MdpInboundBatchException(409, "idempotency key conflict");
- return existing;
- }
- if (string.Equals(existing.Status, "RECEIVED", StringComparison.OrdinalIgnoreCase)
- && existing.UpdateTime > now - InProgressWindow)
- throw new MdpInboundBatchException(409, "in-progress");
- var n = await _db.Ado.ExecuteCommandAsync(
- """
- UPDATE mdp_inbound_request
- SET status='RECEIVED', request_hash=@h, snapshot_id=@snap,
- error_msg=NULL, update_time=@now
- WHERE id=@id AND status=@old
- """,
- new SugarParameter("@h", hash),
- new SugarParameter("@snap", (object)snapshotId ?? DBNull.Value),
- new SugarParameter("@now", now),
- new SugarParameter("@id", existing.Id),
- new SugarParameter("@old", existing.Status));
- if (n != 1)
- throw new MdpInboundBatchException(409, "in-progress");
- existing.Status = "RECEIVED";
- existing.RequestHash = hash;
- existing.SnapshotId = snapshotId;
- existing.UpdateTime = now;
- if (string.IsNullOrWhiteSpace(existing.SyncBatchId) || existing.SyncBatchId == "PENDING")
- {
- existing.SyncBatchId = $"INB_{now:yyyyMMdd}_{existing.Id:D6}";
- await _db.Updateable<MdpInboundRequest>()
- .SetColumns(r => new MdpInboundRequest { SyncBatchId = existing.SyncBatchId })
- .Where(r => r.Id == existing.Id)
- .ExecuteCommandAsync(ct);
- }
- return existing;
- }
- }
- private async Task TryArchiveEnvelopeAsync(
- MdpInboundReceiveArgs args, MdpInboundRequest request, string hash, CancellationToken ct)
- {
- try
- {
- var exists = await _db.Queryable<MdpInboundEnvelope>()
- .Where(e => e.RequestId == request.Id)
- .AnyAsync(ct);
- if (exists)
- return;
- var headers = args.Headers
- .Where(kv => !IsSensitiveHeader(kv.Key))
- .ToDictionary(kv => kv.Key, kv => kv.Value, StringComparer.OrdinalIgnoreCase);
- await _db.Insertable(new MdpInboundEnvelope
- {
- RequestId = request.Id,
- TenantId = args.TenantId,
- AccessKey = args.AccessKey,
- EntityCode = request.EntityCode,
- Path = args.Path,
- HeadersJson = JsonSerializer.Serialize(headers, MdpInboundJson.Options),
- Body = Encoding.UTF8.GetString(args.RawBody),
- BodySha256 = hash,
- CreateTime = DateTime.Now
- }).ExecuteCommandAsync(ct);
- }
- catch (Exception ex)
- {
- _logger.LogError(ex, "inbound envelope archive failed requestId={Id}", request.Id);
- }
- }
- private MdpInboundPreparedRow PrepareRow(
- JsonElement el,
- MdpEntity entity,
- IReadOnlyList<MdpFieldMap> maps,
- MdpInboundGrant grant,
- long boundTenantId,
- List<MdpInboundRejectedRow> rejected,
- out MdpInboundBatchException batchDeny)
- {
- batchDeny = null;
- var incoming = MdpInboundFieldMapper.ElementToDict(el);
- var mapped = _mapper.MapRow(incoming, maps);
- var missing = _mapper.MissingRequired(mapped, maps);
- var previewKey = TryGetString(mapped, "bizKey");
- if (missing.Count > 0)
- {
- rejected.Add(new MdpInboundRejectedRow
- {
- BizKey = previewKey,
- Reason = "missing required field " + string.Join(",", missing)
- });
- return null;
- }
- var sourceUpdatedAt = TryGetString(mapped, "sourceUpdatedAt");
- if (!HasIsoOffset(sourceUpdatedAt))
- {
- rejected.Add(new MdpInboundRejectedRow { BizKey = previewKey, Reason = "sourceUpdatedAt required with timezone offset" });
- return null;
- }
- var sourceVersion = TryGetString(mapped, "sourceVersion");
- if (!string.IsNullOrWhiteSpace(sourceVersion) && !long.TryParse(sourceVersion, out _))
- {
- rejected.Add(new MdpInboundRejectedRow { BizKey = previewKey, Reason = "sourceVersion must be long" });
- return null;
- }
- if (TryGetLong(mapped, "tenant_id", out var rowTenant) || TryGetLong(mapped, "TenantId", out rowTenant))
- {
- if (rowTenant != boundTenantId)
- {
- batchDeny = new MdpInboundBatchException(403, "tenant mismatch");
- return null;
- }
- }
- mapped["tenant_id"] = boundTenantId;
- if (grant.FactoryId is > 0)
- {
- if (!TryGetLong(mapped, "factory_id", out var rowFactory) && !TryGetLong(mapped, "FactoryId", out rowFactory))
- {
- rejected.Add(new MdpInboundRejectedRow { BizKey = previewKey, Reason = "factory mismatch" });
- return null;
- }
- if (rowFactory != grant.FactoryId.Value)
- {
- rejected.Add(new MdpInboundRejectedRow { BizKey = previewKey, Reason = "factory mismatch" });
- return null;
- }
- }
- var bizKey = MdpStagingWriter.BuildBizKey(entity.BizKeyExpr, mapped);
- if (string.IsNullOrWhiteSpace(bizKey))
- {
- rejected.Add(new MdpInboundRejectedRow { BizKey = previewKey, Reason = "missing biz key fields" });
- return null;
- }
- if (!string.IsNullOrWhiteSpace(previewKey) && !string.Equals(previewKey, bizKey, StringComparison.Ordinal))
- {
- rejected.Add(new MdpInboundRejectedRow { BizKey = previewKey, Reason = "bizKey mismatch" });
- return null;
- }
- var op = TryGetString(mapped, "op");
- if (string.IsNullOrWhiteSpace(op))
- op = "upsert";
- if (string.Equals(op, "delete", StringComparison.OrdinalIgnoreCase))
- {
- if (!DeleteSupported(entity))
- {
- rejected.Add(new MdpInboundRejectedRow { BizKey = bizKey, Reason = "op=delete not supported" });
- return null;
- }
- mapped["is_deleted"] = 1;
- }
- var rawJson = JsonSerializer.Serialize(mapped, MdpInboundJson.Options);
- return new MdpInboundPreparedRow
- {
- Dict = mapped,
- RawJson = rawJson,
- BizKey = bizKey,
- SourceUpdatedAt = sourceUpdatedAt,
- SourceVersion = sourceVersion
- };
- }
- private async Task<(bool RejectStale, bool EqualTsOverride, MdpInboundPreparedRow Row)> GuardStaleAsync(
- MdpEntity entity,
- long tenantId,
- string sourceTable,
- MdpInboundPreparedRow incoming,
- string inboundSource,
- CancellationToken ct)
- {
- var existing = await LoadGuardRowAsync(entity.TargetTableName, tenantId, sourceTable, incoming.BizKey, inboundSource, ct);
- if (existing == null)
- return (false, false, incoming);
- var oldHasVer = long.TryParse(existing.SourceVersion, out var oldVer);
- var newHasVer = long.TryParse(incoming.SourceVersion, out var newVer);
- if (oldHasVer && newHasVer)
- {
- if (newVer < oldVer)
- return (true, false, incoming);
- if (newVer > oldVer)
- return (false, false, incoming);
- }
- else
- {
- var oldHasTs = DateTimeOffset.TryParse(existing.SourceUpdatedAt, out var oldTs);
- if (!oldHasTs)
- return (false, false, incoming);
- if (!DateTimeOffset.TryParse(incoming.SourceUpdatedAt, out var newTs))
- return (true, false, incoming);
- if (newTs < oldTs)
- return (true, false, incoming);
- if (newTs > oldTs)
- return (false, false, incoming);
- }
- var same = string.Equals(incoming.RawJson, existing.RawData, StringComparison.Ordinal);
- return (false, !same, incoming);
- }
- private async Task<MdpInboundStgGuardRow> LoadGuardRowAsync(
- string table, long tenantId, string sourceTable, string bizKey, string inboundSource, CancellationToken ct)
- {
- if (string.IsNullOrWhiteSpace(table) || !TableNameRe.IsMatch(table))
- throw new InvalidOperationException("illegal target_table_name");
- var rows = await _db.Ado.SqlQueryAsync<MdpInboundStgGuardRow>(
- $"""
- SELECT id AS Id,
- source_system AS SourceSystem,
- source_row_id AS SourceRowId,
- JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.sourceUpdatedAt')) AS SourceUpdatedAt,
- JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.sourceVersion')) AS SourceVersion,
- raw_data AS RawData
- FROM `{table}`
- WHERE tenant_id=@t AND source_table=@st AND source_biz_key=@bk
- ORDER BY CASE WHEN source_system=@sys THEN 0 ELSE 1 END
- LIMIT 1
- """,
- new SugarParameter("@t", tenantId),
- new SugarParameter("@st", sourceTable),
- new SugarParameter("@bk", bizKey),
- new SugarParameter("@sys", inboundSource));
- return rows?.FirstOrDefault();
- }
- private async Task TransferIfNeededAsync(
- MdpEntity entity, long tenantId, string sourceTable, string bizKey, string inboundSource, CancellationToken ct)
- {
- var existing = await LoadGuardRowAsync(entity.TargetTableName, tenantId, sourceTable, bizKey, inboundSource, ct);
- if (existing == null)
- return;
- if (string.Equals(existing.SourceSystem, inboundSource, StringComparison.OrdinalIgnoreCase))
- return;
- if (string.IsNullOrWhiteSpace(entity.TargetTableName) || !TableNameRe.IsMatch(entity.TargetTableName))
- throw new InvalidOperationException("illegal target_table_name");
- await _db.Ado.ExecuteCommandAsync(
- $"""
- UPDATE `{entity.TargetTableName}`
- SET source_system=@sys, source_row_id=@bk
- WHERE id=@id
- """,
- new SugarParameter("@sys", inboundSource),
- new SugarParameter("@bk", bizKey),
- new SugarParameter("@id", existing.Id));
- }
- private async Task MarkFailedAsync(long requestId, string error, long elapsedMs, CancellationToken ct)
- {
- try
- {
- await _db.Updateable<MdpInboundRequest>()
- .SetColumns(r => new MdpInboundRequest
- {
- Status = "FAILED",
- HttpStatus = 500,
- ResultCode = 500,
- ErrorMsg = error,
- ElapsedMs = (int)elapsedMs,
- UpdateTime = DateTime.Now
- })
- .Where(r => r.Id == requestId)
- .ExecuteCommandAsync(ct);
- }
- catch (Exception ex)
- {
- _logger.LogError(ex, "inbound mark FAILED failed requestId={Id}", requestId);
- }
- }
- private static bool DeleteSupported(MdpEntity entity) =>
- (entity.Remark ?? string.Empty).Contains("inbound_delete=supported", StringComparison.OrdinalIgnoreCase);
- private static bool HasIsoOffset(string value) =>
- !string.IsNullOrWhiteSpace(value)
- && IsoOffsetRe.IsMatch(value)
- && DateTimeOffset.TryParse(value, out _);
- private static bool TryGetPropertyIgnoreCase(JsonElement obj, string name, out JsonElement value)
- {
- foreach (var prop in obj.EnumerateObject())
- {
- if (string.Equals(prop.Name, name, StringComparison.OrdinalIgnoreCase))
- {
- value = prop.Value;
- return true;
- }
- }
- value = default;
- return false;
- }
- private static string TryGetString(IDictionary<string, object?> row, string key)
- {
- foreach (var kv in row)
- {
- if (!string.Equals(kv.Key, key, StringComparison.OrdinalIgnoreCase))
- continue;
- return kv.Value?.ToString();
- }
- return null;
- }
- private static bool TryGetLong(IDictionary<string, object?> row, string key, out long value)
- {
- value = 0;
- var s = TryGetString(row, key);
- return !string.IsNullOrWhiteSpace(s) && long.TryParse(s, out value);
- }
- private static bool IsSensitiveHeader(string name) =>
- string.Equals(name, "X-Signature", StringComparison.OrdinalIgnoreCase)
- || string.Equals(name, "Authorization", StringComparison.OrdinalIgnoreCase)
- || string.Equals(name, "Cookie", StringComparison.OrdinalIgnoreCase);
- private static bool IsDuplicateKey(Exception ex)
- {
- for (var e = ex; e != null; e = e.InnerException)
- {
- var m = e.Message ?? string.Empty;
- if (m.Contains("Duplicate", StringComparison.OrdinalIgnoreCase)
- || m.Contains("1062")
- || m.Contains("uk_inbound_idem", StringComparison.OrdinalIgnoreCase))
- return true;
- }
- return false;
- }
- private static string ToSha256Hex(byte[] raw) =>
- Convert.ToHexString(SHA256.HashData(raw ?? [])).ToLowerInvariant();
- private static string Truncate(string s, int max) =>
- string.IsNullOrEmpty(s) ? s : (s.Length <= max ? s : s[..max]);
- private static object Envelope(int code, string message, object data) =>
- new { code, message, data };
- private static MdpInboundOutcome Fail(int status, string message) =>
- new() { HttpStatus = status, Body = Envelope(status, message, null) };
- private static MdpInboundOutcome Ok(object data) =>
- new() { HttpStatus = 200, Body = Envelope(0, "ok", data) };
- private async Task<string> PatchTransformEnqueuedAsync(long requestId, string receiptJson, CancellationToken ct)
- {
- if (string.IsNullOrWhiteSpace(receiptJson))
- return receiptJson;
- try
- {
- var node = JsonNode.Parse(receiptJson);
- if (node is JsonObject obj && obj["data"] is JsonObject data)
- data["transformEnqueued"] = true;
- var updated = node.ToJsonString(MdpInboundJson.Options);
- await _db.Updateable<MdpInboundRequest>()
- .SetColumns(r => new MdpInboundRequest
- {
- ReceiptJson = updated,
- UpdateTime = DateTime.Now
- })
- .Where(r => r.Id == requestId)
- .ExecuteCommandAsync(ct);
- return updated;
- }
- catch (Exception ex)
- {
- _logger.LogError(ex, "inbound patch transformEnqueued failed requestId={Id}", requestId);
- return receiptJson;
- }
- }
- private static MdpInboundOutcome Receipt(string receiptJson, bool replay)
- {
- if (string.IsNullOrWhiteSpace(receiptJson))
- return new MdpInboundOutcome { HttpStatus = 202, Body = Envelope(0, "ok", null) };
- try
- {
- var node = JsonNode.Parse(receiptJson);
- if (replay && node is JsonObject obj && obj["data"] is JsonObject data)
- data["idempotentReplay"] = true;
- return new MdpInboundOutcome { HttpStatus = 202, Body = node };
- }
- catch (JsonException)
- {
- return new MdpInboundOutcome { HttpStatus = 202, Body = Envelope(0, "ok", receiptJson) };
- }
- }
- }
|