fix(aihr): preserve tenant in upload queue

This commit is contained in:
2026-07-14 06:49:38 +08:00
parent 7eb787995d
commit 6ec1772674
4 changed files with 79 additions and 40 deletions
@@ -37,7 +37,7 @@ import org.dromara.aihr.domain.AihrSopDto.UploadResponse;
import org.dromara.aihr.domain.AihrSopDto.VectorIndexStatusResponse; import org.dromara.aihr.domain.AihrSopDto.VectorIndexStatusResponse;
import org.dromara.aihr.domain.AihrSopDto.VectorizeResponse; import org.dromara.aihr.domain.AihrSopDto.VectorizeResponse;
import org.dromara.common.core.exception.ServiceException; import org.dromara.common.core.exception.ServiceException;
import org.dromara.common.satoken.utils.LoginHelper; import org.dromara.common.tenant.helper.TenantHelper;
import org.apache.tika.metadata.Metadata; import org.apache.tika.metadata.Metadata;
import org.apache.tika.metadata.TikaCoreProperties; import org.apache.tika.metadata.TikaCoreProperties;
import org.apache.tika.parser.AutoDetectParser; import org.apache.tika.parser.AutoDetectParser;
@@ -3245,7 +3245,7 @@ public class AihrSopSeedService {
} }
private String tenantId() { private String tenantId() {
return firstNonBlank(LoginHelper.getTenantId(), TENANT_ID); return firstNonBlank(TenantHelper.getTenantId(), TENANT_ID);
} }
private static String normalizeBaseUrl(String baseUrl) { private static String normalizeBaseUrl(String baseUrl) {
@@ -6,6 +6,8 @@ import org.dromara.aihr.domain.AihrSopDto.UploadEnqueueResponse;
import org.dromara.aihr.domain.AihrSopDto.UploadItemResponse; import org.dromara.aihr.domain.AihrSopDto.UploadItemResponse;
import org.dromara.aihr.domain.AihrSopDto.UploadResponse; import org.dromara.aihr.domain.AihrSopDto.UploadResponse;
import org.dromara.common.core.exception.ServiceException; import org.dromara.common.core.exception.ServiceException;
import org.dromara.common.satoken.utils.LoginHelper;
import org.dromara.common.tenant.helper.TenantHelper;
import org.springframework.beans.factory.annotation.Value; import org.springframework.beans.factory.annotation.Value;
import org.springframework.dao.DataAccessException; import org.springframework.dao.DataAccessException;
import org.springframework.jdbc.core.JdbcTemplate; import org.springframework.jdbc.core.JdbcTemplate;
@@ -108,11 +110,11 @@ public class AihrUploadQueueService {
insert into aihr_knowledge_upload_item insert into aihr_knowledge_upload_item
(tenant_id, batch_id, file_name, category, staging_path, status, create_time, update_time) (tenant_id, batch_id, file_name, category, staging_path, status, create_time, update_time)
values (?, ?, ?, ?, ?, 0, now(), now()) values (?, ?, ?, ?, ?, 0, now(), now())
""", TENANT_ID, batch, fileName, category == null ? "" : category.trim(), staged.toString()); """, currentTenantId(), batch, fileName, category == null ? "" : category.trim(), staged.toString());
// staging_path 含 UUID 全局唯一,用它反查主键,避免 last_insert_id 的跨连接问题 // staging_path 含 UUID 全局唯一,用它反查主键,避免 last_insert_id 的跨连接问题
Long itemId = jdbcTemplate.query(""" Long itemId = jdbcTemplate.query("""
select id from aihr_knowledge_upload_item where tenant_id = ? and staging_path = ? select id from aihr_knowledge_upload_item where tenant_id = ? and staging_path = ?
""", rs -> rs.next() ? rs.getLong(1) : null, TENANT_ID, staged.toString()); """, rs -> rs.next() ? rs.getLong(1) : null, currentTenantId(), staged.toString());
triggerProcessing(); triggerProcessing();
return new UploadEnqueueResponse(itemId, batch, fileName, "queued"); return new UploadEnqueueResponse(itemId, batch, fileName, "queued");
} }
@@ -128,7 +130,7 @@ public class AihrUploadQueueService {
where tenant_id = ? where tenant_id = ?
order by id desc order by id desc
limit %d limit %d
""".formatted(safeLimit), this::mapItem, TENANT_ID); """.formatted(safeLimit), this::mapItem, currentTenantId());
} }
return jdbcTemplate.query(""" return jdbcTemplate.query("""
select id, batch_id, file_name, category, status, error, doc_id, fragment_count, update_time select id, batch_id, file_name, category, status, error, doc_id, fragment_count, update_time
@@ -136,7 +138,7 @@ public class AihrUploadQueueService {
where tenant_id = ? and batch_id = ? where tenant_id = ? and batch_id = ?
order by id asc order by id asc
limit %d limit %d
""".formatted(safeLimit), this::mapItem, TENANT_ID, batch); """.formatted(safeLimit), this::mapItem, currentTenantId(), batch);
} }
public UploadItemResponse retry(Long id) { public UploadItemResponse retry(Long id) {
@@ -144,13 +146,14 @@ public class AihrUploadQueueService {
if (id == null) { if (id == null) {
throw new ServiceException("重试条目不存在"); throw new ServiceException("重试条目不存在");
} }
String tenantId = currentTenantId();
ItemRow row = jdbcTemplate.query(""" ItemRow row = jdbcTemplate.query("""
select id, batch_id, file_name, category, staging_path, status select tenant_id, id, batch_id, file_name, category, staging_path, status
from aihr_knowledge_upload_item from aihr_knowledge_upload_item
where tenant_id = ? and id = ? where tenant_id = ? and id = ?
""", rs -> rs.next() """, rs -> rs.next()
? new ItemRow(rs.getLong("id"), rs.getString("batch_id"), rs.getString("file_name"), rs.getString("category"), rs.getString("staging_path"), rs.getInt("status")) ? new ItemRow(rs.getString("tenant_id"), rs.getLong("id"), rs.getString("batch_id"), rs.getString("file_name"), rs.getString("category"), rs.getString("staging_path"), rs.getInt("status"))
: null, TENANT_ID, id); : null, tenantId, id);
if (row == null) { if (row == null) {
throw new ServiceException("重试条目不存在"); throw new ServiceException("重试条目不存在");
} }
@@ -158,14 +161,14 @@ public class AihrUploadQueueService {
throw new ServiceException("仅失败状态的文件可重试"); throw new ServiceException("仅失败状态的文件可重试");
} }
if (!Files.isRegularFile(Path.of(row.stagingPath()))) { if (!Files.isRegularFile(Path.of(row.stagingPath()))) {
markFailed(id, "暂存文件已丢失,请重新上传该文件"); markFailed(tenantId, id, "暂存文件已丢失,请重新上传该文件");
throw new ServiceException("暂存文件已丢失,请重新上传该文件"); throw new ServiceException("暂存文件已丢失,请重新上传该文件");
} }
int updated = jdbcTemplate.update(""" int updated = jdbcTemplate.update("""
update aihr_knowledge_upload_item update aihr_knowledge_upload_item
set status = 0, error = null, update_time = now() set status = 0, error = null, update_time = now()
where tenant_id = ? and id = ? and status = 3 where tenant_id = ? and id = ? and status = 3
""", TENANT_ID, id); """, tenantId, id);
if (updated == 0) { if (updated == 0) {
throw new ServiceException("仅失败状态的文件可重试"); throw new ServiceException("仅失败状态的文件可重试");
} }
@@ -176,7 +179,7 @@ public class AihrUploadQueueService {
/** 领取待处理条目并在受限并发下加工;每领取一条占用一个许可,处理完释放后继续领取。 */ /** 领取待处理条目并在受限并发下加工;每领取一条占用一个许可,处理完释放后继续领取。 */
private void triggerProcessing() { private void triggerProcessing() {
while (workerPermits.tryAcquire()) { while (workerPermits.tryAcquire()) {
Long claimed = claimNext(); QueueCandidate claimed = claimNext();
if (claimed == null) { if (claimed == null) {
workerPermits.release(); workerPermits.release();
return; return;
@@ -192,31 +195,31 @@ public class AihrUploadQueueService {
} }
} }
private Long claimNext() { private QueueCandidate claimNext() {
try { try {
// 同名文件不并发领取:saveDocument 对同 knowledge+fileName 是替换语义,并发会撞 attach 唯一键。 // 同名文件不并发领取:saveDocument 对同 knowledge+fileName 是替换语义,并发会撞 attach 唯一键。
// 视频最多同时加工 1 个:单个视频要占数分钟的 ffmpeg+ASR+vision,双路并发互相拖慢。 // 视频最多同时加工 1 个:单个视频要占数分钟的 ffmpeg+ASR+vision,双路并发互相拖慢。
boolean videoBusy = hasProcessingVideo(); boolean videoBusy = hasProcessingVideo();
List<QueueCandidate> candidates = jdbcTemplate.query(""" List<QueueCandidate> candidates = jdbcTemplate.query("""
select i.id, i.file_name from aihr_knowledge_upload_item i select i.tenant_id, i.id, i.file_name from aihr_knowledge_upload_item i
where i.tenant_id = ? and i.status = 0 where i.status = 0
and (? = 0 or i.file_name not regexp '\\\\.(mp4|mov|avi|mkv|webm|m4v)$') and (? = 0 or i.file_name not regexp '\\\\.(mp4|mov|avi|mkv|webm|m4v)$')
and not exists ( and not exists (
select 1 from (select file_name from aihr_knowledge_upload_item where tenant_id = ? and status = 1) p select 1 from aihr_knowledge_upload_item p
where p.file_name = i.file_name where p.tenant_id = i.tenant_id and p.status = 1 and p.file_name = i.file_name
) )
order by i.id asc order by i.id asc
limit 10 limit 10
""", (rs, rowNum) -> new QueueCandidate(rs.getLong(1), rs.getString(2)), """, (rs, rowNum) -> new QueueCandidate(rs.getString(1), rs.getLong(2), rs.getString(3)),
TENANT_ID, videoBusy ? 1 : 0, TENANT_ID); videoBusy ? 1 : 0);
for (QueueCandidate candidate : candidates) { for (QueueCandidate candidate : candidates) {
int updated = jdbcTemplate.update(""" int updated = jdbcTemplate.update("""
update aihr_knowledge_upload_item update aihr_knowledge_upload_item
set status = 1, update_time = now() set status = 1, update_time = now()
where tenant_id = ? and id = ? and status = 0 where tenant_id = ? and id = ? and status = 0
""", TENANT_ID, candidate.id()); """, candidate.tenantId(), candidate.id());
if (updated == 1) { if (updated == 1) {
return candidate.id(); return candidate;
} }
} }
return null; return null;
@@ -228,28 +231,30 @@ public class AihrUploadQueueService {
private boolean hasProcessingVideo() { private boolean hasProcessingVideo() {
List<String> processing = jdbcTemplate.query(""" List<String> processing = jdbcTemplate.query("""
select file_name from aihr_knowledge_upload_item where tenant_id = ? and status = 1 select file_name from aihr_knowledge_upload_item where status = 1
""", (rs, rowNum) -> rs.getString(1), TENANT_ID); """, (rs, rowNum) -> rs.getString(1));
return processing.stream().anyMatch(AihrVideoService::videoFile); return processing.stream().anyMatch(AihrVideoService::videoFile);
} }
private record QueueCandidate(Long id, String fileName) { private record QueueCandidate(String tenantId, Long id, String fileName) {
} }
private void processItem(Long id) { private void processItem(QueueCandidate claim) {
String tenantId = claim.tenantId();
Long id = claim.id();
ItemRow row = jdbcTemplate.query(""" ItemRow row = jdbcTemplate.query("""
select id, batch_id, file_name, category, staging_path, status select tenant_id, id, batch_id, file_name, category, staging_path, status
from aihr_knowledge_upload_item from aihr_knowledge_upload_item
where tenant_id = ? and id = ? where tenant_id = ? and id = ?
""", rs -> rs.next() """, rs -> rs.next()
? new ItemRow(rs.getLong("id"), rs.getString("batch_id"), rs.getString("file_name"), rs.getString("category"), rs.getString("staging_path"), rs.getInt("status")) ? new ItemRow(rs.getString("tenant_id"), rs.getLong("id"), rs.getString("batch_id"), rs.getString("file_name"), rs.getString("category"), rs.getString("staging_path"), rs.getInt("status"))
: null, TENANT_ID, id); : null, tenantId, id);
if (row == null) { if (row == null) {
return; return;
} }
Path staged = Path.of(row.stagingPath()); Path staged = Path.of(row.stagingPath());
if (!Files.isRegularFile(staged)) { if (!Files.isRegularFile(staged)) {
markFailed(id, "暂存文件已丢失,请重新上传该文件"); markFailed(tenantId, id, "暂存文件已丢失,请重新上传该文件");
return; return;
} }
try { try {
@@ -259,13 +264,14 @@ public class AihrUploadQueueService {
update aihr_knowledge_upload_item update aihr_knowledge_upload_item
set status = 2, fragment_count = ?, update_time = now() set status = 2, fragment_count = ?, update_time = now()
where tenant_id = ? and id = ? and status = 1 where tenant_id = ? and id = ? and status = 1
""", extracted, TENANT_ID, id); """, extracted, tenantId, id);
if (updated == 1) { if (updated == 1) {
deleteQuietly(staged); deleteQuietly(staged);
} }
return; return;
} }
UploadResponse result = sopSeedService.processStagedDocument(row.fileName(), row.category(), staged); UploadResponse result = TenantHelper.dynamic(tenantId,
() -> sopSeedService.processStagedDocument(row.fileName(), row.category(), staged));
// 0 片段的完成态(如图片待 OCR)把说明写进 error 列,面板可见原因 // 0 片段的完成态(如图片待 OCR)把说明写进 error 列,面板可见原因
String note = result.fragments() != null && result.fragments() == 0 ? truncateError(result.summary()) : null; String note = result.fragments() != null && result.fragments() == 0 ? truncateError(result.summary()) : null;
// CAS:仅当仍是本工人持有的「处理中」才写完成,防止清扫重置后被后来的工人覆盖状态 // CAS:仅当仍是本工人持有的「处理中」才写完成,防止清扫重置后被后来的工人覆盖状态
@@ -273,7 +279,7 @@ public class AihrUploadQueueService {
update aihr_knowledge_upload_item update aihr_knowledge_upload_item
set status = 2, error = ?, doc_id = ?, fragment_count = ?, update_time = now() set status = 2, error = ?, doc_id = ?, fragment_count = ?, update_time = now()
where tenant_id = ? and id = ? and status = 1 where tenant_id = ? and id = ? and status = 1
""", note, result.docId(), result.fragments(), TENANT_ID, id); """, note, result.docId(), result.fragments(), tenantId, id);
if (updated == 1) { if (updated == 1) {
deleteQuietly(staged); deleteQuietly(staged);
} else { } else {
@@ -281,7 +287,7 @@ public class AihrUploadQueueService {
} }
} catch (Exception e) { } catch (Exception e) {
log.warn("upload item {} process failed: {}", id, e.getMessage()); log.warn("upload item {} process failed: {}", id, e.getMessage());
markFailed(id, e.getMessage()); markFailed(tenantId, id, e.getMessage());
} }
} }
@@ -341,12 +347,12 @@ public class AihrUploadQueueService {
insert into aihr_knowledge_upload_item insert into aihr_knowledge_upload_item
(tenant_id, batch_id, file_name, category, staging_path, status, create_time, update_time) (tenant_id, batch_id, file_name, category, staging_path, status, create_time, update_time)
values (?, ?, ?, ?, ?, 0, now(), now()) values (?, ?, ?, ?, ?, 0, now(), now())
""", TENANT_ID, archive.batchId(), file.fileName(), archive.category(), file.stagedPath().toString()); """, archive.tenantId(), archive.batchId(), file.fileName(), archive.category(), file.stagedPath().toString());
} }
} catch (RuntimeException e) { } catch (RuntimeException e) {
for (ExtractedFile file : files) { for (ExtractedFile file : files) {
jdbcTemplate.update("delete from aihr_knowledge_upload_item where tenant_id = ? and staging_path = ?", jdbcTemplate.update("delete from aihr_knowledge_upload_item where tenant_id = ? and staging_path = ?",
TENANT_ID, file.stagedPath().toString()); archive.tenantId(), file.stagedPath().toString());
deleteQuietly(file.stagedPath()); deleteQuietly(file.stagedPath());
} }
throw new IOException("ZIP文件入队失败", e); throw new IOException("ZIP文件入队失败", e);
@@ -373,7 +379,7 @@ public class AihrUploadQueueService {
(update_time < (now() - interval %d minute) and file_name not regexp '\\\\.(mp4|mov|avi|mkv|webm|m4v)$') (update_time < (now() - interval %d minute) and file_name not regexp '\\\\.(mp4|mov|avi|mkv|webm|m4v)$')
or update_time < (now() - interval %d minute) or update_time < (now() - interval %d minute)
) )
""".formatted(STUCK_MINUTES, VIDEO_STUCK_MINUTES), TENANT_ID); """.formatted(STUCK_MINUTES, VIDEO_STUCK_MINUTES));
if (reset > 0) { if (reset > 0) {
log.info("upload queue reset {} stuck items", reset); log.info("upload queue reset {} stuck items", reset);
} }
@@ -405,14 +411,14 @@ public class AihrUploadQueueService {
} }
} }
private void markFailed(Long id, String message) { private void markFailed(String tenantId, Long id, String message) {
String error = truncateError(message == null || message.isBlank() ? "处理失败" : message); String error = truncateError(message == null || message.isBlank() ? "处理失败" : message);
// 不覆盖已被清扫重置(0)或已完成(2)的行:处理中标失败,失败态允许更新错误信息(重试路径) // 不覆盖已被清扫重置(0)或已完成(2)的行:处理中标失败,失败态允许更新错误信息(重试路径)
jdbcTemplate.update(""" jdbcTemplate.update("""
update aihr_knowledge_upload_item update aihr_knowledge_upload_item
set status = 3, error = ?, update_time = now() set status = 3, error = ?, update_time = now()
where tenant_id = ? and id = ? and status in (1, 3) where tenant_id = ? and id = ? and status in (1, 3)
""", error, TENANT_ID, id); """, error, tenantId, id);
} }
private static String truncateError(String value) { private static String truncateError(String value) {
@@ -427,7 +433,7 @@ public class AihrUploadQueueService {
select id, batch_id, file_name, category, status, error, doc_id, fragment_count, update_time select id, batch_id, file_name, category, status, error, doc_id, fragment_count, update_time
from aihr_knowledge_upload_item from aihr_knowledge_upload_item
where tenant_id = ? and id = ? where tenant_id = ? and id = ?
""", this::mapItem, TENANT_ID, id); """, this::mapItem, currentTenantId(), id);
return rows.isEmpty() ? null : rows.get(0); return rows.isEmpty() ? null : rows.get(0);
} }
@@ -490,6 +496,11 @@ public class AihrUploadQueueService {
} }
} }
private String currentTenantId() {
String tenantId = LoginHelper.getTenantId();
return tenantId == null || tenantId.isBlank() ? TENANT_ID : tenantId;
}
private void ensureTable() { private void ensureTable() {
if (tableReady) { if (tableReady) {
return; return;
@@ -521,6 +532,6 @@ public class AihrUploadQueueService {
} }
} }
private record ItemRow(Long id, String batchId, String fileName, String category, String stagingPath, int status) { private record ItemRow(String tenantId, Long id, String batchId, String fileName, String category, String stagingPath, int status) {
} }
} }
@@ -0,0 +1,27 @@
package org.dromara.aihr.service;
import org.junit.jupiter.api.Tag;
import org.junit.jupiter.api.Test;
import java.nio.file.Files;
import java.nio.file.Path;
import static org.junit.jupiter.api.Assertions.assertTrue;
@Tag("dev")
class AihrUploadQueueServiceTest {
@Test
void queueKeepsTenantWhenWorkerLeavesLoginThread() 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);
assertTrue(code.contains("select i.tenant_id, i.id, i.file_name"));
assertTrue(code.contains("p.tenant_id = i.tenant_id"));
assertTrue(code.contains("TenantHelper.dynamic(tenantId"));
assertTrue(code.contains("new ItemRow(rs.getString(\"tenant_id\")"));
}
}
+1
View File
@@ -97,3 +97,4 @@
- 2026-07-14 SOP/RAG 租户边界复核:知识库搜索、知识缺口、SOP 问答评审、答案反馈、异步导入任务、文档/片段/向量写入及 Qdrant tenant filter 改为当前登录租户;无登录上下文时回退本地开发租户 `000000`。定向 AIHR 测试 48/48 通过;正式多租户检索隔离仍需在测试环境做跨租户回归,批量上传队列仍待单独审查。 - 2026-07-14 SOP/RAG 租户边界复核:知识库搜索、知识缺口、SOP 问答评审、答案反馈、异步导入任务、文档/片段/向量写入及 Qdrant tenant filter 改为当前登录租户;无登录上下文时回退本地开发租户 `000000`。定向 AIHR 测试 48/48 通过;正式多租户检索隔离仍需在测试环境做跨租户回归,批量上传队列仍待单独审查。
- 2026-07-14 剩余固定租户审查:`AihrPracticeSeedService` 仍有 79 处运行时 `TENANT_ID`,覆盖场景、训练记录、派题、复盘、导出和 Prompt 管理;这些路径可按同步请求上下文分批迁移。`AihrUploadQueueService` 仍有 19 处固定租户引用,但其后台 worker 会脱离登录线程并调用 SOP 加工,不能机械替换为 `LoginHelper`;下一步需让队列条目携带租户、worker 按条目租户加工,并补跨租户领取/隔离测试。 - 2026-07-14 剩余固定租户审查:`AihrPracticeSeedService` 仍有 79 处运行时 `TENANT_ID`,覆盖场景、训练记录、派题、复盘、导出和 Prompt 管理;这些路径可按同步请求上下文分批迁移。`AihrUploadQueueService` 仍有 19 处固定租户引用,但其后台 worker 会脱离登录线程并调用 SOP 加工,不能机械替换为 `LoginHelper`;下一步需让队列条目携带租户、worker 按条目租户加工,并补跨租户领取/隔离测试。
- 2026-07-14 训练租户边界修复:`AihrPracticeSeedService` 的场景、训练记录、派题、复盘、导出、Prompt、标注和音频明细运行时查询/写入统一使用当前登录租户;无登录上下文回退本地开发租户 `000000`。定向 AIHR 测试 48/48、管理端完整打包均通过。`AihrUploadQueueService` 仍保留固定租户,原因是后台 worker 需要下一批显式携带条目租户后再调用 SOP 加工,避免机械替换造成跨租户队列停摆。 - 2026-07-14 训练租户边界修复:`AihrPracticeSeedService` 的场景、训练记录、派题、复盘、导出、Prompt、标注和音频明细运行时查询/写入统一使用当前登录租户;无登录上下文回退本地开发租户 `000000`。定向 AIHR 测试 48/48、管理端完整打包均通过。`AihrUploadQueueService` 仍保留固定租户,原因是后台 worker 需要下一批显式携带条目租户后再调用 SOP 加工,避免机械替换造成跨租户队列停摆。
- 2026-07-14 异步上传队列租户边界修复:请求入队、查询和重试按当前登录租户;后台领取按队列条目租户跨租户扫描,文件加工通过 `TenantHelper.dynamic` 进入对应租户,ZIP 子条目、失败回写和卡死清扫均沿用条目租户;新增队列结构回归测试。AIHR 全量测试 50/50、`ruoyi-admin` 跳测试打包通过;正式多租户文件上传、向量写入和跨租户隔离仍需测试环境回归。