From 652e48884b581830f8dcdc140886dc10c4554124 Mon Sep 17 00:00:00 2001 From: key Date: Wed, 29 Jul 2026 15:09:01 +0800 Subject: [PATCH] fix(aihr): recover stalled media processing --- .../aihr/controller/AihrSopController.java | 6 + .../aihr/service/AihrSopSeedService.java | 211 +++++++++++--- .../aihr/service/AihrUploadQueueService.java | 269 ++++++++++++++++-- .../service/AihrVisionImagePreprocessor.java | 146 ++++++++++ .../aihr/service/AihrSopSeedServiceTest.java | 19 ++ .../service/AihrUploadQueueServiceTest.java | 22 +- .../AihrVisionImagePreprocessorTest.java | 53 ++++ backend/script/sql/aihr_knowledge_mysql8.sql | 4 +- .../aihr_20260729_media_reprocess_mysql8.sql | 36 +++ docs/API_INTEGRATION.md | 8 +- docs/BRD_PRODUCTION_MIGRATION_RUNBOOK.md | 12 +- frontend/src/api/aihr/processing.ts | 15 + .../processing-media-retry-contract.test.ts | 18 ++ frontend/src/views/knowledge/processing.vue | 46 ++- scripts/release-preflight.sh | 29 ++ scripts/reset-dev-db.sh | 1 + scripts/tests/aihr-schema-migrations.test.sh | 17 ++ scripts/tests/release-preflight-scope.test.sh | 2 + 18 files changed, 838 insertions(+), 76 deletions(-) create mode 100644 backend/ruoyi-modules/ruoyi-aihr/src/main/java/org/dromara/aihr/service/AihrVisionImagePreprocessor.java create mode 100644 backend/ruoyi-modules/ruoyi-aihr/src/test/java/org/dromara/aihr/service/AihrVisionImagePreprocessorTest.java create mode 100644 backend/script/sql/update/aihr_20260729_media_reprocess_mysql8.sql create mode 100644 frontend/src/utils/processing-media-retry-contract.test.ts diff --git a/backend/ruoyi-modules/ruoyi-aihr/src/main/java/org/dromara/aihr/controller/AihrSopController.java b/backend/ruoyi-modules/ruoyi-aihr/src/main/java/org/dromara/aihr/controller/AihrSopController.java index f2847bfa..5af11489 100644 --- a/backend/ruoyi-modules/ruoyi-aihr/src/main/java/org/dromara/aihr/controller/AihrSopController.java +++ b/backend/ruoyi-modules/ruoyi-aihr/src/main/java/org/dromara/aihr/controller/AihrSopController.java @@ -96,6 +96,12 @@ public class AihrSopController { return R.ok(uploadQueueService.retry(id)); } + @SaCheckRole(value = {TenantConstants.SUPER_ADMIN_ROLE_KEY, HR_OPERATOR_ROLE}, mode = SaMode.OR) + @PostMapping("/doc/processing-tasks/{id}/retry") + public R retryProcessingTask(@PathVariable Long id) { + return R.ok(uploadQueueService.retryAttachment(id)); + } + @SaCheckLogin @PostMapping("/search") public R search(@RequestBody SearchRequest request) { diff --git a/backend/ruoyi-modules/ruoyi-aihr/src/main/java/org/dromara/aihr/service/AihrSopSeedService.java b/backend/ruoyi-modules/ruoyi-aihr/src/main/java/org/dromara/aihr/service/AihrSopSeedService.java index e2d0b513..43a2c8d9 100644 --- a/backend/ruoyi-modules/ruoyi-aihr/src/main/java/org/dromara/aihr/service/AihrSopSeedService.java +++ b/backend/ruoyi-modules/ruoyi-aihr/src/main/java/org/dromara/aihr/service/AihrSopSeedService.java @@ -68,6 +68,7 @@ import java.net.URI; import java.net.http.HttpClient; import java.net.http.HttpRequest; import java.net.http.HttpResponse; +import java.net.http.HttpTimeoutException; import java.security.MessageDigest; import java.security.NoSuchAlgorithmException; import java.nio.file.Files; @@ -740,6 +741,100 @@ public class AihrSopSeedService { () -> fileFingerprint(stagedFile)); } + /** + * Reprocesses a media attachment from a queue-staged copy while retaining the original OSS object + * and its existing knowledge-space membership. + */ + @Transactional(rollbackFor = Exception.class) + public UploadResponse reprocessStagedAttachment(Long attachmentId, Path stagedFile) { + if (attachmentId == null || attachmentId <= 0 || stagedFile == null || !Files.isRegularFile(stagedFile)) { + throw new ServiceException("MEDIA_RETRY_SOURCE_MISSING: 待重试资料不存在"); + } + PendingMediaAttachment attachment = jdbcTemplate.query(""" + select a.id, a.knowledge_id, k.name as space_name, + coalesce(k.text_block_size, 800) as block_size, + coalesce(k.overlap_char, 120) as overlap_char, + a.oss_id, a.doc_id, a.name + from aihr_knowledge_attach a + join aihr_knowledge_info k on k.tenant_id = a.tenant_id and k.id = a.knowledge_id + where a.tenant_id = ? and a.id = ? and a.status in (0, 1, 3) + for update + """, rs -> rs.next() ? new PendingMediaAttachment( + rs.getLong("id"), rs.getLong("knowledge_id"), rs.getString("space_name"), + Math.max(200, rs.getInt("block_size")), + Math.max(0, Math.min(rs.getInt("overlap_char"), 200)), + rs.getLong("oss_id"), rs.getString("doc_id"), rs.getString("name") + ) : null, tenantId(), attachmentId); + if (attachment == null || attachment.ossId() <= 0 || isBlank(attachment.fileName()) + || (!imageFile(attachment.fileName()) && !AihrVideoService.videoFile(attachment.fileName()))) { + throw new ServiceException("MEDIA_RETRY_NOT_ALLOWED: 仅可重试待处理的图片或视频"); + } + + String docId = isBlank(attachment.docId()) + ? UUID.randomUUID().toString().replace("-", "") : attachment.docId(); + if (!docId.equals(attachment.docId())) { + jdbcTemplate.update(""" + update aihr_knowledge_attach set doc_id = ?, update_time = now() + where tenant_id = ? and id = ? + """, docId, tenantId(), attachment.id()); + } + String content = readContent(stagedFile, attachment.fileName()); + if (isBlank(content)) { + String reason = AihrVideoService.videoFile(attachment.fileName()) + ? "MEDIA_NO_TEXT:video" : "MEDIA_NO_TEXT:image"; + markAttachStatus(attachment.knowledgeId(), docId, 3, reason); + return new UploadResponse(docId, attachment.ossId(), attachment.fileName(), + attachment.spaceName(), 0, + "MEDIA_NO_TEXT: 未提取到可用于知识检索的文字,可检查源文件或模型配置后重试", + List.of("需重试"), List.of()); + } + + FileFingerprint sourceFingerprint = fileFingerprint(stagedFile); + DocumentFingerprint fingerprint = new DocumentFingerprint( + sourceFingerprint.sizeBytes(), sourceFingerprint.fileSha256(), sourceFingerprint.fileMd5(), textSha256(content)); + DocumentInsight insight = documentInsight(attachment.fileName(), content, attachment.spaceName()); + List fragments = split(content, attachment.blockSize(), attachment.overlap()); + + deleteQdrantDoc(attachment.knowledgeId(), docId); + jdbcTemplate.update(""" + delete from aihr_knowledge_fragment + where tenant_id = ? and knowledge_id = ? and doc_id = ? + """, tenantId(), attachment.knowledgeId(), docId); + for (int i = 0; i < fragments.size(); i++) { + jdbcTemplate.update(""" + insert into aihr_knowledge_fragment + (tenant_id, knowledge_id, idx, doc_id, content, create_time, update_time, remark) + values (?, ?, ?, ?, ?, now(), now(), ?) + """, tenantId(), attachment.knowledgeId(), i + 1, docId, fragments.get(i), + "media-reprocess:" + attachment.fileName()); + } + embedFragments(attachment.knowledgeId(), attachment.spaceName(), docId, fragments); + updateOssInsight(attachment.ossId(), insight, fingerprint); + markAttachStatus(attachment.knowledgeId(), docId, 2, + "media-reprocess:" + firstNonBlank(insight.classifiedBy(), "manual")); + + List snippets = fragments.stream() + .limit(3) + .map(fragment -> new SnippetResponse(attachment.fileName(), fragment)) + .toList(); + return new UploadResponse(docId, attachment.ossId(), attachment.fileName(), attachment.spaceName(), + fragments.size(), insight.summary(), insight.tags(), snippets); + } + + @Transactional + public void markMediaReprocessFailed(Long attachmentId, String reason) { + if (attachmentId == null || attachmentId <= 0) { + return; + } + String safeReason = reason != null && reason.matches("[A-Z][A-Z0-9_]{2,63}:.*") + ? truncate(reason, 180) : "MEDIA_PROCESSING_FAILED: 媒体解析失败"; + jdbcTemplate.update(""" + update aihr_knowledge_attach + set status = 3, remark = ?, update_time = now() + where tenant_id = ? and id = ? and status in (0, 1, 3) + """, safeReason, tenantId(), attachmentId); + } + /** * Removes only one document-to-space membership. The shared OSS object is deleted only after a * cross-tenant reference count confirms that no other knowledge attachment still uses it. @@ -1294,16 +1389,17 @@ public class AihrSopSeedService { for (SpaceConfig space : spaces) { String docId = upsertAttach(space.knowledgeId(), uploadedOss.getOssId(), fileName); docIds.add(docId); - markAttachStatus(space.knowledgeId(), docId, 0, video ? "pending-transcribe" : "pending-ocr"); + markAttachStatus(space.knowledgeId(), docId, 3, + video ? "MEDIA_NO_TEXT:video" : "MEDIA_NO_TEXT:image"); } } catch (RuntimeException error) { deleteUploadedOssQuietly(uploadedOss); throw error; } String hint = video - ? "视频未提取到语音转写或画面文字,已在所选知识空间标记待处理" - : "图片未提取到文字,已在所选知识空间标记待处理"; - List tags = video ? List.of("视频", "待转写") : List.of("图片", "待OCR"); + ? "MEDIA_NO_TEXT: 视频未提取到语音转写或画面文字,可从资料处理中心重试" + : "MEDIA_NO_TEXT: 图片未提取到可用文字,可从资料处理中心重试"; + List tags = video ? List.of("视频", "需重试") : List.of("图片", "需重试"); return new UploadResponse(docIds.get(0), uploadedOss.getOssId(), fileName, String.join("、", spaces.stream().map(SpaceConfig::name).toList()), 0, hint, tags, List.of()); } @@ -1344,8 +1440,8 @@ public class AihrSopSeedService { } /** - * OCR/转写拿不到文字的图片或视频:只落 OSS + attach 标「等待解析」(status 0),0 片段返回; - * 配置视觉/ASR 模型后重新上传同名文件即可替换加工。 + * OCR/转写拿不到文字的图片或视频:保留 OSS,attach 标记为可重试失败(status 3),0 片段返回; + * 管理端可直接从原 OSS 重新入队,不要求用户再次上传。 */ private UploadResponse savePendingImage(String fileName, String category, Supplier ossUploader) { boolean video = AihrVideoService.videoFile(fileName); @@ -1360,16 +1456,12 @@ public class AihrSopSeedService { deleteUploadedOssQuietly(uploadedOss); throw error; } - markAttachStatus(config.knowledgeId(), docId, 0, video ? "pending-transcribe" : "pending-ocr"); - String hint; - if (video) { - hint = "视频未提取到语音转写或画面文字,已入库标记待处理;确认已启用 asr(语音)或 vision(画面)模型后重新上传即可解析"; - } else { - hint = visionRuntime().isPresent() - ? "图片 OCR 未提取到文字(当前模型可能不支持图片输入,或图片中无可识别文字),已入库标记待处理" - : "未配置可用视觉模型,图片已入库标记待处理;在模型管理启用 vision/多模态 chat 模型后重新上传即可解析"; - } - List tags = video ? List.of("视频", "待转写") : List.of("图片", "待OCR"); + markAttachStatus(config.knowledgeId(), docId, 3, + video ? "MEDIA_NO_TEXT:video" : "MEDIA_NO_TEXT:image"); + String hint = video + ? "MEDIA_NO_TEXT: 视频未提取到语音转写或画面文字,可从资料处理中心重试" + : "MEDIA_NO_TEXT: 图片未提取到可用文字,可从资料处理中心重试"; + List tags = video ? List.of("视频", "需重试") : List.of("图片", "需重试"); return new UploadResponse(docId, uploadedOss.getOssId(), fileName, knowledgeName, 0, hint, tags, List.of()); } @@ -1572,7 +1664,9 @@ public class AihrSopSeedService { } private String callVisionOcr(ChatRuntime runtime, byte[] imageBytes, String mimeType, String prompt) throws Exception { - String dataUrl = "data:" + mimeType + ";base64," + Base64.getEncoder().encodeToString(imageBytes); + AihrVisionImagePreprocessor.PreparedImage prepared = AihrVisionImagePreprocessor.prepare(imageBytes); + String dataUrl = "data:" + prepared.mimeType() + ";base64," + + Base64.getEncoder().encodeToString(prepared.bytes()); ObjectNode body = objectMapper.createObjectNode(); body.put("model", runtime.modelName()); @@ -1600,16 +1694,26 @@ public class AihrSopSeedService { if (!isBlank(runtime.apiKey())) { builder.header("Authorization", "Bearer " + runtime.apiKey()); } - HttpResponse response = HttpClient.newBuilder() - .connectTimeout(Duration.ofSeconds(15)) - .build() - .send(builder.build(), HttpResponse.BodyHandlers.ofString()); + HttpResponse response; + try { + response = HttpClient.newBuilder() + .connectTimeout(Duration.ofSeconds(15)) + .build() + .send(builder.build(), HttpResponse.BodyHandlers.ofString()); + } catch (HttpTimeoutException e) { + throw new ServiceException("VISION_TIMEOUT: 视觉模型调用超时"); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new ServiceException("VISION_INTERRUPTED: 视觉模型调用已中断"); + } catch (IOException e) { + throw new ServiceException("VISION_NETWORK_ERROR: 视觉模型网络调用失败"); + } if (!ok(response.statusCode())) { - throw new IllegalStateException("vision OCR HTTP " + response.statusCode()); + throw new ServiceException("VISION_HTTP_" + response.statusCode() + ": 视觉模型拒绝图片请求"); } JsonNode choices = objectMapper.readTree(response.body()).path("choices"); if (!choices.isArray() || choices.isEmpty()) { - throw new IllegalStateException("vision OCR response missing choices"); + throw new ServiceException("VISION_BAD_RESPONSE: 视觉模型返回格式无效"); } return choices.get(0).path("message").path("content").asText(); } @@ -1722,7 +1826,7 @@ public class AihrSopSeedService { long totalBytes = rows.stream().mapToLong(ProcessingTaskRow::sizeBytes).sum(); int completed = (int) rows.stream().filter(row -> row.status() == 2).count(); int processing = (int) rows.stream().filter(row -> row.status() == 1).count(); - int failed = (int) rows.stream().filter(row -> row.status() == 3).count(); + int failed = (int) rows.stream().filter(row -> row.status() == 0 || row.status() == 3).count(); int embedded = rows.stream().mapToInt(ProcessingTaskRow::embeddingCount).sum(); int fragments = rows.stream().mapToInt(ProcessingTaskRow::fragmentCount).sum(); @@ -1730,7 +1834,7 @@ public class AihrSopSeedService { new MetricResponse("total", "资料总量", formatBytes(totalBytes), "已接入 " + rows.size() + " 个文件", "danger"), new MetricResponse("completed", "已完成", String.valueOf(completed), "可引用资料", "success"), new MetricResponse("processing", "处理中", String.valueOf(processing), "进行中任务", "primary"), - new MetricResponse("failed", "失败", String.valueOf(failed), failed == 0 ? "暂无异常" : "待处理异常", "warning") + new MetricResponse("failed", "需处理", String.valueOf(failed), failed == 0 ? "暂无异常" : "可从任务列表重试", "warning") ); Map summaries = new LinkedHashMap<>(); @@ -1825,7 +1929,7 @@ public class AihrSopSeedService { vectorStatus(row.fragmentCount(), row.embeddingCount(), row.status()), vectorType(row.fragmentCount(), row.embeddingCount(), row.status()), formatTime(row.updateTime()), - row.status() == 3 ? "重试" : row.status() == 0 ? "开始" : "查看", + row.status() == 3 || row.status() == 0 ? "重试" : "查看", row.remark(), row.summary(), row.tags(), @@ -1834,8 +1938,9 @@ public class AihrSopSeedService { } private ProcessingEventResponse processingEvent(ProcessingTaskRow row) { - String resultType = row.status() == 3 ? "danger" : row.status() == 1 ? "primary" : "success"; - String detail = row.status() == 3 ? row.remark() : "生成 " + row.fragmentCount() + " 个片段,入库 " + row.embeddingCount() + " 个向量"; + String resultType = row.status() == 3 || row.status() == 0 ? "danger" : row.status() == 1 ? "primary" : "success"; + String detail = row.status() == 3 || row.status() == 0 + ? row.remark() : "生成 " + row.fragmentCount() + " 个片段,入库 " + row.embeddingCount() + " 个向量"; return new ProcessingEventResponse( formatTime(row.updateTime()), "TASK-" + row.id(), @@ -2055,7 +2160,7 @@ public class AihrSopSeedService { private static String statusLabel(int status) { return switch (status) { - case 0 -> "等待解析"; + case 0 -> "待重试"; case 1 -> "解析中"; case 2 -> "已完成"; case 3 -> "失败"; @@ -2077,12 +2182,12 @@ public class AihrSopSeedService { case 1 -> "Tika 解析"; case 2 -> "向量化入库"; case 3 -> "解析失败"; - default -> "等待解析"; + default -> "待重试"; }; } private static String vectorStatus(int fragments, int embeddings, int status) { - if (status == 3) { + if (status == 0 || status == 3) { return "失败原因"; } if (fragments <= 0) { @@ -2098,7 +2203,7 @@ public class AihrSopSeedService { } private static String vectorType(int fragments, int embeddings, int status) { - if (status == 3) { + if (status == 0 || status == 3) { return "danger"; } if (fragments <= 0 || embeddings <= 0) { @@ -3833,12 +3938,17 @@ public class AihrSopSeedService { byte[] bytes = file.getBytes(); String text = callVisionOcr(runtime.get(), bytes, guessMimeType(fileName)); return normalizeExtractedText(text); + } catch (ServiceException e) { + log.warn("image OCR failed file={} model={} reason={}", + AihrSensitiveText.forModel(fileName), runtime.get().modelName(), mediaFailureReason(e)); + throw e; } catch (Exception e) { - log.warn("image ocr failed for {} via {}(处理错误已隐藏)", AihrSensitiveText.forModel(fileName), runtime.get().modelName()); - return ""; + log.warn("image OCR failed file={} model={} reason=VISION_OCR_FAILED", + AihrSensitiveText.forModel(fileName), runtime.get().modelName()); + throw new ServiceException("VISION_OCR_FAILED: 图片文字识别失败"); } } - return ""; + throw new ServiceException("VISION_NOT_CONFIGURED: 未配置可用视觉模型"); } try (InputStream input = file.getInputStream()) { return readKnowledgeDocument(input, fileName); @@ -3856,9 +3966,14 @@ public class AihrSopSeedService { } try { return Optional.of(normalizeExtractedText(callVisionOcr(runtime.get(), bytes, mimeType, FRAME_PROMPT))); + } catch (ServiceException e) { + log.warn("video frame OCR failed file={} reason={}", + AihrSensitiveText.forModel(fileName), mediaFailureReason(e)); + throw e; } catch (Exception e) { - log.warn("video frame ocr failed for {}(处理错误已隐藏)", AihrSensitiveText.forModel(fileName)); - return Optional.empty(); + log.warn("video frame OCR failed file={} reason=VISION_OCR_FAILED", + AihrSensitiveText.forModel(fileName)); + throw new ServiceException("VISION_OCR_FAILED: 视频关键帧识别失败"); } }); } @@ -3872,12 +3987,17 @@ public class AihrSopSeedService { byte[] bytes = Files.readAllBytes(file); String text = callVisionOcr(runtime.get(), bytes, guessMimeType(fileName)); return normalizeExtractedText(text); + } catch (ServiceException e) { + log.warn("image OCR failed file={} model={} reason={}", + AihrSensitiveText.forModel(fileName), runtime.get().modelName(), mediaFailureReason(e)); + throw e; } catch (Exception e) { - log.warn("image ocr failed for {} via {}(处理错误已隐藏)", AihrSensitiveText.forModel(fileName), runtime.get().modelName()); - return ""; + log.warn("image OCR failed file={} model={} reason=VISION_OCR_FAILED", + AihrSensitiveText.forModel(fileName), runtime.get().modelName()); + throw new ServiceException("VISION_OCR_FAILED: 图片文字识别失败"); } } - return ""; + throw new ServiceException("VISION_NOT_CONFIGURED: 未配置可用视觉模型"); } try (InputStream input = Files.newInputStream(file)) { return readKnowledgeDocument(input, fileName); @@ -4085,6 +4205,13 @@ public class AihrSopSeedService { return value.substring(0, maxLength); } + private static String mediaFailureReason(ServiceException error) { + String message = error == null ? "" : error.getMessage(); + int separator = message == null ? -1 : message.indexOf(':'); + String code = separator > 0 ? message.substring(0, separator) : message; + return code != null && code.matches("[A-Z][A-Z0-9_]{2,63}") ? code : "MEDIA_PROCESSING_FAILED"; + } + private String tenantId() { return firstNonBlank(TenantHelper.getTenantId(), TENANT_ID); } @@ -4250,6 +4377,10 @@ public class AihrSopSeedService { private record SpaceConfig(long knowledgeId, String code, String name, int blockSize, int overlap) { } + private record PendingMediaAttachment(long id, long knowledgeId, String spaceName, int blockSize, + int overlap, long ossId, String docId, String fileName) { + } + private record EmbeddingRuntime(String modelName, String baseUrl, String apiKey, int dimension) { } diff --git a/backend/ruoyi-modules/ruoyi-aihr/src/main/java/org/dromara/aihr/service/AihrUploadQueueService.java b/backend/ruoyi-modules/ruoyi-aihr/src/main/java/org/dromara/aihr/service/AihrUploadQueueService.java index b42e8ab7..8b69e147 100644 --- a/backend/ruoyi-modules/ruoyi-aihr/src/main/java/org/dromara/aihr/service/AihrUploadQueueService.java +++ b/backend/ruoyi-modules/ruoyi-aihr/src/main/java/org/dromara/aihr/service/AihrUploadQueueService.java @@ -8,14 +8,23 @@ import org.dromara.aihr.domain.AihrSopDto.UploadEnqueueResponse; import org.dromara.aihr.domain.AihrSopDto.UploadItemResponse; import org.dromara.aihr.domain.AihrSopDto.UploadResponse; import org.dromara.common.core.exception.ServiceException; +import org.dromara.common.oss.core.OssClient; +import org.dromara.common.oss.factory.OssFactory; import org.dromara.common.tenant.helper.TenantHelper; +import org.dromara.system.domain.vo.SysOssVo; +import org.dromara.system.service.ISysOssService; import org.springframework.beans.factory.annotation.Value; import org.springframework.dao.DataAccessException; import org.springframework.jdbc.core.JdbcTemplate; import org.springframework.stereotype.Service; +import org.springframework.transaction.annotation.Transactional; +import org.springframework.transaction.support.TransactionSynchronization; +import org.springframework.transaction.support.TransactionSynchronizationManager; import org.springframework.web.multipart.MultipartFile; import java.io.IOException; +import java.io.InputStream; +import java.io.OutputStream; import java.nio.charset.Charset; import java.nio.charset.StandardCharsets; import java.nio.file.Files; @@ -67,6 +76,7 @@ public class AihrUploadQueueService { private static final ObjectMapper JSON = new ObjectMapper(); private final AihrSopSeedService sopSeedService; + private final ISysOssService ossService; private final JdbcTemplate jdbcTemplate; private final ScheduledExecutorService scheduledExecutorService; private final Path stagingRoot; @@ -78,10 +88,12 @@ public class AihrUploadQueueService { private volatile boolean tableReady; public AihrUploadQueueService(AihrSopSeedService sopSeedService, + ISysOssService ossService, JdbcTemplate jdbcTemplate, ScheduledExecutorService scheduledExecutorService, @Value("${aihr.upload.staging:./.data/staging}") String stagingRootConfig) { this.sopSeedService = sopSeedService; + this.ossService = ossService; this.jdbcTemplate = jdbcTemplate; this.scheduledExecutorService = scheduledExecutorService; this.stagingRoot = Path.of(stagingRootConfig).toAbsolutePath().normalize(); @@ -167,6 +179,67 @@ public class AihrUploadQueueService { """.formatted(safeLimit), this::mapItem, currentTenantId(), batch); } + /** + * Restages a failed/legacy-pending media attachment from its original OSS object and puts it back + * on the same asynchronous worker. The attachment row is claimed with CAS to prevent double retry. + */ + @Transactional(rollbackFor = Exception.class) + public UploadEnqueueResponse retryAttachment(Long attachmentId) { + ensureTable(); + if (attachmentId == null || attachmentId <= 0) { + throw new ServiceException("MEDIA_RETRY_NOT_FOUND: 待重试资料不存在"); + } + String tenantId = currentTenantId(); + SourceAttachment source = jdbcTemplate.query(""" + select a.id, a.oss_id, a.name, a.status, k.code as space_code, k.name as category + from aihr_knowledge_attach a + join aihr_knowledge_info k on k.tenant_id = a.tenant_id and k.id = a.knowledge_id + where a.tenant_id = ? and a.id = ? + """, rs -> rs.next() ? new SourceAttachment( + rs.getLong("id"), rs.getLong("oss_id"), rs.getString("name"), rs.getInt("status"), + rs.getString("space_code"), rs.getString("category") + ) : null, tenantId, attachmentId); + if (source == null || source.ossId() <= 0 || source.fileName() == null + || (!AihrSopSeedService.supportedFile(source.fileName()) + || (!AihrVideoService.videoFile(source.fileName()) && !isImage(source.fileName()))) + || (source.status() != 0 && source.status() != 3)) { + throw new ServiceException("MEDIA_RETRY_NOT_ALLOWED: 仅可重试待处理的图片或视频"); + } + + Path staged = stagingPath(source.fileName()); + try { + int claimed = jdbcTemplate.update(""" + update aihr_knowledge_attach + set status = 1, remark = 'media-retry:queued', update_time = now() + where tenant_id = ? and id = ? and status in (0, 3) + """, tenantId, attachmentId); + if (claimed != 1) { + throw new ServiceException("MEDIA_RETRY_ALREADY_QUEUED: 资料已进入解析队列"); + } + stageSourceObject(source, staged); + String batchId = "retry-" + UUID.randomUUID(); + jdbcTemplate.update(""" + insert into aihr_knowledge_upload_item + (tenant_id, batch_id, file_name, category, space_codes_json, staging_path, + source_attach_id, status, create_time, update_time) + values (?, ?, ?, ?, ?, ?, ?, 0, now(), now()) + """, tenantId, batchId, source.fileName(), source.category(), + writeSpaceCodes(List.of(source.spaceCode())), staged.toString(), attachmentId); + Long itemId = jdbcTemplate.query(""" + select id from aihr_knowledge_upload_item where tenant_id = ? and staging_path = ? + """, rs -> rs.next() ? rs.getLong(1) : null, tenantId, staged.toString()); + triggerAfterCommit(); + return new UploadEnqueueResponse(itemId, batchId, source.fileName(), "queued"); + } catch (RuntimeException | IOException e) { + deleteQuietly(staged); + if (e instanceof ServiceException serviceException) { + throw serviceException; + } + throw new ServiceException("MEDIA_RETRY_SOURCE_UNAVAILABLE: 原始文件读取失败"); + } + } + + @Transactional(rollbackFor = Exception.class) public UploadItemResponse retry(Long id) { ensureTable(); if (id == null) { @@ -174,12 +247,14 @@ public class AihrUploadQueueService { } String tenantId = currentTenantId(); ItemRow row = jdbcTemplate.query(""" - select tenant_id, id, batch_id, file_name, category, space_codes_json, staging_path, status + select tenant_id, id, batch_id, file_name, category, space_codes_json, staging_path, + source_attach_id, status from aihr_knowledge_upload_item where tenant_id = ? and id = ? """, rs -> rs.next() ? new ItemRow(rs.getString("tenant_id"), rs.getLong("id"), rs.getString("batch_id"), rs.getString("file_name"), - rs.getString("category"), readSpaceCodes(rs.getString("space_codes_json")), rs.getString("staging_path"), rs.getInt("status")) + rs.getString("category"), readSpaceCodes(rs.getString("space_codes_json")), + rs.getString("staging_path"), rs.getObject("source_attach_id", Long.class), rs.getInt("status")) : null, tenantId, id); if (row == null) { throw new ServiceException("重试条目不存在"); @@ -199,7 +274,17 @@ public class AihrUploadQueueService { if (updated == 0) { throw new ServiceException("仅失败状态的文件可重试"); } - triggerProcessing(); + if (row.sourceAttachId() != null) { + int claimed = jdbcTemplate.update(""" + update aihr_knowledge_attach + set status = 1, remark = 'media-retry:queued', update_time = now() + where tenant_id = ? and id = ? and status in (0, 3) + """, tenantId, row.sourceAttachId()); + if (claimed != 1) { + throw new ServiceException("MEDIA_RETRY_ALREADY_QUEUED: 资料已进入解析队列"); + } + } + triggerAfterCommit(); return item(id); } @@ -270,19 +355,25 @@ public class AihrUploadQueueService { String tenantId = claim.tenantId(); Long id = claim.id(); ItemRow row = jdbcTemplate.query(""" - select tenant_id, id, batch_id, file_name, category, space_codes_json, staging_path, status + select tenant_id, id, batch_id, file_name, category, space_codes_json, staging_path, + source_attach_id, status from aihr_knowledge_upload_item where tenant_id = ? and id = ? """, rs -> rs.next() ? new ItemRow(rs.getString("tenant_id"), rs.getLong("id"), rs.getString("batch_id"), rs.getString("file_name"), - rs.getString("category"), readSpaceCodes(rs.getString("space_codes_json")), rs.getString("staging_path"), rs.getInt("status")) + rs.getString("category"), readSpaceCodes(rs.getString("space_codes_json")), + rs.getString("staging_path"), rs.getObject("source_attach_id", Long.class), rs.getInt("status")) : null, tenantId, id); if (row == null) { return; } Path staged = Path.of(row.stagingPath()); if (!Files.isRegularFile(staged)) { - markFailed(tenantId, id, "暂存文件已丢失,请重新上传该文件"); + String message = row.sourceAttachId() == null + ? "暂存文件已丢失,请重新上传该文件" + : "MEDIA_RETRY_STAGING_MISSING: 重试暂存文件已丢失,请从原文件重新入队"; + markFailed(tenantId, id, message); + markSourceAttachmentFailed(tenantId, row.sourceAttachId(), message); return; } try { @@ -298,28 +389,36 @@ public class AihrUploadQueueService { } return; } - UploadResponse result = TenantHelper.dynamic(tenantId, () -> row.spaceCodes().isEmpty() - ? sopSeedService.processStagedDocument(row.fileName(), row.category(), staged) - : sopSeedService.processStagedDocument(row.fileName(), row.spaceCodes(), staged)); - // 0 片段的完成态(如图片待 OCR)把说明写进 error 列,面板可见原因 - String note = result.fragments() != null && result.fragments() == 0 ? truncateError(result.summary()) : null; - // CAS:仅当仍是本工人持有的「处理中」才写完成,防止清扫重置后被后来的工人覆盖状态 + UploadResponse result = TenantHelper.dynamic(tenantId, () -> { + if (row.sourceAttachId() != null) { + return sopSeedService.reprocessStagedAttachment(row.sourceAttachId(), staged); + } + return row.spaceCodes().isEmpty() + ? sopSeedService.processStagedDocument(row.fileName(), row.category(), staged) + : sopSeedService.processStagedDocument(row.fileName(), row.spaceCodes(), staged); + }); + boolean completed = result.fragments() != null && result.fragments() > 0; + String note = completed ? null : truncateError(firstNonBlank( + result.summary(), "MEDIA_NO_TEXT: 未提取到可用文字")); + // 0 片段不是完成态:保留暂存文件和可重试失败状态,避免附件永久显示成“等待解析”。 int updated = jdbcTemplate.update(""" update aihr_knowledge_upload_item - set status = 2, error = ?, doc_id = ?, fragment_count = ?, update_time = now() + set status = ?, error = ?, doc_id = ?, fragment_count = ?, update_time = now() where tenant_id = ? and id = ? and status = 1 - """, note, result.docId(), result.fragments(), tenantId, id); - if (updated == 1) { + """, completed ? 2 : 3, note, result.docId(), result.fragments(), tenantId, id); + if (updated == 1 && completed) { deleteQuietly(staged); } else { - log.info("upload item {} finished but row was re-claimed, skip status write", id); + log.info("upload item {} result persisted completed={} updated={}", id, completed, updated); } } catch (ArchiveException e) { log.warn("upload item {} archive rejected: {}", id, e.getMessage(), e); markFailed(tenantId, id, e.getMessage()); } catch (Exception e) { - log.warn("upload item {} process failed", id, e); - markFailed(tenantId, id, "资料加工失败,请重试或联系管理员"); + String message = processingFailureMessage(e); + log.warn("upload item {} process failed reason={}", id, failureCode(message)); + markFailed(tenantId, id, message); + markSourceAttachmentFailed(tenantId, row.sourceAttachId(), message); } } @@ -542,6 +641,46 @@ public class AihrUploadQueueService { """, error, tenantId, id); } + private void markSourceAttachmentFailed(String tenantId, Long attachmentId, String message) { + if (attachmentId == null) { + return; + } + try { + TenantHelper.dynamic(tenantId, () -> { + sopSeedService.markMediaReprocessFailed(attachmentId, message); + return null; + }); + } catch (Exception e) { + log.warn("media source attachment {} failure status update skipped", attachmentId); + } + } + + private static String processingFailureMessage(Exception error) { + if (error instanceof ServiceException) { + String message = error.getMessage(); + if (message != null && message.matches("[A-Z][A-Z0-9_]{2,63}:.*")) { + return truncateError(message); + } + } + return "资料加工失败,请重试或联系管理员"; + } + + private static String failureCode(String message) { + int separator = message == null ? -1 : message.indexOf(':'); + return separator > 0 ? message.substring(0, separator) : "PROCESSING_FAILED"; + } + + private static String firstNonBlank(String... values) { + if (values != null) { + for (String value : values) { + if (value != null && !value.isBlank()) { + return value; + } + } + } + return ""; + } + private static String truncateError(String value) { if (value == null || value.length() <= MAX_ERROR_CHARS) { return value; @@ -586,12 +725,66 @@ public class AihrUploadQueueService { }; } + private void stageSourceObject(SourceAttachment source, Path target) throws IOException { + if (ossService == null) { + throw new ServiceException("MEDIA_RETRY_SOURCE_UNAVAILABLE: 对象存储服务不可用"); + } + SysOssVo object = ossService.getById(source.ossId()); + if (object == null || object.getOssId() == null || !source.ossId().equals(object.getOssId()) + || object.getService() == null || object.getService().isBlank() + || object.getFileName() == null || object.getFileName().isBlank()) { + throw new ServiceException("MEDIA_RETRY_SOURCE_UNAVAILABLE: 原始文件不存在"); + } + long maxBytes = AihrVideoService.videoFile(source.fileName()) ? MAX_VIDEO_BYTES : MAX_FILE_BYTES; + OssClient client = OssFactory.instance(object.getService()); + Files.createDirectories(target.getParent()); + try (InputStream input = client.getObjectContent(object.getFileName()); + OutputStream output = Files.newOutputStream(target)) { + copyWithLimit(input, output, maxBytes); + } catch (RuntimeException | IOException e) { + deleteQuietly(target); + throw e; + } + } + + private static void copyWithLimit(InputStream input, OutputStream output, long maxBytes) throws IOException { + byte[] buffer = new byte[8192]; + long copied = 0; + int read; + while ((read = input.read(buffer)) != -1) { + copied += read; + if (copied > maxBytes) { + throw new ServiceException("MEDIA_RETRY_SOURCE_TOO_LARGE: 原始文件超过处理上限"); + } + output.write(buffer, 0, read); + } + } + + private void triggerAfterCommit() { + if (!TransactionSynchronizationManager.isSynchronizationActive()) { + triggerProcessing(); + return; + } + TransactionSynchronizationManager.registerSynchronization(new TransactionSynchronization() { + @Override + public void afterCommit() { + triggerProcessing(); + } + }); + } + private Path stagingPath(String fileName) { String lower = fileName.toLowerCase(); String extension = lower.contains(".") ? lower.substring(lower.lastIndexOf('.')) : ""; return stagingRoot.resolve(UUID.randomUUID() + extension); } + private static boolean isImage(String fileName) { + String lower = fileName == null ? "" : fileName.toLowerCase(Locale.ROOT); + return lower.endsWith(".jpg") || lower.endsWith(".jpeg") || lower.endsWith(".png") + || lower.endsWith(".gif") || lower.endsWith(".webp") || lower.endsWith(".bmp"); + } + /** 去掉路径分隔与控制字符,防止暂存路径逃逸出 stagingRoot。 */ private static String sanitizeFileName(String original) { String name = original == null || original.isBlank() ? "knowledge.txt" : original.trim(); @@ -666,6 +859,9 @@ public class AihrUploadQueueService { } if (!runtimeSchemaBootstrap) { AihrSchemaMigrationGuard.requireTables(jdbcTemplate, "aihr_knowledge_upload_item"); + AihrSchemaMigrationGuard.requireColumns(jdbcTemplate, "aihr_knowledge_upload_item", "source_attach_id"); + AihrSchemaMigrationGuard.requireIndexes(jdbcTemplate, "aihr_knowledge_upload_item", + "idx_aihr_upload_item_source"); tableReady = true; return; } @@ -678,6 +874,7 @@ public class AihrUploadQueueService { `category` varchar(100) DEFAULT '' COMMENT '目标分类', `space_codes_json` json DEFAULT NULL COMMENT '目标知识空间编码', `staging_path` varchar(500) NOT NULL COMMENT '暂存文件路径', + `source_attach_id` bigint DEFAULT NULL COMMENT '从既有附件重新解析时的附件ID', `status` tinyint NOT NULL DEFAULT 0 COMMENT '0待处理 1处理中 2完成 3失败', `error` varchar(500) DEFAULT NULL COMMENT '失败原因', `doc_id` varchar(80) DEFAULT NULL COMMENT '入库文档ID', @@ -686,14 +883,46 @@ public class AihrUploadQueueService { `update_time` datetime DEFAULT NULL COMMENT '更新时间', PRIMARY KEY (`id`), KEY `idx_aihr_upload_item_batch` (`tenant_id`, `batch_id`, `id`), - KEY `idx_aihr_upload_item_status` (`tenant_id`, `status`, `update_time`) + KEY `idx_aihr_upload_item_status` (`tenant_id`, `status`, `update_time`), + KEY `idx_aihr_upload_item_source` (`tenant_id`, `source_attach_id`, `id`) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci COMMENT='知识库批量上传队列'; """); + ensureDevelopmentRetrySchema(); tableReady = true; } } + private void ensureDevelopmentRetrySchema() { + Integer column = jdbcTemplate.queryForObject(""" + select count(*) from information_schema.columns + where table_schema = database() and table_name = 'aihr_knowledge_upload_item' + and column_name = 'source_attach_id' + """, Integer.class); + if (column == null || column == 0) { + jdbcTemplate.execute(""" + alter table aihr_knowledge_upload_item + add column source_attach_id bigint default null comment '从既有附件重新解析时的附件ID' + after staging_path + """); + } + Integer index = jdbcTemplate.queryForObject(""" + select count(*) from information_schema.statistics + where table_schema = database() and table_name = 'aihr_knowledge_upload_item' + and index_name = 'idx_aihr_upload_item_source' + """, Integer.class); + if (index == null || index == 0) { + jdbcTemplate.execute(""" + alter table aihr_knowledge_upload_item + add key idx_aihr_upload_item_source (tenant_id, source_attach_id, id) + """); + } + } + private record ItemRow(String tenantId, Long id, String batchId, String fileName, String category, - List spaceCodes, String stagingPath, int status) { + List spaceCodes, String stagingPath, Long sourceAttachId, int status) { + } + + private record SourceAttachment(Long id, Long ossId, String fileName, int status, String spaceCode, + String category) { } } diff --git a/backend/ruoyi-modules/ruoyi-aihr/src/main/java/org/dromara/aihr/service/AihrVisionImagePreprocessor.java b/backend/ruoyi-modules/ruoyi-aihr/src/main/java/org/dromara/aihr/service/AihrVisionImagePreprocessor.java new file mode 100644 index 00000000..2102e1ff --- /dev/null +++ b/backend/ruoyi-modules/ruoyi-aihr/src/main/java/org/dromara/aihr/service/AihrVisionImagePreprocessor.java @@ -0,0 +1,146 @@ +package org.dromara.aihr.service; + +import org.dromara.common.core.exception.ServiceException; + +import javax.imageio.IIOImage; +import javax.imageio.ImageIO; +import javax.imageio.ImageReadParam; +import javax.imageio.ImageReader; +import javax.imageio.ImageWriteParam; +import javax.imageio.ImageWriter; +import javax.imageio.stream.ImageInputStream; +import javax.imageio.stream.ImageOutputStream; +import java.awt.Color; +import java.awt.Graphics2D; +import java.awt.RenderingHints; +import java.awt.image.BufferedImage; +import java.io.ByteArrayInputStream; +import java.io.ByteArrayOutputStream; +import java.io.IOException; +import java.util.Iterator; + +/** + * Bounds image dimensions and request size before an image is embedded in a vision-model data URL. + * Metadata is read before decoding and source subsampling is used for oversized images so a compressed + * multi-gigapixel JPEG does not have to be expanded at full resolution in the JVM. + */ +final class AihrVisionImagePreprocessor { + + static final int MAX_LONG_EDGE = 2048; + static final long MAX_PIXELS = 4_194_304L; + static final int MAX_ENCODED_BYTES = 5 * 1024 * 1024; + private static final float JPEG_QUALITY = 0.82F; + + private AihrVisionImagePreprocessor() { + } + + static PreparedImage prepare(byte[] source) { + if (source == null || source.length == 0) { + throw new ServiceException("VISION_IMAGE_EMPTY: 图片内容为空"); + } + try (ImageInputStream input = ImageIO.createImageInputStream(new ByteArrayInputStream(source))) { + if (input == null) { + throw new ServiceException("VISION_IMAGE_INVALID: 无法读取图片"); + } + Iterator readers = ImageIO.getImageReaders(input); + if (!readers.hasNext()) { + throw new ServiceException("VISION_IMAGE_FORMAT_UNSUPPORTED: 图片格式不受支持"); + } + ImageReader reader = readers.next(); + try { + reader.setInput(input, true, true); + int width = reader.getWidth(0); + int height = reader.getHeight(0); + if (width <= 0 || height <= 0) { + throw new ServiceException("VISION_IMAGE_INVALID: 图片尺寸无效"); + } + int subsampling = sourceSubsampling(width, height); + ImageReadParam readParam = reader.getDefaultReadParam(); + if (subsampling > 1) { + readParam.setSourceSubsampling(subsampling, subsampling, 0, 0); + } + BufferedImage decoded = reader.read(0, readParam); + if (decoded == null) { + throw new ServiceException("VISION_IMAGE_INVALID: 图片解码失败"); + } + BufferedImage bounded = boundDimensions(decoded); + BufferedImage processed = bounded; + byte[] encoded = encodeJpeg(processed, JPEG_QUALITY); + if (encoded.length > MAX_ENCODED_BYTES) { + processed = boundDimensions(bounded, 0.75D); + encoded = encodeJpeg(processed, 0.70F); + } + if (encoded.length > MAX_ENCODED_BYTES) { + throw new ServiceException("VISION_IMAGE_TOO_LARGE: 图片压缩后仍超过视觉模型请求上限"); + } + return new PreparedImage(encoded, "image/jpeg", processed.getWidth(), processed.getHeight(), + width, height); + } finally { + reader.dispose(); + } + } catch (ServiceException e) { + throw e; + } catch (IOException | RuntimeException e) { + throw new ServiceException("VISION_IMAGE_INVALID: 图片预处理失败"); + } + } + + private static int sourceSubsampling(int width, int height) { + long pixels = (long) width * height; + double edgeScale = Math.max(width / (double) MAX_LONG_EDGE, height / (double) MAX_LONG_EDGE); + double pixelScale = Math.sqrt(pixels / (double) MAX_PIXELS); + return Math.max(1, (int) Math.ceil(Math.max(edgeScale, pixelScale))); + } + + private static BufferedImage boundDimensions(BufferedImage source) { + long pixels = (long) source.getWidth() * source.getHeight(); + double scale = Math.min(1D, Math.min( + MAX_LONG_EDGE / (double) Math.max(source.getWidth(), source.getHeight()), + Math.sqrt(MAX_PIXELS / (double) pixels) + )); + return boundDimensions(source, scale); + } + + private static BufferedImage boundDimensions(BufferedImage source, double scale) { + int width = Math.max(1, (int) Math.floor(source.getWidth() * Math.min(1D, scale))); + int height = Math.max(1, (int) Math.floor(source.getHeight() * Math.min(1D, scale))); + BufferedImage rgb = new BufferedImage(width, height, BufferedImage.TYPE_INT_RGB); + Graphics2D graphics = rgb.createGraphics(); + try { + graphics.setColor(Color.WHITE); + graphics.fillRect(0, 0, width, height); + graphics.setRenderingHint(RenderingHints.KEY_INTERPOLATION, RenderingHints.VALUE_INTERPOLATION_BICUBIC); + graphics.setRenderingHint(RenderingHints.KEY_RENDERING, RenderingHints.VALUE_RENDER_QUALITY); + graphics.drawImage(source, 0, 0, width, height, null); + } finally { + graphics.dispose(); + } + return rgb; + } + + private static byte[] encodeJpeg(BufferedImage image, float quality) throws IOException { + Iterator writers = ImageIO.getImageWritersByFormatName("jpeg"); + if (!writers.hasNext()) { + throw new IOException("JPEG writer unavailable"); + } + ImageWriter writer = writers.next(); + try (ByteArrayOutputStream output = new ByteArrayOutputStream(); + ImageOutputStream imageOutput = ImageIO.createImageOutputStream(output)) { + writer.setOutput(imageOutput); + ImageWriteParam writeParam = writer.getDefaultWriteParam(); + if (writeParam.canWriteCompressed()) { + writeParam.setCompressionMode(ImageWriteParam.MODE_EXPLICIT); + writeParam.setCompressionQuality(quality); + } + writer.write(null, new IIOImage(image, null, null), writeParam); + imageOutput.flush(); + return output.toByteArray(); + } finally { + writer.dispose(); + } + } + + record PreparedImage(byte[] bytes, String mimeType, int width, int height, + int originalWidth, int originalHeight) { + } +} diff --git a/backend/ruoyi-modules/ruoyi-aihr/src/test/java/org/dromara/aihr/service/AihrSopSeedServiceTest.java b/backend/ruoyi-modules/ruoyi-aihr/src/test/java/org/dromara/aihr/service/AihrSopSeedServiceTest.java index 75ceeb9d..c86f1940 100644 --- a/backend/ruoyi-modules/ruoyi-aihr/src/test/java/org/dromara/aihr/service/AihrSopSeedServiceTest.java +++ b/backend/ruoyi-modules/ruoyi-aihr/src/test/java/org/dromara/aihr/service/AihrSopSeedServiceTest.java @@ -174,6 +174,25 @@ public class AihrSopSeedServiceTest { assertTrue(code.contains("!modelService.visionAllowed()")); } + @Test + @Tag("dev") + public void mediaWithoutTextFailsExplicitlyAndCanBeRestagedFromOss() throws Exception { + Path source = Path.of("src/main/java/org/dromara/aihr/service/AihrSopSeedService.java"); + if (!Files.exists(source)) { + source = Path.of("ruoyi-modules/ruoyi-aihr/src/main/java/org/dromara/aihr/service/AihrSopSeedService.java"); + } + String code = Files.readString(source).replace("\r\n", "\n"); + + assertTrue(code.contains("AihrVisionImagePreprocessor.prepare(imageBytes)")); + assertTrue(code.contains("markAttachStatus(space.knowledgeId(), docId, 3")); + assertTrue(code.contains("markAttachStatus(config.knowledgeId(), docId, 3")); + assertTrue(code.contains("public UploadResponse reprocessStagedAttachment")); + assertTrue(code.contains("MEDIA_NO_TEXT")); + assertTrue(AihrSopSeedService.class + .getMethod("reprocessStagedAttachment", Long.class, Path.class) + .getAnnotation(org.springframework.transaction.annotation.Transactional.class) != null); + } + @Test @Tag("dev") public void documentUploadCleansNewOssWhenAttachWriteFails() throws Exception { diff --git a/backend/ruoyi-modules/ruoyi-aihr/src/test/java/org/dromara/aihr/service/AihrUploadQueueServiceTest.java b/backend/ruoyi-modules/ruoyi-aihr/src/test/java/org/dromara/aihr/service/AihrUploadQueueServiceTest.java index 3945c6ab..682f354b 100644 --- a/backend/ruoyi-modules/ruoyi-aihr/src/test/java/org/dromara/aihr/service/AihrUploadQueueServiceTest.java +++ b/backend/ruoyi-modules/ruoyi-aihr/src/test/java/org/dromara/aihr/service/AihrUploadQueueServiceTest.java @@ -44,7 +44,7 @@ class AihrUploadQueueServiceTest { when(jdbcTemplate.update(anyString(), any(Object[].class))).thenReturn(1); ScheduledExecutorService executor = Executors.newSingleThreadScheduledExecutor(); try { - AihrUploadQueueService service = new AihrUploadQueueService(null, jdbcTemplate, executor, staging.toString()); + AihrUploadQueueService service = new AihrUploadQueueService(null, null, jdbcTemplate, executor, staging.toString()); int extracted = assertDoesNotThrow(() -> expandArchive(service, itemRow(archive))); @@ -83,6 +83,23 @@ class AihrUploadQueueServiceTest { assertFalse(code.contains("where tenant_id = ? and status = 1\n and (")); } + @Test + void zeroFragmentMediaIsRetryableAndExistingAttachmentsKeepTheirSourceLineage() throws Exception { + Path source = Path.of("src/main/java/org/dromara/aihr/service/AihrUploadQueueService.java"); + if (!Files.exists(source)) { + source = Path.of("ruoyi-modules/ruoyi-aihr/src/main/java/org/dromara/aihr/service/AihrUploadQueueService.java"); + } + String code = Files.readString(source).replace("\r\n", "\n"); + + assertTrue(code.contains("boolean completed = result.fragments() != null && result.fragments() > 0")); + assertTrue(code.contains("completed ? 2 : 3")); + assertTrue(code.contains("if (updated == 1 && completed)")); + assertTrue(code.contains("source_attach_id")); + assertTrue(code.contains("sopSeedService.reprocessStagedAttachment")); + assertTrue(code.contains("TransactionSynchronizationManager.registerSynchronization")); + assertTrue(code.contains("MEDIA_RETRY_ALREADY_QUEUED")); + } + private static void writeZip(Path archive, Charset charset, String... entries) throws IOException { try (ZipOutputStream output = new ZipOutputStream(Files.newOutputStream(archive), charset)) { for (String entry : entries) { @@ -97,7 +114,8 @@ class AihrUploadQueueServiceTest { Class rowType = Class.forName("org.dromara.aihr.service.AihrUploadQueueService$ItemRow"); Constructor constructor = rowType.getDeclaredConstructors()[0]; constructor.setAccessible(true); - return constructor.newInstance("000000", 1L, "batch", "legacy.zip", "__auto__", java.util.List.of(), archive.toString(), 1); + return constructor.newInstance("000000", 1L, "batch", "legacy.zip", "__auto__", + java.util.List.of(), archive.toString(), null, 1); } private static int expandArchive(AihrUploadQueueService service, Object itemRow) throws Exception { diff --git a/backend/ruoyi-modules/ruoyi-aihr/src/test/java/org/dromara/aihr/service/AihrVisionImagePreprocessorTest.java b/backend/ruoyi-modules/ruoyi-aihr/src/test/java/org/dromara/aihr/service/AihrVisionImagePreprocessorTest.java new file mode 100644 index 00000000..0ec9860a --- /dev/null +++ b/backend/ruoyi-modules/ruoyi-aihr/src/test/java/org/dromara/aihr/service/AihrVisionImagePreprocessorTest.java @@ -0,0 +1,53 @@ +package org.dromara.aihr.service; + +import org.junit.jupiter.api.Tag; +import org.junit.jupiter.api.Test; + +import javax.imageio.ImageIO; +import java.awt.Color; +import java.awt.Graphics2D; +import java.awt.image.BufferedImage; +import java.io.ByteArrayInputStream; +import java.io.ByteArrayOutputStream; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; + +@Tag("dev") +class AihrVisionImagePreprocessorTest { + + @Test + void downsamplesLargeImagesBeforeVisionRequest() throws Exception { + BufferedImage source = new BufferedImage(3200, 2400, BufferedImage.TYPE_INT_RGB); + Graphics2D graphics = source.createGraphics(); + try { + graphics.setColor(Color.WHITE); + graphics.fillRect(0, 0, source.getWidth(), source.getHeight()); + graphics.setColor(Color.BLACK); + graphics.drawString("AIHR OCR fixture", 120, 160); + } finally { + graphics.dispose(); + } + ByteArrayOutputStream encoded = new ByteArrayOutputStream(); + ImageIO.write(source, "png", encoded); + + AihrVisionImagePreprocessor.PreparedImage prepared = + AihrVisionImagePreprocessor.prepare(encoded.toByteArray()); + BufferedImage decoded = ImageIO.read(new ByteArrayInputStream(prepared.bytes())); + + assertEquals("image/jpeg", prepared.mimeType()); + assertEquals(3200, prepared.originalWidth()); + assertEquals(2400, prepared.originalHeight()); + assertTrue(decoded.getWidth() <= AihrVisionImagePreprocessor.MAX_LONG_EDGE); + assertTrue(decoded.getHeight() <= AihrVisionImagePreprocessor.MAX_LONG_EDGE); + assertTrue((long) decoded.getWidth() * decoded.getHeight() <= AihrVisionImagePreprocessor.MAX_PIXELS); + assertTrue(prepared.bytes().length <= AihrVisionImagePreprocessor.MAX_ENCODED_BYTES); + } + + @Test + void rejectsUnknownImagePayloadWithoutForwardingIt() { + assertThrows(org.dromara.common.core.exception.ServiceException.class, + () -> AihrVisionImagePreprocessor.prepare("not-an-image".getBytes())); + } +} diff --git a/backend/script/sql/aihr_knowledge_mysql8.sql b/backend/script/sql/aihr_knowledge_mysql8.sql index a1f1a999..7a2a3795 100644 --- a/backend/script/sql/aihr_knowledge_mysql8.sql +++ b/backend/script/sql/aihr_knowledge_mysql8.sql @@ -199,6 +199,7 @@ CREATE TABLE IF NOT EXISTS `aihr_knowledge_upload_item` ( `category` varchar(100) DEFAULT '' COMMENT '目标分类', `space_codes_json` json DEFAULT NULL COMMENT '目标知识空间编码数组', `staging_path` varchar(500) NOT NULL COMMENT '暂存文件路径', + `source_attach_id` bigint DEFAULT NULL COMMENT '从既有附件重新解析时的附件ID', `status` tinyint NOT NULL DEFAULT 0 COMMENT '0待处理 1处理中 2完成 3失败', `error` varchar(500) DEFAULT NULL COMMENT '失败原因', `doc_id` varchar(80) DEFAULT NULL COMMENT '入库文档ID', @@ -207,7 +208,8 @@ CREATE TABLE IF NOT EXISTS `aihr_knowledge_upload_item` ( `update_time` datetime DEFAULT NULL COMMENT '更新时间', PRIMARY KEY (`id`), KEY `idx_aihr_upload_item_batch` (`tenant_id`, `batch_id`, `id`), - KEY `idx_aihr_upload_item_status` (`tenant_id`, `status`, `update_time`) + KEY `idx_aihr_upload_item_status` (`tenant_id`, `status`, `update_time`), + KEY `idx_aihr_upload_item_source` (`tenant_id`, `source_attach_id`, `id`) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci COMMENT='知识库批量上传队列'; CREATE TABLE IF NOT EXISTS `aihr_knowledge_space_grant` ( diff --git a/backend/script/sql/update/aihr_20260729_media_reprocess_mysql8.sql b/backend/script/sql/update/aihr_20260729_media_reprocess_mysql8.sql new file mode 100644 index 00000000..90a2bfdc --- /dev/null +++ b/backend/script/sql/update/aihr_20260729_media_reprocess_mysql8.sql @@ -0,0 +1,36 @@ +-- 资料处理媒体重试:记录从既有知识附件重新解析的来源,支持直接从 OSS 原件重新入队。 +-- 可重复执行;依赖 aihr_knowledge_upload_item。 + +SET @has_source_attach_id := ( + SELECT COUNT(*) FROM information_schema.columns + WHERE table_schema = DATABASE() + AND table_name = 'aihr_knowledge_upload_item' + AND column_name = 'source_attach_id' +); +SET @media_reprocess_ddl := IF( + @has_source_attach_id = 0, + 'ALTER TABLE aihr_knowledge_upload_item ADD COLUMN source_attach_id bigint DEFAULT NULL COMMENT ''从既有附件重新解析时的附件ID'' AFTER staging_path', + 'SELECT 1' +); +PREPARE media_reprocess_stmt FROM @media_reprocess_ddl; +EXECUTE media_reprocess_stmt; +DEALLOCATE PREPARE media_reprocess_stmt; + +SET @has_media_source_index := ( + SELECT COUNT(*) FROM information_schema.statistics + WHERE table_schema = DATABASE() + AND table_name = 'aihr_knowledge_upload_item' + AND index_name = 'idx_aihr_upload_item_source' +); +SET @media_reprocess_ddl := IF( + @has_media_source_index = 0, + 'ALTER TABLE aihr_knowledge_upload_item ADD KEY idx_aihr_upload_item_source (tenant_id, source_attach_id, id)', + 'SELECT 1' +); +PREPARE media_reprocess_stmt FROM @media_reprocess_ddl; +EXECUTE media_reprocess_stmt; +DEALLOCATE PREPARE media_reprocess_stmt; + +SET @has_source_attach_id := NULL; +SET @has_media_source_index := NULL; +SET @media_reprocess_ddl := NULL; diff --git a/docs/API_INTEGRATION.md b/docs/API_INTEGRATION.md index b8a4444b..761018fa 100644 --- a/docs/API_INTEGRATION.md +++ b/docs/API_INTEGRATION.md @@ -23,7 +23,7 @@ | 银城大喇叭 | 员工 `GET /api/aihr/broadcast/unread-count`、`GET /api/aihr/broadcast/messages?pageNum=&pageSize=`、`GET /api/aihr/broadcast/messages/{id}`、`POST /api/aihr/broadcast/messages/{id}/read`、`GET /api/aihr/broadcast/attachments/{id}/content`;员工详情可经统一 `POST /api/knowledge/query` 的 `broadcastMessageId` 发起文字追问;管理 `GET /api/aihr/broadcast/admin/messages`、`POST /api/aihr/broadcast/admin/attachments`、`GET /api/aihr/broadcast/admin/attachments/{id}`、`POST /api/aihr/broadcast/messages`、`POST /api/aihr/broadcast/messages/{id}/withdraw` | 发布支持全员或按部门/岗位/人员定向、必读、单个公司文件与异步提炼。文件支持 txt/md/PDF/Word/Excel/PPT,100MB 内;先受控写入 OSS,再异步解析并复用现有模型生成摘要,只有 `READY` 且未绑定、属于当前租户和上传人的文件才可随消息发布。文件提炼完成后异步生成生活顾问、保洁、保安、工程维修、财务、人力、运营、审计风控、管理层中有原文依据的岗位解读;员工详情附件返回 `insightStatus/defaultPerspectiveCode/perspectiveLabels/perspectives`,每项依据仅含原文段号与已校验的短引用,不返回提取全文或原始 OSS 地址。`defaultPerspectiveCode` 由服务端按当前 APP 手机号精确匹配在职岗位,客户端不能指定;所有有权限员工仍可查看并切换全部已生成视角。`PENDING/PARTIAL/FAILED` 均不阻断摘要、原文件下载、消息发布或追问。下载和追问每次重新校验当前 APP 在职身份、租户和消息状态,追问上下文由服务端拼接消息正文与提取文本,客户端仍只传 `broadcastMessageId`。定向目标在界面称“定向提醒”,用于必读与范围提示,不改变全租户公开频道的可见性。发布、阅读和撤回保持既有幂等与审计约束;管理接口只允许 `superadmin` 或 `hr_operator`。当前仍不包含消息修订、撤回后补推或短信/电话强触达。 | | 员工直通车 `/pages/user/direct/index`、管理端 `/content/direct` | 员工 `GET /api/aihr/direct/channels`、`POST /api/aihr/direct/feedback`、`GET /api/aihr/direct/mine`、`GET /api/aihr/direct/mine/{id}`;处理端 `GET /api/aihr/direct/admin/feedback`、`POST /api/aihr/direct/admin/feedback/{id}/reply` | 员工可选择总裁、财务、人力、审计、运营并点对点提交;反馈内容支持语音转文字输入(复用 `/api/ai/asr`;文字为主路径,录音不可用或权限失败仅提示改文字输入)。总裁/审计默认匿名;匿名仅表示业务处理界面不显示提交人,系统仍保存内部账号和姓名快照供本人查询、幂等与审计。`direct_president/direct_finance/direct_hr/direct_audit/direct_operations` 仅处理各自频道,`superadmin` 可处理全部;列表、详情和回复均由服务端按角色收窄。一个反馈只允许一次正式回复,员工可在“我的反馈”查看状态和回复;当前不扩展为工单 SLA、转派或多轮聊天。 | | 成果投稿(当前页面名“工作上报”)`/pages/user/report/index` | `POST /api/aihr/work-report/organize`、`POST /attachment`、`POST /reports`、`GET /reports/mine`;主管/运营另有列表和审核接口 | 只承载 CASE/VIDEO/SOP/KNOWLEDGE 四类投稿。会话式页面、无状态整理、附件、幂等正式提交和历史状态已部署;审核通过不自动入知识库。当前整理服务不读取图片/视频内容,只把附件名称作为不可信元数据;不得与日常工作记录或“今日工作成果”混用 | -| 资料处理 `/knowledge/processing` | `GET /api/knowledge/processing/overview`、`POST /api/knowledge/doc/upload-async`、`GET /api/knowledge/doc/upload-items`、`POST /api/knowledge/doc/upload-items/{id}/retry` | 已接入解析任务状态聚合;页面只保留“批量导入”,接口暂存+入队即秒回,后台 worker(并发 2)逐条解析/归类/向量化;ZIP 在 worker 内安全解压后把支持的子文件继续入同一批次队列,页面按批次轮询进度、失败可单文件重试;不提供浏览器目录选择或服务端目录导入入口 | +| 资料处理 `/knowledge/processing` | `GET /api/knowledge/processing/overview`、`POST /api/knowledge/doc/upload-async`、`GET /api/knowledge/doc/upload-items`、`POST /api/knowledge/doc/upload-items/{id}/retry`、`POST /api/knowledge/doc/processing-tasks/{attachmentId}/retry` | 已接入解析任务状态聚合;页面只保留“批量导入”,接口暂存+入队即秒回,后台 worker(并发 2)逐条解析/归类/向量化;ZIP 在 worker 内安全解压后把支持的子文件继续入同一批次队列,页面按批次轮询进度。零片段媒体不再记为完成或永久“等待解析”,而是保留为可重试失败;管理端可从原 OSS 文件重新入队,无需用户重复上传;不提供浏览器目录选择或服务端目录导入入口 | | 组织人员同步 | `POST /api/aihr/org/sync` | 从开放组织同步系统的 `/api/open/v1/sync/snapshot` 拉取 `company/department/employee/employee_project_assignment` 快照,分页参数使用 `limit`;员工手机号只落 `person_phone` 用于移动端身份映射,不在组织列表响应暴露;岗位识别 `position/job_title/post/job_name/role/title` 等字段。`dryRun` 必须显式传入 `true`(预检)或 `false`(写入),省略或传 `null` 直接拒绝;默认 `replaceExisting=true`,写入前必须先运行 `{"dryRun":true}`。dry-run 不访问本地快照表、不执行 DDL/写库,返回 `phoneLinked/maskedPhone/suspectText/warnings` 且 `syncedCount=0`。非 dry-run 写入(含 `replaceExisting=false`)遇到脱敏手机号一律默认拒绝,避免空值覆盖本地登录身份;覆盖写入遇到员工被跳过、手机号不完整、疑似乱码或同项目重复关系时同样默认拒绝;`allowPartialReplace=true` 只能在身份字段契约一致且异常逐项确认后使用,不能绕过主体身份漂移或重复成员关系,此时脱敏员工的本地既有有效手机号会被保留而非清空。相同员工可由多条有效项目分配展开为多个项目成员行,唯一约束为 `tenant_id + project_code + ext_party_id`。2026-07-25 生产安全状态:既有快照 3417 行、2938 人在职、3392 个手机号映射;姓名已按稳定 `employee_id` 定向补齐,未执行全量覆盖。当前上游返回 3424 名员工、3398 个脱敏手机号和 1 条重复项目成员关系,优先字段 `employee_number` 仅匹配既有主体 `7/3417`,稳定 `employee_id` 匹配 `3417/3417`;全量覆盖仍只允许 dry-run,日常岗位变化使用下方增量接口。 | | 组织人员增量同步 | `POST /api/aihr/org/sync-changes` | 主动拉取上游 `/sync/changes?resource_type=employee`,即使事件没有进入 outbox、`changed_fields` 为空,也会按 `resource_id` 回源 `/employees/{id}`,并结合任职快照只替换受影响员工。主体固定使用稳定 `employee.id`;上游只返回脱敏手机号时保留本地既有有效手机号,上游显式清空手机号时同步清除本地值(对应移动端登录身份失效),手机号字段整体缺失则拒绝写入;范围字段明确不一致视为迁出,写入时删除该员工本地旧记录,范围字段缺失则保守跳过不删。任职快照失败、资源 ID 不一致或项目关系重复时拒绝写入。首次调用必须传 `sinceTime`,后续可传响应的 `nextCursor`;仍须先 `dryRun=true` 再以相同起点执行 `dryRun=false`,且只处理当前租户有效组织绑定范围。 | | 移动端手机号登录 | `GET /resource/sms/code`、`POST /auth/mobile/sms-login` | 已复用 sms4j 阿里云配置 `config1` 和 RuoYi `sms` 授权策略;移动专用接口固定服务端默认租户,客户端不传也不能选择 `tenantId`;验证码按“默认租户 + 手机号”隔离。校验在手机号粒度的分布式锁内完成:仅匹配成功才消费,输错不会作废原验证码。移动端令牌只关联 `app_user`;若同一手机号存在任何后台/系统账号(即使同时存在 `app_user`),一律拒绝登录而不复用或并置身份;手机号完全不存在时才自动注册 `app_user`。`aihr.sms.dev-fixed-code` 非空时不真发短信、验证码固定(dev 默认 `123456`)。prod 默认关闭,试点期只有同时设置 `AIHR_SMS_DEV_FIXED_CODE` 与 `AIHR_SMS_PROD_FIXED_CODE_ENABLED=true` 才启用固定码。 | @@ -115,7 +115,7 @@ SOP 文档上传第三片已经落最小后端边界: | 能力 | 后端接口 | 处理 | |---|---|---| | 文档上传解析 | `POST /api/knowledge/doc/upload` | `multipart/form-data`,字段 `file` 和 `category`;支持 `.txt/.md/.markdown/.pdf/.doc/.docx/.xls/.xlsx/.ppt/.pptx` 及图片 `.jpg/.jpeg/.png/.gif/.webp/.bmp`、100MB 内;先写 `sys_oss`,再绑定 `aihr_knowledge_attach.oss_id` 并切分写入 `aihr_knowledge_fragment`;管理端该接口 timeout 为 180s,避免 PDF 同步解析/归类/向量化接近默认 50s 时被客户端断开 | -| 图片视觉 OCR | 同一上传/导入链路 | 上传图片时,若启用了 `aihr_model_config.category='vision'`(无则回退 `category='chat'`)的模型,会把图片编码为 base64 data URL 调用该供应商 OpenAI-compatible `/chat/completions` 提取文字,识别结果按普通正文切分写入 fragment;无可用视觉模型或 OCR 无结果时按「待处理」落库(attach status=0、0 片段),不算失败,启用视觉模型后重新上传同名文件即可解析。不依赖 Tesseract,识别质量由所配置视觉模型决定 | +| 图片视觉 OCR | 同一上传/导入链路 | 上传图片时,若启用了 `aihr_model_config.category='vision'`(无则回退 `category='chat'`)的模型,服务端先读取图片元数据并通过源采样把长边限制到 2048、像素限制到约 419 万,再统一压缩为不超过 5MB 的 JPEG data URL 调用供应商 OpenAI-compatible `/chat/completions`。模型未配置、HTTP/网络/超时错误分别写入安全原因码;OCR 成功但无可用文字时 attach/queue 均记为可重试失败(status=3、0 片段),由管理端从原 OSS 重试,不再显示成永久“等待解析”。不依赖 Tesseract,识别质量由所配置视觉模型决定 | | 智能归类与标签 | 同一上传/导入链路 | `category=__auto__` 时,解析正文后优先调用已启用的 `aihr_model_config.category='chat'` 模型生成分类、摘要、标签和归类理由,温度固定为 `0`;无模型或调用失败时按文件名/正文关键词兜底。最终分类写入 `aihr_knowledge_info/attach`,摘要和标签写入 `sys_oss.ext1` | | 重复资料处理 | 同一上传/导入链路 | 上传时计算原始文件 `aihrFileSha256`、解析文本 `aihrTextSha256` 和 `md5` 写入 `sys_oss.ext1`;重复判断优先按文件 SHA-256,其次按文本 SHA-256,最后用同名同大小兼容旧数据。命中重复时复用原附件并迁移到最新分类,删除其他重复附件和旧 fragment | | 归类稳定性 | 同一上传/导入链路 | 命中重复资料时优先复用已有 `sys_oss.ext1` 里的分类、摘要和标签;只有旧资料没有存过模型结果时才重新分析,避免同文件因重复上传或更换模型导致标签漂移 | @@ -129,8 +129,8 @@ SOP 文档上传第三片已经落最小后端边界: | 解析状态聚合 | `GET /api/knowledge/processing/overview` | 聚合 `aihr_knowledge_attach.status`、fragment 数、embedding 数、`sys_oss.ext1.fileSize`,生成资料处理页指标、分类、任务、链路和事件列表 | | 服务端目录导入(运维) | `POST /api/knowledge/doc/import-local` | JSON `{ directory, category, limit }`;`directory` 只能是 `AIHR_IMPORT_ROOT` / `aihr.import.root` 下的相对目录,默认根目录为 `./.data/import`;逐文件复用上传解析链路,同步执行,仅保留给运维/调试,不在资料处理页暴露 | | 服务端导入任务(运维) | `POST /api/knowledge/doc/import-local-task`、`GET /api/knowledge/doc/import-tasks`、`POST /api/knowledge/doc/import-tasks/{id}/cancel` | 启动后台目录导入并返回任务;任务写入 `aihr_knowledge_import_task`,调用方可查询总数、成功数、失败数、当前文件和进度,运行中任务可取消;仅保留给运维/调试,不在资料处理页暴露 | -| 批量异步上传 | `POST /api/knowledge/doc/upload-async`(multipart `file`+`category`+`batchId`)、`GET /api/knowledge/doc/upload-items?batchId=`、`POST /api/knowledge/doc/upload-items/{id}/retry` | 上传只做暂存(`aihr.upload.staging`,默认 `./.data/staging`)+ 写入 `aihr_knowledge_upload_item`(0待处理/1处理中/2完成/3失败)即返回;视频和 ZIP 单文件 ≤500MB,其他文件 ≤100MB。生产 Spring/Undertow multipart 上限为单文件 500MB、请求 520MB,防止视频/ZIP 在入队前被拦截;worker 限制最多 1000 个 ZIP 子文件、解压总量 ≤2GB,拒绝嵌套 ZIP/不安全路径并忽略未知格式;ZIP 文件名优先按 UTF-8 解码,旧版中文 Windows ZIP 自动回退 GBK,仍无法解压时返回安全提示;同名子文件自动追加编号并以单条批量写入入队,避免目录扁平化覆盖和半批次入队;完成/失败状态 CAS 写入,卡住 30 分钟(视频 120 分钟)由清扫重置,失败暂存文件保留 72h 供重试;同步接口 `POST /api/knowledge/doc/upload` 不支持 ZIP,保留给 SOP 页单文件即时预览 | -| 视频解析 | 走批量异步上传,支持 `.mp4/.mov/.avi/.mkv/.webm/.m4v`,单文件 ≤500MB、时长 ≤60 分钟 | 依赖服务器 ffmpeg/ffprobe;抽音轨按 5 分钟分段调 asr 模型转写(`[mm:ss]` 时间戳),按 30s 间隔(≤40 帧)抽关键帧调 vision 模型提取画面文字(相邻重复画面去重);两类文本合并后走归类/切片/向量化;无 asr 且无 vision 结果时按「待处理」落库不报错;视频同时只加工 1 个 | +| 批量异步上传 | `POST /api/knowledge/doc/upload-async`(multipart `file`+`category`+`batchId`)、`GET /api/knowledge/doc/upload-items?batchId=`、`POST /api/knowledge/doc/upload-items/{id}/retry`、`POST /api/knowledge/doc/processing-tasks/{attachmentId}/retry` | 上传只做暂存(`aihr.upload.staging`,默认 `./.data/staging`)+ 写入 `aihr_knowledge_upload_item`(0待处理/1处理中/2完成/3失败)即返回;`source_attach_id` 记录从原 OSS 重试的附件。视频和 ZIP 单文件 ≤500MB,其他文件 ≤100MB。生产 Spring/Undertow multipart 上限为单文件 500MB、请求 520MB;worker 限制最多 1000 个 ZIP 子文件、解压总量 ≤2GB,拒绝嵌套 ZIP/不安全路径并忽略未知格式;ZIP 文件名优先按 UTF-8 解码,旧版中文 Windows ZIP 自动回退 GBK;同名子文件自动追加编号并以单条批量写入入队;只有片段数大于 0 才记为完成并删除暂存文件,零片段或安全媒体错误保持失败态供重试;卡住 30 分钟(视频 120 分钟)由清扫重置,失败暂存文件保留 72h;同步接口 `POST /api/knowledge/doc/upload` 不支持 ZIP,保留给 SOP 页单文件即时预览 | +| 视频解析 | 走批量异步上传,支持 `.mp4/.mov/.avi/.mkv/.webm/.m4v`,单文件 ≤500MB、时长 ≤60 分钟 | 依赖服务器 ffmpeg/ffprobe;抽音轨按 5 分钟分段调 asr 模型转写(`[mm:ss]` 时间戳),按 30s 间隔(≤40 帧)抽关键帧调 vision 模型提取画面文字(相邻重复画面去重);两类文本合并后走归类/切片/向量化;无音轨、ASR 无文本且 vision 无结果时以 `MEDIA_NO_TEXT` 进入可重试失败态,不冒充仍在排队;视频同时只加工 1 个 | 模型能力第二片已经落最小后端边界: diff --git a/docs/BRD_PRODUCTION_MIGRATION_RUNBOOK.md b/docs/BRD_PRODUCTION_MIGRATION_RUNBOOK.md index 7a0f66c2..083f21c7 100644 --- a/docs/BRD_PRODUCTION_MIGRATION_RUNBOOK.md +++ b/docs/BRD_PRODUCTION_MIGRATION_RUNBOOK.md @@ -118,6 +118,7 @@ mysql --default-character-set=utf8mb4 "$DB_NAME" < backend/script/sql/update/aih mysql --default-character-set=utf8mb4 "$DB_NAME" < backend/script/sql/update/aihr_20260724_agent_orchestration_mysql8.sql mysql --default-character-set=utf8mb4 "$DB_NAME" < backend/script/sql/update/aihr_20260725_deepseek_v4_model_mysql8.sql mysql --default-character-set=utf8mb4 "$DB_NAME" < backend/script/sql/update/aihr_20260725_sop_answer_partial_evidence_prompt_mysql8.sql +mysql --default-character-set=utf8mb4 "$DB_NAME" < backend/script/sql/update/aihr_20260729_media_reprocess_mysql8.sql mysql --default-character-set=utf8mb4 "$DB_NAME" < backend/script/sql/aihr_personal_knowledge_mysql8.sql ``` @@ -171,7 +172,7 @@ WHERE table_schema = DATABASE() OR (table_name = 'aihr_sop_answer_review' AND column_name IN ('prompt_version', 'requester_ext_party_id')) OR (table_name = 'aihr_knowledge_info' AND column_name IN ('code', 'space_type', 'sensitivity_level', 'status')) OR (table_name = 'aihr_knowledge_attach' AND column_name = 'category_id') - OR (table_name = 'aihr_knowledge_upload_item' AND column_name = 'space_codes_json') + OR (table_name = 'aihr_knowledge_upload_item' AND column_name IN ('space_codes_json', 'source_attach_id')) OR (table_name = 'aihr_org_snapshot' AND column_name = 'hire_date') OR (table_name = 'aihr_candidate_material' AND column_name IN ('reviewer', 'reviewed_time')) OR (table_name = 'aihr_interview_result' AND column_name IN ('questions_json', 'reviewed_score', 'review_status', 'reviewer', 'review_note', 'reviewed_time')) @@ -190,6 +191,13 @@ WHERE table_schema = DATABASE() ) ORDER BY table_name, column_name; +SELECT index_name, GROUP_CONCAT(column_name ORDER BY seq_in_index) AS indexed_columns +FROM information_schema.statistics +WHERE table_schema = DATABASE() + AND table_name = 'aihr_knowledge_upload_item' + AND index_name = 'idx_aihr_upload_item_source' +GROUP BY index_name; + SELECT MIN(non_unique) AS non_unique, GROUP_CONCAT(column_name ORDER BY seq_in_index) AS indexed_columns, SUM(sub_part IS NOT NULL) AS prefix_columns @@ -253,7 +261,7 @@ GROUP BY d.dimension_code ORDER BY d.dimension_code; ``` -本批完整发布要求 `release-preflight` 的 `64/64` 必需表与列契约均存在,新增项包括 `aihr_practice_help_event`、`aihr_direct_feedback`、`aihr_broadcast_attachment`、`aihr_broadcast_message.attachment_id` 和 `aihr_agent_run`。此前 `57/57`、`62/62` 都只代表较窄的历史发布清单。2026-07-25 生产只读预检已达到 `64/64`,后端 JAR/AIHR 模块与发布源一致,并确认公司文件附件契约 `15/15` 及多岗位解读队列索引存在;这只证明迁移与产物已发布,不替代 Agent、训练或试点业务验收。Agent 审计字段应为 `14/14`,且表中不得出现问题、答案或附件列。工作上报幂等字段为 `2/2`,大喇叭发布/撤回审计字段为 `3/3`,其发布幂等索引为 `0 | tenant_id,published_by,publish_request_key | 0`,M1 必读/目标载荷字段契约为 `2/2`,两张目标表字段契约为 `16/16`,定向规则与对象快照唯一约束分别为 `4/4`、`3/3`,公司文件附件字段和消息绑定字段契约为 `15/15`,并应存在 `(insight_status,status,update_time)` 解读队列索引;公司消息追问会话字段为 `1/1`,工作助手新增字段为 `8/8`,会话证据字段为 `8/8`,专项批量/内容快照字段为 `8/8`,候选资料审核字段为 `2/2`,面试复核字段为 `6/6`,字段查询覆盖所有列,内置 Prompt 数量为 `6`。私有资料配置只应返回 `ruoyi-personal`、私有策略和 `configured` 状态,不显示访问密钥。五个岗位各有 `20` 个场景,`daily/special` 各有 `100` 道题;其中 `14` 个未完成正式审核的高风险场景及其 `28` 道题保持禁用。岗位/SOP/任务/资格表及新记忆表为空是允许的;结构、索引和内置内容通过只证明 schema/seed 与发布产物一致,不证明员工直达反馈、文件提炼、多岗位解读、员工确认记忆、完整个人知识库或严格试点已经完成业务验收。 +本批完整发布要求 `release-preflight` 的 `64/64` 必需表与列契约均存在,新增项包括 `aihr_practice_help_event`、`aihr_direct_feedback`、`aihr_broadcast_attachment`、`aihr_broadcast_message.attachment_id` 和 `aihr_agent_run`;资料重试还必须通过 `remote_media_reprocess_contract=1/1`,即 `source_attach_id` 与 `(tenant_id, source_attach_id, id)` 索引同时存在。此前 `57/57`、`62/62` 都只代表较窄的历史发布清单。2026-07-25 生产只读预检已达到 `64/64`,后端 JAR/AIHR 模块与发布源一致,并确认公司文件附件契约 `15/15` 及多岗位解读队列索引存在;这只证明迁移与产物已发布,不替代 Agent、训练或试点业务验收。Agent 审计字段应为 `14/14`,且表中不得出现问题、答案或附件列。工作上报幂等字段为 `2/2`,大喇叭发布/撤回审计字段为 `3/3`,其发布幂等索引为 `0 | tenant_id,published_by,publish_request_key | 0`,M1 必读/目标载荷字段契约为 `2/2`,两张目标表字段契约为 `16/16`,定向规则与对象快照唯一约束分别为 `4/4`、`3/3`,公司文件附件字段和消息绑定字段契约为 `15/15`,并应存在 `(insight_status,status,update_time)` 解读队列索引;公司消息追问会话字段为 `1/1`,工作助手新增字段为 `8/8`,会话证据字段为 `8/8`,专项批量/内容快照字段为 `8/8`,候选资料审核字段为 `2/2`,面试复核字段为 `6/6`,字段查询覆盖所有列,内置 Prompt 数量为 `6`。私有资料配置只应返回 `ruoyi-personal`、私有策略和 `configured` 状态,不显示访问密钥。五个岗位各有 `20` 个场景,`daily/special` 各有 `100` 道题;其中 `14` 个未完成正式审核的高风险场景及其 `28` 道题保持禁用。岗位/SOP/任务/资格表及新记忆表为空是允许的;结构、索引和内置内容通过只证明 schema/seed 与发布产物一致,不证明员工直达反馈、文件提炼、多岗位解读、员工确认记忆、完整个人知识库或严格试点已经完成业务验收。 ## 迁移后应用回归 diff --git a/frontend/src/api/aihr/processing.ts b/frontend/src/api/aihr/processing.ts index 5fd46acd..498a5b21 100644 --- a/frontend/src/api/aihr/processing.ts +++ b/frontend/src/api/aihr/processing.ts @@ -75,6 +75,13 @@ export type ProcessingOverview = { rules: ProcessingRules; }; +export type ProcessingRetryResponse = { + itemId: number; + batchId: string; + fileName: string; + status: string; +}; + export type LocalImportRequest = { directory: string; category: string; @@ -113,6 +120,14 @@ export function getProcessingOverview(): Promise> }); } +export function retryProcessingTask(id: number): Promise> { + return request({ + url: `/api/knowledge/doc/processing-tasks/${id}/retry`, + method: 'post', + headers: { repeatSubmit: false } + }); +} + export function importLocalKnowledgeDocs(data: LocalImportRequest): Promise> { return request({ url: '/api/knowledge/doc/import-local', diff --git a/frontend/src/utils/processing-media-retry-contract.test.ts b/frontend/src/utils/processing-media-retry-contract.test.ts new file mode 100644 index 00000000..3f02a69b --- /dev/null +++ b/frontend/src/utils/processing-media-retry-contract.test.ts @@ -0,0 +1,18 @@ +import { describe, expect, it } from 'vitest'; +import { readFileSync } from 'node:fs'; +import { resolve } from 'node:path'; + +const read = (path: string) => readFileSync(resolve(process.cwd(), path), 'utf8'); + +describe('资料处理媒体重试契约', () => { + it('把零文本媒体展示为待重试并从原文件重新入队', () => { + const api = read('src/api/aihr/processing.ts'); + const page = read('src/views/knowledge/processing.vue'); + + expect(api).toContain('/api/knowledge/doc/processing-tasks/${id}/retry'); + expect(page).toContain("const statusTabs = ['全部', '待重试', '解析中', '已完成', '失败']"); + expect(page).toContain('canRetryProcessingTask(row)'); + expect(page).toContain('await retryProcessingTask(task.id)'); + expect(page).toContain('已从原文件重新排队'); + }); +}); diff --git a/frontend/src/views/knowledge/processing.vue b/frontend/src/views/knowledge/processing.vue index c822cc49..e229bcf0 100644 --- a/frontend/src/views/knowledge/processing.vue +++ b/frontend/src/views/knowledge/processing.vue @@ -149,6 +149,20 @@ + + +