# Prompt - 元数据拉取核心功能 ## 输入引用 引用相关的 docs 文档链接: - [REQ-010-6.md](../requirements/REQ-010-6.md) - 元数据拉取核心功能需求文档 - [REQ-010-1.md](../requirements/REQ-010-1.md) - 数据库表结构设计和创建 - [REQ-010-2.md](../requirements/REQ-010-2.md) - 基础实体类和Mapper创建 - [REQ-010-3.md](../requirements/REQ-010-3.md) - Salesforce组织配置管理 - [REQ-010-4.md](../requirements/REQ-010-4.md) - 元数据任务定义管理 - [REQ-010-5.md](../requirements/REQ-010-5.md) - Metadata API客户端封装 - [0015-metadata-retrieve-core.md](../decisions/adr/0015-metadata-retrieve-core.md) - 元数据拉取核心功能架构决策 ## Context Maps 强制列出本次 Prompt 依赖的 Canvas 文件: - [Authentication.canvas](../Authentication.canvas) - 项目架构视觉化展示 - **相关节点**: [集成核心](node_integration_core) - 提供与Salesforce的各种连接方式 - **相关节点**: [SessionManager](node_session_manager_detail) - 会话管理,提供登录服务 ## 目标 实现 Salesforce 元数据拉取的核心功能,包括手动触发拉取、异步拉取执行、状态监控、拉取历史记录、拉取进度查询、拉取取消功能等。 ## 输出格式 ### 代码示例(语言:Java) #### 1. JobExecutionStatus 枚举类 ```java package com.datai.salesforce.metadata.enums; import com.baomidou.mybatisplus.annotation.EnumValue; import com.fasterxml.jackson.annotation.JsonValue; /** * 作业执行状态枚举 */ public enum JobExecutionStatus { PENDING("pending", "待执行"), PROCESSING("processing", "执行中"), SUCCESS("success", "成功"), FAILED("failed", "失败"), PARTIAL_SUCCESS("partial_success", "部分成功"), CANCELLED("cancelled", "已取消"); @EnumValue private final String code; @JsonValue private final String displayName; JobExecutionStatus(String code, String displayName) { this.code = code; this.displayName = displayName; } public String getCode() { return code; } public String getDisplayName() { return displayName; } public static JobExecutionStatus fromCode(String code) { for (JobExecutionStatus status : JobExecutionStatus.values()) { if (status.getCode().equals(code)) { return status; } } throw new IllegalArgumentException("Invalid job execution status code: " + code); } } ``` #### 2. MetadataRetrieveController 控制器 ```java package com.datai.salesforce.metadata.controller; import com.baomidou.mybatisplus.extension.plugins.pagination.Page; import com.datai.salesforce.metadata.dto.RetrieveRequest; import com.datai.salesforce.metadata.dto.RetrieveResponse; import com.datai.salesforce.metadata.dto.RetrieveProgressResponse; import com.datai.salesforce.metadata.entity.DataiMetaJobExecution; import com.datai.salesforce.metadata.service.IMetadataRetrieveService; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.http.ResponseEntity; import org.springframework.validation.annotation.Validated; import org.springframework.web.bind.annotation.*; import javax.validation.Valid; /** * 元数据拉取控制器 */ @Slf4j @RestController @RequestMapping("/api/metadata/retrieve") @Validated public class MetadataRetrieveController { @Autowired private IMetadataRetrieveService metadataRetrieveService; /** * 手动触发拉取 */ @PostMapping public ResponseEntity triggerRetrieve(@Valid @RequestBody RetrieveRequest request) { log.info("Trigger metadata retrieve, taskId: {}, orgConfigId: {}", request.getTaskId(), request.getOrgConfigId()); RetrieveResponse response = metadataRetrieveService.triggerRetrieve(request); return ResponseEntity.ok(response); } /** * 查询拉取历史记录 */ @GetMapping("/history") public ResponseEntity> getRetrieveHistory( @RequestParam(required = false) Long taskId, @RequestParam(required = false) Long orgConfigId, @RequestParam(defaultValue = "1") int page, @RequestParam(defaultValue = "10") int size) { log.info("Get retrieve history, taskId: {}, orgConfigId: {}, page: {}, size: {}", taskId, orgConfigId, page, size); Page pageResult = metadataRetrieveService.getRetrieveHistory(taskId, orgConfigId, page, size); return ResponseEntity.ok(pageResult); } /** * 查询拉取进度 */ @GetMapping("/progress/{jobId}") public ResponseEntity getRetrieveProgress(@PathVariable String jobId) { log.info("Get retrieve progress, jobId: {}", jobId); RetrieveProgressResponse response = metadataRetrieveService.getRetrieveProgress(jobId); return ResponseEntity.ok(response); } /** * 取消拉取 */ @PostMapping("/cancel/{jobId}") public ResponseEntity cancelRetrieve(@PathVariable String jobId) { log.info("Cancel retrieve, jobId: {}", jobId); metadataRetrieveService.cancelRetrieve(jobId); return ResponseEntity.ok().build(); } } ``` #### 3. RetrieveRequest DTO ```java package com.datai.salesforce.metadata.dto; import lombok.Data; import javax.validation.constraints.NotNull; /** * 拉取请求 */ @Data public class RetrieveRequest { @NotNull(message = "任务ID不能为空") private Long taskId; @NotNull(message = "组织配置ID不能为空") private Long orgConfigId; } ``` #### 4. RetrieveResponse DTO ```java package com.datai.salesforce.metadata.dto; import lombok.AllArgsConstructor; import lombok.Data; import lombok.NoArgsConstructor; /** * 拉取响应 */ @Data @NoArgsConstructor @AllArgsConstructor public class RetrieveResponse { private String jobId; private String message; } ``` #### 5. RetrieveProgressResponse DTO ```java package com.datai.salesforce.metadata.dto; import lombok.AllArgsConstructor; import lombok.Data; import lombok.NoArgsConstructor; /** * 拉取进度响应 */ @Data @NoArgsConstructor @AllArgsConstructor public class RetrieveProgressResponse { private String jobId; private String status; private int progress; private String message; } ``` #### 6. IMetadataRetrieveService 服务接口 ```java package com.datai.salesforce.metadata.service; import com.baomidou.mybatisplus.extension.plugins.pagination.Page; import com.datai.salesforce.metadata.dto.RetrieveProgressResponse; import com.datai.salesforce.metadata.dto.RetrieveRequest; import com.datai.salesforce.metadata.dto.RetrieveResponse; import com.datai.salesforce.metadata.entity.DataiMetaJobExecution; /** * 元数据拉取服务接口 */ public interface IMetadataRetrieveService { /** * 手动触发拉取 */ RetrieveResponse triggerRetrieve(RetrieveRequest request); /** * 查询拉取历史记录 */ Page getRetrieveHistory(Long taskId, Long orgConfigId, int page, int size); /** * 查询拉取进度 */ RetrieveProgressResponse getRetrieveProgress(String jobId); /** * 取消拉取 */ void cancelRetrieve(String jobId); } ``` #### 7. MetadataRetrieveServiceImpl 服务实现 ```java package com.datai.salesforce.metadata.service.impl; import com.baomidou.mybatisplus.core.conditions.query.QueryWrapper; import com.baomidou.mybatisplus.extension.plugins.pagination.Page; import com.datai.salesforce.metadata.client.MetadataApiClient; import com.datai.salesforce.metadata.dto.RetrieveProgressResponse; import com.datai.salesforce.metadata.dto.RetrieveRequest; import com.datai.salesforce.metadata.dto.RetrieveResponse; import com.datai.salesforce.metadata.entity.DataiMetaJobExecution; import com.datai.salesforce.metadata.entity.DataiMetaTask; import com.datai.salesforce.metadata.enums.JobExecutionStatus; import com.datai.salesforce.metadata.enums.ScheduleType; import com.datai.salesforce.metadata.exception.RetrieveException; import com.datai.salesforce.metadata.mapper.DataiMetaJobExecutionMapper; import com.datai.salesforce.metadata.mapper.DataiMetaTaskMapper; import com.datai.salesforce.metadata.model.RetrieveRequest as SfRetrieveRequest; import com.datai.salesforce.metadata.model.RetrieveResult as SfRetrieveResult; import com.datai.salesforce.metadata.service.IMetadataRetrieveService; import com.datai.salesforce.metadata.util.PackageXmlEditor; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.scheduling.annotation.Async; import org.springframework.stereotype.Service; import org.springframework.transaction.annotation.Transactional; import java.util.Date; import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.Future; /** * 元数据拉取服务实现 */ @Slf4j @Service public class MetadataRetrieveServiceImpl implements IMetadataRetrieveService { @Autowired private MetadataApiClient metadataApiClient; @Autowired private DataiMetaTaskMapper taskMapper; @Autowired private DataiMetaJobExecutionMapper jobExecutionMapper; @Autowired private PackageXmlEditor packageXmlEditor; private final ConcurrentHashMap> runningJobs = new ConcurrentHashMap<>(); @Override @Transactional public RetrieveResponse triggerRetrieve(RetrieveRequest request) { DataiMetaTask task = taskMapper.selectById(request.getTaskId()); if (task == null) { throw new RetrieveException("任务不存在"); } String jobId = UUID.randomUUID().toString(); DataiMetaJobExecution jobExecution = new DataiMetaJobExecution(); jobExecution.setJobId(jobId); jobExecution.setTaskId(request.getTaskId()); jobExecution.setOrgConfigId(request.getOrgConfigId()); jobExecution.setJobType("RETRIEVE"); jobExecution.setStatus(JobExecutionStatus.PENDING.getCode()); jobExecution.setStartTime(new Date()); jobExecution.setCreateTime(new Date()); jobExecution.setUpdateTime(new Date()); jobExecutionMapper.insert(jobExecution); executeRetrieveAsync(jobId, request.getTaskId(), request.getOrgConfigId()); return new RetrieveResponse(jobId, "拉取任务已创建"); } @Async("metadataTaskExecutor") public void executeRetrieveAsync(String jobId, Long taskId, Long orgConfigId) { try { updateJobStatus(jobId, JobExecutionStatus.PROCESSING); DataiMetaTask task = taskMapper.selectById(taskId); String packageXml = task.getPackageXml(); SfRetrieveRequest sfRequest = new SfRetrieveRequest(); sfRequest.setOrgConfigId(orgConfigId); sfRequest.setPackageNames(new String[]{}); sfRequest.setSinglePackage(false); sfRequest.setSpecificFiles(new String[]{}); sfRequest.setUnpackaged(packageXmlEditor.parsePackageXml(packageXml)); sfRequest.setTimeout(300000); SfRetrieveResult result = metadataApiClient.retrieve(sfRequest); if (result.isSuccess()) { updateJobStatus(jobId, JobExecutionStatus.SUCCESS); log.info("Retrieve job {} completed successfully", jobId); } else { updateJobStatus(jobId, JobExecutionStatus.FAILED); log.error("Retrieve job {} failed: {}", jobId, result.getMessage()); } } catch (Exception e) { log.error("Retrieve job {} failed", jobId, e); updateJobStatus(jobId, JobExecutionStatus.FAILED); } finally { runningJobs.remove(jobId); } } @Override public Page getRetrieveHistory(Long taskId, Long orgConfigId, int page, int size) { QueryWrapper queryWrapper = new QueryWrapper<>(); queryWrapper.eq("job_type", "RETRIEVE"); if (taskId != null) { queryWrapper.eq("task_id", taskId); } if (orgConfigId != null) { queryWrapper.eq("org_config_id", orgConfigId); } queryWrapper.orderByDesc("create_time"); Page pageParam = new Page<>(page, size); return jobExecutionMapper.selectPage(pageParam, queryWrapper); } @Override public RetrieveProgressResponse getRetrieveProgress(String jobId) { DataiMetaJobExecution jobExecution = jobExecutionMapper.selectOne( new QueryWrapper().eq("job_id", jobId) ); if (jobExecution == null) { throw new RetrieveException("作业不存在"); } int progress = calculateProgress(jobExecution.getStatus()); return new RetrieveProgressResponse( jobId, jobExecution.getStatus(), progress, jobExecution.getMessage() ); } @Override public void cancelRetrieve(String jobId) { Future future = runningJobs.get(jobId); if (future != null) { boolean cancelled = future.cancel(true); if (cancelled) { updateJobStatus(jobId, JobExecutionStatus.CANCELLED); log.info("Retrieve job {} cancelled successfully", jobId); } else { throw new RetrieveException("取消失败,作业可能已完成或无法取消"); } } else { throw new RetrieveException("作业不存在或已完成"); } } private void updateJobStatus(String jobId, JobExecutionStatus status) { DataiMetaJobExecution jobExecution = new DataiMetaJobExecution(); jobExecution.setJobId(jobId); jobExecution.setStatus(status.getCode()); jobExecution.setUpdateTime(new Date()); if (status == JobExecutionStatus.SUCCESS || status == JobExecutionStatus.FAILED || status == JobExecutionStatus.CANCELLED) { jobExecution.setEndTime(new Date()); } jobExecutionMapper.update(jobExecution, new QueryWrapper().eq("job_id", jobId)); } private int calculateProgress(String status) { switch (JobExecutionStatus.fromCode(status)) { case PENDING: return 0; case PROCESSING: return 50; case SUCCESS: case PARTIAL_SUCCESS: return 100; case FAILED: case CANCELLED: return 0; default: return 0; } } } ``` #### 8. RetrieveException 异常类 ```java package com.datai.salesforce.metadata.exception; /** * 拉取异常 */ public class RetrieveException extends RuntimeException { public RetrieveException(String message) { super(message); } public RetrieveException(String message, Throwable cause) { super(message, cause); } } ``` #### 9. 单元测试 ```java package com.datai.salesforce.metadata.service.impl; import com.baomidou.mybatisplus.extension.plugins.pagination.Page; import com.datai.salesforce.metadata.dto.RetrieveProgressResponse; import com.datai.salesforce.metadata.dto.RetrieveRequest; import com.datai.salesforce.metadata.dto.RetrieveResponse; import com.datai.salesforce.metadata.entity.DataiMetaJobExecution; import com.datai.salesforce.metadata.enums.JobExecutionStatus; import com.datai.salesforce.metadata.exception.RetrieveException; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.extension.ExtendWith; import org.mockito.InjectMocks; import org.mockito.Mock; import org.mockito.junit.jupiter.MockitoExtension; import static org.junit.jupiter.api.Assertions.*; import static org.mockito.ArgumentMatchers.any; import static org.mockito.Mockito.*; @ExtendWith(MockitoExtension.class) class MetadataRetrieveServiceImplTest { @Mock private MetadataApiClient metadataApiClient; @Mock private DataiMetaTaskMapper taskMapper; @Mock private DataiMetaJobExecutionMapper jobExecutionMapper; @InjectMocks private MetadataRetrieveServiceImpl metadataRetrieveService; @Test void testTriggerRetrieve() { RetrieveRequest request = new RetrieveRequest(); request.setTaskId(1L); request.setOrgConfigId(1L); when(taskMapper.selectById(any())).thenReturn(new com.datai.salesforce.metadata.entity.DataiMetaTask()); when(jobExecutionMapper.insert(any())).thenReturn(1); RetrieveResponse response = metadataRetrieveService.triggerRetrieve(request); assertNotNull(response); assertNotNull(response.getJobId()); } @Test void testGetRetrieveHistory() { Page page = metadataRetrieveService.getRetrieveHistory(1L, 1L, 1, 10); assertNotNull(page); } @Test void testGetRetrieveProgress() { DataiMetaJobExecution jobExecution = new DataiMetaJobExecution(); jobExecution.setJobId("test-job-id"); jobExecution.setStatus(JobExecutionStatus.PROCESSING.getCode()); when(jobExecutionMapper.selectOne(any())).thenReturn(jobExecution); RetrieveProgressResponse response = metadataRetrieveService.getRetrieveProgress("test-job-id"); assertNotNull(response); assertEquals("test-job-id", response.getJobId()); assertEquals(JobExecutionStatus.PROCESSING.getCode(), response.getStatus()); assertEquals(50, response.getProgress()); } @Test void testCancelRetrieve() { assertThrows(RetrieveException.class, () -> { metadataRetrieveService.cancelRetrieve("test-job-id"); }); } } ``` ## 约束 列出使用此提示词时的约束条件,例如: - **技术栈限制**: 必须使用 Spring Boot 3、MyBatis Plus、Redis - **架构约束**: 必须遵循 Authentication.canvas 中定义的架构和调用关系 - **模块约束**: 必须在 datai-salesforce-metadata 模块下实现 - **认证约束**: 必须使用 SessionManager 进行会话管理和自动重新登录 - **API约束**: 必须使用现有的集成核心功能进行 API 调用 - **异步约束**: 必须使用异步线程池执行长时间任务 - **依赖约束**: 必须依赖于 REQ-010-1, REQ-010-2, REQ-010-3, REQ-010-4, REQ-010-5 ## Rule Set "请严格参考 @Authentication.canvas 中的状态机转移逻辑,不要自行发挥。" **具体规则**: - 必须使用 Canvas 中定义的类名和方法名 - 必须遵循 Canvas 中定义的调用关系 - 必须参考 Canvas 中的流程图逻辑 - 必须使用 SessionManager 进行会话管理和自动重新登录 - 必须使用现有的认证模块进行 OAuth 认证 - 必须使用现有的集成核心功能进行 API 调用 - 必须遵循现有的异常处理机制 - 必须遵循现有的日志记录规范 ## 验收标准 定义验证输出质量的具体标准,例如: - **功能完整性**: 所有拉取功能能够正常工作,异步执行机制正常 - **代码规范性**: 代码符合项目编码规范,有清晰的注释 - **性能要求**: 异步执行不影响系统响应,状态轮询频率合理 - **可维护性**: 代码结构清晰,易于扩展和维护 - **可测试性**: 代码易于单元测试和集成测试 ## 风险 识别使用此提示词可能带来的风险,例如: - **异步执行风险**: 异步执行机制复杂可能导致状态管理困难 - **状态轮询风险**: 状态轮询频率不当可能导致 API 限流 - **拉取取消风险**: 拉取取消功能复杂可能导致资源泄漏 - **历史记录风险**: 拉取历史记录过多可能影响查询性能 - **进度查询风险**: 进度查询不准确可能导致用户体验差 ## 使用记录 | 日期 | 使用场景 | 输入参数 | 输出结果 | 反馈 | 改进措施 | |------|---------|---------|---------|------|----------| | 2026-01-19 | 元数据拉取核心功能实现 | REQ-010-6 需求文档、ADR 文档 | JobExecutionStatus 枚举类、MetadataRetrieveController 控制器、RetrieveRequest DTO、RetrieveResponse DTO、RetrieveProgressResponse DTO、IMetadataRetrieveService 服务接口、MetadataRetrieveServiceImpl 服务实现、RetrieveException 异常类、单元测试 | 待反馈 | 待改进 |