MdpInboundReceiveService.cs 37 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955
  1. using System.Security.Cryptography;
  2. using System.Text;
  3. using System.Text.Json;
  4. using System.Text.Json.Nodes;
  5. using System.Text.RegularExpressions;
  6. using Admin.NET.Plugin.AiDOP.DataPlatform.Executors;
  7. using Admin.NET.Plugin.AiDOP.Entity.DataPlatform;
  8. using Microsoft.Extensions.Logging;
  9. namespace Admin.NET.Plugin.AiDOP.DataPlatform.Inbound;
  10. /// <summary>入站接收主流水线。步骤顺序按任务书 §6.3 定死,勿换序。</summary>
  11. public sealed class MdpInboundReceiveService : ITransient
  12. {
  13. private static readonly Regex TableNameRe = new(@"^[A-Za-z0-9_]+$", RegexOptions.Compiled);
  14. private static readonly Regex IsoOffsetRe = new(@"[Zz]|[+\-]\d{2}:\d{2}", RegexOptions.Compiled);
  15. private const int MaxRowBytes = 64 * 1024;
  16. private static readonly TimeSpan InProgressWindow = TimeSpan.FromMinutes(10);
  17. private readonly ISqlSugarClient _db;
  18. private readonly MdpInboundAuthService _auth;
  19. private readonly MdpInboundFieldMapper _mapper;
  20. private readonly MdpStagingWriter _writer;
  21. private readonly MdpInboundModuleTrigger _trigger;
  22. private readonly MdpInboundSnapshotService _snapshots;
  23. private readonly MdmMirrorUpsertService _mirror;
  24. private readonly ILogger<MdpInboundReceiveService> _logger;
  25. public MdpInboundReceiveService(
  26. ISqlSugarClient db,
  27. MdpInboundAuthService auth,
  28. MdpInboundFieldMapper mapper,
  29. MdpStagingWriter writer,
  30. MdpInboundModuleTrigger trigger,
  31. MdpInboundSnapshotService snapshots,
  32. MdmMirrorUpsertService mirror,
  33. ILogger<MdpInboundReceiveService> logger)
  34. {
  35. _db = db;
  36. _auth = auth;
  37. _mapper = mapper;
  38. _writer = writer;
  39. _trigger = trigger;
  40. _snapshots = snapshots;
  41. _mirror = mirror;
  42. _logger = logger;
  43. }
  44. public async Task<MdpInboundOutcome> ReceiveAsync(MdpInboundReceiveArgs args, CancellationToken ct)
  45. {
  46. var sw = System.Diagnostics.Stopwatch.StartNew();
  47. var entityCode = (args.EntityCode ?? string.Empty).Trim().ToUpperInvariant();
  48. var hash = ToSha256Hex(args.RawBody);
  49. MdpEntity entity;
  50. MdpInboundGrant grant;
  51. try
  52. {
  53. entity = await LoadInboundEntityAsync(entityCode, ct)
  54. ?? throw new MdpInboundBatchException(403, "entity not authorized");
  55. var denied = await _auth.CheckAsync(args.AccessKey, args.TenantId, entityCode, args.ClientIp, ct);
  56. if (denied != null)
  57. return Fail(denied.Value.Status, denied.Value.Message);
  58. grant = await _auth.GetGrantAsync(args.AccessKey, entityCode, ct)
  59. ?? throw new MdpInboundBatchException(403, "entity not authorized");
  60. }
  61. catch (MdpInboundBatchException ex)
  62. {
  63. return Fail(ex.HttpStatus, ex.Message);
  64. }
  65. MdpInboundParsedBatch batch;
  66. try
  67. {
  68. batch = ParseBatch(args.RawBody, grant, entity);
  69. }
  70. catch (MdpInboundBatchException ex)
  71. {
  72. return Fail(ex.HttpStatus, ex.Message);
  73. }
  74. catch (JsonException)
  75. {
  76. return Fail(400, "invalid json");
  77. }
  78. var source = await _db.Queryable<MdpSource>()
  79. .Where(s => s.SourceCode == grant.SourceCode && s.Status == 1)
  80. .FirstAsync(ct);
  81. if (source == null)
  82. return Fail(500, "inbound source not configured");
  83. var sourceTable = string.IsNullOrWhiteSpace(entity.SourceTableName)
  84. ? entity.EntityCode
  85. : entity.SourceTableName!;
  86. try
  87. {
  88. await _snapshots.EnsureAcceptAsync(
  89. batch.SnapshotId, args.AccessKey, args.TenantId, entity.EntityCode, batch.Seq, ct);
  90. }
  91. catch (MdpInboundBatchException ex)
  92. {
  93. return Fail(ex.HttpStatus, ex.Message);
  94. }
  95. MdpInboundRequest request;
  96. try
  97. {
  98. request = await BeginRequestAsync(args, entity, hash, batch.SnapshotId, ct);
  99. }
  100. catch (MdpInboundBatchException ex)
  101. {
  102. return Fail(ex.HttpStatus, ex.Message);
  103. }
  104. if (request.Status == "COMMITTED" &&
  105. string.Equals(request.RequestHash, hash, StringComparison.OrdinalIgnoreCase) &&
  106. !string.IsNullOrWhiteSpace(request.ReceiptJson))
  107. {
  108. return Receipt(request.ReceiptJson, replay: true);
  109. }
  110. await TryArchiveEnvelopeAsync(args, request, hash, ct);
  111. try
  112. {
  113. var maps = await _mapper.LoadAsync(args.TenantId, entity.EntityCode, args.AccessKey, ct);
  114. var accepted = new List<MdpInboundPreparedRow>();
  115. var rejected = new List<MdpInboundRejectedRow>();
  116. var staleRejected = 0;
  117. var equalTsOverride = 0;
  118. foreach (var el in batch.Rows)
  119. {
  120. var prepared = PrepareRow(el, entity, maps, grant, args.TenantId, rejected, out var batchDeny);
  121. if (batchDeny != null)
  122. {
  123. await MarkFailedAsync(request.Id, batchDeny.Message, sw.ElapsedMilliseconds, ct);
  124. return Fail(batchDeny.HttpStatus, batchDeny.Message);
  125. }
  126. if (prepared == null)
  127. continue;
  128. if (entity.StaleGuardEnabled == 1)
  129. {
  130. var decision = await GuardStaleAsync(
  131. entity, args.TenantId, sourceTable, prepared, grant.SourceCode, ct);
  132. if (decision.RejectStale)
  133. {
  134. staleRejected++;
  135. rejected.Add(new MdpInboundRejectedRow { BizKey = prepared.BizKey, Reason = "stale" });
  136. continue;
  137. }
  138. if (decision.EqualTsOverride)
  139. equalTsOverride++;
  140. prepared = decision.Row ?? prepared;
  141. }
  142. accepted.Add(prepared);
  143. }
  144. var ctx = new MdpPullContext { TenantId = args.TenantId, BatchId = request.SyncBatchId };
  145. var written = new List<MdpInboundPreparedRow>();
  146. var tran = await _db.Ado.UseTranAsync(async () =>
  147. {
  148. foreach (var row in accepted)
  149. {
  150. await TransferIfNeededAsync(
  151. entity, args.TenantId, sourceTable, row.BizKey, grant.SourceCode, ct);
  152. var n = await _writer.UpsertAsync(
  153. source, entity, sourceTable, row.Dict, row.RawJson, row.BizKey, ctx);
  154. if (n == 0)
  155. {
  156. rejected.Add(new MdpInboundRejectedRow
  157. {
  158. BizKey = row.BizKey,
  159. Reason = "tenant_unresolved"
  160. });
  161. }
  162. else
  163. {
  164. written.Add(row);
  165. }
  166. }
  167. var data = new MdpInboundReceiptData
  168. {
  169. EntityCode = entity.EntityCode,
  170. Accepted = written.Count,
  171. Rejected = rejected.Count,
  172. RejectedRows = rejected,
  173. IdempotentReplay = false,
  174. SyncBatchId = request.SyncBatchId,
  175. RequestId = request.Id,
  176. TransformEnqueued = false,
  177. StaleRejected = staleRejected,
  178. EqualTsOverride = equalTsOverride
  179. };
  180. var resultCode = rejected.Count > 0 ? 2 : 0;
  181. var body = Envelope(resultCode, resultCode == 2 ? "partial" : "ok", data);
  182. var receiptJson = JsonSerializer.Serialize(body, MdpInboundJson.Options);
  183. var now = DateTime.Now;
  184. await _db.Updateable<MdpInboundRequest>()
  185. .SetColumns(r => new MdpInboundRequest
  186. {
  187. Status = "COMMITTED",
  188. HttpStatus = 202,
  189. ResultCode = resultCode,
  190. Accepted = written.Count,
  191. Rejected = rejected.Count,
  192. StaleRejected = staleRejected,
  193. EqualTsOverride = equalTsOverride,
  194. ReceiptJson = receiptJson,
  195. ErrorMsg = null,
  196. ElapsedMs = (int)sw.ElapsedMilliseconds,
  197. UpdateTime = now
  198. })
  199. .Where(r => r.Id == request.Id)
  200. .ExecuteCommandAsync(ct);
  201. if (!string.IsNullOrWhiteSpace(batch.SnapshotId) && batch.Seq is > 0)
  202. await _snapshots.AdvanceAsync(batch.SnapshotId, batch.Seq.Value, written.Count, ct);
  203. });
  204. if (!tran.IsSuccess)
  205. {
  206. var err = tran.ErrorException?.Message ?? "transaction failed";
  207. await MarkFailedAsync(request.Id, Truncate(err, 1000), sw.ElapsedMilliseconds, ct);
  208. _logger.LogError(tran.ErrorException, "inbound write transaction failed requestId={Id}", request.Id);
  209. return Fail(500, "inbound write failed");
  210. }
  211. var committed = await _db.Queryable<MdpInboundRequest>()
  212. .Where(r => r.Id == request.Id)
  213. .FirstAsync(ct);
  214. var receiptJson = committed?.ReceiptJson;
  215. try
  216. {
  217. await _mirror.MirrorCommittedAsync(
  218. entity.EntityCode, args.TenantId, grant.FactoryId, grant.SourceCode, written, ct);
  219. }
  220. catch (Exception ex)
  221. {
  222. _logger.LogError(ex, "inbound mirror after commit failed requestId={Id}", request.Id);
  223. }
  224. var enqueued = false;
  225. try
  226. {
  227. enqueued = await _trigger.EnqueueForEntityAsync(
  228. entity.EntityCode, args.TenantId, grant.FactoryId, ct);
  229. }
  230. catch (Exception ex)
  231. {
  232. _logger.LogError(ex, "inbound trigger after commit failed requestId={Id}", request.Id);
  233. }
  234. if (enqueued)
  235. receiptJson = await PatchTransformEnqueuedAsync(request.Id, receiptJson, ct);
  236. return Receipt(receiptJson, replay: false);
  237. }
  238. catch (Exception ex)
  239. {
  240. await MarkFailedAsync(request.Id, Truncate(ex.Message, 1000), sw.ElapsedMilliseconds, ct);
  241. _logger.LogError(ex, "inbound receive failed requestId={Id}", request.Id);
  242. return Fail(500, "inbound receive failed");
  243. }
  244. }
  245. public async Task<MdpInboundOutcome> SchemaAsync(
  246. string entityCode, string accessKey, long tenantId, string clientIp, CancellationToken ct)
  247. {
  248. var code = (entityCode ?? string.Empty).Trim().ToUpperInvariant();
  249. var denied = await _auth.CheckAsync(accessKey, tenantId, code, clientIp, ct);
  250. if (denied != null)
  251. return Fail(denied.Value.Status, denied.Value.Message);
  252. var entity = await LoadInboundEntityAsync(code, ct);
  253. if (entity == null)
  254. return Fail(403, "entity not authorized");
  255. var maps = await _mapper.LoadAsync(tenantId, entity.EntityCode, accessKey, ct);
  256. return Ok(_mapper.BuildSchema(entity, maps));
  257. }
  258. public async Task<MdpInboundOutcome> ReceiptAsync(
  259. string syncBatchId, string accessKey, CancellationToken ct)
  260. {
  261. var row = await _db.Queryable<MdpInboundRequest>()
  262. .Where(r => r.SyncBatchId == syncBatchId && r.AccessKey == accessKey)
  263. .FirstAsync(ct);
  264. if (row == null)
  265. return Fail(404, "receipt not found");
  266. object receipt = null;
  267. if (!string.IsNullOrWhiteSpace(row.ReceiptJson))
  268. {
  269. try
  270. {
  271. receipt = JsonSerializer.Deserialize<object>(row.ReceiptJson);
  272. }
  273. catch (JsonException)
  274. {
  275. receipt = row.ReceiptJson;
  276. }
  277. }
  278. return Ok(new
  279. {
  280. status = row.Status,
  281. accepted = row.Accepted,
  282. rejected = row.Rejected,
  283. staleRejected = row.StaleRejected,
  284. receiptJson = receipt
  285. });
  286. }
  287. /// <summary>NDJSON 分块:每行须为合法接收报文,并带根级 idempotencyKey。</summary>
  288. public async Task<MdpInboundOutcome> BulkAsync(MdpInboundReceiveArgs args, CancellationToken ct)
  289. {
  290. var text = Encoding.UTF8.GetString(args.RawBody ?? []);
  291. var lines = text.Replace("\r\n", "\n").Split('\n');
  292. var nonempty = lines.Count(l => !string.IsNullOrWhiteSpace(l));
  293. if (nonempty == 0)
  294. return Fail(400, "empty bulk");
  295. if (nonempty > 100)
  296. return Fail(413, "bulk line limit exceeded");
  297. var chunks = new List<object>();
  298. var lineNo = 0;
  299. foreach (var rawLine in lines)
  300. {
  301. var line = rawLine.Trim();
  302. if (string.IsNullOrEmpty(line))
  303. continue;
  304. lineNo++;
  305. string idem;
  306. try
  307. {
  308. using var doc = JsonDocument.Parse(line);
  309. if (doc.RootElement.ValueKind != JsonValueKind.Object
  310. || !TryGetPropertyIgnoreCase(doc.RootElement, "idempotencyKey", out var keyEl)
  311. || keyEl.ValueKind != JsonValueKind.String
  312. || string.IsNullOrWhiteSpace(keyEl.GetString()))
  313. {
  314. chunks.Add(new { line = lineNo, httpStatus = 400, message = "missing line idempotencyKey" });
  315. continue;
  316. }
  317. idem = keyEl.GetString();
  318. }
  319. catch (JsonException)
  320. {
  321. chunks.Add(new { line = lineNo, httpStatus = 400, message = "invalid json" });
  322. continue;
  323. }
  324. var outcome = await ReceiveAsync(new MdpInboundReceiveArgs
  325. {
  326. EntityCode = args.EntityCode,
  327. RawBody = Encoding.UTF8.GetBytes(line),
  328. AccessKey = args.AccessKey,
  329. TenantId = args.TenantId,
  330. IdempotencyKey = idem,
  331. ClientIp = args.ClientIp,
  332. Path = args.Path,
  333. Headers = args.Headers
  334. }, ct);
  335. chunks.Add(new { line = lineNo, httpStatus = outcome.HttpStatus, body = outcome.Body });
  336. }
  337. return new MdpInboundOutcome
  338. {
  339. HttpStatus = 202,
  340. Body = Envelope(0, "ok", new { chunks })
  341. };
  342. }
  343. /// <summary>日终对账:扫当日请求 + 按 sync_batch_id 取 stg 键,字典序拼接后 sha256。首期不加新索引。</summary>
  344. public async Task<MdpInboundOutcome> DigestAsync(
  345. string entityCode, string accessKey, long tenantId, string clientIp, string dateRaw, CancellationToken ct)
  346. {
  347. var code = (entityCode ?? string.Empty).Trim().ToUpperInvariant();
  348. var denied = await _auth.CheckAsync(accessKey, tenantId, code, clientIp, ct);
  349. if (denied != null)
  350. return Fail(denied.Value.Status, denied.Value.Message);
  351. if (!DateTime.TryParse(dateRaw, out var parsed))
  352. return Fail(400, "date required (yyyy-MM-dd)");
  353. var from = parsed.Date;
  354. var to = from.AddDays(1);
  355. var entity = await LoadInboundEntityAsync(code, ct);
  356. if (entity == null)
  357. return Fail(403, "entity not authorized");
  358. var reqs = await _db.Queryable<MdpInboundRequest>()
  359. .Where(r => r.AccessKey == accessKey && r.CreateTime >= from && r.CreateTime < to)
  360. .Where("UPPER(entity_code) = @code", new SugarParameter("@code", code))
  361. .ToListAsync(ct);
  362. var accepted = reqs.Where(r => r.Status == "COMMITTED").Sum(r => r.Accepted);
  363. var rejected = reqs.Sum(r => r.Rejected);
  364. var batchIds = reqs
  365. .Where(r => r.Status == "COMMITTED" && !string.IsNullOrWhiteSpace(r.SyncBatchId))
  366. .Select(r => r.SyncBatchId)
  367. .Distinct()
  368. .ToList();
  369. var keys = new List<string>();
  370. if (batchIds.Count > 0
  371. && !string.IsNullOrWhiteSpace(entity.TargetTableName)
  372. && TableNameRe.IsMatch(entity.TargetTableName))
  373. {
  374. keys = await _db.Ado.SqlQueryAsync<string>(
  375. $"""
  376. SELECT source_biz_key FROM `{entity.TargetTableName}`
  377. WHERE tenant_id=@t AND sync_batch_id IN ({string.Join(",", batchIds.Select((_, i) => "@b" + i))})
  378. """,
  379. new[] { new SugarParameter("@t", tenantId) }
  380. .Concat(batchIds.Select((b, i) => new SugarParameter("@b" + i, b)))
  381. .ToArray());
  382. }
  383. var sorted = keys
  384. .Where(k => !string.IsNullOrWhiteSpace(k))
  385. .Distinct(StringComparer.Ordinal)
  386. .OrderBy(k => k, StringComparer.Ordinal)
  387. .ToList();
  388. var digest = ToSha256Hex(Encoding.UTF8.GetBytes(string.Join("\n", sorted)));
  389. return Ok(new
  390. {
  391. entityCode = code,
  392. date = from.ToString("yyyy-MM-dd"),
  393. requestCount = reqs.Count,
  394. accepted,
  395. rejected,
  396. committedBatches = batchIds.Count,
  397. bizKeyCount = sorted.Count,
  398. bizKeyDigest = digest,
  399. sort = "ordinal-asc join \\n then sha256-hex"
  400. });
  401. }
  402. private async Task<MdpEntity> LoadInboundEntityAsync(string entityCode, CancellationToken ct)
  403. {
  404. return await _db.Queryable<MdpEntity>()
  405. .Where(e => e.InboundEnabled == 1)
  406. .Where("UPPER(entity_code) = @code", new SugarParameter("@code", entityCode))
  407. .FirstAsync(ct);
  408. }
  409. private static MdpInboundParsedBatch ParseBatch(byte[] rawBody, MdpInboundGrant grant, MdpEntity entity)
  410. {
  411. using var doc = JsonDocument.Parse(Encoding.UTF8.GetString(rawBody));
  412. var root = doc.RootElement;
  413. JsonElement listEl;
  414. string snapshotId = null;
  415. int? parsedSeq = null;
  416. if (root.ValueKind == JsonValueKind.Array)
  417. {
  418. listEl = root;
  419. }
  420. else if (root.ValueKind == JsonValueKind.Object
  421. && TryGetPropertyIgnoreCase(root, "data", out var data)
  422. && data.ValueKind == JsonValueKind.Object
  423. && TryGetPropertyIgnoreCase(data, "list", out var list)
  424. && list.ValueKind == JsonValueKind.Array)
  425. {
  426. listEl = list;
  427. if (TryGetPropertyIgnoreCase(data, "snapshotId", out var snap) && snap.ValueKind == JsonValueKind.String)
  428. snapshotId = snap.GetString();
  429. if (TryGetPropertyIgnoreCase(data, "seq", out var seqEl)
  430. && seqEl.ValueKind == JsonValueKind.Number
  431. && seqEl.TryGetInt32(out var seqVal))
  432. parsedSeq = seqVal;
  433. }
  434. else
  435. {
  436. throw new MdpInboundBatchException(400, "invalid payload: expect data.list or root array");
  437. }
  438. var rows = new List<JsonElement>();
  439. foreach (var item in listEl.EnumerateArray())
  440. {
  441. if (item.ValueKind != JsonValueKind.Object)
  442. throw new MdpInboundBatchException(400, "invalid payload: row must be object");
  443. AssertNoCaseDuplicateKeys(item);
  444. var raw = item.GetRawText();
  445. if (Encoding.UTF8.GetByteCount(raw) > MaxRowBytes)
  446. throw new MdpInboundBatchException(413, "row too large");
  447. rows.Add(item.Clone());
  448. }
  449. var limit = Math.Min(
  450. grant.BatchRowLimit > 0 ? grant.BatchRowLimit : 500,
  451. entity.InboundMaxRows > 0 ? entity.InboundMaxRows : 500);
  452. if (rows.Count > limit)
  453. throw new MdpInboundBatchException(413, "row limit exceeded");
  454. return new MdpInboundParsedBatch { Rows = rows, SnapshotId = snapshotId, Seq = parsedSeq };
  455. }
  456. private static void AssertNoCaseDuplicateKeys(JsonElement obj)
  457. {
  458. var seen = new HashSet<string>(StringComparer.OrdinalIgnoreCase);
  459. foreach (var prop in obj.EnumerateObject())
  460. {
  461. if (!seen.Add(prop.Name))
  462. throw new MdpInboundBatchException(400, "duplicate key");
  463. }
  464. }
  465. private async Task<MdpInboundRequest> BeginRequestAsync(
  466. MdpInboundReceiveArgs args, MdpEntity entity, string hash, string snapshotId, CancellationToken ct)
  467. {
  468. var now = DateTime.Now;
  469. var draft = new MdpInboundRequest
  470. {
  471. TenantId = args.TenantId,
  472. AccessKey = args.AccessKey,
  473. EntityCode = entity.EntityCode,
  474. IdempotencyKey = args.IdempotencyKey,
  475. RequestHash = hash,
  476. SyncBatchId = "PENDING",
  477. SnapshotId = snapshotId,
  478. Status = "RECEIVED",
  479. ClientIp = args.ClientIp,
  480. CreateTime = now,
  481. UpdateTime = now
  482. };
  483. try
  484. {
  485. var id = await _db.Insertable(draft).ExecuteReturnBigIdentityAsync();
  486. draft.Id = id;
  487. draft.SyncBatchId = $"INB_{now:yyyyMMdd}_{id:D6}";
  488. await _db.Updateable<MdpInboundRequest>()
  489. .SetColumns(r => new MdpInboundRequest { SyncBatchId = draft.SyncBatchId, UpdateTime = now })
  490. .Where(r => r.Id == id)
  491. .ExecuteCommandAsync(ct);
  492. return draft;
  493. }
  494. catch (Exception ex) when (IsDuplicateKey(ex))
  495. {
  496. var existing = await _db.Queryable<MdpInboundRequest>()
  497. .Where(r => r.TenantId == args.TenantId
  498. && r.AccessKey == args.AccessKey
  499. && r.IdempotencyKey == args.IdempotencyKey)
  500. .Where("UPPER(entity_code) = @code", new SugarParameter("@code", entity.EntityCode.ToUpperInvariant()))
  501. .FirstAsync(ct)
  502. ?? throw new MdpInboundBatchException(409, "in-progress");
  503. if (string.Equals(existing.Status, "COMMITTED", StringComparison.OrdinalIgnoreCase))
  504. {
  505. if (!string.Equals(existing.RequestHash, hash, StringComparison.OrdinalIgnoreCase))
  506. throw new MdpInboundBatchException(409, "idempotency key conflict");
  507. return existing;
  508. }
  509. if (string.Equals(existing.Status, "RECEIVED", StringComparison.OrdinalIgnoreCase)
  510. && existing.UpdateTime > now - InProgressWindow)
  511. throw new MdpInboundBatchException(409, "in-progress");
  512. var n = await _db.Ado.ExecuteCommandAsync(
  513. """
  514. UPDATE mdp_inbound_request
  515. SET status='RECEIVED', request_hash=@h, snapshot_id=@snap,
  516. error_msg=NULL, update_time=@now
  517. WHERE id=@id AND status=@old
  518. """,
  519. new SugarParameter("@h", hash),
  520. new SugarParameter("@snap", (object)snapshotId ?? DBNull.Value),
  521. new SugarParameter("@now", now),
  522. new SugarParameter("@id", existing.Id),
  523. new SugarParameter("@old", existing.Status));
  524. if (n != 1)
  525. throw new MdpInboundBatchException(409, "in-progress");
  526. existing.Status = "RECEIVED";
  527. existing.RequestHash = hash;
  528. existing.SnapshotId = snapshotId;
  529. existing.UpdateTime = now;
  530. if (string.IsNullOrWhiteSpace(existing.SyncBatchId) || existing.SyncBatchId == "PENDING")
  531. {
  532. existing.SyncBatchId = $"INB_{now:yyyyMMdd}_{existing.Id:D6}";
  533. await _db.Updateable<MdpInboundRequest>()
  534. .SetColumns(r => new MdpInboundRequest { SyncBatchId = existing.SyncBatchId })
  535. .Where(r => r.Id == existing.Id)
  536. .ExecuteCommandAsync(ct);
  537. }
  538. return existing;
  539. }
  540. }
  541. private async Task TryArchiveEnvelopeAsync(
  542. MdpInboundReceiveArgs args, MdpInboundRequest request, string hash, CancellationToken ct)
  543. {
  544. try
  545. {
  546. var exists = await _db.Queryable<MdpInboundEnvelope>()
  547. .Where(e => e.RequestId == request.Id)
  548. .AnyAsync(ct);
  549. if (exists)
  550. return;
  551. var headers = args.Headers
  552. .Where(kv => !IsSensitiveHeader(kv.Key))
  553. .ToDictionary(kv => kv.Key, kv => kv.Value, StringComparer.OrdinalIgnoreCase);
  554. await _db.Insertable(new MdpInboundEnvelope
  555. {
  556. RequestId = request.Id,
  557. TenantId = args.TenantId,
  558. AccessKey = args.AccessKey,
  559. EntityCode = request.EntityCode,
  560. Path = args.Path,
  561. HeadersJson = JsonSerializer.Serialize(headers, MdpInboundJson.Options),
  562. Body = Encoding.UTF8.GetString(args.RawBody),
  563. BodySha256 = hash,
  564. CreateTime = DateTime.Now
  565. }).ExecuteCommandAsync(ct);
  566. }
  567. catch (Exception ex)
  568. {
  569. _logger.LogError(ex, "inbound envelope archive failed requestId={Id}", request.Id);
  570. }
  571. }
  572. private MdpInboundPreparedRow PrepareRow(
  573. JsonElement el,
  574. MdpEntity entity,
  575. IReadOnlyList<MdpFieldMap> maps,
  576. MdpInboundGrant grant,
  577. long boundTenantId,
  578. List<MdpInboundRejectedRow> rejected,
  579. out MdpInboundBatchException batchDeny)
  580. {
  581. batchDeny = null;
  582. var incoming = MdpInboundFieldMapper.ElementToDict(el);
  583. var mapped = _mapper.MapRow(incoming, maps);
  584. var missing = _mapper.MissingRequired(mapped, maps);
  585. var previewKey = TryGetString(mapped, "bizKey");
  586. if (missing.Count > 0)
  587. {
  588. rejected.Add(new MdpInboundRejectedRow
  589. {
  590. BizKey = previewKey,
  591. Reason = "missing required field " + string.Join(",", missing)
  592. });
  593. return null;
  594. }
  595. var sourceUpdatedAt = TryGetString(mapped, "sourceUpdatedAt");
  596. if (!HasIsoOffset(sourceUpdatedAt))
  597. {
  598. rejected.Add(new MdpInboundRejectedRow { BizKey = previewKey, Reason = "sourceUpdatedAt required with timezone offset" });
  599. return null;
  600. }
  601. var sourceVersion = TryGetString(mapped, "sourceVersion");
  602. if (!string.IsNullOrWhiteSpace(sourceVersion) && !long.TryParse(sourceVersion, out _))
  603. {
  604. rejected.Add(new MdpInboundRejectedRow { BizKey = previewKey, Reason = "sourceVersion must be long" });
  605. return null;
  606. }
  607. if (TryGetLong(mapped, "tenant_id", out var rowTenant) || TryGetLong(mapped, "TenantId", out rowTenant))
  608. {
  609. if (rowTenant != boundTenantId)
  610. {
  611. batchDeny = new MdpInboundBatchException(403, "tenant mismatch");
  612. return null;
  613. }
  614. }
  615. mapped["tenant_id"] = boundTenantId;
  616. if (grant.FactoryId is > 0)
  617. {
  618. if (!TryGetLong(mapped, "factory_id", out var rowFactory) && !TryGetLong(mapped, "FactoryId", out rowFactory))
  619. {
  620. rejected.Add(new MdpInboundRejectedRow { BizKey = previewKey, Reason = "factory mismatch" });
  621. return null;
  622. }
  623. if (rowFactory != grant.FactoryId.Value)
  624. {
  625. rejected.Add(new MdpInboundRejectedRow { BizKey = previewKey, Reason = "factory mismatch" });
  626. return null;
  627. }
  628. }
  629. var bizKey = MdpStagingWriter.BuildBizKey(entity.BizKeyExpr, mapped);
  630. if (string.IsNullOrWhiteSpace(bizKey))
  631. {
  632. rejected.Add(new MdpInboundRejectedRow { BizKey = previewKey, Reason = "missing biz key fields" });
  633. return null;
  634. }
  635. if (!string.IsNullOrWhiteSpace(previewKey) && !string.Equals(previewKey, bizKey, StringComparison.Ordinal))
  636. {
  637. rejected.Add(new MdpInboundRejectedRow { BizKey = previewKey, Reason = "bizKey mismatch" });
  638. return null;
  639. }
  640. var op = TryGetString(mapped, "op");
  641. if (string.IsNullOrWhiteSpace(op))
  642. op = "upsert";
  643. if (string.Equals(op, "delete", StringComparison.OrdinalIgnoreCase))
  644. {
  645. if (!DeleteSupported(entity))
  646. {
  647. rejected.Add(new MdpInboundRejectedRow { BizKey = bizKey, Reason = "op=delete not supported" });
  648. return null;
  649. }
  650. mapped["is_deleted"] = 1;
  651. }
  652. var rawJson = JsonSerializer.Serialize(mapped, MdpInboundJson.Options);
  653. return new MdpInboundPreparedRow
  654. {
  655. Dict = mapped,
  656. RawJson = rawJson,
  657. BizKey = bizKey,
  658. SourceUpdatedAt = sourceUpdatedAt,
  659. SourceVersion = sourceVersion
  660. };
  661. }
  662. private async Task<(bool RejectStale, bool EqualTsOverride, MdpInboundPreparedRow Row)> GuardStaleAsync(
  663. MdpEntity entity,
  664. long tenantId,
  665. string sourceTable,
  666. MdpInboundPreparedRow incoming,
  667. string inboundSource,
  668. CancellationToken ct)
  669. {
  670. var existing = await LoadGuardRowAsync(entity.TargetTableName, tenantId, sourceTable, incoming.BizKey, inboundSource, ct);
  671. if (existing == null)
  672. return (false, false, incoming);
  673. var oldHasVer = long.TryParse(existing.SourceVersion, out var oldVer);
  674. var newHasVer = long.TryParse(incoming.SourceVersion, out var newVer);
  675. if (oldHasVer && newHasVer)
  676. {
  677. if (newVer < oldVer)
  678. return (true, false, incoming);
  679. if (newVer > oldVer)
  680. return (false, false, incoming);
  681. }
  682. else
  683. {
  684. var oldHasTs = DateTimeOffset.TryParse(existing.SourceUpdatedAt, out var oldTs);
  685. if (!oldHasTs)
  686. return (false, false, incoming);
  687. if (!DateTimeOffset.TryParse(incoming.SourceUpdatedAt, out var newTs))
  688. return (true, false, incoming);
  689. if (newTs < oldTs)
  690. return (true, false, incoming);
  691. if (newTs > oldTs)
  692. return (false, false, incoming);
  693. }
  694. var same = string.Equals(incoming.RawJson, existing.RawData, StringComparison.Ordinal);
  695. return (false, !same, incoming);
  696. }
  697. private async Task<MdpInboundStgGuardRow> LoadGuardRowAsync(
  698. string table, long tenantId, string sourceTable, string bizKey, string inboundSource, CancellationToken ct)
  699. {
  700. if (string.IsNullOrWhiteSpace(table) || !TableNameRe.IsMatch(table))
  701. throw new InvalidOperationException("illegal target_table_name");
  702. var rows = await _db.Ado.SqlQueryAsync<MdpInboundStgGuardRow>(
  703. $"""
  704. SELECT id AS Id,
  705. source_system AS SourceSystem,
  706. source_row_id AS SourceRowId,
  707. JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.sourceUpdatedAt')) AS SourceUpdatedAt,
  708. JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.sourceVersion')) AS SourceVersion,
  709. raw_data AS RawData
  710. FROM `{table}`
  711. WHERE tenant_id=@t AND source_table=@st AND source_biz_key=@bk
  712. ORDER BY CASE WHEN source_system=@sys THEN 0 ELSE 1 END
  713. LIMIT 1
  714. """,
  715. new SugarParameter("@t", tenantId),
  716. new SugarParameter("@st", sourceTable),
  717. new SugarParameter("@bk", bizKey),
  718. new SugarParameter("@sys", inboundSource));
  719. return rows?.FirstOrDefault();
  720. }
  721. private async Task TransferIfNeededAsync(
  722. MdpEntity entity, long tenantId, string sourceTable, string bizKey, string inboundSource, CancellationToken ct)
  723. {
  724. var existing = await LoadGuardRowAsync(entity.TargetTableName, tenantId, sourceTable, bizKey, inboundSource, ct);
  725. if (existing == null)
  726. return;
  727. if (string.Equals(existing.SourceSystem, inboundSource, StringComparison.OrdinalIgnoreCase))
  728. return;
  729. if (string.IsNullOrWhiteSpace(entity.TargetTableName) || !TableNameRe.IsMatch(entity.TargetTableName))
  730. throw new InvalidOperationException("illegal target_table_name");
  731. await _db.Ado.ExecuteCommandAsync(
  732. $"""
  733. UPDATE `{entity.TargetTableName}`
  734. SET source_system=@sys, source_row_id=@bk
  735. WHERE id=@id
  736. """,
  737. new SugarParameter("@sys", inboundSource),
  738. new SugarParameter("@bk", bizKey),
  739. new SugarParameter("@id", existing.Id));
  740. }
  741. private async Task MarkFailedAsync(long requestId, string error, long elapsedMs, CancellationToken ct)
  742. {
  743. try
  744. {
  745. await _db.Updateable<MdpInboundRequest>()
  746. .SetColumns(r => new MdpInboundRequest
  747. {
  748. Status = "FAILED",
  749. HttpStatus = 500,
  750. ResultCode = 500,
  751. ErrorMsg = error,
  752. ElapsedMs = (int)elapsedMs,
  753. UpdateTime = DateTime.Now
  754. })
  755. .Where(r => r.Id == requestId)
  756. .ExecuteCommandAsync(ct);
  757. }
  758. catch (Exception ex)
  759. {
  760. _logger.LogError(ex, "inbound mark FAILED failed requestId={Id}", requestId);
  761. }
  762. }
  763. private static bool DeleteSupported(MdpEntity entity) =>
  764. (entity.Remark ?? string.Empty).Contains("inbound_delete=supported", StringComparison.OrdinalIgnoreCase);
  765. private static bool HasIsoOffset(string value) =>
  766. !string.IsNullOrWhiteSpace(value)
  767. && IsoOffsetRe.IsMatch(value)
  768. && DateTimeOffset.TryParse(value, out _);
  769. private static bool TryGetPropertyIgnoreCase(JsonElement obj, string name, out JsonElement value)
  770. {
  771. foreach (var prop in obj.EnumerateObject())
  772. {
  773. if (string.Equals(prop.Name, name, StringComparison.OrdinalIgnoreCase))
  774. {
  775. value = prop.Value;
  776. return true;
  777. }
  778. }
  779. value = default;
  780. return false;
  781. }
  782. private static string TryGetString(IDictionary<string, object?> row, string key)
  783. {
  784. foreach (var kv in row)
  785. {
  786. if (!string.Equals(kv.Key, key, StringComparison.OrdinalIgnoreCase))
  787. continue;
  788. return kv.Value?.ToString();
  789. }
  790. return null;
  791. }
  792. private static bool TryGetLong(IDictionary<string, object?> row, string key, out long value)
  793. {
  794. value = 0;
  795. var s = TryGetString(row, key);
  796. return !string.IsNullOrWhiteSpace(s) && long.TryParse(s, out value);
  797. }
  798. private static bool IsSensitiveHeader(string name) =>
  799. string.Equals(name, "X-Signature", StringComparison.OrdinalIgnoreCase)
  800. || string.Equals(name, "Authorization", StringComparison.OrdinalIgnoreCase)
  801. || string.Equals(name, "Cookie", StringComparison.OrdinalIgnoreCase);
  802. private static bool IsDuplicateKey(Exception ex)
  803. {
  804. for (var e = ex; e != null; e = e.InnerException)
  805. {
  806. var m = e.Message ?? string.Empty;
  807. if (m.Contains("Duplicate", StringComparison.OrdinalIgnoreCase)
  808. || m.Contains("1062")
  809. || m.Contains("uk_inbound_idem", StringComparison.OrdinalIgnoreCase))
  810. return true;
  811. }
  812. return false;
  813. }
  814. private static string ToSha256Hex(byte[] raw) =>
  815. Convert.ToHexString(SHA256.HashData(raw ?? [])).ToLowerInvariant();
  816. private static string Truncate(string s, int max) =>
  817. string.IsNullOrEmpty(s) ? s : (s.Length <= max ? s : s[..max]);
  818. private static object Envelope(int code, string message, object data) =>
  819. new { code, message, data };
  820. private static MdpInboundOutcome Fail(int status, string message) =>
  821. new() { HttpStatus = status, Body = Envelope(status, message, null) };
  822. private static MdpInboundOutcome Ok(object data) =>
  823. new() { HttpStatus = 200, Body = Envelope(0, "ok", data) };
  824. private async Task<string> PatchTransformEnqueuedAsync(long requestId, string receiptJson, CancellationToken ct)
  825. {
  826. if (string.IsNullOrWhiteSpace(receiptJson))
  827. return receiptJson;
  828. try
  829. {
  830. var node = JsonNode.Parse(receiptJson);
  831. if (node is JsonObject obj && obj["data"] is JsonObject data)
  832. data["transformEnqueued"] = true;
  833. var updated = node.ToJsonString(MdpInboundJson.Options);
  834. await _db.Updateable<MdpInboundRequest>()
  835. .SetColumns(r => new MdpInboundRequest
  836. {
  837. ReceiptJson = updated,
  838. UpdateTime = DateTime.Now
  839. })
  840. .Where(r => r.Id == requestId)
  841. .ExecuteCommandAsync(ct);
  842. return updated;
  843. }
  844. catch (Exception ex)
  845. {
  846. _logger.LogError(ex, "inbound patch transformEnqueued failed requestId={Id}", requestId);
  847. return receiptJson;
  848. }
  849. }
  850. private static MdpInboundOutcome Receipt(string receiptJson, bool replay)
  851. {
  852. if (string.IsNullOrWhiteSpace(receiptJson))
  853. return new MdpInboundOutcome { HttpStatus = 202, Body = Envelope(0, "ok", null) };
  854. try
  855. {
  856. var node = JsonNode.Parse(receiptJson);
  857. if (replay && node is JsonObject obj && obj["data"] is JsonObject data)
  858. data["idempotentReplay"] = true;
  859. return new MdpInboundOutcome { HttpStatus = 202, Body = node };
  860. }
  861. catch (JsonException)
  862. {
  863. return new MdpInboundOutcome { HttpStatus = 202, Body = Envelope(0, "ok", receiptJson) };
  864. }
  865. }
  866. }