MdpApiPushExecutor.cs 5.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130
  1. using System.Net.Http.Headers;
  2. using System.Text;
  3. using System.Text.Json;
  4. using Admin.NET.Plugin.AiDOP.Entity.DataPlatform;
  5. namespace Admin.NET.Plugin.AiDOP.DataPlatform.Executors;
  6. /// <summary>
  7. /// 出站回写:扫描 mdp_outbox 待推记录,按 target_source_code 调第三方写接口。
  8. /// </summary>
  9. public sealed class MdpApiPushExecutor : ITransient
  10. {
  11. private const int MaxRetry = 3;
  12. private readonly IHttpClientFactory _httpClientFactory;
  13. private readonly ISqlSugarClient _db;
  14. public MdpApiPushExecutor(IHttpClientFactory httpClientFactory, ISqlSugarClient db)
  15. {
  16. _httpClientFactory = httpClientFactory;
  17. _db = db;
  18. }
  19. /// <summary>处理一批待推记录,返回成功/失败计数。</summary>
  20. public async Task<(int success, int failed, int skipped)> PushPendingAsync(int take = 50, CancellationToken cancellationToken = default)
  21. {
  22. var pending = await _db.Queryable<MdpOutbox>()
  23. .Where(x => x.Status == 0 && x.RetryCount < MaxRetry)
  24. .OrderBy(x => x.Id)
  25. .Take(take)
  26. .ToListAsync(cancellationToken);
  27. var success = 0;
  28. var failed = 0;
  29. var skipped = 0;
  30. foreach (var item in pending)
  31. {
  32. cancellationToken.ThrowIfCancellationRequested();
  33. try
  34. {
  35. var source = await _db.Queryable<MdpSource>()
  36. .Where(x => x.SourceCode == item.TargetSourceCode && x.Status == 1)
  37. .FirstAsync(cancellationToken);
  38. if (source == null || string.IsNullOrWhiteSpace(source.ApiBaseUrl))
  39. {
  40. await MarkAsync(item, 2, null, $"目标源 {item.TargetSourceCode} 未配置 API", cancellationToken);
  41. failed++;
  42. continue;
  43. }
  44. // payload 约定:{ "path": "/api/xxx", "method": "POST", "body": {...} }
  45. using var payloadDoc = JsonDocument.Parse(string.IsNullOrWhiteSpace(item.PayloadJson) ? "{}" : item.PayloadJson!);
  46. var root = payloadDoc.RootElement;
  47. var path = root.TryGetProperty("path", out var p) ? p.GetString() : $"/outbox/{item.ActionCode}";
  48. var method = root.TryGetProperty("method", out var m) ? (m.GetString() ?? "POST") : "POST";
  49. var bodyEl = root.TryGetProperty("body", out var b) ? b : root;
  50. var client = _httpClientFactory.CreateClient("MdpApiPush");
  51. client.Timeout = TimeSpan.FromSeconds(60);
  52. var url = source.ApiBaseUrl!.TrimEnd('/') + (path!.StartsWith('/') ? path : "/" + path);
  53. using var request = new HttpRequestMessage(new HttpMethod(method), url);
  54. ApplyAuth(request, source);
  55. request.Content = new StringContent(bodyEl.GetRawText(), Encoding.UTF8, "application/json");
  56. using var response = await client.SendAsync(request, cancellationToken);
  57. var respBody = await response.Content.ReadAsStringAsync(cancellationToken);
  58. if (response.IsSuccessStatusCode)
  59. {
  60. await MarkAsync(item, 1, respBody, null, cancellationToken);
  61. success++;
  62. }
  63. else
  64. {
  65. item.RetryCount++;
  66. var status = item.RetryCount >= MaxRetry ? 2 : 0;
  67. await MarkAsync(item, status, respBody, $"HTTP {(int)response.StatusCode}", cancellationToken);
  68. if (status == 2) failed++; else skipped++;
  69. }
  70. }
  71. catch (Exception ex)
  72. {
  73. item.RetryCount++;
  74. var status = item.RetryCount >= MaxRetry ? 2 : 0;
  75. await MarkAsync(item, status, null, Truncate(ex.Message, 900), cancellationToken);
  76. if (status == 2) failed++; else skipped++;
  77. }
  78. }
  79. return (success, failed, skipped);
  80. }
  81. private async Task MarkAsync(MdpOutbox item, int status, string? responseJson, string? error, CancellationToken ct)
  82. {
  83. item.Status = status;
  84. item.ResponseJson = Truncate(responseJson, 4000);
  85. item.ErrorMsg = error;
  86. item.UpdateTime = DateTime.Now;
  87. await _db.Updateable(item)
  88. .UpdateColumns(x => new { x.Status, x.RetryCount, x.ResponseJson, x.ErrorMsg, x.UpdateTime })
  89. .ExecuteCommandAsync(ct);
  90. }
  91. private static void ApplyAuth(HttpRequestMessage request, MdpSource source)
  92. {
  93. var authType = (source.ApiAuthType ?? "NONE").Trim().ToUpperInvariant();
  94. if (authType is "NONE" or "") return;
  95. Dictionary<string, string>? cfg = null;
  96. if (!string.IsNullOrWhiteSpace(source.ApiAuthConfig))
  97. {
  98. try { cfg = JsonSerializer.Deserialize<Dictionary<string, string>>(source.ApiAuthConfig!); }
  99. catch { /* ignore */ }
  100. }
  101. cfg ??= new Dictionary<string, string>();
  102. if ((authType is "BEARER" or "TOKEN") && (cfg.TryGetValue("token", out var token) || cfg.TryGetValue("access_token", out token)))
  103. request.Headers.Authorization = new AuthenticationHeaderValue("Bearer", token);
  104. else if (authType == "BASIC" && cfg.TryGetValue("username", out var user) && cfg.TryGetValue("password", out var pwd))
  105. request.Headers.Authorization = new AuthenticationHeaderValue("Basic", Convert.ToBase64String(Encoding.UTF8.GetBytes($"{user}:{pwd}")));
  106. else if (authType == "APIKEY")
  107. {
  108. var header = cfg.GetValueOrDefault("header") ?? "X-API-Key";
  109. if (cfg.TryGetValue("apiKey", out var key) || cfg.TryGetValue("key", out key))
  110. request.Headers.TryAddWithoutValidation(header, key);
  111. }
  112. }
  113. private static string? Truncate(string? s, int max) =>
  114. string.IsNullOrEmpty(s) ? s : (s.Length <= max ? s : s[..max]);
  115. }