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(row.ReceiptJson); } catch (JsonException) { receipt = row.ReceiptJson; } } return Ok(new { status = row.Status, accepted = row.Accepted, rejected = row.Rejected, staleRejected = row.StaleRejected, receiptJson = receipt }); } /// NDJSON 分块:每行须为合法接收报文,并带根级 idempotencyKey。 public async Task 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(); 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 }) }; } /// 日终对账:扫当日请求 + 按 sync_batch_id 取 stg 键,字典序拼接后 sha256。首期不加新索引。 public async Task 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() .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(); if (batchIds.Count > 0 && !string.IsNullOrWhiteSpace(entity.TargetTableName) && TableNameRe.IsMatch(entity.TargetTableName)) { keys = await _db.Ado.SqlQueryAsync( $""" 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 LoadInboundEntityAsync(string entityCode, CancellationToken ct) { return await _db.Queryable() .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(); 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(StringComparer.OrdinalIgnoreCase); foreach (var prop in obj.EnumerateObject()) { if (!seen.Add(prop.Name)) throw new MdpInboundBatchException(400, "duplicate key"); } } private async Task 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() .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() .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() .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() .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 maps, MdpInboundGrant grant, long boundTenantId, List 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 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( $""" 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() .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 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 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 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() .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) }; } } }