fix(aihr): recover stalled media processing
This commit is contained in:
+6
@@ -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<UploadEnqueueResponse> retryProcessingTask(@PathVariable Long id) {
|
||||
return R.ok(uploadQueueService.retryAttachment(id));
|
||||
}
|
||||
|
||||
@SaCheckLogin
|
||||
@PostMapping("/search")
|
||||
public R<SearchResponse> search(@RequestBody SearchRequest request) {
|
||||
|
||||
+171
-40
@@ -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<String> 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<SnippetResponse> 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<String> tags = video ? List.of("视频", "待转写") : List.of("图片", "待OCR");
|
||||
? "MEDIA_NO_TEXT: 视频未提取到语音转写或画面文字,可从资料处理中心重试"
|
||||
: "MEDIA_NO_TEXT: 图片未提取到可用文字,可从资料处理中心重试";
|
||||
List<String> 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<SysOssVo> 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<String> 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<String> 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<String> response = HttpClient.newBuilder()
|
||||
.connectTimeout(Duration.ofSeconds(15))
|
||||
.build()
|
||||
.send(builder.build(), HttpResponse.BodyHandlers.ofString());
|
||||
HttpResponse<String> 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<String, CategorySummary> 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) {
|
||||
}
|
||||
|
||||
|
||||
+249
-20
@@ -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<String> spaceCodes, String stagingPath, int status) {
|
||||
List<String> spaceCodes, String stagingPath, Long sourceAttachId, int status) {
|
||||
}
|
||||
|
||||
private record SourceAttachment(Long id, Long ossId, String fileName, int status, String spaceCode,
|
||||
String category) {
|
||||
}
|
||||
}
|
||||
|
||||
+146
@@ -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<ImageReader> 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<ImageWriter> 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) {
|
||||
}
|
||||
}
|
||||
+19
@@ -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 {
|
||||
|
||||
+20
-2
@@ -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 {
|
||||
|
||||
+53
@@ -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()));
|
||||
}
|
||||
}
|
||||
@@ -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` (
|
||||
|
||||
@@ -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;
|
||||
Reference in New Issue
Block a user