后台管理接口修复 后端BUG修复

This commit is contained in:
super
2026-05-12 21:03:42 +08:00
parent 5c990e651e
commit 0d78d63437
30 changed files with 3444 additions and 394 deletions
@@ -355,6 +355,9 @@ public class AppearancePatentCozeClient {
text(firstNonNull(node.get("result"),
firstNonNull(node.get("conclusion"),
firstNonNull(itemNode.get("result"), itemNode.get("conclusion"))))),
text(firstNonNull(
firstNonNull(node.get("status"), firstNonNull(node.get("row_status"), node.get("rowStatus"))),
firstNonNull(itemNode.get("status"), firstNonNull(itemNode.get("row_status"), itemNode.get("rowStatus"))))),
text(firstNonNull(node.get("title_reason"),
firstNonNull(node.get("titleReason"),
firstNonNull(itemNode.get("title_reason"), itemNode.get("titleReason"))))),
@@ -465,6 +468,7 @@ public class AppearancePatentCozeClient {
row.setAppearanceRisk(result.appearance());
row.setPatentRisk(result.patent());
row.setConclusion(result.result());
row.setStatus(result.status());
row.setTitleReason(result.titleReason());
row.setAppearanceReason(result.appearanceReason());
row.setPatentReason(result.patentReason());
@@ -490,6 +494,7 @@ public class AppearancePatentCozeClient {
row.setTitle(source.getTitle());
row.setError(source.getError());
row.setDone(source.getDone());
row.setStatus(source.getStatus());
row.setTitleRisk(source.getTitleRisk());
row.setAppearanceRisk(source.getAppearanceRisk());
row.setPatentRisk(source.getPatentRisk());
@@ -942,6 +947,7 @@ public class AppearancePatentCozeClient {
String appearance,
String patent,
String result,
String status,
String titleReason,
String appearanceReason,
String patentReason
@@ -48,6 +48,10 @@ public class AppearancePatentResultRowDto {
@Schema(description = "单行完成标记。当前主要使用请求体顶层 done 控制任务收尾,该字段仅作兼容。", example = "true")
private Boolean done;
@JsonAlias({"row_status", "rowStatus", "Status"})
@Schema(description = "Coze row status", example = "success", accessMode = Schema.AccessMode.READ_ONLY)
private String status;
@Schema(description = "Java 调用 Coze 后生成的标题维度检测结果,对应最终 xlsx 的“标题维度(商标)”列。Python 回传请求中不要传该字段;即使传入,后端也会以 Java/Coze 处理结果为准。", example = "标题未发现明显商标侵权风险。", accessMode = Schema.AccessMode.READ_ONLY)
private String titleRisk;
@@ -118,7 +118,8 @@ public class AppearancePatentTaskService {
"标题维度(商标)",
"外观维度(外观设计专利)",
"专利维度(发明/实用新型专利)",
"结论"
"结论",
"status"
);
private final LocalFileStorageService localFileStorageService;
@@ -424,7 +425,7 @@ public class AppearancePatentTaskService {
completeSubmittedChunk(context);
return null;
});
submitCozeForSubmittedChunk(context);
scheduleCozePipelineForSubmittedChunk(context);
return;
}
FileTaskEntity task = fileTaskMapper.selectById(taskId);
@@ -510,7 +511,7 @@ public class AppearancePatentTaskService {
task.setUpdatedAt(LocalDateTime.now());
fileTaskMapper.updateById(task);
}
submitCozeForSubmittedChunk(context);
scheduleCozePipelineForSubmittedChunk(context);
}
@Transactional
@@ -560,14 +561,9 @@ public class AppearancePatentTaskService {
if (transactionManager != null) {
LocalDateTime threshold = LocalDateTime.now().minusMinutes(Math.max(1, properties.getStaleTimeoutMinutes()));
long thresholdMillis = threshold.atZone(ZoneId.systemDefault()).toInstant().toEpochMilli();
List<FileTaskEntity> tasks = fileTaskMapper.selectList(new LambdaQueryWrapper<FileTaskEntity>()
.eq(FileTaskEntity::getModuleType, MODULE_TYPE)
.eq(FileTaskEntity::getStatus, STATUS_RUNNING)
.lt(FileTaskEntity::getUpdatedAt, threshold)
.last("limit 50"));
List<FileTaskEntity> tasks = listStaleFinalizeCandidates(threshold);
for (FileTaskEntity task : tasks) {
if (isJavaSideProcessing(task.getId())) {
touchJavaSideTaskActivity(task.getId());
if (!isOwnerCurrent(ownerFromTask(task))) {
continue;
}
long heartbeatMillis = taskCacheService.getTaskHeartbeatMillis(task.getId());
@@ -589,14 +585,9 @@ public class AppearancePatentTaskService {
}
LocalDateTime threshold = LocalDateTime.now().minusMinutes(Math.max(1, properties.getStaleTimeoutMinutes()));
long thresholdMillis = threshold.atZone(ZoneId.systemDefault()).toInstant().toEpochMilli();
List<FileTaskEntity> tasks = fileTaskMapper.selectList(new LambdaQueryWrapper<FileTaskEntity>()
.eq(FileTaskEntity::getModuleType, MODULE_TYPE)
.eq(FileTaskEntity::getStatus, STATUS_RUNNING)
.lt(FileTaskEntity::getUpdatedAt, threshold)
.last("limit 50"));
List<FileTaskEntity> tasks = listStaleFinalizeCandidates(threshold);
for (FileTaskEntity task : tasks) {
if (isJavaSideProcessing(task.getId())) {
touchJavaSideTaskActivity(task.getId());
if (!isOwnerCurrent(ownerFromTask(task))) {
continue;
}
long heartbeatMillis = taskCacheService.getTaskHeartbeatMillis(task.getId());
@@ -608,14 +599,46 @@ public class AppearancePatentTaskService {
continue;
}
try (lockHandle) {
FileTaskEntity latestTask = fileTaskMapper.selectById(task.getId());
if (latestTask != null) {
finalizeTask(latestTask, "Python interrupted before uploading final appearance patent result", allRowCount(latestTask), true);
}
finalizeStaleTask(task.getId(), "Python interrupted before uploading final appearance patent result");
}
}
}
public void debugFinalizeStaleTask(Long taskId) {
if (taskId == null || taskId <= 0) {
return;
}
if (transactionManager != null) {
TaskDistributedLockService.LockHandle lockHandle = acquireTaskLock(taskId, 0L);
if (lockHandle == null) {
return;
}
try (lockHandle) {
inNewTransaction(() -> {
finalizeStaleTask(taskId, "Python interrupted before uploading final appearance patent result");
return null;
});
}
return;
}
TaskDistributedLockService.LockHandle lockHandle = acquireTaskLock(taskId, 0L);
if (lockHandle == null) {
return;
}
try (lockHandle) {
finalizeStaleTask(taskId, "Python interrupted before uploading final appearance patent result");
}
}
private List<FileTaskEntity> listStaleFinalizeCandidates(LocalDateTime threshold) {
return fileTaskMapper.selectList(new LambdaQueryWrapper<FileTaskEntity>()
.eq(FileTaskEntity::getModuleType, MODULE_TYPE)
.eq(FileTaskEntity::getStatus, STATUS_RUNNING)
.lt(FileTaskEntity::getCreatedAt, threshold)
.orderByAsc(FileTaskEntity::getCreatedAt)
.last("limit 200"));
}
private SubmitContext persistSubmittedChunk(Long taskId, AppearancePatentSubmitResultRequest request) {
FileTaskEntity task = fileTaskMapper.selectById(taskId);
if (task == null || !MODULE_TYPE.equals(task.getModuleType())) {
@@ -699,17 +722,24 @@ public class AppearancePatentTaskService {
}
List<TaskChunkEntity> chunks = loadSubmittedChunks(task.getId());
if (chunks.isEmpty()) {
log.info("[appearance-patent] skip stale recovery coze submission because no submitted chunks remain taskId={}",
task.getId());
return;
}
FileResultEntity result = findOrCreateResultRecordForAssembly(task, allRowCount(task));
if (result == null) {
log.warn("[appearance-patent] stale recovery could not create result record taskId={}", task.getId());
return;
}
TaskFileJobEntity job = taskFileJobService.enqueueAssembleResult(
task.getId(), MODULE_TYPE, result.getId(), buildTaskOwnerScopeKey(task.getId()));
task.getId(), MODULE_TYPE, result.getId(), buildTaskOwnerScopeKey(task));
if (job == null || "SUCCESS".equals(job.getStatus())) {
log.info("[appearance-patent] stale recovery skipped coze submission because result job unavailable taskId={} resultId={} jobStatus={}",
task.getId(), result.getId(), job == null ? null : job.getStatus());
return;
}
log.info("[appearance-patent] stale recovery created assemble job taskId={} resultId={} jobId={} jobStatus={} chunkCount={}",
task.getId(), result.getId(), job.getId(), job.getStatus(), chunks.size());
Map<String, List<AppearancePatentParsedRowVo>> allRowsByBaseId = loadAllRowsByBaseId(task);
boolean pendingCoze = submitCozeBatches(task, result, job, chunks, allRowsByBaseId);
saveCozePipelineProgress(task, job);
@@ -722,14 +752,127 @@ public class AppearancePatentTaskService {
}
}
private void scheduleCozePipelineForSubmittedChunk(SubmitContext context) {
if (context == null || context.task() == null || context.task().getId() == null) {
return;
}
FileTaskEntity task = fileTaskMapper.selectById(context.task().getId());
if (task == null || !MODULE_TYPE.equals(task.getModuleType()) || STATUS_SUCCESS.equals(task.getStatus())) {
return;
}
List<TaskChunkEntity> chunks = loadSubmittedChunks(task.getId());
if (chunks.isEmpty()) {
return;
}
FileResultEntity result = findOrCreateResultRecordForAssembly(task, allRowCount(task));
if (result == null) {
return;
}
TaskFileJobEntity job = taskFileJobService.enqueueAssembleResult(
task.getId(), MODULE_TYPE, result.getId(), buildTaskOwnerScopeKey(task));
if (job == null || "SUCCESS".equals(job.getStatus())) {
return;
}
taskFileJobService.requeue(job.getId(), "Appearance patent result uploaded, scheduling Coze/file assembly");
touchJavaSideTaskActivity(task.getId());
}
private void finalizeStaleTask(Long taskId, String error) {
FileTaskEntity task = fileTaskMapper.selectById(taskId);
if (task == null || !MODULE_TYPE.equals(task.getModuleType()) || !STATUS_RUNNING.equals(task.getStatus())) {
return;
}
if (tryRecoverTimedOutPythonTask(task)) {
return;
}
finalizeTask(task, error, allRowCount(task), true);
}
private boolean tryRecoverTimedOutPythonTask(FileTaskEntity task) {
if (task == null || task.getId() == null) {
return false;
}
Long taskId = task.getId();
boolean uploadComplete = isResultSubmissionComplete(taskId);
long pendingCozeStates = countPendingCozeStates(taskId);
long activeAssembleJobs = taskFileJobService.countActiveAssembleJobs(taskId, MODULE_TYPE);
log.info("[appearance-patent] stale recovery probe taskId={} uploadComplete={} pendingCozeStates={} activeAssembleJobs={} persistedRows={}",
taskId, uploadComplete, pendingCozeStates, activeAssembleJobs, hasPersistedResultRows(taskId));
if (uploadComplete && (pendingCozeStates > 0 || activeAssembleJobs > 0)) {
touchJavaSideTaskActivity(taskId);
return true;
}
if (!hasPersistedResultRows(taskId)) {
log.info("[appearance-patent] stale recovery aborted because no persisted rows taskId={}", taskId);
return false;
}
if (!uploadComplete) {
int forcedScopes = markSubmissionCompleteOnPythonTimeout(taskId);
if (forcedScopes <= 0) {
log.warn("[appearance-patent] stale recovery aborted because submission could not be marked complete taskId={}",
taskId);
return false;
}
int pendingRows = collectPendingCozeCandidates(task, loadSubmittedChunks(taskId)).size();
log.warn("[appearance-patent] python heartbeat timed out, forcing coze flush taskId={} pendingRows={} pendingCozeStates={} activeAssembleJobs={} forcedScopes={}",
taskId, pendingRows, pendingCozeStates, activeAssembleJobs, forcedScopes);
} else {
log.warn("[appearance-patent] stale running task resuming coze/file assembly after python timeout taskId={} pendingCozeStates={} activeAssembleJobs={}",
taskId, pendingCozeStates, activeAssembleJobs);
}
submitCozeForSubmittedChunk(new SubmitContext(task, null, null, 0, true, null));
touchJavaSideTaskActivity(taskId);
return true;
}
private int markSubmissionCompleteOnPythonTimeout(Long taskId) {
if (taskId == null || taskId <= 0) {
return 0;
}
List<TaskScopeStateEntity> inputStates = taskScopeStateMapper.selectList(new LambdaQueryWrapper<TaskScopeStateEntity>()
.eq(TaskScopeStateEntity::getTaskId, taskId)
.eq(TaskScopeStateEntity::getModuleType, MODULE_TYPE)
.isNull(TaskScopeStateEntity::getCozeStatus)
.orderByDesc(TaskScopeStateEntity::getUpdatedAt));
if (inputStates == null || inputStates.isEmpty()) {
TaskChunkEntity latestChunk = taskChunkMapper.selectOne(new LambdaQueryWrapper<TaskChunkEntity>()
.eq(TaskChunkEntity::getTaskId, taskId)
.eq(TaskChunkEntity::getModuleType, MODULE_TYPE)
.orderByDesc(TaskChunkEntity::getUpdatedAt)
.last("limit 1"));
if (latestChunk == null || latestChunk.getScopeHash() == null || latestChunk.getScopeHash().isBlank()) {
log.warn("[appearance-patent] stale recovery found no latest chunk to complete taskId={}", taskId);
return 0;
}
log.info("[appearance-patent] stale recovery recreating input scope from chunk taskId={} scopeKey={} scopeHash={} chunkTotal={}",
taskId, latestChunk.getScopeKey(), latestChunk.getScopeHash(), latestChunk.getChunkTotal());
upsertScopeState(taskId,
firstNonBlank(latestChunk.getScopeKey(), "task:" + taskId),
latestChunk.getScopeHash(),
latestChunk.getChunkTotal(),
null,
true,
false);
return 1;
}
LocalDateTime now = LocalDateTime.now();
int updated = 0;
for (TaskScopeStateEntity state : inputStates) {
if (state == null || state.getId() == null) {
continue;
}
state.setCompleted(1);
if (state.getLastChunkAt() == null) {
state.setLastChunkAt(now);
}
state.setUpdatedAt(now);
state.setStateJson("{\"phase\":\"RECEIVED\",\"coze\":\"PENDING\"}");
taskScopeStateMapper.updateById(state);
updated++;
}
return updated;
}
private void upsertScopeState(Long taskId,
String scopeKey,
String scopeHash,
@@ -765,6 +908,8 @@ public class AppearancePatentTaskService {
if (scope.getId() == null) {
try {
taskScopeStateMapper.insert(scope);
log.info("[appearance-patent] scope state inserted taskId={} scope={} scopeHash={} completed={} cozeDone={}",
taskId, scopeKey, scopeHash, completed, cozeDone);
return;
} catch (DuplicateKeyException ex) {
log.info("[appearance-patent] duplicate scope state inserted concurrently taskId={} scope={}", taskId, scopeKey);
@@ -792,6 +937,8 @@ public class AppearancePatentTaskService {
}
}
taskScopeStateMapper.updateById(scope);
log.info("[appearance-patent] scope state updated taskId={} scope={} scopeHash={} completed={} cozeDone={}",
taskId, scopeKey, scopeHash, completed, cozeDone);
}
private <T> T inNewTransaction(Supplier<T> action) {
@@ -1085,6 +1232,7 @@ public class AppearancePatentTaskService {
row.setUrl(firstNonBlank(representative.getUrl(), sibling.getUrl()));
row.setTitle(firstNonBlank(representative.getTitle(), sibling.getTitle()));
row.setError(representative.getError());
row.setStatus(representative.getStatus());
row.setTitleRisk(representative.getTitleRisk());
row.setAppearanceRisk(representative.getAppearanceRisk());
row.setPatentRisk(representative.getPatentRisk());
@@ -1185,7 +1333,7 @@ public class AppearancePatentTaskService {
return null;
}
TaskFileJobEntity job = taskFileJobService.enqueueAssembleResult(
task.getId(), MODULE_TYPE, result.getId(), buildTaskOwnerScopeKey(task.getId()));
task.getId(), MODULE_TYPE, result.getId(), buildTaskOwnerScopeKey(task));
if (dispatchWhenIdle
&& job != null
&& "RUNNING".equals(job.getStatus())
@@ -1220,11 +1368,6 @@ public class AppearancePatentTaskService {
if (job == null || job.getTaskId() == null || job.getResultId() == null) {
throw new BusinessException("result file job arguments are incomplete");
}
if (!isJobOwnedByCurrentInstance(job)) {
log.info("[appearance-patent] skip result file job because owner is another instance jobId={} taskId={} owner={} current={}",
job.getId(), job.getTaskId(), ownerFromScopeKey(job.getScopeKey()), currentInstanceId());
return false;
}
TaskDistributedLockService.LockHandle lockHandle = acquireTaskLock(job.getTaskId(), TASK_LOCK_WAIT_MILLIS);
if (lockHandle == null) {
taskFileJobService.requeue(job.getId(), "Task is busy, waiting for appearance patent result merge");
@@ -1241,6 +1384,13 @@ public class AppearancePatentTaskService {
if (task == null || !MODULE_TYPE.equals(task.getModuleType())) {
throw new BusinessException("task not found");
}
ensureTaskOwnedByCurrentInstance(task, "assemble result file");
job.setScopeKey(buildTaskOwnerScopeKey(task));
if (!isJobOwnedByCurrentInstance(job)) {
log.info("[appearance-patent] skip result file job because owner is another instance jobId={} taskId={} owner={} current={}",
job.getId(), job.getTaskId(), ownerFromScopeKey(job.getScopeKey()), currentInstanceId());
return false;
}
FileResultEntity result = fileResultMapper.selectById(job.getResultId());
if (result == null || !MODULE_TYPE.equals(result.getModuleType())) {
throw new BusinessException("result record not found");
@@ -1250,7 +1400,8 @@ public class AppearancePatentTaskService {
.eq(TaskChunkEntity::getModuleType, MODULE_TYPE)
.orderByAsc(TaskChunkEntity::getChunkIndex));
int cozeWorkUnits = countCozeWorkUnits(chunks, Math.max(1, properties.getCozeBatchSize()));
int totalProgressUnits = Math.max(3, cozeWorkUnits + 3);
int plannedCozeUnits = Math.max(cozeWorkUnits, countAllCozeStates(task.getId()));
int totalProgressUnits = Math.max(3, plannedCozeUnits + 3);
if (countPendingCozeStates(task.getId()) > 0) {
taskFileJobService.touchRunning(job.getId());
touchJavaSideTaskActivity(task.getId());
@@ -1269,10 +1420,10 @@ public class AppearancePatentTaskService {
if (STATUS_RUNNING.equals(task.getStatus()) && !isResultSubmissionComplete(task.getId())) {
taskFileJobService.touchRunning(job.getId());
touchJavaSideTaskActivity(task.getId());
saveFileBuildProgress(task, job, totalProgressUnits, Math.max(1, cozeWorkUnits), "等待 Python 继续回传数据");
saveFileBuildProgress(task, job, totalProgressUnits, Math.max(1, plannedCozeUnits), "等待 Python 继续回传数据");
return false;
}
completeCozeFileJob(task, result, job, totalProgressUnits, cozeWorkUnits);
completeCozeFileJob(task, result, job, totalProgressUnits);
return true;
}
@@ -1555,10 +1706,10 @@ public class AppearancePatentTaskService {
markCozeStateTerminal(state, COZE_STATUS_FAILED, "Coze batch context missing");
return;
}
if (!isOwnerCurrent(context.ownerInstanceId())) {
log.info("[appearance-patent] coze poll skipped after context refresh because owner is another instance taskId={} stateId={} owner={} current={}",
state.getTaskId(), state.getId(), context.ownerInstanceId(), currentInstanceId());
return;
if (!isOwnerCurrent(context.ownerInstanceId())) {
log.info("[appearance-patent] coze poll skipped after context refresh because owner is another instance taskId={} stateId={} owner={} current={}",
state.getTaskId(), state.getId(), context.ownerInstanceId(), currentInstanceId());
return;
}
taskFileJobService.touchRunning(context.jobId());
log.info("[appearance-patent] coze poll start taskId={} stateId={} executeId={} jobId={} chunk={} batch={}/{}",
@@ -1924,6 +2075,9 @@ public class AppearancePatentTaskService {
if (taskId == null || context == null || countPendingCozeStates(taskId) > 0) {
return;
}
if (!isOwnerCurrent(context.ownerInstanceId())) {
return;
}
TaskDistributedLockService.LockHandle taskLockHandle = acquireTaskLock(taskId, 0L);
if (taskLockHandle == null) {
return;
@@ -1974,12 +2128,11 @@ public class AppearancePatentTaskService {
private void completeCozeFileJob(FileTaskEntity task,
FileResultEntity result,
TaskFileJobEntity job,
int totalProgressUnits,
int cozeWorkUnits) {
int totalProgressUnits) {
if (countPendingCozeStates(task.getId()) > 0) {
throw new BusinessException("Coze 结果仍在处理中,暂不能生成结果文件");
}
int assembleProgress = Math.max(1, Math.min(totalProgressUnits - 2, cozeWorkUnits));
int assembleProgress = Math.max(1, totalProgressUnits - 2);
saveFileBuildProgress(task, job, totalProgressUnits, assembleProgress, "正在组装 xlsx");
assembleResultWorkbook(task, result);
saveFileBuildProgress(task, job, totalProgressUnits, totalProgressUnits - 1, "正在上传结果文件");
@@ -2216,8 +2369,9 @@ public class AppearancePatentTaskService {
return "coze:task:" + taskId + ":rows:" + DigestUtil.sha256Hex(rowKeys.toString());
}
private String buildTaskOwnerScopeKey(Long taskId) {
return "task:" + taskId + ":owner:" + currentInstanceId();
private String buildTaskOwnerScopeKey(FileTaskEntity task) {
Long taskId = task == null ? null : task.getId();
return "task:" + taskId + ":owner:" + firstNonBlank(ownerFromTask(task), currentInstanceId());
}
private String currentInstanceId() {
@@ -2577,7 +2731,8 @@ public class AppearancePatentTaskService {
row.createCell(col++).setCellValue(resultRow == null ? missingReason : userFacingCozeCellValue(resultRow, resultRow.getTitleRisk()));
row.createCell(col++).setCellValue(resultRow == null ? missingReason : userFacingCozeCellValue(resultRow, resultRow.getAppearanceRisk()));
row.createCell(col++).setCellValue(resultRow == null ? missingReason : userFacingCozeCellValue(resultRow, resultRow.getPatentRisk()));
row.createCell(col).setCellValue(resultRow == null ? "未送检" : userFacingConclusion(resultRow));
row.createCell(col++).setCellValue(resultRow == null ? "未送检" : userFacingConclusion(resultRow));
row.createCell(col).setCellValue(resultRow == null ? "" : firstNonBlank(resultRow.getStatus(), ""));
}
writeReasonSheet(workbook, headerStyle, rowsToWrite, resultMap);
workbook.write(fos);
@@ -2962,7 +3117,7 @@ public class AppearancePatentTaskService {
}
vo.setFileJobId(job.getId());
vo.setFileStatus(job.getStatus());
vo.setFileError(job.getErrorMessage());
vo.setFileError(STATUS_FAILED.equals(job.getStatus()) ? firstNonBlank(job.getErrorMessage(), null) : null);
attachFileProgress(vo, row, job);
}