新增业务外部接口日志管理

This commit is contained in:
wkc
2026-07-20 15:35:54 +08:00
parent 0a0313af7c
commit 2d9828ea5f
39 changed files with 1603 additions and 497 deletions

View File

@@ -15,6 +15,7 @@ import com.ruoyi.common.core.page.TableDataInfo;
import com.ruoyi.common.core.page.TableSupport;
import com.ruoyi.common.utils.SecurityUtils;
import com.ruoyi.lsfx.constants.LsfxConstants;
import com.ruoyi.lsfx.domain.CallerContext;
import io.swagger.v3.oas.annotations.Operation;
import io.swagger.v3.oas.annotations.tags.Tag;
import jakarta.annotation.Resource;
@@ -77,8 +78,8 @@ public class CcdiFileUploadController extends BaseController {
}
try {
String username = SecurityUtils.getUsername();
String batchId = fileUploadService.batchUploadFiles(projectId, files, username);
CallerContext caller = CallerContext.from(SecurityUtils.getLoginUser());
String batchId = fileUploadService.batchUploadFiles(projectId, files, caller);
return AjaxResult.success("上传任务已提交", batchId);
} catch (RejectedExecutionException e) {
log.warn("线程池已满,拒绝上传请求: projectId={}, fileCount={}", projectId, files.length);
@@ -130,16 +131,14 @@ public class CcdiFileUploadController extends BaseController {
return AjaxResult.error("开始日期和结束日期不能为空");
}
Long userId = SecurityUtils.getUserId();
String username = SecurityUtils.getUsername();
CallerContext caller = CallerContext.from(SecurityUtils.getLoginUser());
String batchId = fileUploadService.submitPullBankInfo(
dto.getProjectId(),
dto.getIdCards(),
dataChannelCode,
dto.getStartDate(),
dto.getEndDate(),
userId,
username
caller
);
return AjaxResult.success("拉取任务已提交", batchId);
}

View File

