MdpInboundReceiveService.cs 41 KB

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