ExternalApiController.java 26 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576
  1. package com.agent.management.external;
  2. import com.agent.management.config.ExternalApiAuthFilter;
  3. import com.agent.management.config.ExternalApiProperties;
  4. import com.agent.management.engine.AgentIgnoreFilter;
  5. import com.agent.management.engine.WorkflowEngine;
  6. import com.agent.management.engine.WorkflowRunDirManager;
  7. import com.agent.management.engine.WorkflowRunEvent;
  8. import com.agent.management.external.dto.ApiError;
  9. import com.agent.management.external.dto.ExternalDtos.CreateRunRequest;
  10. import com.agent.management.external.dto.ExternalDtos.CreateRunResponse;
  11. import com.agent.management.external.dto.ExternalDtos.FieldDef;
  12. import com.agent.management.external.dto.ExternalDtos.Links;
  13. import com.agent.management.external.dto.ExternalDtos.LogEntry;
  14. import com.agent.management.external.dto.ExternalDtos.NodeLogsResponse;
  15. import com.agent.management.external.dto.ExternalDtos.NodeResult;
  16. import com.agent.management.external.dto.ExternalDtos.RunResultResponse;
  17. import com.agent.management.external.dto.ExternalDtos.StartRunResponse;
  18. import com.agent.management.external.dto.ExternalDtos.WorkflowListResponse;
  19. import com.agent.management.external.dto.ExternalDtos.WorkflowSummary;
  20. import com.agent.management.external.dto.ExternalDtos.WorkspaceInfo;
  21. import com.agent.management.model.entity.Workflow;
  22. import com.agent.management.model.entity.WorkflowRun;
  23. import com.agent.management.model.entity.WorkflowRunNode;
  24. import com.agent.management.repository.WorkflowRepository;
  25. import com.agent.management.repository.WorkflowRunNodeRepository;
  26. import com.agent.management.repository.WorkflowRunRepository;
  27. import com.fasterxml.jackson.core.type.TypeReference;
  28. import com.fasterxml.jackson.databind.JsonNode;
  29. import com.fasterxml.jackson.databind.ObjectMapper;
  30. import com.fasterxml.jackson.datatype.jsr310.JavaTimeModule;
  31. import lombok.RequiredArgsConstructor;
  32. import lombok.extern.slf4j.Slf4j;
  33. import org.springframework.core.io.Resource;
  34. import org.springframework.core.io.UrlResource;
  35. import org.springframework.http.HttpHeaders;
  36. import org.springframework.http.MediaType;
  37. import org.springframework.http.ResponseEntity;
  38. import org.springframework.web.bind.annotation.*;
  39. import org.springframework.web.multipart.MultipartFile;
  40. import org.springframework.web.servlet.mvc.method.annotation.SseEmitter;
  41. import java.io.IOException;
  42. import java.net.URLEncoder;
  43. import java.nio.charset.StandardCharsets;
  44. import java.nio.file.Files;
  45. import java.nio.file.Path;
  46. import java.time.Instant;
  47. import java.time.LocalDateTime;
  48. import java.time.ZoneOffset;
  49. import java.util.*;
  50. import java.util.stream.Collectors;
  51. import java.util.zip.ZipEntry;
  52. import java.util.zip.ZipOutputStream;
  53. /**
  54. * 外部 API(/api/v1/**)入口:第三方系统通过此 Controller 调用平台工作流。
  55. *
  56. * <p>URL 路径使用 kebab-case 唯一标识 {@code workflowName}(与 Skill 一致)。
  57. * 旧调用方使用数字 id 的 URL 仍可通过 by-id 兜底逻辑访问(用户手动改名为非数字前都自动兼容)。
  58. *
  59. * 流程:
  60. * - POST /workflows/{workflowName}/runs 创建运行
  61. * - POST /workflows/{workflowName}/runs/{runId}/start 启动执行
  62. * - GET /workflows/{workflowName}/runs/{runId}/stream 订阅 SSE 流
  63. * - GET /workflows/{workflowName}/runs/{runId} 查询运行结果
  64. * - GET /workflows/{workflowName}/runs/{runId}/workspace 下载工作空间 zip
  65. * - GET /workflows/{workflowName}/runs/{runId}/nodes/{nodeId}/logs 查询节点日志
  66. * - GET /workflows 列出当前 API Key 可调用的工作流
  67. */
  68. @Slf4j
  69. @RestController
  70. @RequestMapping("/api/v1")
  71. @RequiredArgsConstructor
  72. public class ExternalApiController {
  73. private final ExternalApiProperties properties;
  74. private final ExternalRunRegistry runRegistry;
  75. private final SseEventBus eventBus;
  76. private final WorkflowEngine workflowEngine;
  77. private final WorkflowRepository workflowRepository;
  78. private final WorkflowRunRepository workflowRunRepository;
  79. private final WorkflowRunNodeRepository workflowRunNodeRepository;
  80. private final WorkflowRunDirManager runDirManager;
  81. private final AgentIgnoreFilter agentIgnoreFilter;
  82. private static final ObjectMapper MAPPER = new ObjectMapper().registerModule(new JavaTimeModule());
  83. private static final int DEFAULT_TTL_HOURS = 720;
  84. private static final long SSE_TIMEOUT_MS = 1_800_000L;
  85. // ============================================================
  86. // POST /workflows/{workflowName}/runs — 创建运行
  87. // ============================================================
  88. @PostMapping(value = "/workflows/{workflowName}/runs", consumes = MediaType.APPLICATION_JSON_VALUE)
  89. public ResponseEntity<?> createRunJson(@PathVariable String workflowName,
  90. @RequestAttribute(ExternalApiAuthFilter.ATTR_KEY_ENTRY) ExternalApiProperties.KeyEntry keyEntry,
  91. @RequestBody(required = false) CreateRunRequest req) {
  92. return doCreateRun(workflowName, req, null, null);
  93. }
  94. @PostMapping(value = "/workflows/{workflowName}/runs", consumes = MediaType.MULTIPART_FORM_DATA_VALUE)
  95. public ResponseEntity<?> createRunMultipart(@PathVariable String workflowName,
  96. @RequestAttribute(ExternalApiAuthFilter.ATTR_KEY_ENTRY) ExternalApiProperties.KeyEntry keyEntry,
  97. @RequestParam(value = "payload", required = false) String payloadJson,
  98. @RequestParam(value = "files", required = false) MultipartFile[] files,
  99. @RequestHeader(value = "X-Relative-Path", required = false) List<String> relativePaths) {
  100. CreateRunRequest req = parsePayload(payloadJson);
  101. return doCreateRun(workflowName, req, files, relativePaths);
  102. }
  103. private ResponseEntity<?> doCreateRun(String workflowName, CreateRunRequest req,
  104. MultipartFile[] files, List<String> relativePaths) {
  105. // 1. 工作流解析
  106. Workflow wf = resolveWorkflow(workflowName);
  107. if (wf == null) {
  108. return error(404, "WORKFLOW_NOT_FOUND", "工作流不存在: " + workflowName);
  109. }
  110. if (wf.getGraphData() == null || wf.getGraphData().isBlank()
  111. || wf.getGraphData().equals("{\"nodes\":[],\"edges\":[]}")) {
  112. return error(422, "GRAPH_INVALID", "工作流图为空,无法执行");
  113. }
  114. // 2. 准备 inputs
  115. Map<String, Object> inputs = new HashMap<>();
  116. if (req != null && req.getVariables() != null) {
  117. inputs.putAll(req.getVariables());
  118. }
  119. // 3. 生成 runId + 创建工作目录 + 写入上传文件
  120. String runId;
  121. try {
  122. runId = generateRunId();
  123. Path runDir = runDirManager.createRunDir(wf.getId(), wf.getName(), runId);
  124. List<String> uploadedFiles = new ArrayList<>();
  125. if (files != null) {
  126. for (int i = 0; i < files.length; i++) {
  127. MultipartFile file = files[i];
  128. if (file == null || file.isEmpty()) continue;
  129. String rel = (relativePaths != null && i < relativePaths.size() && !relativePaths.get(i).isBlank())
  130. ? relativePaths.get(i) : file.getOriginalFilename();
  131. if (rel == null || rel.isBlank()) continue;
  132. Path dest = runDir.resolve(rel).normalize();
  133. if (!dest.startsWith(runDir)) {
  134. return error(400, "VALIDATION_FAILED", "非法上传路径: " + rel);
  135. }
  136. Files.createDirectories(dest.getParent());
  137. file.transferTo(dest.toFile());
  138. uploadedFiles.add(runDir.relativize(dest).toString().replace('\\', '/'));
  139. }
  140. }
  141. if (!uploadedFiles.isEmpty()) {
  142. inputs.put("uploadedFiles", uploadedFiles);
  143. }
  144. } catch (IllegalArgumentException e) {
  145. return error(400, "VALIDATION_FAILED", e.getMessage());
  146. } catch (IOException e) {
  147. log.error("[ExternalApi] 创建运行失败 workflowName={}", workflowName, e);
  148. return error(500, "INTERNAL_ERROR", "工作目录创建失败");
  149. }
  150. // 4. 注册 CREATED 状态
  151. int ttlHours = (req != null && req.getTtlHours() != null && req.getTtlHours() > 0)
  152. ? req.getTtlHours() : DEFAULT_TTL_HOURS;
  153. runRegistry.register(runId, wf.getId(), wf.getName(), inputs, ttlHours);
  154. eventBus.register(runId);
  155. // 5. async=false:立即启动并直接进入 SSE 流
  156. boolean async = req == null || req.getAsync() == null || req.getAsync();
  157. if (!async) {
  158. return startRun(wf.getName(), wf.getId(), runId, inputs);
  159. }
  160. // 6. async=true:返回 runId + links(URL 使用 workflowName)
  161. CreateRunResponse body = new CreateRunResponse();
  162. body.setRunId(runId);
  163. body.setWorkflowId(wf.getId());
  164. body.setWorkflowName(wf.getName());
  165. body.setDisplayName(wf.getDisplayName());
  166. body.setStatus("CREATED");
  167. body.setCreatedAt(Instant.now());
  168. Links links = new Links();
  169. String base = "/api/v1/workflows/" + wf.getName() + "/runs/" + runId;
  170. links.setStart(base + "/start");
  171. links.setStream(base + "/stream");
  172. links.setResult(base);
  173. links.setWorkspace(base + "/workspace");
  174. body.setLinks(links);
  175. return ResponseEntity.ok(ApiError.ok(body));
  176. }
  177. // ============================================================
  178. // POST /workflows/{workflowName}/runs/{runId}/start — 启动运行
  179. // ============================================================
  180. @PostMapping("/workflows/{workflowName}/runs/{runId}/start")
  181. public ResponseEntity<?> startRun(@PathVariable String workflowName,
  182. @PathVariable String runId) {
  183. ExternalRunRegistry.CreatedRun created = runRegistry.consume(runId);
  184. if (created == null) {
  185. return error(404, "RUN_NOT_FOUND", "运行不存在或已启动: " + runId);
  186. }
  187. // 校验 pathVariable 与 created 一致(兼容数字字符串)
  188. if (!matchesWorkflow(created, workflowName)) {
  189. return error(400, "VALIDATION_FAILED", "runId 与 workflowName 不匹配");
  190. }
  191. return startRun(created.workflowName(), created.workflowId(), runId, created.inputs());
  192. }
  193. private ResponseEntity<?> startRun(String workflowName, Long workflowId, String runId, Map<String, Object> inputs) {
  194. workflowEngine.executeForExternal(workflowId, runId, inputs);
  195. StartRunResponse body = new StartRunResponse();
  196. body.setRunId(runId);
  197. body.setStatus("RUNNING");
  198. body.setStreamUrl("/api/v1/workflows/" + workflowName + "/runs/" + runId + "/stream");
  199. return ResponseEntity.ok(ApiError.ok(body));
  200. }
  201. // ============================================================
  202. // GET /workflows/{workflowName}/runs/{runId}/stream — SSE 流
  203. // ============================================================
  204. @GetMapping(value = "/workflows/{workflowName}/runs/{runId}/stream", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
  205. public ResponseEntity<SseEmitter> streamRun(@PathVariable String workflowName,
  206. @PathVariable String runId,
  207. @RequestHeader(value = "Last-Event-ID", required = false) Long lastEventId) {
  208. if (!eventBus.exists(runId)) {
  209. WorkflowRun record = findRunRecord(runId);
  210. if (record == null) {
  211. return ResponseEntity.notFound().build();
  212. }
  213. eventBus.register(runId);
  214. eventBus.publish(runId,
  215. record.getStatus().equals("FAILED") ? "workflow_error" : "workflow_complete",
  216. WorkflowRunEvent.workflowComplete(runId, parseOutputs(record)));
  217. }
  218. SseEmitter emitter = eventBus.subscribe(runId, lastEventId);
  219. if (emitter == null) {
  220. return ResponseEntity.notFound().build();
  221. }
  222. return ResponseEntity.ok()
  223. .header("X-Accel-Buffering", "no")
  224. .header("Cache-Control", "no-cache")
  225. .body(emitter);
  226. }
  227. // ============================================================
  228. // GET /workflows/{workflowName}/runs/{runId} — 查询运行结果
  229. // ============================================================
  230. @GetMapping("/workflows/{workflowName}/runs/{runId}")
  231. public ResponseEntity<?> getRunResult(@PathVariable String workflowName,
  232. @PathVariable String runId) {
  233. Workflow wf = resolveWorkflow(workflowName);
  234. if (wf == null) {
  235. return error(404, "WORKFLOW_NOT_FOUND", "工作流不存在: " + workflowName);
  236. }
  237. RunResultResponse body = new RunResultResponse();
  238. body.setRunId(runId);
  239. body.setWorkflowId(wf.getId());
  240. body.setWorkflowName(wf.getName());
  241. body.setDisplayName(wf.getDisplayName());
  242. SseEventBus.RunStatus live = eventBus.getStatus(runId);
  243. if (live != null) {
  244. body.setStatus(live.status());
  245. }
  246. WorkflowRun record = findRunRecord(runId);
  247. if (record == null) {
  248. ExternalRunRegistry.CreatedRun created = runRegistry.peek(runId);
  249. if (created != null) {
  250. body.setStatus("CREATED");
  251. return ResponseEntity.ok(ApiError.ok(body));
  252. }
  253. return error(404, "RUN_NOT_FOUND", "运行不存在: " + runId);
  254. }
  255. body.setStatus(record.getStatus());
  256. if (record.getStartedAt() != null) {
  257. body.setStartedAt(record.getStartedAt().toInstant(ZoneOffset.UTC));
  258. }
  259. if (record.getCompletedAt() != null) {
  260. body.setCompletedAt(record.getCompletedAt().toInstant(ZoneOffset.UTC));
  261. }
  262. body.setError(record.getError());
  263. body.setOutputs(parseOutputs(record));
  264. List<WorkflowRunNode> nodes = workflowRunNodeRepository.findByRunIdOrderBySortOrderAsc(record.getId());
  265. List<NodeResult> nodeDtos = nodes.stream().map(n -> {
  266. NodeResult dto = new NodeResult();
  267. dto.setNodeId(n.getNodeId());
  268. dto.setNodeType(n.getNodeType());
  269. dto.setLabel(n.getLabel());
  270. dto.setStatus(n.getStatus());
  271. dto.setOutput(parseJsonToMap(n.getOutput()));
  272. dto.setError(n.getError());
  273. dto.setLogsCount(countLogs(n.getLogs()));
  274. dto.setLogsUrl("/api/v1/workflows/" + wf.getName() + "/runs/" + runId + "/nodes/" + n.getNodeId() + "/logs");
  275. return dto;
  276. }).collect(Collectors.toList());
  277. body.setNodes(nodeDtos);
  278. Path runDir = runDirManager.getRunDir(wf.getId(), wf.getName(), runId);
  279. if (runDir != null) {
  280. try {
  281. final int[] fileCount = {0};
  282. final long[] totalBytes = {0};
  283. Files.walk(runDir)
  284. .filter(p -> !Files.isDirectory(p) && shouldIncludeInZip(runDir, p))
  285. .forEach(p -> {
  286. fileCount[0]++;
  287. try { totalBytes[0] += Files.size(p); } catch (IOException ignore) {}
  288. });
  289. WorkspaceInfo ws = new WorkspaceInfo();
  290. ws.setFileCount(fileCount[0]);
  291. ws.setTotalBytes(totalBytes[0]);
  292. ws.setDownloadUrl("/api/v1/workflows/" + wf.getName() + "/runs/" + runId + "/workspace");
  293. body.setWorkspace(ws);
  294. } catch (IOException e) {
  295. log.warn("[ExternalApi] 统计工作目录失败 runId={}: {}", runId, e.getMessage());
  296. }
  297. }
  298. return ResponseEntity.ok(ApiError.ok(body));
  299. }
  300. // ============================================================
  301. // GET /workflows/{workflowName}/runs/{runId}/workspace
  302. // ============================================================
  303. @GetMapping("/workflows/{workflowName}/runs/{runId}/workspace")
  304. public ResponseEntity<Resource> downloadWorkspace(@PathVariable String workflowName,
  305. @PathVariable String runId) throws IOException {
  306. Workflow wf = resolveWorkflow(workflowName);
  307. if (wf == null) {
  308. return ResponseEntity.notFound().build();
  309. }
  310. Path runDir = runDirManager.getRunDir(wf.getId(), wf.getName(), runId);
  311. if (runDir == null) {
  312. return ResponseEntity.notFound().build();
  313. }
  314. Path zipPath = Files.createTempFile("external-run-" + runId, ".zip");
  315. try (ZipOutputStream zos = new ZipOutputStream(Files.newOutputStream(zipPath))) {
  316. Files.walk(runDir)
  317. .filter(p -> !Files.isDirectory(p) && shouldIncludeInZip(runDir, p))
  318. .forEach(p -> {
  319. ZipEntry entry = new ZipEntry(runDir.relativize(p).toString().replace('\\', '/'));
  320. try {
  321. zos.putNextEntry(entry);
  322. Files.copy(p, zos);
  323. zos.closeEntry();
  324. } catch (IOException e) {
  325. log.warn("[ExternalApi] zip 条目写入失败: {}", p, e);
  326. }
  327. });
  328. }
  329. Resource resource = new UrlResource(zipPath.toUri());
  330. String filename = URLEncoder.encode("workflow-" + wf.getName() + "-run-" + runId + ".zip", StandardCharsets.UTF_8);
  331. return ResponseEntity.ok()
  332. .header(HttpHeaders.CONTENT_DISPOSITION, "attachment; filename*=UTF-8''" + filename)
  333. .contentType(MediaType.APPLICATION_OCTET_STREAM)
  334. .body(resource);
  335. }
  336. // ============================================================
  337. // GET /workflows/{workflowName}/runs/{runId}/nodes/{nodeId}/logs
  338. // ============================================================
  339. @GetMapping("/workflows/{workflowName}/runs/{runId}/nodes/{nodeId}/logs")
  340. public ResponseEntity<?> getNodeLogs(@PathVariable String workflowName,
  341. @PathVariable String runId,
  342. @PathVariable String nodeId) {
  343. WorkflowRun record = findRunRecord(runId);
  344. if (record == null) {
  345. return error(404, "RUN_NOT_FOUND", "运行不存在: " + runId);
  346. }
  347. List<WorkflowRunNode> nodes = workflowRunNodeRepository.findByRunIdOrderBySortOrderAsc(record.getId());
  348. WorkflowRunNode target = nodes.stream()
  349. .filter(n -> nodeId.equals(n.getNodeId()))
  350. .findFirst().orElse(null);
  351. if (target == null) {
  352. return error(404, "NODE_NOT_FOUND", "节点不存在: " + nodeId);
  353. }
  354. List<LogEntry> logs = new ArrayList<>();
  355. if (target.getLogs() != null && !target.getLogs().isBlank()) {
  356. try {
  357. List<JsonNode> raw = MAPPER.readValue(target.getLogs(), new TypeReference<List<JsonNode>>() {});
  358. for (JsonNode n : raw) {
  359. LogEntry e = new LogEntry();
  360. e.setType(n.path("type").asText(""));
  361. e.setMessage(n.path("message").asText(""));
  362. e.setDetail(n.path("detail").isNull() ? null : n.path("detail").asText(""));
  363. String ts = n.path("timestamp").asText("");
  364. if (!ts.isEmpty()) {
  365. try { e.setTimestamp(Instant.parse(ts)); } catch (Exception ignore) {}
  366. }
  367. logs.add(e);
  368. }
  369. } catch (Exception e) {
  370. log.warn("[ExternalApi] 解析节点日志失败 nodeId={}: {}", nodeId, e.getMessage());
  371. }
  372. }
  373. NodeLogsResponse body = new NodeLogsResponse();
  374. body.setRunId(runId);
  375. body.setNodeId(nodeId);
  376. body.setLogs(logs);
  377. return ResponseEntity.ok(ApiError.ok(body));
  378. }
  379. // ============================================================
  380. // GET /workflows — 列出可调用工作流
  381. // ============================================================
  382. @GetMapping("/workflows")
  383. public ResponseEntity<?> listWorkflows(@RequestAttribute(ExternalApiAuthFilter.ATTR_KEY_ENTRY) ExternalApiProperties.KeyEntry keyEntry) {
  384. List<Workflow> all = workflowRepository.findAllByOrderByUpdatedAtDesc();
  385. List<Workflow> accessible = all.stream()
  386. .filter(w -> properties.canAccess(keyEntry, w.getName(), w.getId()))
  387. .filter(w -> w.getGraphData() != null && !w.getGraphData().isBlank())
  388. .collect(Collectors.toList());
  389. List<WorkflowSummary> items = new ArrayList<>();
  390. for (Workflow w : accessible) {
  391. WorkflowSummary s = new WorkflowSummary();
  392. s.setId(w.getId());
  393. s.setName(w.getName());
  394. s.setDisplayName(w.getDisplayName());
  395. s.setDescription(w.getDescription());
  396. try {
  397. JsonNode root = MAPPER.readTree(w.getGraphData());
  398. s.setInputs(extractFields(root, "userInput"));
  399. s.setOutputs(extractFields(root, "output"));
  400. } catch (Exception e) {
  401. s.setInputs(List.of());
  402. s.setOutputs(List.of());
  403. }
  404. items.add(s);
  405. }
  406. WorkflowListResponse body = new WorkflowListResponse();
  407. body.setTotal(items.size());
  408. body.setItems(items);
  409. return ResponseEntity.ok(ApiError.ok(body));
  410. }
  411. // ============================================================
  412. // 工具方法
  413. // ============================================================
  414. /**
  415. * 解析工作流:优先 findByName,回退 findById(兼容数字 id URL)。
  416. */
  417. private Workflow resolveWorkflow(String nameOrId) {
  418. if (nameOrId == null || nameOrId.isEmpty()) return null;
  419. return workflowRepository.findByName(nameOrId)
  420. .orElseGet(() -> {
  421. if (nameOrId.matches("\\d+")) {
  422. return workflowRepository.findById(Long.valueOf(nameOrId)).orElse(null);
  423. }
  424. return null;
  425. });
  426. }
  427. /** 判断 created run 是否与 pathVariable workflowName 匹配(name 或数字 id 兼容) */
  428. private boolean matchesWorkflow(ExternalRunRegistry.CreatedRun created, String workflowName) {
  429. if (created == null || workflowName == null) return false;
  430. if (workflowName.equals(created.workflowName())) return true;
  431. if (workflowName.matches("\\d+") && Long.valueOf(workflowName).equals(created.workflowId())) return true;
  432. return false;
  433. }
  434. private ResponseEntity<ApiError> error(int httpStatus, String code, String message) {
  435. return ResponseEntity.status(httpStatus).body(ApiError.error(code, message));
  436. }
  437. private String generateRunId() {
  438. return UUID.randomUUID().toString().substring(0, 8);
  439. }
  440. private CreateRunRequest parsePayload(String payloadJson) {
  441. if (payloadJson == null || payloadJson.isBlank()) return new CreateRunRequest();
  442. try {
  443. return MAPPER.readValue(payloadJson, CreateRunRequest.class);
  444. } catch (Exception e) {
  445. log.warn("[ExternalApi] 解析 payload 失败,使用默认值: {}", e.getMessage());
  446. return new CreateRunRequest();
  447. }
  448. }
  449. private WorkflowRun findRunRecord(String runId) {
  450. return workflowRunRepository.findAllByOrderByStartedAtDesc().stream()
  451. .filter(r -> runId.equals(r.getRunId()))
  452. .findFirst().orElse(null);
  453. }
  454. @SuppressWarnings("unchecked")
  455. private Map<String, Object> parseOutputs(WorkflowRun record) {
  456. if (record.getOutputs() == null || record.getOutputs().isBlank()) return Map.of();
  457. try {
  458. return MAPPER.readValue(record.getOutputs(), Map.class);
  459. } catch (Exception e) {
  460. return Map.of();
  461. }
  462. }
  463. @SuppressWarnings("unchecked")
  464. private Map<String, Object> parseJsonToMap(String json) {
  465. if (json == null || json.isBlank()) return null;
  466. try {
  467. return MAPPER.readValue(json, Map.class);
  468. } catch (Exception e) {
  469. return null;
  470. }
  471. }
  472. private int countLogs(String logsJson) {
  473. if (logsJson == null || logsJson.isBlank()) return 0;
  474. try {
  475. return MAPPER.readTree(logsJson).size();
  476. } catch (Exception e) {
  477. return 0;
  478. }
  479. }
  480. private List<FieldDef> extractFields(JsonNode root, String nodeType) {
  481. List<FieldDef> result = new ArrayList<>();
  482. JsonNode nodes = root.path("nodes");
  483. if (!nodes.isArray()) return result;
  484. for (JsonNode node : nodes) {
  485. if (!nodeType.equals(node.path("type").asText(""))) continue;
  486. JsonNode data = node.path("data");
  487. JsonNode varArr = data.path("variables");
  488. if (varArr.isArray()) {
  489. for (JsonNode v : varArr) {
  490. result.add(toFieldDef(v));
  491. }
  492. }
  493. JsonNode fieldArr = data.path("fields");
  494. if (fieldArr.isArray()) {
  495. for (JsonNode f : fieldArr) {
  496. result.add(toFieldDef(f));
  497. }
  498. }
  499. }
  500. return result;
  501. }
  502. private FieldDef toFieldDef(JsonNode n) {
  503. FieldDef f = new FieldDef();
  504. f.setName(n.path("name").asText(""));
  505. f.setType(n.path("type").asText("string"));
  506. f.setRequired(n.path("required").asBoolean(false));
  507. f.setDescription(n.path("description").asText(""));
  508. return f;
  509. }
  510. private boolean shouldIncludeInZip(Path runDir, Path p) {
  511. Path rel = runDir.relativize(p);
  512. String relStr = rel.toString().replace('\\', '/');
  513. if (agentIgnoreFilter.shouldIgnore(relStr, false)) return false;
  514. Path cur = rel;
  515. while (cur != null && cur.getNameCount() > 1) {
  516. cur = cur.getParent();
  517. if (cur == null) break;
  518. if (agentIgnoreFilter.shouldIgnore(cur.toString().replace('\\', '/'), true)) return false;
  519. }
  520. return true;
  521. }
  522. }