@@ -13,6 +13,7 @@ import com.ruoyi.ccdi.project.domain.vo.CcdiProjectHistoryListItemVO;
import com.ruoyi.ccdi.project.domain.vo.CcdiProjectStatusCountsVO;
import com.ruoyi.ccdi.project.domain.vo.CcdiProjectVO;
import com.ruoyi.ccdi.project.service.ICcdiProjectService;
import com.ruoyi.lsfx.domain.CallerContext;
import io.swagger.v3.oas.annotations.Operation;
import io.swagger.v3.oas.annotations.tags.Tag;
import jakarta.annotation.Resource;
@@ -43,7 +44,7 @@ public class CcdiProjectController extends BaseController {
@Operation(summary = "创建项目")
@PreAuthorize("@ss.hasPermi('ccdi:project:add')")
public AjaxResult createProject(@Validated @RequestBody CcdiProjectSaveDTO dto) {
CcdiProjectVO vo = projectService.createProject(dto);
CcdiProjectVO vo = projectService.createProject(dto, CallerContext.from(SecurityUtils.getLoginUser()));
return AjaxResult.success("项目创建成功", vo);
}
@@ -130,7 +131,7 @@ public class CcdiProjectController extends BaseController {
@Operation(summary = "导入历史项目")
@PreAuthorize("@ss.hasPermi('ccdi:project:add')")
public AjaxResult importFromHistory(@Validated @RequestBody CcdiProjectImportHistoryDTO dto) {
CcdiProjectVO vo = projectService.importFromHistory(dto, SecurityUtils.getUsername());
CcdiProjectVO vo = projectService.importFromHistory(dto, CallerContext.from(SecurityUtils.getLoginUser()));
return AjaxResult.success("项目创建成功", vo);
}

View File

@@ -4,6 +4,7 @@ import com.baomidou.mybatisplus.extension.plugins.pagination.Page;
import com.ruoyi.ccdi.project.domain.dto.CcdiFileUploadQueryDTO;
import com.ruoyi.ccdi.project.domain.entity.CcdiFileUploadRecord;
import com.ruoyi.ccdi.project.domain.vo.CcdiFileUploadStatisticsVO;
import com.ruoyi.lsfx.domain.CallerContext;
import org.springframework.web.multipart.MultipartFile;
import java.util.List;
@@ -24,7 +25,7 @@ public interface ICcdiFileUploadService {
* @param username 上传人
* @return 批次ID
*/
String batchUploadFiles(Long projectId, MultipartFile[] files, String username);
String batchUploadFiles(Long projectId, MultipartFile[] files, CallerContext caller);
/**
* 解析身份证文件
@@ -51,8 +52,7 @@ public interface ICcdiFileUploadService {
String dataChannelCode,
String startDate,
String endDate,
Long userId,
String username);
CallerContext caller);
/**
* 删除上传记录并清理关联数据

View File

@@ -7,6 +7,7 @@ import com.ruoyi.ccdi.project.domain.dto.CcdiProjectSaveDTO;
import com.ruoyi.ccdi.project.domain.vo.CcdiProjectHistoryListItemVO;
import com.ruoyi.ccdi.project.domain.vo.CcdiProjectStatusCountsVO;
import com.ruoyi.ccdi.project.domain.vo.CcdiProjectVO;
import com.ruoyi.lsfx.domain.CallerContext;
import java.util.List;
@@ -22,7 +23,7 @@ public interface ICcdiProjectService {
* @param dto 项目保存DTO
* @return 项目VO
*/
CcdiProjectVO createProject(CcdiProjectSaveDTO dto);
CcdiProjectVO createProject(CcdiProjectSaveDTO dto, CallerContext caller);
/**
* 更新项目
@@ -81,7 +82,7 @@ public interface ICcdiProjectService {
* @param operator 操作人
* @return 新建项目
*/
CcdiProjectVO importFromHistory(CcdiProjectImportHistoryDTO dto, String operator);
CcdiProjectVO importFromHistory(CcdiProjectImportHistoryDTO dto, CallerContext caller);
/**
* 查询各状态的项目总数(不受搜索条件影响)

View File

@@ -19,6 +19,7 @@ import com.ruoyi.ccdi.project.service.ICcdiProjectService;
import com.ruoyi.common.exception.ServiceException;
import com.ruoyi.lsfx.client.LsfxAnalysisClient;
import com.ruoyi.lsfx.constants.LsfxConstants;
import com.ruoyi.lsfx.domain.CallerContext;
import com.ruoyi.lsfx.domain.request.FetchInnerFlowRequest;
import com.ruoyi.lsfx.domain.request.GetBankStatementRequest;
import com.ruoyi.lsfx.domain.request.GetFileUploadStatusRequest;
@@ -157,8 +158,7 @@ public class CcdiFileUploadServiceImpl implements ICcdiFileUploadService {
String dataChannelCode,
String startDate,
String endDate,
Long userId,
String username) {
CallerContext caller) {
if (projectId == null) {
throw new IllegalArgumentException("项目ID不能为空");
}
@@ -209,7 +209,7 @@ public class CcdiFileUploadServiceImpl implements ICcdiFileUploadService {
record.setFileStatus("uploading");
record.setAccountNos(normalized);
record.setUploadTime(now);
record.setUploadUser(username);
record.setUploadUser(caller.username());
records.add(record);
}
if (records.isEmpty()) {
@@ -223,7 +223,7 @@ public class CcdiFileUploadServiceImpl implements ICcdiFileUploadService {
public void afterCommit() {
CompletableFuture.runAsync(() -> submitPullBankInfoTasks(
projectId, lsfxProjectId, records, normalizedIdCards,
normalizedDataChannelCode, startDate, endDate, batchId
normalizedDataChannelCode, startDate, endDate, batchId, caller
));
}
});
@@ -345,9 +345,9 @@ public class CcdiFileUploadServiceImpl implements ICcdiFileUploadService {
@Transactional
@Override
public String batchUploadFiles(Long projectId, MultipartFile[] files, String username) {
public String batchUploadFiles(Long projectId, MultipartFile[] files, CallerContext caller) {
log.info("【文件上传】开始批量上传: projectId={}, 文件数量={}, username={}",
projectId, files.length, username);
projectId, files.length, caller.username());
projectService.ensureProjectNotArchived(projectId, "已归档项目暂不允许上传或拉取数据");
projectService.ensureProjectWritable(projectId, "当前项目正在进行银行流水打标,暂不允许上传或拉取数据");
@@ -406,7 +406,7 @@ public class CcdiFileUploadServiceImpl implements ICcdiFileUploadService {
record.setFileSize(file.getSize());
record.setFileStatus("uploading");
record.setUploadTime(now);
record.setUploadUser(username);
record.setUploadUser(caller.username());
records.add(record);
}
} catch (IOException e) {
@@ -438,7 +438,7 @@ public class CcdiFileUploadServiceImpl implements ICcdiFileUploadService {
public void afterCommit() {
log.info("【文件上传】事务已提交,启动异步任务");
CompletableFuture.runAsync(() -> {
submitTasksAsync(projectId, finalLsfxProjectId, tempFilePaths, records, batchId);
submitTasksAsync(projectId, finalLsfxProjectId, tempFilePaths, records, batchId, caller);
});
}
});
@@ -496,7 +496,8 @@ public class CcdiFileUploadServiceImpl implements ICcdiFileUploadService {
private void submitTasksAsync(Long projectId, Integer lsfxProjectId,
List<String> tempFilePaths,
List<CcdiFileUploadRecord> records,
String batchId) {
String batchId,
CallerContext caller) {
log.info("【文件上传】调度线程启动: projectId={}, batchId={}", projectId, batchId);
List<CompletableFuture<Boolean>> futures = new ArrayList<>();
@@ -519,7 +520,7 @@ public class CcdiFileUploadServiceImpl implements ICcdiFileUploadService {
try {
// 尝试提交异步任务
CompletableFuture<Boolean> future = CompletableFuture.supplyAsync(
() -> processFileAsync(projectId, lsfxProjectId, tempFilePath, record.getId(), batchId, record),
() -> processFileAsync(projectId, lsfxProjectId, tempFilePath, record.getId(), batchId, record, caller),
fileUploadExecutor
);
futures.add(future);
@@ -600,7 +601,8 @@ public class CcdiFileUploadServiceImpl implements ICcdiFileUploadService {
String dataChannelCode,
String startDate,
String endDate,
String batchId) {
String batchId,
CallerContext caller) {
log.info("【拉取本行信息】调度线程启动: projectId={}, batchId={}", projectId, batchId);
List<CompletableFuture<Boolean>> futures = new ArrayList<>();
@@ -619,7 +621,8 @@ public class CcdiFileUploadServiceImpl implements ICcdiFileUploadService {
while (!submitted && retryCount < 2) {
try {
CompletableFuture<Boolean> future = CompletableFuture.supplyAsync(
() -> processPullBankInfoAsync(projectId, lsfxProjectId, record, idCard, dataChannelCode, startDate, endDate),
() -> processPullBankInfoAsync(projectId, lsfxProjectId, record, idCard,
dataChannelCode, startDate, endDate, caller),
fileUploadExecutor
);
futures.add(future);
@@ -660,7 +663,8 @@ public class CcdiFileUploadServiceImpl implements ICcdiFileUploadService {
String idCard,
String dataChannelCode,
String startDate,
String endDate ) {
String endDate,
CallerContext caller) {
try {
String normalizedDataChannelCode = normalizePullBankInfoDataChannelCode(dataChannelCode);
FetchInnerFlowRequest request = new FetchInnerFlowRequest();
@@ -677,7 +681,7 @@ public class CcdiFileUploadServiceImpl implements ICcdiFileUploadService {
}
request.setUploadUserId(LsfxConstants.DEFAULT_USER_ID);
FetchInnerFlowResponse response = lsfxClient.fetchInnerFlow(request);
FetchInnerFlowResponse response = lsfxClient.fetchInnerFlow(caller, request);
if (response == null || response.getData() == null || response.getData().isEmpty()) {
throw new RuntimeException("拉取本行信息失败: 未返回logId");
}
@@ -687,7 +691,7 @@ public class CcdiFileUploadServiceImpl implements ICcdiFileUploadService {
throw new RuntimeException("拉取本行信息失败: 未返回logId");
}
processRecordAfterLogIdReady(projectId, lsfxProjectId, record, logId);
processRecordAfterLogIdReady(projectId, lsfxProjectId, record, logId, caller);
return true;
} catch (Exception e) {
log.error("【拉取本行信息】处理失败: idCard={}, recordId={}", idCard, record.getId(), e);
@@ -709,7 +713,8 @@ public class CcdiFileUploadServiceImpl implements ICcdiFileUploadService {
*/
@Async("fileUploadExecutor")
public boolean processFileAsync(Long projectId, Integer lsfxProjectId, String tempFilePath,
Long recordId, String batchId, CcdiFileUploadRecord record) {
Long recordId, String batchId, CcdiFileUploadRecord record,
CallerContext caller) {
log.info("【文件上传】开始处理文件: fileName={}, recordId={}, tempPath={}",
record.getFileName(), recordId, tempFilePath);
@@ -730,7 +735,7 @@ public class CcdiFileUploadServiceImpl implements ICcdiFileUploadService {
throw new RuntimeException("临时文件不存在: " + tempFilePath);
}
UploadFileResponse uploadResponse = lsfxClient.uploadFile(lsfxProjectId, file, record.getFileName());
UploadFileResponse uploadResponse = lsfxClient.uploadFile(caller, lsfxProjectId, file, record.getFileName());
if (uploadResponse == null || uploadResponse.getData() == null
|| uploadResponse.getData().getUploadLogList() == null
|| uploadResponse.getData().getUploadLogList().isEmpty()) {
@@ -744,7 +749,7 @@ public class CcdiFileUploadServiceImpl implements ICcdiFileUploadService {
}
log.info("【文件上传】文件上传成功: logId={}", logId);
processRecordAfterLogIdReady(projectId, lsfxProjectId, record, logId, true);
processRecordAfterLogIdReady(projectId, lsfxProjectId, record, logId, true, caller);
log.info("【文件上传】处理完成: fileName={}", record.getFileName());
return true;
@@ -780,22 +785,24 @@ public class CcdiFileUploadServiceImpl implements ICcdiFileUploadService {
private void processRecordAfterLogIdReady(Long projectId,
Integer lsfxProjectId,
CcdiFileUploadRecord record,
Integer logId) {
processRecordAfterLogIdReady(projectId, lsfxProjectId, record, logId, false);
Integer logId,
CallerContext caller) {
processRecordAfterLogIdReady(projectId, lsfxProjectId, record, logId, false, caller);
}
private void processRecordAfterLogIdReady(Long projectId,
Integer lsfxProjectId,
CcdiFileUploadRecord record,
Integer logId,
boolean preserveRecordFileName) {
boolean preserveRecordFileName,
CallerContext caller) {
log.info("【文件上传】步骤3: 更新状态为解析中, logId={}", logId);
record.setLogId(logId);
record.setFileStatus("parsing");
recordMapper.updateById(record);
log.info("【文件上传】步骤4: 开始轮询解析状态");
boolean parsingComplete = waitForParsingComplete(lsfxProjectId, logId.toString());
boolean parsingComplete = waitForParsingComplete(caller, lsfxProjectId, logId.toString());
if (!parsingComplete) {
throw new RuntimeException("解析超时(超过10分钟),请检查文件格式是否正确");
}
@@ -805,7 +812,7 @@ public class CcdiFileUploadServiceImpl implements ICcdiFileUploadService {
statusRequest.setGroupId(lsfxProjectId);
statusRequest.setLogId(logId);
GetFileUploadStatusResponse statusResponse = lsfxClient.getFileUploadStatus(statusRequest);
GetFileUploadStatusResponse statusResponse = lsfxClient.getFileUploadStatus(caller, statusRequest);
if (statusResponse == null || statusResponse.getData() == null
|| statusResponse.getData().getLogs() == null
|| statusResponse.getData().getLogs().isEmpty()) {
@@ -846,8 +853,8 @@ public class CcdiFileUploadServiceImpl implements ICcdiFileUploadService {
log.info("【文件上传】步骤7: 获取流水数据");
String fallbackCretNo = extractIdCardFromFileName(record.getFileName());
FetchBankStatementResult fetchResult = fetchAndSaveBankStatements(projectId, lsfxProjectId, logId,
fallbackCretNo);
FetchBankStatementResult fetchResult = fetchAndSaveBankStatements(caller, projectId, lsfxProjectId,
logId, fallbackCretNo);
if (!fetchResult.isSuccess()) {
updateFailedRecord(record, fetchResult.getErrorMessage());
return;
@@ -867,7 +874,7 @@ public class CcdiFileUploadServiceImpl implements ICcdiFileUploadService {
* @param logId 文件ID
* @return true=解析完成false=超时未完成
*/
private boolean waitForParsingComplete(Integer groupId, String logId) {
private boolean waitForParsingComplete(CallerContext caller, Integer groupId, String logId) {
log.info("【文件上传】开始轮询解析状态: groupId={}, logId={}", groupId, logId);
int maxRetries = 300;
@@ -876,7 +883,7 @@ public class CcdiFileUploadServiceImpl implements ICcdiFileUploadService {
for (int i = 1; i <= maxRetries; i++) {
try {
// 调用检查解析状态接口
CheckParseStatusResponse response = lsfxClient.checkParseStatus(groupId, logId);
CheckParseStatusResponse response = lsfxClient.checkParseStatus(caller, groupId, logId);
if (response == null || response.getData() == null) {
log.warn("【文件上传】轮询第{}次: 响应数据为空", i);
@@ -919,9 +926,8 @@ public class CcdiFileUploadServiceImpl implements ICcdiFileUploadService {
* @param groupId 流水分析平台项目ID
* @param logId 文件ID
*/
private FetchBankStatementResult fetchAndSaveBankStatements(Long projectId, Integer groupId,
Integer logId,
String fallbackCretNo) {
private FetchBankStatementResult fetchAndSaveBankStatements(CallerContext caller, Long projectId, Integer groupId,
Integer logId, String fallbackCretNo) {
log.info("【文件上传】开始获取流水数据: projectId={}, groupId={}, logId={}",
projectId, groupId, logId);
@@ -934,7 +940,7 @@ public class CcdiFileUploadServiceImpl implements ICcdiFileUploadService {
firstRequest.setPageNow(1);
firstRequest.setPageSize(1);
GetBankStatementResponse firstResponse = lsfxClient.getBankStatement(firstRequest);
GetBankStatementResponse firstResponse = lsfxClient.getBankStatement(caller, firstRequest);
if (firstResponse == null || firstResponse.getData() == null) {
result.setSuccess(false);
result.setErrorMessage("获取流水数据失败: 响应数据为空");
@@ -968,7 +974,7 @@ public class CcdiFileUploadServiceImpl implements ICcdiFileUploadService {
request.setPageNow(pageNow);
request.setPageSize(pageSize);
GetBankStatementResponse response = lsfxClient.getBankStatement(request);
GetBankStatementResponse response = lsfxClient.getBankStatement(caller, request);
if (response == null || response.getData() == null
|| response.getData().getBankStatementList() == null) {
result.setSuccess(false);

View File

@@ -19,6 +19,7 @@ import com.ruoyi.ccdi.project.service.CcdiProjectAccessService;
import com.ruoyi.ccdi.project.service.ICcdiProjectService;
import com.ruoyi.common.exception.ServiceException;
import com.ruoyi.lsfx.client.LsfxAnalysisClient;
import com.ruoyi.lsfx.domain.CallerContext;
import com.ruoyi.lsfx.domain.request.GetTokenRequest;
import com.ruoyi.lsfx.domain.response.GetTokenResponse;
import jakarta.annotation.Resource;
@@ -61,9 +62,9 @@ public class CcdiProjectServiceImpl implements ICcdiProjectService {
@Override
@Transactional(rollbackFor = Exception.class)
public CcdiProjectVO createProject(CcdiProjectSaveDTO dto) {
public CcdiProjectVO createProject(CcdiProjectSaveDTO dto, CallerContext caller) {
// 1. 调用流水分析平台获取projectId
Integer lsfxProjectId = callLsfxPlatform(dto.getProjectName());
Integer lsfxProjectId = callLsfxPlatform(dto.getProjectName(), caller);
// 2. 创建项目实体
CcdiProject project = new CcdiProject();
@@ -163,18 +164,18 @@ public class CcdiProjectServiceImpl implements ICcdiProjectService {
@Override
@Transactional(rollbackFor = Exception.class)
public CcdiProjectVO importFromHistory(CcdiProjectImportHistoryDTO dto, String operator) {
public CcdiProjectVO importFromHistory(CcdiProjectImportHistoryDTO dto, CallerContext caller) {
projectAccessService.assertSourceProjectsReadable(dto.getSourceProjectIds());
CcdiProjectSaveDTO saveDTO = new CcdiProjectSaveDTO();
saveDTO.setProjectName(dto.getProjectName());
saveDTO.setDescription(dto.getDescription());
saveDTO.setConfigType("default");
CcdiProjectVO project = createProject(saveDTO);
CcdiProjectVO project = createProject(saveDTO, caller);
TransactionSynchronizationManager.registerSynchronization(new TransactionSynchronization() {
@Override
public void afterCommit() {
applicationEventPublisher.publishEvent(
new CcdiProjectHistoryImportSubmittedEvent(project.getProjectId(), project.getLsfxProjectId(), dto, operator)
new CcdiProjectHistoryImportSubmittedEvent(project.getProjectId(), project.getLsfxProjectId(), dto, caller.username())
);
}
});
@@ -380,7 +381,7 @@ public class CcdiProjectServiceImpl implements ICcdiProjectService {
* @return 流水分析平台项目ID
* @throws ServiceException 调用失败或响应无效时抛出
*/
private Integer callLsfxPlatform(String projectName) {
private Integer callLsfxPlatform(String projectName, CallerContext caller) {
// 构建请求参数
GetTokenRequest request = new GetTokenRequest();
request.setProjectNo("902000_" + System.currentTimeMillis());
@@ -393,7 +394,7 @@ public class CcdiProjectServiceImpl implements ICcdiProjectService {
request.setDepartmentCode("902000");
// 调用流水分析平台(异常处理和日志已在 LsfxAnalysisClient 中完成)
GetTokenResponse response = lsfxAnalysisClient.getToken(request);
GetTokenResponse response = lsfxAnalysisClient.getToken(caller, request);
// 业务层校验:确保响应有效
if (response == null || response.getData() == null) {

View File

@@ -6,6 +6,7 @@ import com.ruoyi.ccdi.project.service.ICcdiFileUploadService;
import com.ruoyi.common.core.domain.entity.SysUser;
import com.ruoyi.common.core.domain.model.LoginUser;
import com.ruoyi.common.core.domain.AjaxResult;
import com.ruoyi.lsfx.domain.CallerContext;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
@@ -30,6 +31,7 @@ import static org.mockito.Mockito.when;
class CcdiFileUploadControllerTest {
private static final Long PROJECT_ID = 100L;
private static final CallerContext CALLER = CallerContext.of(9527L, "admin");
@InjectMocks
private CcdiFileUploadController controller;
@@ -72,7 +74,7 @@ class CcdiFileUploadControllerTest {
};
setLoginUser(9527L, "admin");
when(fileUploadService.batchUploadFiles(PROJECT_ID, files, "admin"))
when(fileUploadService.batchUploadFiles(PROJECT_ID, files, CALLER))
.thenReturn("batch-1");
AjaxResult result = controller.batchUpload(PROJECT_ID, files);
@@ -80,7 +82,7 @@ class CcdiFileUploadControllerTest {
assertEquals(200, result.get("code"));
assertEquals("batch-1", result.get("data"));
verify(projectAccessService).assertCanOperate(PROJECT_ID);
verify(fileUploadService).batchUploadFiles(PROJECT_ID, files, "admin");
verify(fileUploadService).batchUploadFiles(PROJECT_ID, files, CALLER);
}
@Test
@@ -93,7 +95,8 @@ class CcdiFileUploadControllerTest {
dto.setEndDate("2026-03-10");
setLoginUser(9527L, "admin");
when(fileUploadService.submitPullBankInfo(PROJECT_ID, dto.getIdCards(), "ZJRCU", "2026-03-01", "2026-03-10", 9527L, "admin"))
when(fileUploadService.submitPullBankInfo(PROJECT_ID, dto.getIdCards(), "ZJRCU",
"2026-03-01", "2026-03-10", CALLER))
.thenReturn("batch-1");
AjaxResult result = controller.pullBankInfo(dto);
@@ -109,7 +112,7 @@ class CcdiFileUploadControllerTest {
dto.setDataChannelCode("JZL");
setLoginUser(9527L, "admin");
when(fileUploadService.submitPullBankInfo(PROJECT_ID, dto.getIdCards(), "JZL", null, null, 9527L, "admin"))
when(fileUploadService.submitPullBankInfo(PROJECT_ID, dto.getIdCards(), "JZL", null, null, CALLER))
.thenReturn("batch-1");
AjaxResult result = controller.pullBankInfo(dto);

View File

@@ -17,6 +17,7 @@ import com.ruoyi.ccdi.project.service.ICcdiProjectService;
import com.ruoyi.common.exception.ServiceException;
import com.ruoyi.lsfx.client.LsfxAnalysisClient;
import com.ruoyi.lsfx.constants.LsfxConstants;
import com.ruoyi.lsfx.domain.CallerContext;
import com.ruoyi.lsfx.domain.request.FetchInnerFlowRequest;
import com.ruoyi.lsfx.domain.request.GetBankStatementRequest;
import com.ruoyi.lsfx.domain.response.CheckParseStatusResponse;
@@ -34,6 +35,7 @@ import org.mockito.junit.jupiter.MockitoExtension;
import org.slf4j.LoggerFactory;
import org.springframework.mock.web.MockMultipartFile;
import org.springframework.test.util.ReflectionTestUtils;
import org.springframework.security.core.context.SecurityContextHolder;
import org.springframework.transaction.support.TransactionSynchronizationManager;
import org.springframework.web.multipart.MultipartFile;
@@ -63,6 +65,7 @@ import static org.mockito.Mockito.doAnswer;
import static org.mockito.Mockito.never;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
import org.mockito.ArgumentCaptor;
@ExtendWith(MockitoExtension.class)
class CcdiFileUploadServiceImplTest {
@@ -71,6 +74,7 @@ class CcdiFileUploadServiceImplTest {
private static final Integer LSFX_PROJECT_ID = 200;
private static final Long RECORD_ID = 300L;
private static final Integer LOG_ID = 400;
private static final CallerContext CALLER = CallerContext.of(9527L, "admin");
private static final int MAX_ERROR_MESSAGE_LENGTH = 2000;
@InjectMocks
@@ -149,8 +153,7 @@ class CcdiFileUploadServiceImplTest {
LsfxConstants.DATA_CHANNEL_ZJRCU,
"2026-03-01",
"2026-03-10",
9527L,
"admin"
CALLER
);
assertNotNull(batchId);
@@ -177,8 +180,7 @@ class CcdiFileUploadServiceImplTest {
LsfxConstants.DATA_CHANNEL_ZJRCU,
"2026-01-01",
"2026-01-31",
1L,
"tester"
CALLER
));
}
@@ -195,7 +197,7 @@ class CcdiFileUploadServiceImplTest {
);
assertThrows(ServiceException.class,
() -> service.batchUploadFiles(PROJECT_ID, new MultipartFile[]{file}, "tester"));
() -> service.batchUploadFiles(PROJECT_ID, new MultipartFile[]{file}, CALLER));
}
@Test
@@ -222,7 +224,7 @@ class CcdiFileUploadServiceImplTest {
TransactionSynchronizationManager.initSynchronization();
try {
String batchId = service.batchUploadFiles(PROJECT_ID, new MultipartFile[]{file}, "tester");
String batchId = service.batchUploadFiles(PROJECT_ID, new MultipartFile[]{file}, CALLER);
assertNotNull(batchId);
assertNotNull(inserted.get());
@@ -251,12 +253,12 @@ class CcdiFileUploadServiceImplTest {
);
IllegalArgumentException exception = assertThrows(IllegalArgumentException.class,
() -> service.batchUploadFiles(PROJECT_ID, new MultipartFile[]{file}, "tester"));
() -> service.batchUploadFiles(PROJECT_ID, new MultipartFile[]{file}, CALLER));
assertTrue(exception.getMessage().contains("身份证"));
assertFalse(Files.exists(tempDir.resolve("temp")));
verify(recordMapper, never()).insertBatch(any());
verify(lsfxClient, never()).uploadFile(any(), org.mockito.ArgumentMatchers.<java.io.File>any(), any());
verify(lsfxClient, never()).uploadFile(any(), any(), org.mockito.ArgumentMatchers.<java.io.File>any(), any());
}
@Test
@@ -272,12 +274,12 @@ class CcdiFileUploadServiceImplTest {
);
IllegalArgumentException exception = assertThrows(IllegalArgumentException.class,
() -> service.batchUploadFiles(PROJECT_ID, new MultipartFile[]{file}, "tester"));
() -> service.batchUploadFiles(PROJECT_ID, new MultipartFile[]{file}, CALLER));
assertTrue(exception.getMessage().contains("文件名不能为空"));
assertFalse(Files.exists(tempDir.resolve("temp")));
verify(recordMapper, never()).insertBatch(any());
verify(lsfxClient, never()).uploadFile(any(), org.mockito.ArgumentMatchers.<java.io.File>any(), any());
verify(lsfxClient, never()).uploadFile(any(), any(), org.mockito.ArgumentMatchers.<java.io.File>any(), any());
}
@Test
@@ -293,6 +295,29 @@ class CcdiFileUploadServiceImplTest {
assertFalse(Files.exists(batchLogDir));
}
@Test
void submitTasksAsync_shouldKeepEachCapturedCallerAfterSecurityContextIsCleared() throws Exception {
setField("fileUploadExecutor", (Executor) Runnable::run);
CallerContext callerA = CallerContext.of(101L, "userA");
CallerContext callerB = CallerContext.of(202L, "userB");
Path fileA = createTempFile();
Path fileB = createTempFile();
CcdiFileUploadRecord recordA = buildRecord();
CcdiFileUploadRecord recordB = buildRecord();
recordB.setId(RECORD_ID + 1);
when(lsfxClient.uploadFile(any(CallerContext.class), eq(LSFX_PROJECT_ID), any(), any()))
.thenThrow(new RuntimeException("stop after caller capture"));
SecurityContextHolder.clearContext();
invokeSubmitTasksAsync(List.of(fileA.toString()), List.of(recordA), "batch-a", callerA);
invokeSubmitTasksAsync(List.of(fileB.toString()), List.of(recordB), "batch-b", callerB);
ArgumentCaptor<CallerContext> callerCaptor = ArgumentCaptor.forClass(CallerContext.class);
verify(lsfxClient, org.mockito.Mockito.times(2)).uploadFile(
callerCaptor.capture(), eq(LSFX_PROJECT_ID), any(), any());
assertEquals(List.of(callerA, callerB), callerCaptor.getAllValues());
}
@Test
void handleTagRebuildAfterBatchCompletion_shouldLogSkipWhenAllRecordsFailed() {
Logger logger = (Logger) LoggerFactory.getLogger(CcdiFileUploadServiceImpl.class);
@@ -352,18 +377,18 @@ class CcdiFileUploadServiceImplTest {
AtomicInteger sequence = new AtomicInteger();
captureRecordStatus(events, sequence);
when(lsfxClient.uploadFile(eq(LSFX_PROJECT_ID), any(), org.mockito.ArgumentMatchers.anyString()))
when(lsfxClient.uploadFile(eq(CALLER), eq(LSFX_PROJECT_ID), any(), org.mockito.ArgumentMatchers.anyString()))
.thenReturn(buildUploadResponse());
when(lsfxClient.checkParseStatus(LSFX_PROJECT_ID, String.valueOf(LOG_ID)))
when(lsfxClient.checkParseStatus(CALLER, LSFX_PROJECT_ID, String.valueOf(LOG_ID)))
.thenReturn(buildCheckParseStatusResponse(false));
when(lsfxClient.getFileUploadStatus(any())).thenReturn(buildParsedSuccessStatusResponse());
when(lsfxClient.getBankStatement(any(GetBankStatementRequest.class)))
when(lsfxClient.getFileUploadStatus(eq(CALLER), any())).thenReturn(buildParsedSuccessStatusResponse());
when(lsfxClient.getBankStatement(eq(CALLER), any(GetBankStatementRequest.class)))
.thenThrow(new RuntimeException("bank statement fetch failed"));
CcdiFileUploadRecord record = buildRecord();
Path tempFile = createTempFile();
service.processFileAsync(PROJECT_ID, LSFX_PROJECT_ID, tempFile.toString(), RECORD_ID, "batch-1", record);
service.processFileAsync(PROJECT_ID, LSFX_PROJECT_ID, tempFile.toString(), RECORD_ID, "batch-1", record, CALLER);
assertTrue(events.stream().anyMatch(event -> event.endsWith("record:parsed_failed")));
assertFalse(events.stream().anyMatch(event -> event.endsWith("record:parsed_success")));
@@ -379,12 +404,12 @@ class CcdiFileUploadServiceImplTest {
when(projectMapper.selectById(PROJECT_ID)).thenReturn(project);
when(bankStatementMapper.countMatchedStaffCountByProjectId(PROJECT_ID)).thenReturn(1);
when(lsfxClient.uploadFile(eq(LSFX_PROJECT_ID), any(), org.mockito.ArgumentMatchers.anyString()))
when(lsfxClient.uploadFile(eq(CALLER), eq(LSFX_PROJECT_ID), any(), org.mockito.ArgumentMatchers.anyString()))
.thenReturn(buildUploadResponse());
when(lsfxClient.checkParseStatus(LSFX_PROJECT_ID, String.valueOf(LOG_ID)))
when(lsfxClient.checkParseStatus(CALLER, LSFX_PROJECT_ID, String.valueOf(LOG_ID)))
.thenReturn(buildCheckParseStatusResponse(false));
when(lsfxClient.getFileUploadStatus(any())).thenReturn(buildParsedSuccessStatusResponse());
when(lsfxClient.getBankStatement(any(GetBankStatementRequest.class)))
when(lsfxClient.getFileUploadStatus(eq(CALLER), any())).thenReturn(buildParsedSuccessStatusResponse());
when(lsfxClient.getBankStatement(eq(CALLER), any(GetBankStatementRequest.class)))
.thenAnswer(invocation -> {
events.add(sequence.incrementAndGet() + ":bank-fetch");
return buildEmptyBankStatementResponse();
@@ -393,7 +418,7 @@ class CcdiFileUploadServiceImplTest {
CcdiFileUploadRecord record = buildRecord();
Path tempFile = createTempFile();
service.processFileAsync(PROJECT_ID, LSFX_PROJECT_ID, tempFile.toString(), RECORD_ID, "batch-1", record);
service.processFileAsync(PROJECT_ID, LSFX_PROJECT_ID, tempFile.toString(), RECORD_ID, "batch-1", record, CALLER);
int fetchIndex = findEventIndex(events, "bank-fetch");
int successIndex = findEventIndex(events, "record:parsed_success");
@@ -406,18 +431,18 @@ class CcdiFileUploadServiceImplTest {
@Test
void processFileAsync_shouldCleanupInsertedStatementsWhenFetchFails() throws IOException {
when(lsfxClient.uploadFile(eq(LSFX_PROJECT_ID), any(), org.mockito.ArgumentMatchers.anyString()))
when(lsfxClient.uploadFile(eq(CALLER), eq(LSFX_PROJECT_ID), any(), org.mockito.ArgumentMatchers.anyString()))
.thenReturn(buildUploadResponse());
when(lsfxClient.checkParseStatus(LSFX_PROJECT_ID, String.valueOf(LOG_ID)))
when(lsfxClient.checkParseStatus(CALLER, LSFX_PROJECT_ID, String.valueOf(LOG_ID)))
.thenReturn(buildCheckParseStatusResponse(false));
when(lsfxClient.getFileUploadStatus(any())).thenReturn(buildParsedSuccessStatusResponse());
when(lsfxClient.getBankStatement(any(GetBankStatementRequest.class)))
when(lsfxClient.getFileUploadStatus(eq(CALLER), any())).thenReturn(buildParsedSuccessStatusResponse());
when(lsfxClient.getBankStatement(eq(CALLER), any(GetBankStatementRequest.class)))
.thenThrow(new RuntimeException("bank statement fetch failed"));
CcdiFileUploadRecord record = buildRecord();
Path tempFile = createTempFile();
service.processFileAsync(PROJECT_ID, LSFX_PROJECT_ID, tempFile.toString(), RECORD_ID, "batch-1", record);
service.processFileAsync(PROJECT_ID, LSFX_PROJECT_ID, tempFile.toString(), RECORD_ID, "batch-1", record, CALLER);
verify(bankStatementMapper).deleteByProjectIdAndBatchId(PROJECT_ID, LOG_ID);
}
@@ -427,11 +452,11 @@ class CcdiFileUploadServiceImplTest {
GetFileUploadStatusResponse statusResponse = buildParsedSuccessStatusResponse("XX身份证.xlsx");
statusResponse.getData().getLogs().get(0).setFileSize(2048L);
when(lsfxClient.fetchInnerFlow(any())).thenReturn(buildFetchInnerFlowResponse(LOG_ID));
when(lsfxClient.checkParseStatus(LSFX_PROJECT_ID, String.valueOf(LOG_ID)))
when(lsfxClient.fetchInnerFlow(eq(CALLER), any())).thenReturn(buildFetchInnerFlowResponse(LOG_ID));
when(lsfxClient.checkParseStatus(CALLER, LSFX_PROJECT_ID, String.valueOf(LOG_ID)))
.thenReturn(buildCheckParseStatusResponse(false));
when(lsfxClient.getFileUploadStatus(any())).thenReturn(statusResponse);
when(lsfxClient.getBankStatement(any(GetBankStatementRequest.class)))
when(lsfxClient.getFileUploadStatus(eq(CALLER), any())).thenReturn(statusResponse);
when(lsfxClient.getBankStatement(eq(CALLER), any(GetBankStatementRequest.class)))
.thenReturn(buildEmptyBankStatementResponse());
CcdiFileUploadRecord record = buildRecord();
@@ -444,7 +469,8 @@ class CcdiFileUploadServiceImplTest {
"110101199001018888",
LsfxConstants.DATA_CHANNEL_ZJRCU,
"2026-03-01",
"2026-03-10"
"2026-03-10",
CALLER
);
verify(recordMapper, org.mockito.Mockito.atLeastOnce()).updateById(
@@ -457,11 +483,11 @@ class CcdiFileUploadServiceImplTest {
@Test
void processPullBankInfoAsync_shouldFetchJzlWithZeroDateRange() {
when(lsfxClient.fetchInnerFlow(any())).thenReturn(buildFetchInnerFlowResponse(LOG_ID));
when(lsfxClient.checkParseStatus(LSFX_PROJECT_ID, String.valueOf(LOG_ID)))
when(lsfxClient.fetchInnerFlow(eq(CALLER), any())).thenReturn(buildFetchInnerFlowResponse(LOG_ID));
when(lsfxClient.checkParseStatus(CALLER, LSFX_PROJECT_ID, String.valueOf(LOG_ID)))
.thenReturn(buildCheckParseStatusResponse(false));
when(lsfxClient.getFileUploadStatus(any())).thenReturn(buildParsedSuccessStatusResponse());
when(lsfxClient.getBankStatement(any(GetBankStatementRequest.class)))
when(lsfxClient.getFileUploadStatus(eq(CALLER), any())).thenReturn(buildParsedSuccessStatusResponse());
when(lsfxClient.getBankStatement(eq(CALLER), any(GetBankStatementRequest.class)))
.thenReturn(buildEmptyBankStatementResponse());
CcdiFileUploadRecord record = buildRecord();
@@ -473,10 +499,11 @@ class CcdiFileUploadServiceImplTest {
"110101199001018888",
LsfxConstants.DATA_CHANNEL_JZL,
null,
null
null,
CALLER
);
verify(lsfxClient).fetchInnerFlow(argThat((FetchInnerFlowRequest request) ->
verify(lsfxClient).fetchInnerFlow(eq(CALLER), argThat((FetchInnerFlowRequest request) ->
LsfxConstants.DATA_CHANNEL_JZL.equals(request.getDataChannelCode())
&& Integer.valueOf(0).equals(request.getDataStartDateId())
&& Integer.valueOf(0).equals(request.getDataEndDateId())
@@ -485,21 +512,21 @@ class CcdiFileUploadServiceImplTest {
@Test
void processFileAsync_shouldUploadToLsfxWithOriginalRecordFileName() throws IOException {
when(lsfxClient.uploadFile(eq(LSFX_PROJECT_ID), any(), eq("原始流水.xlsx")))
when(lsfxClient.uploadFile(eq(CALLER), eq(LSFX_PROJECT_ID), any(), eq("原始流水.xlsx")))
.thenReturn(buildUploadResponse());
when(lsfxClient.checkParseStatus(LSFX_PROJECT_ID, String.valueOf(LOG_ID)))
when(lsfxClient.checkParseStatus(CALLER, LSFX_PROJECT_ID, String.valueOf(LOG_ID)))
.thenReturn(buildCheckParseStatusResponse(false));
when(lsfxClient.getFileUploadStatus(any())).thenReturn(buildParsedSuccessStatusResponse());
when(lsfxClient.getBankStatement(any(GetBankStatementRequest.class)))
when(lsfxClient.getFileUploadStatus(eq(CALLER), any())).thenReturn(buildParsedSuccessStatusResponse());
when(lsfxClient.getBankStatement(eq(CALLER), any(GetBankStatementRequest.class)))
.thenReturn(buildEmptyBankStatementResponse());
CcdiFileUploadRecord record = buildRecord();
record.setFileName("原始流水.xlsx");
Path tempFile = createTempFile();
service.processFileAsync(PROJECT_ID, LSFX_PROJECT_ID, tempFile.toString(), RECORD_ID, "batch-1", record);
service.processFileAsync(PROJECT_ID, LSFX_PROJECT_ID, tempFile.toString(), RECORD_ID, "batch-1", record, CALLER);
verify(lsfxClient).uploadFile(eq(LSFX_PROJECT_ID), argThat(file ->
verify(lsfxClient).uploadFile(eq(CALLER), eq(LSFX_PROJECT_ID), argThat(file ->
file.getName().startsWith("upload-") && file.getName().endsWith(".xlsx")
), eq("原始流水.xlsx"));
}
@@ -517,14 +544,15 @@ class CcdiFileUploadServiceImplTest {
project.setProjectId(PROJECT_ID);
when(projectMapper.selectById(PROJECT_ID)).thenReturn(project);
when(bankStatementMapper.countMatchedStaffCountByProjectId(PROJECT_ID)).thenReturn(1);
when(lsfxClient.uploadFile(eq(LSFX_PROJECT_ID), any(), eq("张三_330101199001010011_流水.xlsx")))
when(lsfxClient.uploadFile(eq(CALLER), eq(LSFX_PROJECT_ID), any(),
eq("张三_330101199001010011_流水.xlsx")))
.thenReturn(buildUploadResponse());
when(lsfxClient.checkParseStatus(LSFX_PROJECT_ID, String.valueOf(LOG_ID)))
when(lsfxClient.checkParseStatus(CALLER, LSFX_PROJECT_ID, String.valueOf(LOG_ID)))
.thenReturn(buildCheckParseStatusResponse(false));
when(lsfxClient.getFileUploadStatus(any())).thenReturn(buildParsedSuccessStatusResponse());
when(lsfxClient.getBankStatement(any(GetBankStatementRequest.class)))
when(lsfxClient.getFileUploadStatus(eq(CALLER), any())).thenReturn(buildParsedSuccessStatusResponse());
when(lsfxClient.getBankStatement(eq(CALLER), any(GetBankStatementRequest.class)))
.thenAnswer(invocation -> {
GetBankStatementRequest request = invocation.getArgument(0);
GetBankStatementRequest request = invocation.getArgument(1);
if (Integer.valueOf(1).equals(request.getPageSize())) {
return buildBankStatementCountResponse(1);
}
@@ -535,7 +563,8 @@ class CcdiFileUploadServiceImplTest {
record.setFileName("张三_330101199001010011_流水.xlsx");
Path tempFile = createTempFile();
service.processFileAsync(PROJECT_ID, LSFX_PROJECT_ID, tempFile.toString(), RECORD_ID, "batch-1", record);
service.processFileAsync(PROJECT_ID, LSFX_PROJECT_ID, tempFile.toString(), RECORD_ID, "batch-1", record,
CALLER);
assertNotNull(insertedStatements.get());
assertEquals(1, insertedStatements.get().size());
@@ -544,20 +573,20 @@ class CcdiFileUploadServiceImplTest {
@Test
void processFileAsync_shouldKeepOriginalFileNameWhenStatusReturnsDifferentName() throws IOException {
when(lsfxClient.uploadFile(eq(LSFX_PROJECT_ID), any(), org.mockito.ArgumentMatchers.anyString()))
when(lsfxClient.uploadFile(eq(CALLER), eq(LSFX_PROJECT_ID), any(), org.mockito.ArgumentMatchers.anyString()))
.thenReturn(buildUploadResponse());
when(lsfxClient.checkParseStatus(LSFX_PROJECT_ID, String.valueOf(LOG_ID)))
when(lsfxClient.checkParseStatus(CALLER, LSFX_PROJECT_ID, String.valueOf(LOG_ID)))
.thenReturn(buildCheckParseStatusResponse(false));
when(lsfxClient.getFileUploadStatus(any()))
when(lsfxClient.getFileUploadStatus(eq(CALLER), any()))
.thenReturn(buildParsedSuccessStatusResponse("平台返回文件名.xlsx"));
when(lsfxClient.getBankStatement(any(GetBankStatementRequest.class)))
when(lsfxClient.getBankStatement(eq(CALLER), any(GetBankStatementRequest.class)))
.thenReturn(buildEmptyBankStatementResponse());
CcdiFileUploadRecord record = buildRecord();
record.setFileName("原始流水.xlsx");
Path tempFile = createTempFile();
service.processFileAsync(PROJECT_ID, LSFX_PROJECT_ID, tempFile.toString(), RECORD_ID, "batch-1", record);
service.processFileAsync(PROJECT_ID, LSFX_PROJECT_ID, tempFile.toString(), RECORD_ID, "batch-1", record, CALLER);
verify(recordMapper, org.mockito.Mockito.atLeastOnce()).updateById(
org.mockito.ArgumentMatchers.<CcdiFileUploadRecord>argThat(item ->
@@ -573,17 +602,17 @@ class CcdiFileUploadServiceImplTest {
logItem.setStatus(-1);
logItem.setUploadStatusDesc("parse.failed");
when(lsfxClient.uploadFile(eq(LSFX_PROJECT_ID), any(), org.mockito.ArgumentMatchers.anyString()))
when(lsfxClient.uploadFile(eq(CALLER), eq(LSFX_PROJECT_ID), any(), org.mockito.ArgumentMatchers.anyString()))
.thenReturn(buildUploadResponse());
when(lsfxClient.checkParseStatus(LSFX_PROJECT_ID, String.valueOf(LOG_ID)))
when(lsfxClient.checkParseStatus(CALLER, LSFX_PROJECT_ID, String.valueOf(LOG_ID)))
.thenReturn(buildCheckParseStatusResponse(false));
when(lsfxClient.getFileUploadStatus(any())).thenReturn(statusResponse);
when(lsfxClient.getFileUploadStatus(eq(CALLER), any())).thenReturn(statusResponse);
CcdiFileUploadRecord record = buildRecord();
record.setFileName("原始流水.xlsx");
Path tempFile = createTempFile();
service.processFileAsync(PROJECT_ID, LSFX_PROJECT_ID, tempFile.toString(), RECORD_ID, "batch-1", record);
service.processFileAsync(PROJECT_ID, LSFX_PROJECT_ID, tempFile.toString(), RECORD_ID, "batch-1", record, CALLER);
verify(recordMapper, org.mockito.Mockito.atLeastOnce()).updateById(
org.mockito.ArgumentMatchers.<CcdiFileUploadRecord>argThat(item ->
@@ -610,7 +639,7 @@ class CcdiFileUploadServiceImplTest {
String result = service.deleteFileUploadRecord(RECORD_ID, 9527L);
assertEquals("删除成功已开始项目重新打标", result);
verify(lsfxClient, never()).deleteFiles(any());
verify(lsfxClient, never()).deleteFiles(any(), any());
verify(bankStatementMapper).deleteByProjectIdAndBatchId(PROJECT_ID, LOG_ID);
verify(recordMapper).updateById(org.mockito.ArgumentMatchers.<CcdiFileUploadRecord>argThat(item ->
RECORD_ID.equals(item.getId()) && "deleted".equals(item.getFileStatus())
@@ -644,7 +673,7 @@ class CcdiFileUploadServiceImplTest {
() -> service.deleteFileUploadRecord(RECORD_ID, 9527L));
assertTrue(exception.getMessage().contains("历史导入文件不支持删除"));
verify(lsfxClient, never()).deleteFiles(any());
verify(lsfxClient, never()).deleteFiles(any(), any());
}
@Test
@@ -659,7 +688,7 @@ class CcdiFileUploadServiceImplTest {
String result = service.deleteFileUploadRecord(RECORD_ID, 9527L);
assertEquals("删除成功已开始项目重新打标", result);
verify(lsfxClient, never()).deleteFiles(any());
verify(lsfxClient, never()).deleteFiles(any(), any());
verify(bankStatementMapper).deleteByProjectIdAndBatchId(PROJECT_ID, LOG_ID);
verify(recordMapper).updateById(org.mockito.ArgumentMatchers.<CcdiFileUploadRecord>argThat(item ->
"deleted".equals(item.getFileStatus())
@@ -669,7 +698,7 @@ class CcdiFileUploadServiceImplTest {
// @Test
// void processPullBankInfoAsync_shouldMarkParsedFailedWhenFetchInnerFlowThrows() {
// when(lsfxClient.fetchInnerFlow(any())).thenThrow(new RuntimeException("fetch inner flow failed"));
// when(lsfxClient.fetchInnerFlow(eq(CALLER), any())).thenThrow(new RuntimeException("fetch inner flow failed"));
//
// CcdiFileUploadRecord record = buildRecord();
// service.processPullBankInfoAsync(
@@ -692,19 +721,19 @@ class CcdiFileUploadServiceImplTest {
AtomicInteger sequence = new AtomicInteger();
captureRecordStatus(events, sequence);
when(lsfxClient.uploadFile(eq(LSFX_PROJECT_ID), any(), org.mockito.ArgumentMatchers.anyString()))
when(lsfxClient.uploadFile(eq(CALLER), eq(LSFX_PROJECT_ID), any(), org.mockito.ArgumentMatchers.anyString()))
.thenReturn(buildUploadResponse());
when(lsfxClient.checkParseStatus(LSFX_PROJECT_ID, String.valueOf(LOG_ID)))
when(lsfxClient.checkParseStatus(CALLER, LSFX_PROJECT_ID, String.valueOf(LOG_ID)))
.thenReturn(buildCheckParseStatusResponse(false));
when(lsfxClient.getFileUploadStatus(any())).thenReturn(buildParsedSuccessStatusResponse());
when(lsfxClient.getBankStatement(any(GetBankStatementRequest.class)))
when(lsfxClient.getFileUploadStatus(eq(CALLER), any())).thenReturn(buildParsedSuccessStatusResponse());
when(lsfxClient.getBankStatement(eq(CALLER), any(GetBankStatementRequest.class)))
.thenReturn(buildBankStatementResponseWithTotalCount(1))
.thenThrow(new RuntimeException("paged fetch failed"));
CcdiFileUploadRecord record = buildRecord();
Path tempFile = createTempFile();
service.processFileAsync(PROJECT_ID, LSFX_PROJECT_ID, tempFile.toString(), RECORD_ID, "batch-1", record);
service.processFileAsync(PROJECT_ID, LSFX_PROJECT_ID, tempFile.toString(), RECORD_ID, "batch-1", record, CALLER);
assertTrue(events.stream().anyMatch(event -> event.endsWith("record:parsed_failed")));
assertFalse(events.stream().anyMatch(event -> event.endsWith("record:parsed_success")));
@@ -716,18 +745,18 @@ class CcdiFileUploadServiceImplTest {
List<CcdiFileUploadRecord> updates = new ArrayList<>();
captureUpdatedRecords(updates);
when(lsfxClient.uploadFile(eq(LSFX_PROJECT_ID), any(), org.mockito.ArgumentMatchers.anyString()))
when(lsfxClient.uploadFile(eq(CALLER), eq(LSFX_PROJECT_ID), any(), org.mockito.ArgumentMatchers.anyString()))
.thenReturn(buildUploadResponse());
when(lsfxClient.checkParseStatus(LSFX_PROJECT_ID, String.valueOf(LOG_ID)))
when(lsfxClient.checkParseStatus(CALLER, LSFX_PROJECT_ID, String.valueOf(LOG_ID)))
.thenReturn(buildCheckParseStatusResponse(false));
when(lsfxClient.getFileUploadStatus(any())).thenReturn(buildParsedSuccessStatusResponse());
when(lsfxClient.getBankStatement(any(GetBankStatementRequest.class)))
when(lsfxClient.getFileUploadStatus(eq(CALLER), any())).thenReturn(buildParsedSuccessStatusResponse());
when(lsfxClient.getBankStatement(eq(CALLER), any(GetBankStatementRequest.class)))
.thenThrow(new RuntimeException("bank statement fetch failed:" + "x".repeat(3000)));
CcdiFileUploadRecord record = buildRecord();
Path tempFile = createTempFile();
service.processFileAsync(PROJECT_ID, LSFX_PROJECT_ID, tempFile.toString(), RECORD_ID, "batch-1", record);
service.processFileAsync(PROJECT_ID, LSFX_PROJECT_ID, tempFile.toString(), RECORD_ID, "batch-1", record, CALLER);
CcdiFileUploadRecord failedRecord = findLastUpdatedRecordByStatus(updates, "parsed_failed");
assertTrue(failedRecord.getErrorMessage().length() <= MAX_ERROR_MESSAGE_LENGTH);
@@ -738,13 +767,13 @@ class CcdiFileUploadServiceImplTest {
List<CcdiFileUploadRecord> updates = new ArrayList<>();
captureUpdatedRecords(updates);
when(lsfxClient.uploadFile(eq(LSFX_PROJECT_ID), any(), org.mockito.ArgumentMatchers.anyString()))
when(lsfxClient.uploadFile(eq(CALLER), eq(LSFX_PROJECT_ID), any(), org.mockito.ArgumentMatchers.anyString()))
.thenThrow(new RuntimeException("upload failed:" + "x".repeat(3000)));
CcdiFileUploadRecord record = buildRecord();
Path tempFile = createTempFile();
service.processFileAsync(PROJECT_ID, LSFX_PROJECT_ID, tempFile.toString(), RECORD_ID, "batch-1", record);
service.processFileAsync(PROJECT_ID, LSFX_PROJECT_ID, tempFile.toString(), RECORD_ID, "batch-1", record, CALLER);
CcdiFileUploadRecord failedRecord = findLastUpdatedRecordByStatus(updates, "parsed_failed");
assertTrue(failedRecord.getErrorMessage().length() <= MAX_ERROR_MESSAGE_LENGTH);
@@ -752,19 +781,19 @@ class CcdiFileUploadServiceImplTest {
@Test
void fetchAndSaveBankStatements_shouldTrimLeAccountNoBeforeInsert() throws IOException {
when(lsfxClient.uploadFile(eq(LSFX_PROJECT_ID), any(), org.mockito.ArgumentMatchers.anyString()))
when(lsfxClient.uploadFile(eq(CALLER), eq(LSFX_PROJECT_ID), any(), org.mockito.ArgumentMatchers.anyString()))
.thenReturn(buildUploadResponse());
when(lsfxClient.checkParseStatus(LSFX_PROJECT_ID, String.valueOf(LOG_ID)))
when(lsfxClient.checkParseStatus(CALLER, LSFX_PROJECT_ID, String.valueOf(LOG_ID)))
.thenReturn(buildCheckParseStatusResponse(false));
when(lsfxClient.getFileUploadStatus(any())).thenReturn(buildParsedSuccessStatusResponse());
when(lsfxClient.getBankStatement(any(GetBankStatementRequest.class)))
when(lsfxClient.getFileUploadStatus(eq(CALLER), any())).thenReturn(buildParsedSuccessStatusResponse());
when(lsfxClient.getBankStatement(eq(CALLER), any(GetBankStatementRequest.class)))
.thenReturn(buildBankStatementResponseWithItems(1, List.of(buildBankStatementItem(" 62220001 "))))
.thenReturn(buildBankStatementResponseWithItems(1, List.of(buildBankStatementItem(" 62220001 "))));
CcdiFileUploadRecord record = buildRecord();
Path tempFile = createTempFile();
service.processFileAsync(PROJECT_ID, LSFX_PROJECT_ID, tempFile.toString(), RECORD_ID, "batch-1", record);
service.processFileAsync(PROJECT_ID, LSFX_PROJECT_ID, tempFile.toString(), RECORD_ID, "batch-1", record, CALLER);
verify(bankStatementMapper).insertBatch(any());
verify(bankStatementMapper).insertBatch(org.mockito.ArgumentMatchers.argThat(list ->
@@ -773,7 +802,7 @@ class CcdiFileUploadServiceImplTest {
@Test
void fetchAndSaveBankStatements_shouldLogConservativeCountsWhenAffectedRowsAreAmbiguous() {
when(lsfxClient.getBankStatement(any(GetBankStatementRequest.class)))
when(lsfxClient.getBankStatement(eq(CALLER), any(GetBankStatementRequest.class)))
.thenReturn(buildBankStatementResponseWithItems(1, List.of(buildBankStatementItem("62220001"))))
.thenReturn(buildBankStatementResponseWithItems(1, List.of(buildBankStatementItem("62220001"))));
when(bankStatementMapper.insertBatch(any())).thenReturn(1);
@@ -787,6 +816,7 @@ class CcdiFileUploadServiceImplTest {
Object result = ReflectionTestUtils.invokeMethod(
service,
"fetchAndSaveBankStatements",
CALLER,
PROJECT_ID,
LSFX_PROJECT_ID,
LOG_ID,
@@ -825,12 +855,12 @@ class CcdiFileUploadServiceImplTest {
AtomicInteger sequence = new AtomicInteger();
captureRecordStatus(events, sequence);
when(lsfxClient.uploadFile(eq(LSFX_PROJECT_ID), any(), org.mockito.ArgumentMatchers.anyString()))
when(lsfxClient.uploadFile(eq(CALLER), eq(LSFX_PROJECT_ID), any(), org.mockito.ArgumentMatchers.anyString()))
.thenReturn(buildUploadResponse());
when(lsfxClient.checkParseStatus(LSFX_PROJECT_ID, String.valueOf(LOG_ID)))
when(lsfxClient.checkParseStatus(CALLER, LSFX_PROJECT_ID, String.valueOf(LOG_ID)))
.thenReturn(buildCheckParseStatusResponse(false));
when(lsfxClient.getFileUploadStatus(any())).thenReturn(buildParsedSuccessStatusResponse());
when(lsfxClient.getBankStatement(any(GetBankStatementRequest.class)))
when(lsfxClient.getFileUploadStatus(eq(CALLER), any())).thenReturn(buildParsedSuccessStatusResponse());
when(lsfxClient.getBankStatement(eq(CALLER), any(GetBankStatementRequest.class)))
.thenReturn(buildBankStatementResponseWithItems(1, List.of(buildBankStatementItem("62220001"))))
.thenReturn(buildBankStatementResponseWithItems(1, List.of(buildBankStatementItem("62220001"))));
when(bankStatementMapper.insertBatch(any()))
@@ -839,7 +869,7 @@ class CcdiFileUploadServiceImplTest {
CcdiFileUploadRecord record = buildRecord();
Path tempFile = createTempFile();
service.processFileAsync(PROJECT_ID, LSFX_PROJECT_ID, tempFile.toString(), RECORD_ID, "batch-1", record);
service.processFileAsync(PROJECT_ID, LSFX_PROJECT_ID, tempFile.toString(), RECORD_ID, "batch-1", record, CALLER);
assertTrue(events.stream().anyMatch(event -> event.endsWith("record:parsed_failed")));
assertFalse(events.stream().anyMatch(event -> event.endsWith("record:parsed_success")));
@@ -1026,10 +1056,17 @@ class CcdiFileUploadServiceImplTest {
private void invokeSubmitTasksAsync(List<String> tempFilePaths,
List<CcdiFileUploadRecord> records,
String batchId) throws Exception {
invokeSubmitTasksAsync(tempFilePaths, records, batchId, CALLER);
}
private void invokeSubmitTasksAsync(List<String> tempFilePaths,
List<CcdiFileUploadRecord> records,
String batchId,
CallerContext caller) throws Exception {
Method method = CcdiFileUploadServiceImpl.class.getDeclaredMethod("submitTasksAsync",
Long.class, Integer.class, List.class, List.class, String.class);
Long.class, Integer.class, List.class, List.class, String.class, CallerContext.class);
method.setAccessible(true);
method.invoke(service, PROJECT_ID, LSFX_PROJECT_ID, tempFilePaths, records, batchId);
method.invoke(service, PROJECT_ID, LSFX_PROJECT_ID, tempFilePaths, records, batchId, caller);
}
private void setField(String fieldName, Object value) throws Exception {

View File

@@ -18,6 +18,7 @@ import com.ruoyi.ccdi.project.mapper.CcdiProjectMapper;
import com.ruoyi.ccdi.project.service.CcdiProjectAccessService;
import com.ruoyi.common.exception.ServiceException;
import com.ruoyi.lsfx.client.LsfxAnalysisClient;
import com.ruoyi.lsfx.domain.CallerContext;
import com.ruoyi.lsfx.domain.response.GetTokenResponse;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
@@ -48,6 +49,8 @@ import static org.mockito.Mockito.when;
@ExtendWith(MockitoExtension.class)
class CcdiProjectServiceImplTest {
private static final CallerContext CALLER = CallerContext.of(7L, "tester");
@InjectMocks
private CcdiProjectServiceImpl service;
@@ -282,7 +285,7 @@ class CcdiProjectServiceImplTest {
dto.setStartDate("2026-01-01");
dto.setEndDate("2026-01-31");
when(lsfxAnalysisClient.getToken(any())).thenReturn(buildTokenResponse(3001));
when(lsfxAnalysisClient.getToken(any(CallerContext.class), any())).thenReturn(buildTokenResponse(3001));
doAnswer(invocation -> {
CcdiProject project = invocation.getArgument(0);
project.setProjectId(90L);
@@ -291,7 +294,7 @@ class CcdiProjectServiceImplTest {
TransactionSynchronizationManager.initSynchronization();
try {
CcdiProjectVO project = service.importFromHistory(dto, "tester");
CcdiProjectVO project = service.importFromHistory(dto, CALLER);
assertNotNull(project);
assertEquals(90L, project.getProjectId());
@@ -320,7 +323,7 @@ class CcdiProjectServiceImplTest {
dto.setDescription("测试项目");
dto.setConfigType("default");
when(lsfxAnalysisClient.getToken(any())).thenReturn(buildTokenResponse(2001));
when(lsfxAnalysisClient.getToken(any(CallerContext.class), any())).thenReturn(buildTokenResponse(2001));
doAnswer(invocation -> {
CcdiProject project = invocation.getArgument(0);
project.setProjectId(88L);
@@ -333,7 +336,7 @@ class CcdiProjectServiceImplTest {
logger.addAppender(logAppender);
try {
service.createProject(dto);
service.createProject(dto, CALLER);
assertTrue(logAppender.list.stream().map(ILoggingEvent::getFormattedMessage)
.anyMatch(message -> message.contains("项目状态初始化")