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 b4e58ace..261bce86 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 @@ -7,10 +7,13 @@ import org.dromara.aihr.domain.AihrSopDto.LocalImportTaskResponse; import org.dromara.aihr.domain.AihrSopDto.ProcessingOverviewResponse; import org.dromara.aihr.domain.AihrSopDto.SearchRequest; import org.dromara.aihr.domain.AihrSopDto.SearchResponse; +import org.dromara.aihr.domain.AihrSopDto.UploadEnqueueResponse; +import org.dromara.aihr.domain.AihrSopDto.UploadItemResponse; import org.dromara.aihr.domain.AihrSopDto.UploadResponse; import org.dromara.aihr.domain.AihrSopDto.VectorIndexStatusResponse; import org.dromara.aihr.domain.AihrSopDto.VectorizeResponse; import org.dromara.aihr.service.AihrSopSeedService; +import org.dromara.aihr.service.AihrUploadQueueService; import org.dromara.common.core.domain.R; import org.springframework.web.bind.annotation.GetMapping; import org.springframework.web.bind.annotation.PathVariable; @@ -25,7 +28,7 @@ import org.springframework.web.multipart.MultipartFile; import java.util.List; /** - * SOP 知识库 seed API。 + * SOP 知识库 API。 */ @RequiredArgsConstructor @RestController @@ -33,6 +36,25 @@ import java.util.List; public class AihrSopController { private final AihrSopSeedService sopSeedService; + private final AihrUploadQueueService uploadQueueService; + + @PostMapping("/doc/upload-async") + public R uploadAsync(@RequestPart("file") MultipartFile file, + @RequestParam(value = "category", required = false) String category, + @RequestParam(value = "batchId", required = false) String batchId) { + return R.ok(uploadQueueService.enqueue(file, category, batchId)); + } + + @GetMapping("/doc/upload-items") + public R> uploadItems(@RequestParam(value = "batchId", required = false) String batchId, + @RequestParam(value = "limit", required = false, defaultValue = "50") int limit) { + return R.ok(uploadQueueService.items(batchId, limit)); + } + + @PostMapping("/doc/upload-items/{id}/retry") + public R retryUploadItem(@PathVariable Long id) { + return R.ok(uploadQueueService.retry(id)); + } @PostMapping("/search") public R search(@RequestBody SearchRequest request) { diff --git a/backend/ruoyi-modules/ruoyi-aihr/src/main/java/org/dromara/aihr/domain/AihrSopDto.java b/backend/ruoyi-modules/ruoyi-aihr/src/main/java/org/dromara/aihr/domain/AihrSopDto.java index 09378fd1..a40f67c1 100644 --- a/backend/ruoyi-modules/ruoyi-aihr/src/main/java/org/dromara/aihr/domain/AihrSopDto.java +++ b/backend/ruoyi-modules/ruoyi-aihr/src/main/java/org/dromara/aihr/domain/AihrSopDto.java @@ -129,4 +129,21 @@ public final class AihrSopDto { Integer failedTasks ) { } + + public record UploadEnqueueResponse(Long itemId, String batchId, String fileName, String status) { + } + + public record UploadItemResponse( + Long id, + String batchId, + String fileName, + String category, + Integer status, + String statusLabel, + String error, + String docId, + Integer fragmentCount, + String updatedAt + ) { + } } 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 3d8b9950..0fb83984 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 @@ -23,6 +23,7 @@ import org.dromara.aihr.domain.AihrSopDto.SnippetResponse; import org.dromara.aihr.domain.AihrSopDto.UploadResponse; import org.dromara.aihr.domain.AihrSopDto.VectorIndexStatusResponse; import org.dromara.aihr.domain.AihrSopDto.VectorizeResponse; +import org.dromara.common.core.exception.ServiceException; import org.apache.tika.metadata.Metadata; import org.apache.tika.metadata.TikaCoreProperties; import org.apache.tika.parser.AutoDetectParser; @@ -122,19 +123,30 @@ public class AihrSopSeedService { public UploadResponse uploadDoc(MultipartFile file, String category) { if (file == null || file.isEmpty()) { - throw new IllegalArgumentException("上传文件不能为空"); + throw new ServiceException("上传文件不能为空"); } String fileName = Optional.ofNullable(file.getOriginalFilename()).orElse("knowledge.txt").trim(); if (!supportedFile(fileName)) { - throw new IllegalArgumentException("仅支持 txt/md/markdown/pdf/doc/docx/xls/xlsx/ppt/pptx 文件"); + throw new ServiceException("仅支持 txt/md/markdown/pdf/doc/docx/xls/xlsx/ppt/pptx/图片文件"); } if (file.getSize() > 100 * 1024 * 1024) { - throw new IllegalArgumentException("文件不能超过 100MB"); + throw new ServiceException("文件不能超过 100MB"); } return saveDocument(fileName, category, () -> ossService.upload(file), () -> readContent(file, fileName), () -> fileFingerprint(file)); } + /** + * 供后台上传队列(AihrUploadQueueService)复用的完整加工链路:OSS 上传、解析、归类、切片、向量化。 + * 与目录导入同源,入参为暂存文件路径。 + */ + public UploadResponse processStagedDocument(String fileName, String category, Path stagedFile) { + return saveDocument(fileName, category, + () -> ossService.upload(stagedFile.toFile()), + () -> readContent(stagedFile, fileName), + () -> fileFingerprint(stagedFile)); + } + public VectorizeResponse vectorizeMissing() { List targets = jdbcTemplate.query(""" select f.knowledge_id, @@ -454,7 +466,7 @@ public class AihrSopSeedService { try { content = contentReader.get(); if (content.isBlank()) { - throw new IllegalArgumentException("文件内容不能为空"); + throw new ServiceException("文件内容不能为空"); } } catch (RuntimeException e) { throw e; @@ -1967,7 +1979,7 @@ public class AihrSopSeedService { new AutoDetectParser().parse(input, handler, metadata, new ParseContext()); return normalizeExtractedText(handler.toString()); } catch (Exception e) { - throw new IllegalArgumentException("文件解析失败"); + throw new ServiceException("文件解析失败"); } } @@ -1995,7 +2007,7 @@ public class AihrSopSeedService { new AutoDetectParser().parse(input, handler, metadata, new ParseContext()); return normalizeExtractedText(handler.toString()); } catch (Exception e) { - throw new IllegalArgumentException("文件解析失败"); + throw new ServiceException("文件解析失败"); } } @@ -2003,7 +2015,7 @@ public class AihrSopSeedService { try (InputStream input = file.getInputStream()) { return fileFingerprint(input); } catch (IOException e) { - throw new IllegalArgumentException("文件指纹计算失败"); + throw new ServiceException("文件指纹计算失败"); } } @@ -2011,7 +2023,7 @@ public class AihrSopSeedService { try (InputStream input = Files.newInputStream(file)) { return fileFingerprint(input); } catch (IOException e) { - throw new IllegalArgumentException("文件指纹计算失败"); + throw new ServiceException("文件指纹计算失败"); } } @@ -2054,7 +2066,7 @@ public class AihrSopSeedService { try { return new String(file.getBytes(), StandardCharsets.UTF_8); } catch (Exception e) { - throw new IllegalArgumentException("文件读取失败"); + throw new ServiceException("文件读取失败"); } } @@ -2062,7 +2074,7 @@ public class AihrSopSeedService { try { return Files.readString(file, StandardCharsets.UTF_8); } catch (Exception e) { - throw new IllegalArgumentException("文件读取失败"); + throw new ServiceException("文件读取失败"); } } @@ -2140,7 +2152,7 @@ public class AihrSopSeedService { return AUTO_CATEGORY.equals(value) || "智能归类".equals(value) || "自动归类".equals(value); } - private static boolean supportedFile(String fileName) { + static boolean supportedFile(String fileName) { String lower = fileName.toLowerCase(); return textFile(lower) || imageFile(lower) 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 new file mode 100644 index 00000000..d7c3931c --- /dev/null +++ b/backend/ruoyi-modules/ruoyi-aihr/src/main/java/org/dromara/aihr/service/AihrUploadQueueService.java @@ -0,0 +1,395 @@ +package org.dromara.aihr.service; + +import jakarta.annotation.PostConstruct; +import lombok.extern.slf4j.Slf4j; +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.springframework.beans.factory.annotation.Value; +import org.springframework.dao.DataAccessException; +import org.springframework.jdbc.core.JdbcTemplate; +import org.springframework.stereotype.Service; +import org.springframework.web.multipart.MultipartFile; + +import java.io.IOException; +import java.nio.file.Files; +import java.nio.file.Path; +import java.nio.file.StandardCopyOption; +import java.sql.ResultSet; +import java.sql.SQLException; +import java.time.format.DateTimeFormatter; +import java.util.List; +import java.util.UUID; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.Semaphore; +import java.util.concurrent.TimeUnit; + +/** + * 批量上传队列:上传请求只做「暂存文件 + 入队 + 秒回」,解析/归类/向量化由本服务的后台 worker + * 复用 {@link AihrSopSeedService#processStagedDocument} 逐条加工。单文件失败不影响批次,支持按条重试。 + */ +@Service +@Slf4j +public class AihrUploadQueueService { + + private static final String TENANT_ID = "000000"; + private static final DateTimeFormatter TIME_FORMATTER = DateTimeFormatter.ofPattern("MM-dd HH:mm:ss"); + private static final long MAX_FILE_BYTES = 100L * 1024 * 1024; + private static final int MAX_ERROR_CHARS = 480; + /** 加工并发上限:链路含 LLM 归类与 embedding 外呼,并发过高会互相争抢配额并拖慢单文件耗时。 */ + private static final int WORKER_PERMITS = 2; + /** 处理中超过该时长视为孤儿(服务重启/线程异常丢失),由清扫任务重置回待处理。需大于最慢单文件加工时长。 */ + private static final int STUCK_MINUTES = 30; + /** 失败条目的暂存文件保留时长,超时由清扫回收磁盘(之后重试会提示重新上传)。 */ + private static final int STAGING_RETENTION_HOURS = 72; + + private final AihrSopSeedService sopSeedService; + private final JdbcTemplate jdbcTemplate; + private final ScheduledExecutorService scheduledExecutorService; + private final Path stagingRoot; + private final Semaphore workerPermits = new Semaphore(WORKER_PERMITS); + private volatile boolean tableReady; + + public AihrUploadQueueService(AihrSopSeedService sopSeedService, + JdbcTemplate jdbcTemplate, + ScheduledExecutorService scheduledExecutorService, + @Value("${aihr.upload.staging:./.data/staging}") String stagingRootConfig) { + this.sopSeedService = sopSeedService; + this.jdbcTemplate = jdbcTemplate; + this.scheduledExecutorService = scheduledExecutorService; + this.stagingRoot = Path.of(stagingRootConfig).toAbsolutePath().normalize(); + } + + @PostConstruct + void scheduleSweep() { + scheduledExecutorService.scheduleWithFixedDelay(this::sweepAndTrigger, 20, 15, TimeUnit.SECONDS); + } + + public UploadEnqueueResponse enqueue(MultipartFile file, String category, String batchId) { + if (file == null || file.isEmpty()) { + throw new ServiceException("上传文件不能为空"); + } + String fileName = sanitizeFileName(file.getOriginalFilename()); + if (!AihrSopSeedService.supportedFile(fileName)) { + throw new ServiceException("仅支持 txt/md/markdown/pdf/doc/docx/xls/xlsx/ppt/pptx/图片文件"); + } + if (file.getSize() > MAX_FILE_BYTES) { + throw new ServiceException("文件不能超过 100MB"); + } + ensureTable(); + String batch = sanitizeBatchId(batchId); + Path staged = stagingPath(fileName); + try (var input = file.getInputStream()) { + Files.createDirectories(staged.getParent()); + Files.copy(input, staged, StandardCopyOption.REPLACE_EXISTING); + } catch (IOException e) { + throw new ServiceException("暂存文件写入失败:" + e.getMessage()); + } + jdbcTemplate.update(""" + insert into aihr_knowledge_upload_item + (tenant_id, batch_id, file_name, category, staging_path, status, create_time, update_time) + values (?, ?, ?, ?, ?, 0, now(), now()) + """, TENANT_ID, batch, fileName, category == null ? "" : category.trim(), staged.toString()); + // staging_path 含 UUID 全局唯一,用它反查主键,避免 last_insert_id 的跨连接问题 + Long itemId = jdbcTemplate.query(""" + select id from aihr_knowledge_upload_item where tenant_id = ? and staging_path = ? + """, rs -> rs.next() ? rs.getLong(1) : null, TENANT_ID, staged.toString()); + triggerProcessing(); + return new UploadEnqueueResponse(itemId, batch, fileName, "queued"); + } + + public List items(String batchId, int limit) { + ensureTable(); + int safeLimit = Math.max(1, Math.min(limit <= 0 ? 50 : limit, 200)); + String batch = sanitizeBatchId(batchId); + if (batchId == null || batchId.isBlank()) { + return jdbcTemplate.query(""" + select id, batch_id, file_name, category, status, error, doc_id, fragment_count, update_time + from aihr_knowledge_upload_item + where tenant_id = ? + order by id desc + limit %d + """.formatted(safeLimit), this::mapItem, TENANT_ID); + } + return jdbcTemplate.query(""" + select id, batch_id, file_name, category, status, error, doc_id, fragment_count, update_time + from aihr_knowledge_upload_item + where tenant_id = ? and batch_id = ? + order by id asc + limit %d + """.formatted(safeLimit), this::mapItem, TENANT_ID, batch); + } + + public UploadItemResponse retry(Long id) { + ensureTable(); + if (id == null) { + throw new ServiceException("重试条目不存在"); + } + ItemRow row = jdbcTemplate.query(""" + select id, file_name, category, staging_path, status + from aihr_knowledge_upload_item + where tenant_id = ? and id = ? + """, rs -> rs.next() + ? new ItemRow(rs.getLong("id"), rs.getString("file_name"), rs.getString("category"), rs.getString("staging_path"), rs.getInt("status")) + : null, TENANT_ID, id); + if (row == null) { + throw new ServiceException("重试条目不存在"); + } + if (row.status() != 3) { + throw new ServiceException("仅失败状态的文件可重试"); + } + if (!Files.isRegularFile(Path.of(row.stagingPath()))) { + markFailed(id, "暂存文件已丢失,请重新上传该文件"); + throw new ServiceException("暂存文件已丢失,请重新上传该文件"); + } + int updated = jdbcTemplate.update(""" + update aihr_knowledge_upload_item + set status = 0, error = null, update_time = now() + where tenant_id = ? and id = ? and status = 3 + """, TENANT_ID, id); + if (updated == 0) { + throw new ServiceException("仅失败状态的文件可重试"); + } + triggerProcessing(); + return item(id); + } + + /** 领取待处理条目并在受限并发下加工;每领取一条占用一个许可,处理完释放后继续领取。 */ + private void triggerProcessing() { + while (workerPermits.tryAcquire()) { + Long claimed = claimNext(); + if (claimed == null) { + workerPermits.release(); + return; + } + CompletableFuture.runAsync(() -> { + try { + processItem(claimed); + } finally { + workerPermits.release(); + triggerProcessing(); + } + }, scheduledExecutorService); + } + } + + private Long claimNext() { + try { + // 同名文件不并发领取:saveDocument 对同 knowledge+fileName 是替换语义,并发会撞 attach 唯一键 + List candidates = jdbcTemplate.query(""" + select i.id from aihr_knowledge_upload_item i + where i.tenant_id = ? and i.status = 0 + and not exists ( + select 1 from (select file_name from aihr_knowledge_upload_item where tenant_id = ? and status = 1) p + where p.file_name = i.file_name + ) + order by i.id asc + limit 5 + """, (rs, rowNum) -> rs.getLong(1), TENANT_ID, TENANT_ID); + for (Long id : candidates) { + int updated = jdbcTemplate.update(""" + update aihr_knowledge_upload_item + set status = 1, update_time = now() + where tenant_id = ? and id = ? and status = 0 + """, TENANT_ID, id); + if (updated == 1) { + return id; + } + } + return null; + } catch (DataAccessException e) { + log.warn("upload queue claim failed: {}", e.getMessage()); + return null; + } + } + + private void processItem(Long id) { + ItemRow row = jdbcTemplate.query(""" + select id, file_name, category, staging_path, status + from aihr_knowledge_upload_item + where tenant_id = ? and id = ? + """, rs -> rs.next() + ? new ItemRow(rs.getLong("id"), rs.getString("file_name"), rs.getString("category"), rs.getString("staging_path"), rs.getInt("status")) + : null, TENANT_ID, id); + if (row == null) { + return; + } + Path staged = Path.of(row.stagingPath()); + if (!Files.isRegularFile(staged)) { + markFailed(id, "暂存文件已丢失,请重新上传该文件"); + return; + } + try { + UploadResponse result = sopSeedService.processStagedDocument(row.fileName(), row.category(), staged); + // CAS:仅当仍是本工人持有的「处理中」才写完成,防止清扫重置后被后来的工人覆盖状态 + int updated = jdbcTemplate.update(""" + update aihr_knowledge_upload_item + set status = 2, error = null, doc_id = ?, fragment_count = ?, update_time = now() + where tenant_id = ? and id = ? and status = 1 + """, result.docId(), result.fragments(), TENANT_ID, id); + if (updated == 1) { + deleteQuietly(staged); + } else { + log.info("upload item {} finished but row was re-claimed, skip status write", id); + } + } catch (Exception e) { + log.warn("upload item {} process failed: {}", id, e.getMessage()); + markFailed(id, e.getMessage()); + } + } + + /** 清扫卡在「处理中」的孤儿条目(重启/异常导致),重置回待处理后再触发一轮加工;顺带回收过期暂存文件。 */ + private void sweepAndTrigger() { + try { + ensureTable(); + int reset = jdbcTemplate.update(""" + update aihr_knowledge_upload_item + set status = 0, update_time = now() + where tenant_id = ? and status = 1 and update_time < (now() - interval %d minute) + """.formatted(STUCK_MINUTES), TENANT_ID); + if (reset > 0) { + log.info("upload queue reset {} stuck items", reset); + } + cleanExpiredStaging(); + triggerProcessing(); + } catch (Exception e) { + log.debug("upload queue sweep skipped: {}", e.getMessage()); + } + } + + /** 回收超过保留期的暂存文件,防失败条目/入库异常导致磁盘无限增长;活跃条目远不会到期。 */ + private void cleanExpiredStaging() { + if (!Files.isDirectory(stagingRoot)) { + return; + } + long cutoffMillis = System.currentTimeMillis() - STAGING_RETENTION_HOURS * 3600_000L; + try (var stream = Files.list(stagingRoot)) { + stream.filter(Files::isRegularFile) + .filter(path -> { + try { + return Files.getLastModifiedTime(path).toMillis() < cutoffMillis; + } catch (IOException e) { + return false; + } + }) + .forEach(AihrUploadQueueService::deleteQuietly); + } catch (IOException e) { + log.debug("staging cleanup skipped: {}", e.getMessage()); + } + } + + private void markFailed(Long id, String message) { + String error = message == null || message.isBlank() ? "处理失败" : message; + if (error.length() > MAX_ERROR_CHARS) { + error = error.substring(0, MAX_ERROR_CHARS); + } + // 不覆盖已被清扫重置(0)或已完成(2)的行:处理中标失败,失败态允许更新错误信息(重试路径) + jdbcTemplate.update(""" + update aihr_knowledge_upload_item + set status = 3, error = ?, update_time = now() + where tenant_id = ? and id = ? and status in (1, 3) + """, error, TENANT_ID, id); + } + + private UploadItemResponse item(Long id) { + List rows = jdbcTemplate.query(""" + select id, batch_id, file_name, category, status, error, doc_id, fragment_count, update_time + from aihr_knowledge_upload_item + where tenant_id = ? and id = ? + """, this::mapItem, TENANT_ID, id); + return rows.isEmpty() ? null : rows.get(0); + } + + private UploadItemResponse mapItem(ResultSet rs, int rowNum) throws SQLException { + int status = rs.getInt("status"); + java.sql.Timestamp updateTime = rs.getTimestamp("update_time"); + return new UploadItemResponse( + rs.getLong("id"), + rs.getString("batch_id"), + rs.getString("file_name"), + rs.getString("category"), + status, + statusLabel(status), + rs.getString("error"), + rs.getString("doc_id"), + rs.getObject("fragment_count", Integer.class), + updateTime == null ? "" : updateTime.toLocalDateTime().format(TIME_FORMATTER) + ); + } + + private static String statusLabel(int status) { + return switch (status) { + case 0 -> "排队中"; + case 1 -> "加工中"; + case 2 -> "已完成"; + case 3 -> "失败"; + default -> "未知"; + }; + } + + private Path stagingPath(String fileName) { + return stagingRoot.resolve(UUID.randomUUID() + "_" + fileName); + } + + /** 去掉路径分隔与控制字符,防止暂存路径逃逸出 stagingRoot。 */ + private static String sanitizeFileName(String original) { + String name = original == null || original.isBlank() ? "knowledge.txt" : original.trim(); + name = name.replaceAll("[\\\\/\\r\\n\\t]", "_"); + if (name.length() > 200) { + name = name.substring(name.length() - 200); + } + return name; + } + + private static String sanitizeBatchId(String batchId) { + if (batchId == null || batchId.isBlank()) { + return "b-" + UUID.randomUUID(); + } + String cleaned = batchId.trim().replaceAll("[^A-Za-z0-9_-]", ""); + return cleaned.isEmpty() || cleaned.length() > 64 ? "b-" + UUID.randomUUID() : cleaned; + } + + private static void deleteQuietly(Path path) { + try { + Files.deleteIfExists(path); + } catch (IOException ignored) { + // 暂存清理失败不影响业务,留给下次清扫。 + } + } + + private void ensureTable() { + if (tableReady) { + return; + } + synchronized (this) { + if (tableReady) { + return; + } + jdbcTemplate.execute(""" + CREATE TABLE IF NOT EXISTS `aihr_knowledge_upload_item` ( + `id` bigint NOT NULL AUTO_INCREMENT COMMENT '主键', + `tenant_id` varchar(20) DEFAULT '000000' COMMENT '租户编号', + `batch_id` varchar(64) NOT NULL COMMENT '上传批次', + `file_name` varchar(255) NOT NULL COMMENT '文件名', + `category` varchar(100) DEFAULT '' COMMENT '目标分类', + `staging_path` varchar(500) NOT NULL COMMENT '暂存文件路径', + `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', + `fragment_count` int DEFAULT NULL COMMENT '片段数', + `create_time` datetime DEFAULT NULL COMMENT '创建时间', + `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`) + ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci COMMENT='知识库批量上传队列'; + """); + tableReady = true; + } + } + + private record ItemRow(Long id, String fileName, String category, String stagingPath, int status) { + } +} diff --git a/backend/script/sql/aihr_knowledge_mysql8.sql b/backend/script/sql/aihr_knowledge_mysql8.sql index e22b54e1..25217ef4 100644 --- a/backend/script/sql/aihr_knowledge_mysql8.sql +++ b/backend/script/sql/aihr_knowledge_mysql8.sql @@ -124,3 +124,22 @@ VALUES ON DUPLICATE KEY UPDATE `content` = VALUES(`content`), `update_time` = NOW(); + +-- 知识库批量上传队列:上传秒回后由后台 worker 加工,支持单文件重试 +CREATE TABLE IF NOT EXISTS `aihr_knowledge_upload_item` ( + `id` bigint NOT NULL AUTO_INCREMENT COMMENT '主键', + `tenant_id` varchar(20) DEFAULT '000000' COMMENT '租户编号', + `batch_id` varchar(64) NOT NULL COMMENT '上传批次', + `file_name` varchar(255) NOT NULL COMMENT '文件名', + `category` varchar(100) DEFAULT '' COMMENT '目标分类', + `staging_path` varchar(500) NOT NULL COMMENT '暂存文件路径', + `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', + `fragment_count` int DEFAULT NULL COMMENT '片段数', + `create_time` datetime DEFAULT NULL COMMENT '创建时间', + `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`) +) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci COMMENT='知识库批量上传队列'; diff --git a/docs/API_INTEGRATION.md b/docs/API_INTEGRATION.md index c7e9c195..d5ab313a 100644 --- a/docs/API_INTEGRATION.md +++ b/docs/API_INTEGRATION.md @@ -6,14 +6,14 @@ | 页面 | 后端接口 | 处理 | |---|---|---| -| AI面试 `/recruit/interview` | `POST /api/recruit/interview/start`、`/answer`、`/finish` | 已接入 seed API | +| AI面试 `/recruit/interview` | `POST /api/recruit/interview/start`、`/answer`、`/finish` | 已接入真实模型优先链路:配置 chat 模型后 `/start` 动态生成面试题,`/finish` 按前端提交的真实回答做结构化评分;未配置模型时使用本地 Rubric 估分,不返回固定候选人分数 | | 三角色对练 `/train/practice` | `POST /api/train/practice/start`、`/turn`、`/finish` | 已接入编排 API;数据库启用 chat 模型后,`/turn` 客户回复按人设走真 LLM 生成(seed 剧本作剧情锚点),`/finish` 走单次 temperature=0 结构化评分(4 维分+导师改写+点评);模型未配置或调用失败自动回退 seed,契约不变 | | 对练语音 | `POST /api/ai/asr`(multipart 字段 `file`,≤5MB)、`POST /api/ai/tts`(JSON `{text≤300字, voice}`,返回 `{audioUrl}` base64 dataURL) | 走 OpenAI-compatible audio 接口(如硅基流动 SenseVoice/CosyVoice2);模型管理需启用 `category=asr/tts` 配置;未配置或失败返回 fail,前端降级文本 | -| 案例沉淀 `/knowledge/cases` | `POST /api/knowledge/case/upload`、`/organize`、`/curate` | 已接入 seed API | +| 案例沉淀 `/knowledge/cases` | `POST /api/knowledge/case/upload`、`/organize`、`/curate` | `/upload` 改为 multipart 真实语音上传并走 ASR;`/organize` 用真实转写调 chat 模型整理案例,未配置模型时按真实 transcript 本地结构化;不再用固定样例转写冒充成功 | | SOP知识库 `/knowledge/sop` | `POST /api/knowledge/search`、`POST /api/knowledge/doc/upload` | 已接入 MySQL Fulltext + Qdrant 混合召回、OSS-first 文档上传、txt/md/PDF/Word/Excel/PPT 解析和 embedding 写入,失败回退 seed | -| 资料处理 `/knowledge/processing` | `GET /api/knowledge/processing/overview`、`POST /api/knowledge/doc/import-local-task`、`GET /api/knowledge/doc/import-tasks`、`POST /api/knowledge/doc/import-tasks/{id}/cancel`、复用 `POST /api/knowledge/doc/upload` | 已接入解析任务状态聚合;页面支持多文件/目录选择、服务端后台目录导入、进度轮询和运行中任务取消,失败回退 seed | +| 资料处理 `/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/import-local-task`、`GET /api/knowledge/doc/import-tasks`、`POST /api/knowledge/doc/import-tasks/{id}/cancel` | 已接入解析任务状态聚合;**批量上传走异步队列**:接口只暂存+入队即秒回,后台 worker(并发 2)逐条解析/归类/向量化,页面按批次轮询进度、失败可单文件重试;服务端目录导入、进度轮询和任务取消保留,失败回退 seed | | 移动端手机号登录 | `GET /resource/sms/code`、`POST /auth/mobile/sms-login` | 已复用 sms4j 阿里云配置 `config1` 和 RuoYi `sms` 授权策略;短信发送成功后才写 Redis 验证码;手机号不存在时自动注册 `app_user`;`aihr.sms.dev-fixed-code` 非空时不真发短信、验证码固定(dev 默认 `123456`,prod profile 代码级强制失效) | -| 移动端三端首页 `/h5/user`、`/h5/candidate`、`/h5/supervisor` | `GET /api/aihr/mobile/home/{role}` | 已接入员工、候选人、主管首页 seed API;移动端本地 fallback 保演示 | +| 移动端三端首页 `/h5/user`、`/h5/candidate`、`/h5/supervisor` | `GET /api/aihr/mobile/home/{role}` | 已接入员工、候选人、主管首页公开只读 API;移动端本地 fallback 保演示 | | 移动端员工训练闭环 | 复用 `POST /api/train/practice/start`、`/turn`、`/finish`;查询 `GET /api/aihr/mobile/practice/history`、`/practice/reviews`、`/practice/reviews/{id}`、`/profile`;标记 `POST /api/aihr/mobile/practice/reviews/{id}/reviewed` | 员工端登录后带 `Authorization` 与 `clientid` 调用;`mode=mobile` 完成后写入 `aihr_practice_session`,主管端首页完训率、待复盘列表、复盘详情、员工训练历史和能力画像同步变化 | ## 后端落点 @@ -24,7 +24,7 @@ - 包名建议:`org.dromara.aihr` - Controller 返回统一用 `org.dromara.common.core.domain.R` -当前管理端页面流和移动端首页仍保留 seed fallback。知识库、模型能力、文档解析、RAG、chat 按 [ruoyi-ai 能力分片迁移计划](RUOYI_AI_INCREMENTAL_MIGRATION.md) 逐片引入;知识库 DDL 与住宅类 SOP seed 在 `backend/script/sql/aihr_knowledge_mysql8.sql`,模型 DDL 在 `backend/script/sql/aihr_model_mysql8.sql`,训练记录 DDL 在 `backend/script/sql/aihr_practice_mysql8.sql`,组织人员静态快照(2 个住宅项目 22 人,支撑按项目看人数)在 `backend/script/sql/aihr_org_snapshot_mysql8.sql`。 +当前管理端 AI 面试和案例沉淀已改为真实模型/ASR 优先;移动端首页和部分演示态数据仍保留 fallback。知识库、模型能力、文档解析、RAG、chat 按 [ruoyi-ai 能力分片迁移计划](RUOYI_AI_INCREMENTAL_MIGRATION.md) 逐片引入;知识库 DDL 与住宅类 SOP seed 在 `backend/script/sql/aihr_knowledge_mysql8.sql`,模型 DDL 在 `backend/script/sql/aihr_model_mysql8.sql`,训练记录 DDL 在 `backend/script/sql/aihr_practice_mysql8.sql`,组织人员静态快照(2 个住宅项目 22 人,支撑按项目看人数)在 `backend/script/sql/aihr_org_snapshot_mysql8.sql`。 直接打后端 `/api/**` 通常需要登录后的 `Authorization: Bearer `;浏览器内通过已登录前端和 `/dev-api` 代理访问。移动端登录接口为 `POST /auth/mobile/sms-login`,请求 `{ phonenumber, smsCode, tenantId }`,内部固定使用 app 客户端 `428a8310cd442757ae699df5d894f051` 和 `sms` grant;验证码通过后若手机号不存在,会创建 `app_user`,用户名为手机号,备注为“移动端短信自动注册”。移动端 MVP 首页接口 `GET /api/aihr/mobile/home/{role}` 目前仍是 `@SaIgnore` 的公开只读 seed 接口,避免 H5 首屏被后台管理登录态阻断;后续接小程序登录后再收紧为移动端 token。 @@ -55,6 +55,7 @@ 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失败)即返回;后台 worker 并发 2 复用同一解析/归类/向量化链路,同名文件串行防撞唯一键;完成/失败状态 CAS 写入,卡住 30 分钟由清扫重置,失败暂存文件保留 72h 供重试;同步接口 `POST /api/knowledge/doc/upload` 保留给 SOP 页单文件即时预览 | 模型能力第二片已经落最小后端边界: @@ -91,8 +92,8 @@ Qdrant 本地默认值可不配;需要覆盖时用 JVM property 或环境变 - 模型配置 API 文件:`frontend/src/api/aihr/model.ts` - 模型配置页面:`frontend/src/views/system/model/index.vue` - 移动端独立工程:`mobile/src/App.vue`,路由 `/h5/user`、`/h5/candidate`、`/h5/supervisor`;当前三端首页已按高保真原型实现,未登录先走手机号短信登录,登录后 API 优先请求 `/api/aihr/mobile/home/{role}`,不复用后台页面路由。 -- 员工端训练:任务卡“开始训练”复用管理端三角色对练 seed API,完成两轮后展示评分和导师改写;记录写入 `aihr_practice_session`,员工端展示训练历史和能力画像,主管端展示待复盘列表,点进单条可看评分、话术、导师改写,并可标记“已复盘”。 -- 保留本地 seed fallback:接口失败时仍能演示,不让现场演示被后端状态拖死。 +- 员工端训练:任务卡“开始训练”复用管理端三角色对练 API,完成两轮后展示评分和导师改写;记录写入 `aihr_practice_session`,员工端展示训练历史和能力画像,主管端展示待复盘列表,点进单条可看评分、话术、导师改写,并可标记“已复盘”。 +- 演示 fallback 只作为未配置模型或外部接口异常时的降级,不作为真实测试通过证据。 ## 验收 diff --git a/docs/DEV_SETUP.md b/docs/DEV_SETUP.md index a849f67f..9b381f2c 100644 --- a/docs/DEV_SETUP.md +++ b/docs/DEV_SETUP.md @@ -109,11 +109,11 @@ curl -fsS http://127.0.0.1:8080/api/knowledge/doc/vector-index-status -H "Author 当前已跑通管理端五个本地演示流和移动端员工训练闭环: -- AI面试:进入 `/recruit/interview`,点击“生成题目” → “填满 seed 回答” → “完成评分”,应看到“已完成闭环”“建议复试”和新增面试记录。 -- 三角色对练:进入 `/train/practice`,点击“开始对练” → 两次“填入 seed 回复 / 继续一轮” → “结束并评分”,应看到“已完成闭环”“导师改写”和新增对练记录。 -- 案例沉淀:进入 `/knowledge/cases`,点击“选择样例录音” → “AI 整理” → “送审” → “入库”,应看到“已完成闭环”和新增 `seed入库` 记录。 +- AI面试:进入 `/recruit/interview`,点击“生成题目” → 输入真实回答或“填满参考回答” → “完成评分”,应看到新增面试记录;配置 chat 模型时题目和评分都来自真实模型,未配置时使用本地 Rubric。 +- 三角色对练:进入 `/train/practice`,点击“开始对练” → 完成两轮真实或参考回复 → “结束并评分”,应看到“已完成闭环”“导师改写”和新增对练记录。 +- 案例沉淀:进入 `/knowledge/cases`,上传真实语音 → “AI 整理” → “送审” → “入库”,应看到“已完成闭环”和新增案例记录;上传必须先完成 ASR 转写。 - SOP知识库:进入 `/knowledge/sop`,可上传 txt/md/PDF/Word/Excel/PPT 文档入库;点击“检索” → “生成训练题”,应看到“已完成闭环”、命中数据库 SOP 原文片段和训练题;数据库不可用时页面回退 seed。 -- 资料处理:进入 `/knowledge/processing`,应看到资料总量、解析任务表、处理链路、规则与风险;可用少量文件验证“批量导入/选择目录/服务端导入”。服务端导入读取 `./.data/import` 下的相对目录,启动后台任务并在页面显示进度,运行中任务可点“取消”;当前重试粒度是同目录重新导入,不是单失败文件重试。若浏览器上传 PDF 约 50s 后后端出现 `Broken pipe`,先确认前端是否已使用 `uploadKnowledgeDoc` 的 180s timeout。 +- 资料处理:进入 `/knowledge/processing`,应看到资料总量、解析任务表、处理链路、规则与风险;可用少量文件验证“批量导入/选择目录/服务端导入”。**批量导入走异步队列**:提交即返回,页面出现“本次批量上传”进度面板(排队/加工中/完成/失败 + 单条重试),后台 worker 并发 2 逐条解析入库;暂存目录默认 `./.data/staging`(`aihr.upload.staging` 覆盖)。服务端导入读取 `./.data/import` 下的相对目录,启动后台任务并在页面显示进度,运行中任务可点“取消”;目录导入的重试粒度是同目录重新导入,批量上传的重试粒度是单文件。 - 移动端员工训练:进入 `http://127.0.0.1:5174/h5/user`,手机号登录(dev 验证码固定 `123456`)→ “开始训练” → 两次“填入建议回复 / 提交本轮”,应看到评分、导师改写、训练历史和能力画像;切到主管端后完训人数与“待复盘对练”计数增加,并展示待复盘列表;点进单条可看评分、话术、导师改写,并可标记已复盘。训练页有“语音输入”和“播报”按钮,需启用 asr/tts 模型后生效,否则降级文本。 - 真 LLM 激活:在 `/system/model` 给供应商填 api_host/api_key 并启用 `category=chat` 模型后,三角色对练的客户回复与评分即为真实 LLM 生成;再启用 `category=asr/tts`(如硅基流动 SenseVoice/CosyVoice2)语音路径生效。未配置时全链路自动回退 seed。 diff --git a/frontend/src/api/aihr/sop.ts b/frontend/src/api/aihr/sop.ts index 46af609d..cd6395d0 100644 --- a/frontend/src/api/aihr/sop.ts +++ b/frontend/src/api/aihr/sop.ts @@ -67,3 +67,51 @@ export function uploadKnowledgeDoc(data: FormData): Promise> { + return request({ + url: '/api/knowledge/doc/upload-async', + method: 'post', + data, + timeout: 60000, + headers: { repeatSubmit: false } + }); +} + +export function listUploadItems(batchId: string): Promise> { + return request({ + url: '/api/knowledge/doc/upload-items', + method: 'get', + params: { batchId }, + headers: { repeatSubmit: false } + }); +} + +export function retryUploadItem(id: number): Promise> { + return request({ + url: `/api/knowledge/doc/upload-items/${id}/retry`, + method: 'post', + headers: { repeatSubmit: false } + }); +} diff --git a/frontend/src/views/knowledge/processing.vue b/frontend/src/views/knowledge/processing.vue index b5d4dc54..7665f03b 100644 --- a/frontend/src/views/knowledge/processing.vue +++ b/frontend/src/views/knowledge/processing.vue @@ -20,6 +20,45 @@ +
+
+

本次批量上传

+
+ 排队 {{ batchCounts.queued }} + 加工中 {{ batchCounts.processing }} + 完成 {{ batchCounts.done }} + 失败 {{ batchCounts.failed }} + 收起 +
+
+ +
    +
  • 提交失败:{{ text }}
  • +
+ + + + + + + + + + + + +
+
@@ -148,7 +187,11 @@ 图片 OCR 即将支持 - +