WorkflowRunDirManager.java 4.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129
  1. package com.agent.management.engine;
  2. import lombok.extern.slf4j.Slf4j;
  3. import org.springframework.beans.factory.annotation.Value;
  4. import org.springframework.stereotype.Component;
  5. import java.io.IOException;
  6. import java.nio.file.Files;
  7. import java.nio.file.Path;
  8. import java.nio.file.Paths;
  9. import java.util.regex.Pattern;
  10. /**
  11. * 工作流运行目录管理器
  12. * 为每次工作流执行创建独立的工作目录,永久保留
  13. *
  14. * <p>目录结构:
  15. * <ul>
  16. * <li>新建:{@code {basePath}/{workflowName}/{runId}/}(workflowName 为 kebab-case 唯一标识)</li>
  17. * <li>兼容:旧目录 {@code {basePath}/{workflowId}/{runId}/} 仍可读取(用户改名前的历史 run)</li>
  18. * </ul>
  19. */
  20. @Slf4j
  21. @Component
  22. public class WorkflowRunDirManager {
  23. /** runId 白名单:UUID 截取的 8~16 位十六进制 */
  24. private static final Pattern RUN_ID_PATTERN = Pattern.compile("^[a-fA-F0-9]{1,16}$");
  25. private final Path basePath;
  26. public WorkflowRunDirManager(
  27. @Value("${workflow.run-dir:./data/workflow-runs}") String basePath) {
  28. this.basePath = Paths.get(basePath).toAbsolutePath().normalize();
  29. try {
  30. Files.createDirectories(this.basePath);
  31. log.info("[RunDirManager] 工作目录基路径: {}", this.basePath);
  32. } catch (IOException e) {
  33. throw new IllegalStateException("无法创建工作流运行目录: " + this.basePath, e);
  34. }
  35. }
  36. /**
  37. * 为指定工作流执行创建工作目录(新建使用 workflowName 作为目录段)。
  38. *
  39. * @param workflowId 工作流内部 ID(兼容旧目录兜底用)
  40. * @param workflowName 工作流唯一标识(kebab-case),新建目录使用此字段
  41. * @param runId 运行 ID(必须为 1~16 位十六进制)
  42. */
  43. public Path createRunDir(Long workflowId, String workflowName, String runId) {
  44. validateRunId(runId);
  45. Path runDir = resolveByName(workflowName, runId);
  46. try {
  47. Files.createDirectories(runDir);
  48. log.info("[RunDirManager] 创建运行目录: {}", runDir);
  49. return runDir;
  50. } catch (IOException e) {
  51. throw new RuntimeException("创建运行目录失败: " + runDir, e);
  52. }
  53. }
  54. /**
  55. * 获取指定运行的工作目录(已存在时返回,不存在返回 null)。
  56. * 查询顺序:先按 name,再按 id 兜底。
  57. */
  58. public Path getRunDir(Long workflowId, String workflowName, String runId) {
  59. try {
  60. validateRunId(runId);
  61. } catch (IllegalArgumentException e) {
  62. return null;
  63. }
  64. // 1. 优先按 name 查找
  65. if (workflowName != null && !workflowName.isEmpty()) {
  66. Path byName = resolveByName(workflowName, runId);
  67. if (Files.isDirectory(byName)) return byName;
  68. }
  69. // 2. 回退按 id 查找(兼容旧目录)
  70. if (workflowId != null) {
  71. Path byId = resolveById(workflowId, runId);
  72. if (Files.isDirectory(byId)) return byId;
  73. }
  74. return null;
  75. }
  76. /**
  77. * 兼容旧调用:仅按 id 创建/查询。新代码应使用带 name 的重载。
  78. */
  79. @Deprecated
  80. public Path createRunDir(Long workflowId, String runId) {
  81. return createRunDir(workflowId, String.valueOf(workflowId), runId);
  82. }
  83. @Deprecated
  84. public Path getRunDir(Long workflowId, String runId) {
  85. return getRunDir(workflowId, null, runId);
  86. }
  87. /**
  88. * 获取基路径
  89. */
  90. public Path getBasePath() {
  91. return basePath;
  92. }
  93. /**
  94. * 校验 runId 格式,防止路径穿越
  95. */
  96. private void validateRunId(String runId) {
  97. if (runId == null || !RUN_ID_PATTERN.matcher(runId).matches()) {
  98. throw new IllegalArgumentException("非法 runId: " + runId);
  99. }
  100. }
  101. private Path resolveByName(String workflowName, String runId) {
  102. Path runDir = basePath.resolve(workflowName).resolve(runId).normalize();
  103. if (!runDir.startsWith(basePath)) {
  104. throw new IllegalArgumentException("解析后的 runId 路径越界: " + runDir);
  105. }
  106. return runDir;
  107. }
  108. private Path resolveById(Long workflowId, String runId) {
  109. Path runDir = basePath.resolve(String.valueOf(workflowId)).resolve(runId).normalize();
  110. if (!runDir.startsWith(basePath)) {
  111. throw new IllegalArgumentException("解析后的 runId 路径越界: " + runDir);
  112. }
  113. return runDir;
  114. }
  115. }