From e4818a8ee26560ebc7e8f0a222a0e551fe30a8a0 Mon Sep 17 00:00:00 2001 From: matevip Date: Sat, 25 Apr 2026 09:56:19 +0800 Subject: [PATCH] feat(wiki): PR-3 eager pipeline honors per-step model config --- .../wiki/job/WikiModelRoutingService.java | 15 +++++ .../wiki/service/WikiProcessingService.java | 63 ++++++++++++++++--- 2 files changed, 70 insertions(+), 8 deletions(-) diff --git a/mateclaw-server/src/main/java/vip/mate/wiki/job/WikiModelRoutingService.java b/mateclaw-server/src/main/java/vip/mate/wiki/job/WikiModelRoutingService.java index a61a0753..c69b18ea 100644 --- a/mateclaw-server/src/main/java/vip/mate/wiki/job/WikiModelRoutingService.java +++ b/mateclaw-server/src/main/java/vip/mate/wiki/job/WikiModelRoutingService.java @@ -70,6 +70,21 @@ public class WikiModelRoutingService { throw new WikiModelUnavailableException("No available model for step: " + step); } + /** + * RFC-051 PR-3: KB-scoped overload. Synthesizes a minimal job context so + * legacy call sites in {@code WikiProcessingService} (which don't always + * have the {@code WikiProcessingJobEntity} on hand) can still benefit + * from per-step model overrides. {@code jobType} matches the keys the UI + * writes into KB config — {@code "heavy_ingest"} for the eager pipeline, + * {@code "light_enrich"} for wikilink enrichment. + */ + public Long selectModelId(Long kbId, String jobType, WikiJobStep step) { + WikiProcessingJobEntity synthetic = new WikiProcessingJobEntity(); + synthetic.setKbId(kbId); + synthetic.setJobType(jobType); + return selectModelId(synthetic, step); + } + /** * Select a fallback model after a failure. */ diff --git a/mateclaw-server/src/main/java/vip/mate/wiki/service/WikiProcessingService.java b/mateclaw-server/src/main/java/vip/mate/wiki/service/WikiProcessingService.java index 5e67491d..3861e3d1 100644 --- a/mateclaw-server/src/main/java/vip/mate/wiki/service/WikiProcessingService.java +++ b/mateclaw-server/src/main/java/vip/mate/wiki/service/WikiProcessingService.java @@ -78,6 +78,15 @@ public class WikiProcessingService { @org.springframework.beans.factory.annotation.Autowired(required = false) private WikiScaffoldService scaffoldService; + /** + * RFC-051 PR-3: optional model routing service. When wired, route / + * create_page / merge_page LLM calls inside the eager pipeline ask the + * routing chain (stepModels[step] -> wikiDefaultModelId -> system + * default) for a model rather than always pulling the system default. + */ + @org.springframework.beans.factory.annotation.Autowired(required = false) + private vip.mate.wiki.job.WikiModelRoutingService modelRoutingService; + /** Parallel chunk / material processing executor (JDK 21 virtual threads) */ public static final ExecutorService WIKI_EXECUTOR = Executors.newVirtualThreadPerTaskExecutor(); @@ -577,7 +586,8 @@ public class WikiProcessingService { new SystemMessage(systemPrompt), new UserMessage(userPrompt) )); - 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); return applyLlmResponse(kb.getId(), raw.getId(), llmResponse); } @@ -631,7 +641,8 @@ public class WikiProcessingService { new SystemMessage(routeSystem), new UserMessage(routeUser) )); - 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); JsonNode routeJson = parseJsonResponse(routeResponse); if (routeJson == null) { log.warn("[Wiki] Route phase: failed to parse JSON for kbId={}, rawId={}, responseLen={}, first200={}", @@ -829,7 +840,8 @@ public class WikiProcessingService { String batchResponse = callLlmWithResilientRetry(batchPrompt, "batch-create " + subBatch.size() + " pages of raw=" + rawId - + " subBatch=" + (bStart / batchSize + 1)); + + " subBatch=" + (bStart / batchSize + 1), + kbId, vip.mate.wiki.job.WikiJobStep.CREATE_PAGE); List parsedPages = batchParser.parse(batchResponse); @@ -1017,7 +1029,8 @@ public class WikiProcessingService { new SystemMessage(createSystem), new UserMessage(createUser) )); - 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); } /** @@ -1146,7 +1159,8 @@ public class WikiProcessingService { new UserMessage(mergeUser) )); 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); JsonNode mergeJson = parseJsonResponse(response); if (mergeJson == null) { log.warn("[Wiki] Phase B merge page slug='{}' returned unparseable JSON, skipping", slug); @@ -1266,7 +1280,8 @@ public class WikiProcessingService { new UserMessage(user) )); try { - 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); JsonNode json = parseJsonResponse(response); if (json != null) { log.info("[Wiki] Document analysis done for raw={}: topics={}, concepts={}", @@ -1318,6 +1333,28 @@ public class WikiProcessingService { return agentGraphBuilder.buildRuntimeChatModel(defaultModel, WIKI_NO_RETRY); } + /** + * RFC-051 PR-3: build a {@link ChatModel} for an eager-pipeline step, + * honoring the KB-level routing chain when {@link #modelRoutingService} + * is available. Falls back to the system default on any lookup failure + * so a misconfigured KB never blocks ingest. + */ + private ChatModel buildChatModelFor(Long kbId, vip.mate.wiki.job.WikiJobStep step) { + if (modelRoutingService != null && kbId != null && step != null) { + try { + Long modelId = modelRoutingService.selectModelId(kbId, "heavy_ingest", step); + ModelConfigEntity model = modelConfigService.getModel(modelId); + if (model != null) { + return agentGraphBuilder.buildRuntimeChatModel(model, WIKI_NO_RETRY); + } + } catch (Exception e) { + log.warn("[Wiki] Model routing failed for kbId={} step={}, falling back to default: {}", + kbId, step, e.getMessage()); + } + } + return buildChatModel(); + } + /** * 调用 LLM,带"任务完成或模型不可用才终止"的重试策略。 *

@@ -1332,12 +1369,21 @@ public class WikiProcessingService { * 反复瞬时错误把单 chunk 卡到永远;buildChatModel 提到循环外,所有重试复用同一实例。 */ private String callLlmWithResilientRetry(Prompt prompt, String ctx) { + return callLlmWithResilientRetry(prompt, ctx, null, null); + } + + /** + * RFC-051 PR-3: step-aware LLM retry helper. {@code kbId} and {@code step} + * pick the routed chat model; passing {@code null} for either reproduces + * the legacy behavior (system default model). + */ + private String callLlmWithResilientRetry(Prompt prompt, String ctx, Long kbId, vip.mate.wiki.job.WikiJobStep step) { long backoffMs = 1000; final long maxBackoffMs = 60_000; final int maxAttempts = Math.max(1, properties.getLlmMaxAttempts()); final long maxTotalDurationMs = Math.max(1_000L, properties.getLlmMaxTotalDurationMs()); final long startNanos = System.nanoTime(); - final ChatModel chatModel = buildChatModel(); + final ChatModel chatModel = buildChatModelFor(kbId, step); int attempt = 0; while (true) { attempt++; @@ -1599,7 +1645,8 @@ public class WikiProcessingService { new SystemMessage(createSystem), new UserMessage(createUser) )); - 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); com.fasterxml.jackson.databind.JsonNode pageJson = parseJsonResponse(response); if (pageJson == null) return;