MdpInboundReceiveService.cs 40 KB

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