MdpApiPullExecutor.cs 8.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211
  1. using System.Net.Http.Headers;
  2. using System.Text;
  3. using System.Text.Json;
  4. using Admin.NET.Plugin.AiDOP.Entity.DataPlatform;
  5. using SqlSugar;
  6. namespace Admin.NET.Plugin.AiDOP.DataPlatform.Executors;
  7. /// <summary>
  8. /// 方式乙:按 mdp_source API 段 + mdp_entity.source_api_path 拉取 JSON → target_table_name(贴源)。
  9. /// </summary>
  10. public sealed class MdpApiPullExecutor : IMdpSourcePullExecutor, ITransient
  11. {
  12. public string SupportedType => "API_PULL";
  13. private readonly IHttpClientFactory _httpClientFactory;
  14. private readonly ISqlSugarClient _db;
  15. private readonly MdpStagingWriter _writer;
  16. public MdpApiPullExecutor(IHttpClientFactory httpClientFactory, ISqlSugarClient db, MdpStagingWriter writer)
  17. {
  18. _httpClientFactory = httpClientFactory;
  19. _db = db;
  20. _writer = writer;
  21. }
  22. public async Task<MdpPullResult> PullAsync(MdpSource source, MdpEntity entity, MdpPullContext ctx, CancellationToken cancellationToken = default)
  23. {
  24. if (string.IsNullOrWhiteSpace(source.ApiBaseUrl))
  25. throw new InvalidOperationException($"源 {source.SourceCode} 未配置 api_base_url");
  26. if (string.IsNullOrWhiteSpace(entity.SourceApiPath))
  27. throw new InvalidOperationException($"实体 {entity.EntityCode} 未配置 source_api_path");
  28. if (string.IsNullOrWhiteSpace(entity.TargetTableName))
  29. throw new InvalidOperationException($"实体 {entity.EntityCode} 未配置 target_table_name");
  30. var client = _httpClientFactory.CreateClient("MdpApiPull");
  31. client.Timeout = TimeSpan.FromSeconds(120);
  32. var url = CombineUrl(source.ApiBaseUrl!, entity.SourceApiPath!);
  33. if (!string.IsNullOrWhiteSpace(entity.LastCursor) && !ctx.FullRefresh)
  34. url += (url.Contains('?') ? "&" : "?") + "cursor=" + Uri.EscapeDataString(entity.LastCursor);
  35. using var request = new HttpRequestMessage(HttpMethod.Get, url);
  36. ApplyAuth(request, source);
  37. using var response = await client.SendAsync(request, cancellationToken);
  38. var body = await response.Content.ReadAsStringAsync(cancellationToken);
  39. if (!response.IsSuccessStatusCode)
  40. throw new InvalidOperationException($"API_PULL 失败 HTTP {(int)response.StatusCode}: {Truncate(body, 500)}");
  41. using var doc = JsonDocument.Parse(string.IsNullOrWhiteSpace(body) ? "[]" : body);
  42. var items = ResolveArray(doc.RootElement, entity.ResponseDataPath);
  43. var written = 0;
  44. var now = DateTime.Now;
  45. string? newCursor = entity.LastCursor;
  46. var seen = new HashSet<string>(StringComparer.Ordinal);
  47. var pending = new List<(IDictionary<string, object?> Row, string RawJson, string SourceRowId)>();
  48. foreach (var item in items)
  49. {
  50. cancellationToken.ThrowIfCancellationRequested();
  51. var dedup = ResolvePath(item, entity.DedupKeyPath) ?? Guid.NewGuid().ToString("N");
  52. if (!seen.Add(dedup)) continue;
  53. var rawJson = item.GetRawText();
  54. var dict = JsonElementToDict(item);
  55. pending.Add((dict, rawJson, dedup));
  56. newCursor = dedup;
  57. }
  58. if (pending.Count > 0)
  59. {
  60. // 贴源 source_table 优先用逻辑源表名(与 DB 路线/std 过滤一致);未配置时回落 API path
  61. var stagingTable = !string.IsNullOrWhiteSpace(entity.SourceTableName)
  62. ? entity.SourceTableName!
  63. : entity.SourceApiPath!;
  64. var results = await _writer.UpsertBatchAsync(source, entity, stagingTable, pending, ctx);
  65. written += results.Sum();
  66. }
  67. if (!string.IsNullOrEmpty(newCursor) && newCursor != entity.LastCursor)
  68. {
  69. await _db.Updateable<MdpEntity>()
  70. .SetColumns(x => new MdpEntity
  71. {
  72. LastCursor = newCursor,
  73. LastSyncTo = now,
  74. UpdateTime = now
  75. })
  76. .Where(x => x.Id == entity.Id)
  77. .ExecuteCommandAsync(cancellationToken);
  78. }
  79. return new MdpPullResult
  80. {
  81. RowsPulled = items.Count,
  82. RowsWritten = written,
  83. NewCursor = newCursor,
  84. Message = "OK"
  85. };
  86. }
  87. private static void ApplyAuth(HttpRequestMessage request, MdpSource source)
  88. {
  89. var authType = (source.ApiAuthType ?? "NONE").Trim().ToUpperInvariant();
  90. if (authType is "NONE" or "") return;
  91. Dictionary<string, string>? cfg = null;
  92. if (!string.IsNullOrWhiteSpace(source.ApiAuthConfig))
  93. {
  94. try
  95. {
  96. cfg = JsonSerializer.Deserialize<Dictionary<string, string>>(source.ApiAuthConfig!);
  97. }
  98. catch { /* ignore */ }
  99. }
  100. cfg ??= new Dictionary<string, string>();
  101. switch (authType)
  102. {
  103. case "BEARER":
  104. case "TOKEN":
  105. if (cfg.TryGetValue("token", out var token) || cfg.TryGetValue("access_token", out token))
  106. request.Headers.Authorization = new AuthenticationHeaderValue("Bearer", token);
  107. break;
  108. case "BASIC":
  109. if (cfg.TryGetValue("username", out var user) && cfg.TryGetValue("password", out var pwd))
  110. {
  111. var bytes = Encoding.UTF8.GetBytes($"{user}:{pwd}");
  112. request.Headers.Authorization = new AuthenticationHeaderValue("Basic", Convert.ToBase64String(bytes));
  113. }
  114. break;
  115. case "APIKEY":
  116. var header = cfg.GetValueOrDefault("header") ?? "X-API-Key";
  117. if (cfg.TryGetValue("apiKey", out var key) || cfg.TryGetValue("key", out key))
  118. request.Headers.TryAddWithoutValidation(header, key);
  119. break;
  120. case "OAUTH2":
  121. if (cfg.TryGetValue("access_token", out var oauth))
  122. request.Headers.Authorization = new AuthenticationHeaderValue("Bearer", oauth);
  123. break;
  124. }
  125. }
  126. private static List<JsonElement> ResolveArray(JsonElement root, string? path)
  127. {
  128. var el = string.IsNullOrWhiteSpace(path) ? root : ResolveElement(root, path!) ?? root;
  129. if (el.ValueKind == JsonValueKind.Array)
  130. return el.EnumerateArray().Select(x => x.Clone()).ToList();
  131. if (el.ValueKind == JsonValueKind.Object)
  132. return new List<JsonElement> { el.Clone() };
  133. return new List<JsonElement>();
  134. }
  135. private static JsonElement? ResolveElement(JsonElement root, string path)
  136. {
  137. var cur = root;
  138. foreach (var part in path.Split('.', StringSplitOptions.RemoveEmptyEntries | StringSplitOptions.TrimEntries))
  139. {
  140. if (cur.ValueKind != JsonValueKind.Object || !cur.TryGetProperty(part, out cur))
  141. return null;
  142. }
  143. return cur;
  144. }
  145. private static string? ResolvePath(JsonElement el, string? path)
  146. {
  147. if (string.IsNullOrWhiteSpace(path)) return null;
  148. var found = ResolveElement(el, path);
  149. if (found == null) return null;
  150. return found.Value.ValueKind switch
  151. {
  152. JsonValueKind.String => found.Value.GetString(),
  153. JsonValueKind.Number => found.Value.ToString(),
  154. JsonValueKind.True => "true",
  155. JsonValueKind.False => "false",
  156. _ => found.Value.ToString()
  157. };
  158. }
  159. private static string CombineUrl(string baseUrl, string path)
  160. {
  161. baseUrl = baseUrl.TrimEnd('/');
  162. path = path.StartsWith('/') ? path : "/" + path;
  163. return baseUrl + path;
  164. }
  165. private static string Truncate(string s, int max) =>
  166. string.IsNullOrEmpty(s) ? "" : (s.Length <= max ? s : s[..max]);
  167. private static Dictionary<string, object?> JsonElementToDict(JsonElement el)
  168. {
  169. var dict = new Dictionary<string, object?>(StringComparer.OrdinalIgnoreCase);
  170. if (el.ValueKind != JsonValueKind.Object) return dict;
  171. foreach (var prop in el.EnumerateObject())
  172. {
  173. dict[prop.Name] = prop.Value.ValueKind switch
  174. {
  175. JsonValueKind.Null => null,
  176. JsonValueKind.String => prop.Value.GetString(),
  177. JsonValueKind.Number => prop.Value.TryGetInt64(out var l) ? l
  178. : prop.Value.TryGetDecimal(out var d) ? d
  179. : prop.Value.GetDouble(),
  180. JsonValueKind.True => true,
  181. JsonValueKind.False => false,
  182. _ => prop.Value.GetRawText()
  183. };
  184. }
  185. return dict;
  186. }
  187. }