MdpApiPullExecutor.cs 8.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206
  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. foreach (var item in items)
  48. {
  49. cancellationToken.ThrowIfCancellationRequested();
  50. var dedup = ResolvePath(item, entity.DedupKeyPath) ?? Guid.NewGuid().ToString("N");
  51. if (!seen.Add(dedup)) continue;
  52. var rawJson = item.GetRawText();
  53. var dict = JsonElementToDict(item);
  54. // 贴源 source_table 优先用逻辑源表名(与 DB 路线/std 过滤一致);未配置时回落 API path
  55. var stagingTable = !string.IsNullOrWhiteSpace(entity.SourceTableName)
  56. ? entity.SourceTableName!
  57. : entity.SourceApiPath!;
  58. written += await _writer.UpsertAsync(
  59. source, entity, stagingTable, dict, rawJson, dedup, ctx);
  60. newCursor = dedup;
  61. }
  62. if (!string.IsNullOrEmpty(newCursor) && newCursor != entity.LastCursor)
  63. {
  64. await _db.Updateable<MdpEntity>()
  65. .SetColumns(x => new MdpEntity
  66. {
  67. LastCursor = newCursor,
  68. LastSyncTo = now,
  69. UpdateTime = now
  70. })
  71. .Where(x => x.Id == entity.Id)
  72. .ExecuteCommandAsync(cancellationToken);
  73. }
  74. return new MdpPullResult
  75. {
  76. RowsPulled = items.Count,
  77. RowsWritten = written,
  78. NewCursor = newCursor,
  79. Message = "OK"
  80. };
  81. }
  82. private static void ApplyAuth(HttpRequestMessage request, MdpSource source)
  83. {
  84. var authType = (source.ApiAuthType ?? "NONE").Trim().ToUpperInvariant();
  85. if (authType is "NONE" or "") return;
  86. Dictionary<string, string>? cfg = null;
  87. if (!string.IsNullOrWhiteSpace(source.ApiAuthConfig))
  88. {
  89. try
  90. {
  91. cfg = JsonSerializer.Deserialize<Dictionary<string, string>>(source.ApiAuthConfig!);
  92. }
  93. catch { /* ignore */ }
  94. }
  95. cfg ??= new Dictionary<string, string>();
  96. switch (authType)
  97. {
  98. case "BEARER":
  99. case "TOKEN":
  100. if (cfg.TryGetValue("token", out var token) || cfg.TryGetValue("access_token", out token))
  101. request.Headers.Authorization = new AuthenticationHeaderValue("Bearer", token);
  102. break;
  103. case "BASIC":
  104. if (cfg.TryGetValue("username", out var user) && cfg.TryGetValue("password", out var pwd))
  105. {
  106. var bytes = Encoding.UTF8.GetBytes($"{user}:{pwd}");
  107. request.Headers.Authorization = new AuthenticationHeaderValue("Basic", Convert.ToBase64String(bytes));
  108. }
  109. break;
  110. case "APIKEY":
  111. var header = cfg.GetValueOrDefault("header") ?? "X-API-Key";
  112. if (cfg.TryGetValue("apiKey", out var key) || cfg.TryGetValue("key", out key))
  113. request.Headers.TryAddWithoutValidation(header, key);
  114. break;
  115. case "OAUTH2":
  116. if (cfg.TryGetValue("access_token", out var oauth))
  117. request.Headers.Authorization = new AuthenticationHeaderValue("Bearer", oauth);
  118. break;
  119. }
  120. }
  121. private static List<JsonElement> ResolveArray(JsonElement root, string? path)
  122. {
  123. var el = string.IsNullOrWhiteSpace(path) ? root : ResolveElement(root, path!) ?? root;
  124. if (el.ValueKind == JsonValueKind.Array)
  125. return el.EnumerateArray().Select(x => x.Clone()).ToList();
  126. if (el.ValueKind == JsonValueKind.Object)
  127. return new List<JsonElement> { el.Clone() };
  128. return new List<JsonElement>();
  129. }
  130. private static JsonElement? ResolveElement(JsonElement root, string path)
  131. {
  132. var cur = root;
  133. foreach (var part in path.Split('.', StringSplitOptions.RemoveEmptyEntries | StringSplitOptions.TrimEntries))
  134. {
  135. if (cur.ValueKind != JsonValueKind.Object || !cur.TryGetProperty(part, out cur))
  136. return null;
  137. }
  138. return cur;
  139. }
  140. private static string? ResolvePath(JsonElement el, string? path)
  141. {
  142. if (string.IsNullOrWhiteSpace(path)) return null;
  143. var found = ResolveElement(el, path);
  144. if (found == null) return null;
  145. return found.Value.ValueKind switch
  146. {
  147. JsonValueKind.String => found.Value.GetString(),
  148. JsonValueKind.Number => found.Value.ToString(),
  149. JsonValueKind.True => "true",
  150. JsonValueKind.False => "false",
  151. _ => found.Value.ToString()
  152. };
  153. }
  154. private static string CombineUrl(string baseUrl, string path)
  155. {
  156. baseUrl = baseUrl.TrimEnd('/');
  157. path = path.StartsWith('/') ? path : "/" + path;
  158. return baseUrl + path;
  159. }
  160. private static string Truncate(string s, int max) =>
  161. string.IsNullOrEmpty(s) ? "" : (s.Length <= max ? s : s[..max]);
  162. private static Dictionary<string, object?> JsonElementToDict(JsonElement el)
  163. {
  164. var dict = new Dictionary<string, object?>(StringComparer.OrdinalIgnoreCase);
  165. if (el.ValueKind != JsonValueKind.Object) return dict;
  166. foreach (var prop in el.EnumerateObject())
  167. {
  168. dict[prop.Name] = prop.Value.ValueKind switch
  169. {
  170. JsonValueKind.Null => null,
  171. JsonValueKind.String => prop.Value.GetString(),
  172. JsonValueKind.Number => prop.Value.TryGetInt64(out var l) ? l
  173. : prop.Value.TryGetDecimal(out var d) ? d
  174. : prop.Value.GetDouble(),
  175. JsonValueKind.True => true,
  176. JsonValueKind.False => false,
  177. _ => prop.Value.GetRawText()
  178. };
  179. }
  180. return dict;
  181. }
  182. }