fix(wiki): guard page reclassification against concurrent re-trigger

A second POST to /reclassify on the same KB spawned an independent pass over
the same pages, doubling LLM spend and racing the first pass's page-type
writes. Add a per-KB in-flight guard that rejects a concurrent run with a
friendly message (409 via R.fail rather than a generic 500), released in a
finally once the async pass completes. Also count and broadcast per-page
failures so an all-failing run is visible instead of reporting changed=0, and
type the api modelId param as string|number per the snowflake ID convention.
This commit is contained in:
matevip 2026-06-08 21:59:10 +08:00
parent 71ddade3bd
commit 21798d6be5
3 changed files with 46 additions and 14 deletions

View File

@ -285,7 +285,14 @@ public class WikiController {
if (body != null && body.get("modelId") != null) {
modelId = Long.valueOf(String.valueOf(body.get("modelId")));
}
int queued = processingService.reclassifyKB(id, modelId);
int queued;
try {
queued = processingService.reclassifyKB(id, modelId);
} catch (IllegalStateException e) {
// A reclassification is already running for this KB surface a
// friendly message rather than a generic 500.
return R.fail(e.getMessage());
}
Map<String, Object> out = new LinkedHashMap<>();
out.put("queued", queued);
return R.ok(out);

View File

@ -181,6 +181,10 @@ public class WikiProcessingService {
private final ConcurrentHashMap<Long, ProgressCounter> progressCounters = new ConcurrentHashMap<>();
/** KBs with a reclassify pass currently running, used to reject concurrent
* re-triggers (which would double LLM spend and race page-type writes). */
private final java.util.Set<Long> reclassifyInFlight = ConcurrentHashMap.newKeySet();
/**
* Process one raw material.
*/
@ -591,23 +595,38 @@ public class WikiProcessingService {
return 0;
}
// Resolve the classifying model once up front. An explicit modelId wins;
// otherwise route as a CREATE_PAGE step, falling back to the default.
final ChatModel chatModel;
if (modelId != null && modelRoutingService != null) {
chatModel = modelRoutingService.buildChatModel(modelId);
} else {
chatModel = resolveChatModel(kbId, vip.mate.wiki.job.WikiJobStep.CREATE_PAGE).chatModel;
// Reject a concurrent re-trigger on the same KB: two parallel passes would
// double the LLM spend and race each other's updatePageType writes.
if (!reclassifyInFlight.add(kbId)) {
throw new IllegalStateException("A reclassification is already running for this knowledge base");
}
final String systemPrompt = PromptLoader.loadPrompt("wiki/classify-page-system")
.replace("{allowed_page_types}", pageTypeProfileService.describeForPrompt(kbId));
final String userTemplate = PromptLoader.loadPrompt("wiki/classify-page-user");
final ChatModel chatModel;
final String systemPrompt;
final String userTemplate;
try {
// Resolve the classifying model once up front. An explicit modelId wins;
// otherwise route as a CREATE_PAGE step, falling back to the default.
if (modelId != null && modelRoutingService != null) {
chatModel = modelRoutingService.buildChatModel(modelId);
} else {
chatModel = resolveChatModel(kbId, vip.mate.wiki.job.WikiJobStep.CREATE_PAGE).chatModel;
}
systemPrompt = PromptLoader.loadPrompt("wiki/classify-page-system")
.replace("{allowed_page_types}", pageTypeProfileService.describeForPrompt(kbId));
userTemplate = PromptLoader.loadPrompt("wiki/classify-page-user");
} catch (RuntimeException e) {
// Setup failed before any async work was queued release the guard.
reclassifyInFlight.remove(kbId);
throw e;
}
final int total = pages.size();
WIKI_EXECUTOR.submit(() -> {
int done = 0;
int changed = 0;
int failed = 0;
try {
for (WikiPageEntity page : pages) {
done++;
try {
@ -635,6 +654,7 @@ public class WikiProcessingService {
changed++;
}
} catch (Exception e) {
failed++;
log.warn("[Wiki] reclassify failed pageId={} kbId={}: {}",
page.getId(), kbId, e.getMessage());
} finally {
@ -642,9 +662,14 @@ public class WikiProcessingService {
java.util.Map.of("kind", "reclassify", "done", done, "total", total));
}
}
log.info("[Wiki] reclassifyKB done kbId={} pages={} changed={}", kbId, total, changed);
log.info("[Wiki] reclassifyKB done kbId={} pages={} changed={} failed={}",
kbId, total, changed, failed);
progressBus.broadcast(kbId, WikiProgressBus.EVENT_RAW_COMPLETED,
java.util.Map.of("kind", "reclassify", "done", total, "total", total, "changed", changed));
java.util.Map.of("kind", "reclassify", "done", total, "total", total,
"changed", changed, "failed", failed));
} finally {
reclassifyInFlight.remove(kbId);
}
});
log.info("[Wiki] reclassifyKB queued {} page(s) for kbId={} (modelId={})", total, kbId, modelId);

View File

@ -923,7 +923,7 @@ export const wikiApi = {
http.post(`/wiki/knowledge-bases/${kbId}/page-type-profile/validate`, { config }),
resetPageTypeProfile: (kbId: string | number) =>
http.post(`/wiki/knowledge-bases/${kbId}/page-type-profile/reset-default`),
reclassifyKB: (kbId: string | number, modelId?: number | null) =>
reclassifyKB: (kbId: string | number, modelId?: string | number | null) =>
http.post(`/wiki/knowledge-bases/${kbId}/reclassify`, modelId != null ? { modelId } : {}),
// ---- Agent pageType permissions (REQ-3) ----