mirror of
https://gitee.com/mateos/mateclaw.git
synced 2026-09-13 03:13:41 +08:00
feat(wiki): PR-3 eager pipeline honors per-step model config
This commit is contained in:
parent
1281153aa8
commit
e4818a8ee2
@ -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.
|
||||
*/
|
||||
|
||||
@ -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<WikiBatchCreateParser.ParsedPage> 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,带"任务完成或模型不可用才终止"的重试策略。
|
||||
* <p>
|
||||
@ -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;
|
||||
|
||||
|
||||
Loading…
Reference in New Issue
Block a user