datai/docs/archive/prompts/006-metadata-retrieve-core.md

20 KiB
Raw Permalink Blame History

Prompt - 元数据拉取核心功能

输入引用

引用相关的 docs 文档链接:

Context Maps

强制列出本次 Prompt 依赖的 Canvas 文件:

目标

实现 Salesforce 元数据拉取的核心功能,包括手动触发拉取、异步拉取执行、状态监控、拉取历史记录、拉取进度查询、拉取取消功能等。

输出格式

代码示例语言Java

1. JobExecutionStatus 枚举类

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 控制器

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<RetrieveResponse> 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<Page<DataiMetaJobExecution>> 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<DataiMetaJobExecution> pageResult = metadataRetrieveService.getRetrieveHistory(taskId, orgConfigId, page, size);
        return ResponseEntity.ok(pageResult);
    }

    /**
     * 查询拉取进度
     */
    @GetMapping("/progress/{jobId}")
    public ResponseEntity<RetrieveProgressResponse> 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<Void> cancelRetrieve(@PathVariable String jobId) {
        log.info("Cancel retrieve, jobId: {}", jobId);
        metadataRetrieveService.cancelRetrieve(jobId);
        return ResponseEntity.ok().build();
    }
}

3. RetrieveRequest DTO

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

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

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 服务接口

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<DataiMetaJobExecution> getRetrieveHistory(Long taskId, Long orgConfigId, int page, int size);

    /**
     * 查询拉取进度
     */
    RetrieveProgressResponse getRetrieveProgress(String jobId);

    /**
     * 取消拉取
     */
    void cancelRetrieve(String jobId);
}

7. MetadataRetrieveServiceImpl 服务实现

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<String, Future<?>> 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<DataiMetaJobExecution> getRetrieveHistory(Long taskId, Long orgConfigId, int page, int size) {
        QueryWrapper<DataiMetaJobExecution> 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<DataiMetaJobExecution> pageParam = new Page<>(page, size);
        return jobExecutionMapper.selectPage(pageParam, queryWrapper);
    }

    @Override
    public RetrieveProgressResponse getRetrieveProgress(String jobId) {
        DataiMetaJobExecution jobExecution = jobExecutionMapper.selectOne(
            new QueryWrapper<DataiMetaJobExecution>().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<DataiMetaJobExecution>().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 异常类

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. 单元测试

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<DataiMetaJobExecution> 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 异常类、单元测试 待反馈 待改进