mirror of
https://gitee.com/mateos/mateclaw.git
synced 2026-09-15 20:08:18 +08:00
fix(wiki): cancel in-flight LLM work when raw is deleted
This commit is contained in:
parent
4ed88ebbf9
commit
3f25064ef9
@ -484,6 +484,13 @@ public class WikiProcessingService {
|
|||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
try {
|
try {
|
||||||
|
// RFC-051 PR-9: skip remaining chunks if the user deleted the raw
|
||||||
|
// while earlier chunks were still in flight. Counts as a "failed chunk"
|
||||||
|
// for terminal-status accounting (not actually failed, just abandoned).
|
||||||
|
if (isAborted(raw.getId(), "chunk " + (chunkIndex + 1) + "/" + totalChunks)) {
|
||||||
|
failedChunks.incrementAndGet();
|
||||||
|
return;
|
||||||
|
}
|
||||||
log.info("[Wiki] Processing chunk {}/{}: {} chars", chunkIndex + 1, totalChunks, chunk.length());
|
log.info("[Wiki] Processing chunk {}/{}: {} chars", chunkIndex + 1, totalChunks, chunk.length());
|
||||||
int pages = processChunk(kb, raw, chunk, existingPagesIndex, documentMap);
|
int pages = processChunk(kb, raw, chunk, existingPagesIndex, documentMap);
|
||||||
totalPages.addAndGet(pages);
|
totalPages.addAndGet(pages);
|
||||||
@ -606,6 +613,7 @@ public class WikiProcessingService {
|
|||||||
new SystemMessage(systemPrompt),
|
new SystemMessage(systemPrompt),
|
||||||
new UserMessage(userPrompt)
|
new UserMessage(userPrompt)
|
||||||
));
|
));
|
||||||
|
if (isAborted(raw.getId(), "single-chunk legacy")) return 0;
|
||||||
String llmResponse = callLlmWithResilientRetry(prompt, "chunk of raw=" + raw.getId(),
|
String llmResponse = callLlmWithResilientRetry(prompt, "chunk of raw=" + raw.getId(),
|
||||||
kb.getId(), vip.mate.wiki.job.WikiJobStep.CREATE_PAGE);
|
kb.getId(), vip.mate.wiki.job.WikiJobStep.CREATE_PAGE);
|
||||||
|
|
||||||
@ -672,6 +680,7 @@ public class WikiProcessingService {
|
|||||||
new SystemMessage(routeSystem),
|
new SystemMessage(routeSystem),
|
||||||
new UserMessage(routeUser)
|
new UserMessage(routeUser)
|
||||||
));
|
));
|
||||||
|
if (isAborted(rawId, "route phase")) return 0;
|
||||||
String routeResponse = callLlmWithResilientRetry(routePrompt, "route chunk of raw=" + rawId,
|
String routeResponse = callLlmWithResilientRetry(routePrompt, "route chunk of raw=" + rawId,
|
||||||
kbId, vip.mate.wiki.job.WikiJobStep.ROUTE);
|
kbId, vip.mate.wiki.job.WikiJobStep.ROUTE);
|
||||||
|
|
||||||
@ -900,6 +909,10 @@ public class WikiProcessingService {
|
|||||||
new UserMessage(batchUser)
|
new UserMessage(batchUser)
|
||||||
));
|
));
|
||||||
|
|
||||||
|
if (isAborted(rawId, "batch-create sub-batch " + (bStart / batchSize + 1))) {
|
||||||
|
// Abandon remaining sub-batches; return however many pages we already created.
|
||||||
|
return totalCreated;
|
||||||
|
}
|
||||||
String batchResponse = callLlmWithResilientRetry(batchPrompt,
|
String batchResponse = callLlmWithResilientRetry(batchPrompt,
|
||||||
"batch-create " + subBatch.size() + " pages of raw=" + rawId
|
"batch-create " + subBatch.size() + " pages of raw=" + rawId
|
||||||
+ " subBatch=" + (bStart / batchSize + 1),
|
+ " subBatch=" + (bStart / batchSize + 1),
|
||||||
@ -1091,6 +1104,7 @@ public class WikiProcessingService {
|
|||||||
new SystemMessage(createSystem),
|
new SystemMessage(createSystem),
|
||||||
new UserMessage(createUser)
|
new UserMessage(createUser)
|
||||||
));
|
));
|
||||||
|
if (isAborted(raw.getId(), "retry-create slug=" + slug)) return null;
|
||||||
return callLlmWithResilientRetry(prompt, "retry-create slug=" + slug + " of raw=" + raw.getId(),
|
return callLlmWithResilientRetry(prompt, "retry-create slug=" + slug + " of raw=" + raw.getId(),
|
||||||
kb.getId(), vip.mate.wiki.job.WikiJobStep.CREATE_PAGE);
|
kb.getId(), vip.mate.wiki.job.WikiJobStep.CREATE_PAGE);
|
||||||
}
|
}
|
||||||
@ -1113,6 +1127,11 @@ public class WikiProcessingService {
|
|||||||
Long kbId = kb.getId();
|
Long kbId = kb.getId();
|
||||||
Long rawId = raw.getId();
|
Long rawId = raw.getId();
|
||||||
|
|
||||||
|
// RFC-051 PR-9: refuse to materialize a page (or merge into an existing one) tied
|
||||||
|
// to a raw the user just deleted. Prevents zombie pages whose source_raw_ids point
|
||||||
|
// at a tombstoned row.
|
||||||
|
if (isAborted(rawId, "savePageContent slug=" + slug)) return false;
|
||||||
|
|
||||||
// Fallback 0: cross-spelling canonical match (DB has same concept under different slug)
|
// Fallback 0: cross-spelling canonical match (DB has same concept under different slug)
|
||||||
WikiPageEntity existingByCanonical = pageService.findByCanonicalSlug(kbId, slug);
|
WikiPageEntity existingByCanonical = pageService.findByCanonicalSlug(kbId, slug);
|
||||||
if (existingByCanonical != null && !existingByCanonical.getSlug().equals(slug)) {
|
if (existingByCanonical != null && !existingByCanonical.getSlug().equals(slug)) {
|
||||||
@ -1220,6 +1239,7 @@ public class WikiProcessingService {
|
|||||||
new SystemMessage(mergeSystem),
|
new SystemMessage(mergeSystem),
|
||||||
new UserMessage(mergeUser)
|
new UserMessage(mergeUser)
|
||||||
));
|
));
|
||||||
|
if (isAborted(rawId, "merge slug=" + slug)) return false;
|
||||||
String response = callLlmWithResilientRetry(prompt,
|
String response = callLlmWithResilientRetry(prompt,
|
||||||
"merge page slug=" + slug + " of raw=" + rawId,
|
"merge page slug=" + slug + " of raw=" + rawId,
|
||||||
kbId, vip.mate.wiki.job.WikiJobStep.MERGE_PAGE);
|
kbId, vip.mate.wiki.job.WikiJobStep.MERGE_PAGE);
|
||||||
@ -1342,6 +1362,7 @@ public class WikiProcessingService {
|
|||||||
new UserMessage(user)
|
new UserMessage(user)
|
||||||
));
|
));
|
||||||
try {
|
try {
|
||||||
|
if (isAborted(raw.getId(), "doc analysis")) return "";
|
||||||
String response = callLlmWithResilientRetry(prompt, "analyze doc raw=" + raw.getId(),
|
String response = callLlmWithResilientRetry(prompt, "analyze doc raw=" + raw.getId(),
|
||||||
kb.getId(), vip.mate.wiki.job.WikiJobStep.ROUTE);
|
kb.getId(), vip.mate.wiki.job.WikiJobStep.ROUTE);
|
||||||
JsonNode json = parseJsonResponse(response);
|
JsonNode json = parseJsonResponse(response);
|
||||||
@ -1707,6 +1728,7 @@ public class WikiProcessingService {
|
|||||||
new SystemMessage(createSystem),
|
new SystemMessage(createSystem),
|
||||||
new UserMessage(createUser)
|
new UserMessage(createUser)
|
||||||
));
|
));
|
||||||
|
if (isAborted(raw.getId(), "repair page=" + page.getSlug())) return;
|
||||||
String response = callLlmWithResilientRetry(prompt, "repair page=" + page.getSlug(),
|
String response = callLlmWithResilientRetry(prompt, "repair page=" + page.getSlug(),
|
||||||
kb.getId(), vip.mate.wiki.job.WikiJobStep.MERGE_PAGE);
|
kb.getId(), vip.mate.wiki.job.WikiJobStep.MERGE_PAGE);
|
||||||
com.fasterxml.jackson.databind.JsonNode pageJson = parseJsonResponse(response);
|
com.fasterxml.jackson.databind.JsonNode pageJson = parseJsonResponse(response);
|
||||||
@ -1783,6 +1805,28 @@ public class WikiProcessingService {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* RFC-051 PR-9: returns {@code true} when the caller should bail out of an
|
||||||
|
* in-flight processing path because the raw material has been deleted.
|
||||||
|
* <p>
|
||||||
|
* {@link WikiRawMaterialService#delete(Long)} is a logical delete (the
|
||||||
|
* {@code @TableLogic} column flips to 1), so {@code selectById} returns
|
||||||
|
* {@code null} as soon as the deletion commits. Sprinkling this check
|
||||||
|
* right before each LLM call keeps token spend bounded by a single
|
||||||
|
* in-flight call after the user clicks delete.
|
||||||
|
*
|
||||||
|
* @param rawId the raw material id this processing path is about
|
||||||
|
* @param ctx short string used in the log line
|
||||||
|
* @return {@code true} if the raw is gone; caller should stop work
|
||||||
|
*/
|
||||||
|
private boolean isAborted(Long rawId, String ctx) {
|
||||||
|
if (rawService.getById(rawId) == null) {
|
||||||
|
log.info("[Wiki] Aborting {} for raw={}: raw was deleted mid-processing", ctx, rawId);
|
||||||
|
return true;
|
||||||
|
}
|
||||||
|
return false;
|
||||||
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* RFC-051 PR-1c: bridge from {@link DocumentPreprocessService.Chunker} to
|
* RFC-051 PR-1c: bridge from {@link DocumentPreprocessService.Chunker} to
|
||||||
* the existing sentence-boundary chunker. Returns {@code [start, end]}
|
* the existing sentence-boundary chunker. Returns {@code [start, end]}
|
||||||
@ -1838,6 +1882,14 @@ public class WikiProcessingService {
|
|||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// RFC-051 PR-9: text extraction can take many seconds on large binaries.
|
||||||
|
// If the user deleted the raw during that window, persisting chunks for a
|
||||||
|
// tombstoned row is wasted work that the cascade-cleanup already covered.
|
||||||
|
if (isAborted(rawId, "lazy ingest after extract")) {
|
||||||
|
kbService.updateStatus(kbId, "active");
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
int totalChunks;
|
int totalChunks;
|
||||||
// PR-1c: when the preprocessor is on the classpath, normalize +
|
// PR-1c: when the preprocessor is on the classpath, normalize +
|
||||||
// attach metadata; otherwise fall back to the legacy chunker.
|
// attach metadata; otherwise fall back to the legacy chunker.
|
||||||
|
|||||||
Loading…
Reference in New Issue
Block a user