package vip.mate.wiki.service; import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.ObjectMapper; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.ai.chat.messages.SystemMessage; import org.springframework.ai.chat.messages.UserMessage; import org.springframework.ai.chat.model.ChatModel; import org.springframework.ai.chat.model.ChatResponse; import org.springframework.ai.chat.prompt.Prompt; import org.springframework.retry.support.RetryTemplate; import org.springframework.stereotype.Service; import vip.mate.agent.AgentGraphBuilder; import vip.mate.agent.prompt.PromptLoader; import vip.mate.llm.model.ModelConfigEntity; import vip.mate.llm.service.ModelConfigService; import vip.mate.wiki.WikiProperties; import vip.mate.wiki.dto.WikiChunkDraft; import vip.mate.wiki.job.WikiKbConfig; import vip.mate.wiki.model.WikiKnowledgeBaseEntity; import vip.mate.wiki.model.WikiPageEntity; import vip.mate.wiki.model.WikiRawMaterialEntity; import vip.mate.wiki.sse.WikiProgressBus; import java.util.ArrayList; import java.util.List; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.Semaphore; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; /** * Wiki 处理服务 *
* 核心管线:将原始材料通过 LLM 消化为结构化 Wiki 页面。 * * @author MateClaw Team */ @Slf4j @Service @RequiredArgsConstructor public class WikiProcessingService { private final WikiKnowledgeBaseService kbService; private final WikiRawMaterialService rawService; private final WikiPageService pageService; private final WikiChunkService chunkService; private final WikiEmbeddingService embeddingService; private final WikiProperties properties; private final ModelConfigService modelConfigService; private final AgentGraphBuilder agentGraphBuilder; private final ObjectMapper objectMapper; private final WikiProgressBus progressBus; private final WikiCitationService citationService; private final org.springframework.context.ApplicationEventPublisher eventPublisher; /** * Read-the-failover-chain handle. Optional so the existing constructors and * lazy-mode tests don't have to thread a new dependency. When null, the * fallback hop iterates {@code listEnabledModels} in DB order — same * behavior as before this PR. */ @org.springframework.beans.factory.annotation.Autowired(required = false) private vip.mate.llm.service.ModelProviderService modelProviderService; /** * Per-provider failure counter / cooldown bookkeeping. Optional for the * same reason. When wired, fatal errors mark the provider down so other * code paths (chat agent, fallback chain) skip it during cooldown; on a * successful call we clear the failure counter for the provider that * actually responded. */ @org.springframework.beans.factory.annotation.Autowired(required = false) private vip.mate.llm.failover.ProviderHealthTracker providerHealthTracker; @org.springframework.beans.factory.annotation.Autowired(required = false) @org.springframework.context.annotation.Lazy private vip.mate.wiki.job.WikiProcessingJobService wikiJobService; /** * RFC-051 PR-1c: optional preprocessor that fills chunk metadata * (page_number / token_count / header_breadcrumb / source_section). * Marked optional so unit tests that construct this service directly * (without Spring) can opt out without exploding. */ @org.springframework.beans.factory.annotation.Autowired(required = false) private DocumentPreprocessService preprocessService; /** * RFC-051 PR-2: ensures system-page scaffold (overview / log) exists for * the KB before each ingest. Optional so the older lazy-only unit tests * don't need to wire it. */ @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; /** RFC-051 PR-2b/2c: optional overview rebuilder + log appender. */ @org.springframework.beans.factory.annotation.Autowired(required = false) private WikiOverviewService overviewService; @org.springframework.beans.factory.annotation.Autowired(required = false) private WikiLogService logService; /** Parallel chunk / material processing executor (JDK 21 virtual threads) */ public static final ExecutorService WIKI_EXECUTOR = Executors.newVirtualThreadPerTaskExecutor(); /** * RFC-012 M2 v2 UI v2:单 raw 的进度计数器,多个并行 chunk 的 {@code processChunkTwoPhase} * 共享同一份 atomic 计数,避免 6 个 chunk 各写各的 progress 字段时互相覆盖(导致 UI 永远 preparing)。 *
* 生命周期:{@code processRawMaterial} 入口 put,try/finally 出口 remove。 */ private static final class ProgressCounter { final AtomicInteger total = new AtomicInteger(0); final AtomicInteger done = new AtomicInteger(0); /** Page-level 失败计数(chunk 内单页 create / merge 抛异常)。 * 注:DuplicateKeyException 触发的 fallback-to-update 不算 failure, * 内容仍合入了同 slug page。仅 LLM 调用爆炸、JSON 解析失败、内容为空等真失败才递增。 */ final AtomicInteger failed = new AtomicInteger(0); final AtomicBoolean phaseBStarted = new AtomicBoolean(false); /** * 跨 chunk slug 抢占表:canonical slug → 第一个声明该概念的实际 slug。 *
* 解决 LLM 在并行 chunk 中给同一概念起不同 slug 拼写(按词分组 vs 按字分隔)的问题。
* 使用 {@link ConcurrentHashMap#computeIfAbsent} 实现原子抢占:先到的 chunk 把自己的
* slug 注册为 winner,后到的 chunk 看到 winner 后会把内容写入 winner 对应的 page。
*/
final ConcurrentHashMap
* Merge is skipped (not just decremented from count) when a slug is already present.
* Uses ConcurrentHashMap as a concurrent set via putIfAbsent.
*/
final ConcurrentHashMap
* RFC-012 Change 1:材料级并行,受 {@link WikiProperties#getMaxParallelRawMaterials()} 约束。
*/
public void processAllPending(Long kbId) {
List
* 阶段 A(route):一次 LLM 调用决定要 create 哪些新页 + 要 update 哪些已有页(仅 slug 列表)。
* 输入小、输出短,单次稳定在 30s 内返回。
*
* 阶段 B(merge):对 update 列表里的每个 slug 单独发 LLM 调用,输入只塞这一页的现有正文 + 当前
* chunk 文本,输出该页 merge 后的完整内容。每次调用单页规模,远不会触发 nginx 60s 超时。
*
* 新建页直接落库;merge 页因互不依赖,可在当前 chunk 的 virtual thread 内顺序处理(chunk 之间
* 已通过 maxParallelChunks Semaphore 拿到了并行度)。
*
* @return 创建+更新的页面数
*/
private int processChunkTwoPhase(WikiKnowledgeBaseEntity kb, WikiRawMaterialEntity raw,
String textContent, String existingPagesIndex, String documentMap) {
Long kbId = kb.getId();
Long rawId = raw.getId();
String configContent = kb.getConfigContent() != null ? kb.getConfigContent() : "";
String rawTitle = raw.getTitle();
// RFC-012 M2 v2 UI v2:取共享进度计数器(processRawMaterial 入口已 put)。
// 多 chunk 并行时所有 chunk 共享同一份 atomic 计数,避免互相覆盖把 UI 拉回 preparing。
ProgressCounter pc = progressCounters.get(rawId);
// ─── 阶段 A:路由 ───
// Rebuild existingPagesIndex fresh at route time so pages created by earlier chunks
// in this run are visible. This prevents the route from scheduling "create" for a
// concept that was already created by a previous chunk (which would be caught by
// savePageContent and converted to update, but wastes a merge LLM call).
// listSummaries uses a 5-min TTL cache that is evicted on every create/update, so
// this picks up changes from sequential chunks without an extra DB hit when nothing changed.
String freshIndex = buildExistingPagesIndex(kbId);
String routeSystem = PromptLoader.loadPrompt("wiki/route-system");
String routeUserTemplate = PromptLoader.loadPrompt("wiki/route-user");
String documentMapSection = (documentMap != null && !documentMap.isBlank())
? "## 文档全局概念地图(预分析结果,供路由参考)\n\n```json\n" + documentMap + "\n```\n"
: "";
String routeUser = routeUserTemplate
.replace("{config}", configContent)
.replace("{document_map_section}", documentMapSection)
.replace("{existing_pages}", freshIndex)
.replace("{raw_title}", rawTitle)
.replace("{raw_content}", textContent);
// RFC-051 PR-6b: optionally inject Spring AI's structured-output hint so the
// LLM produces strict RouteResult JSON. KB config wins; falls back to global
// mate.wiki.use-structured-route default when the KB hasn't expressed a preference.
boolean useStructured = resolveStructuredRouteFlag(kb);
org.springframework.ai.converter.BeanOutputConverter
* Splits {@code createMetas} into sub-batches of {@code batchCreatePageSize}; after each
* sub-batch the saved pages are appended to {@code liveIndex} so subsequent sub-batches
* can link to them. Returns the count of actually-created (not updated) pages.
*/
private int batchCreatePages(WikiKnowledgeBaseEntity kb, WikiRawMaterialEntity raw,
String chunkText, String existingPagesIndex,
List This consolidates what used to be three near-identical inline retry
* blocks (omitted slug / blank content / unparseable JSON) into one path
* with consistent attempt count, log lines, and progress accounting.
*
* Side effects on success: {@code created} is incremented (only when a
* brand-new page is persisted, not on dedupe), {@code liveIndex} grows so
* subsequent batched pages can wikilink to this one, and {@code pc} ticks
* one {@code done} regardless of outcome with {@code failed} on exhaustion.
*
* @param slugForLog logical slug used in log messages — the actual write
* slug comes from the LLM response or retryMeta.slug
* @return 1 if a new page was persisted, 0 if a parseable response landed
* but the slug already existed (dedup), -1 if all attempts failed
*/
private int tryCreateOneWithRetry(WikiKnowledgeBaseEntity kb,
WikiRawMaterialEntity raw,
String chunkText,
StringBuilder liveIndex,
JsonNode retryMeta,
int maxAttempts,
AtomicInteger created,
ProgressCounter pc,
Long rawId,
Long kbId,
String slugForLog) {
String metaSlug = retryMeta.path("slug").asText("");
String fallbackTitle = retryMeta.path("title").asText("");
String fallbackSummary = retryMeta.path("summary").asText("");
int delta = -1;
for (int attempt = 1; attempt <= maxAttempts; attempt++) {
try {
String retryResult = retrySingleCreate(kb, raw, chunkText, liveIndex.toString(), retryMeta);
if (retryResult == null) {
log.warn("[Wiki] retry-create slug='{}' attempt {}/{}: null response",
slugForLog, attempt, maxAttempts);
continue;
}
JsonNode retryJson = parseJsonResponse(retryResult);
if (retryJson == null) {
log.warn("[Wiki] retry-create slug='{}' attempt {}/{}: unparseable JSON",
slugForLog, attempt, maxAttempts);
continue;
}
String resolvedSlug = retryJson.path("slug").asText(metaSlug);
if (resolvedSlug.isBlank()) resolvedSlug = metaSlug;
String content = retryJson.path("content").asText("");
if (content.isBlank()) {
log.warn("[Wiki] retry-create slug='{}' attempt {}/{}: blank content",
slugForLog, attempt, maxAttempts);
continue;
}
String title = retryJson.path("title").asText(fallbackTitle);
String summary = retryJson.path("summary").asText(fallbackSummary);
boolean wasCreated = savePageContent(kb, raw, resolvedSlug, title, content, summary);
if (wasCreated) {
created.incrementAndGet();
String brief = summary.length() > 100 ? summary.substring(0, 100) : summary;
liveIndex.append("\n- ").append(resolvedSlug).append(": ").append(brief);
delta = 1;
} else {
delta = 0;
}
log.info("[Wiki] retry-create slug='{}' succeeded on attempt {}/{} (newPage={})",
slugForLog, attempt, maxAttempts, wasCreated);
break;
} catch (Exception e) {
log.warn("[Wiki] retry-create slug='{}' attempt {}/{} threw: {}",
slugForLog, attempt, maxAttempts, e.getMessage());
}
}
if (delta < 0) {
log.warn("[Wiki] retry-create slug='{}' exhausted {} attempts — giving up",
slugForLog, maxAttempts);
}
if (pc != null) {
int d = pc.done.incrementAndGet();
if (delta < 0) pc.failed.incrementAndGet();
rawService.updateProgress(rawId, "phase-b", d, pc.total.get());
progressBus.broadcast(kbId, WikiProgressBus.EVENT_CHUNK_DONE,
java.util.Map.of("rawId", rawId, "kind", "create-retry",
"ok", delta >= 0, "done", d, "total", pc.total.get()));
}
return delta;
}
/**
* RFC-047 P1: Shared DB save logic for a new page, extracted for use by both
* {@link #createOnePage} and {@link #batchCreatePages}.
* Handles canonical-slug matching, in-flight slug-claim arbitration, and DuplicateKey fallback.
*
* @return true if a new row was inserted; false if an existing page was updated instead
*/
private boolean savePageContent(WikiKnowledgeBaseEntity kb, WikiRawMaterialEntity raw,
String slug, String title, String content, String pageSummary) {
return savePageContent(kb, raw, slug, title, content, pageSummary, null);
}
private boolean savePageContent(WikiKnowledgeBaseEntity kb, WikiRawMaterialEntity raw,
String slug, String title, String content, String pageSummary,
String pageType) {
Long kbId = kb.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)
WikiPageEntity existingByCanonical = pageService.findByCanonicalSlug(kbId, slug);
if (existingByCanonical != null && !existingByCanonical.getSlug().equals(slug)) {
String actualSlug = existingByCanonical.getSlug();
pageService.updatePageByAi(kbId, actualSlug, content, pageSummary, rawId);
pageService.mergeSourceLineage(existingByCanonical.getId(), rawId, raw.getTitle());
log.info("[Wiki] Phase B create slug='{}' canonical-matches existing '{}', updated",
slug, actualSlug);
return false;
}
// Fallback 0.5: in-flight slug-claim arbitration across parallel chunks
ProgressCounter pcLocal = progressCounters.get(rawId);
String canonical = WikiPageService.canonicalSlug(slug);
if (pcLocal != null && !canonical.isEmpty()) {
final String routedSlug = slug;
String winnerSlug = pcLocal.slugClaims.computeIfAbsent(canonical, k -> routedSlug);
if (!winnerSlug.equals(slug)) {
WikiPageEntity winner = pageService.getBySlug(kbId, winnerSlug);
if (winner != null) {
pageService.updatePageByAi(kbId, winnerSlug, content, pageSummary, rawId);
pageService.mergeSourceLineage(winner.getId(), rawId, raw.getTitle());
log.info("[Wiki] Phase B create slug='{}' lost slug-claim race to '{}', updated",
slug, winnerSlug);
return false;
}
log.info("[Wiki] Phase B create slug='{}' redirects to in-flight winner '{}'",
slug, winnerSlug);
slug = winnerSlug;
}
}
// Fallback 1: slug already in DB (route misclassified or prior run)
WikiPageEntity existing = pageService.getBySlug(kbId, slug);
if (existing != null) {
pageService.updatePageByAi(kbId, slug, content, pageSummary, rawId);
pageService.mergeSourceLineage(existing.getId(), rawId, raw.getTitle());
log.info("[Wiki] Phase B create page slug='{}' done (updated existing)", slug);
return false;
}
String sourceRawIds = "[" + rawId + "]";
try {
WikiPageEntity created = pageService.createPage(kbId, slug, title, content, pageSummary, sourceRawIds, pageType);
pageService.mergeSourceLineage(created.getId(), rawId, raw.getTitle());
log.info("[Wiki] Phase B create page slug='{}' done (created)", slug);
citationService.buildCitationsAsync(created.getId(), kbId);
return true;
} catch (org.springframework.dao.DuplicateKeyException e) {
// Fallback 2: concurrent INSERT race — degrade to update
pageService.updatePageByAi(kbId, slug, content, pageSummary, rawId);
log.info("[Wiki] Phase B create page slug='{}' lost INSERT race -> updated existing", slug);
return false;
}
}
/**
* RFC-012 M2 v2 — 阶段 B 单页 merge:把 chunk 文本合并进一个已有页面。
*
* 输入仅几 KB(该页现有 content + chunk 主题片段),输出仅一页 markdown。
*
* @return true 表示成功 update 一页
*/
private boolean mergeOnePage(WikiKnowledgeBaseEntity kb, WikiRawMaterialEntity raw,
String chunkText, String slug) {
Long kbId = kb.getId();
Long rawId = raw.getId();
WikiPageEntity existing = pageService.getBySlug(kbId, slug);
if (existing == null) {
// 兜底:跨拼写 canonical 匹配——LLM 给的 slug 在 DB 里找不到,
// 但 canonical 形式(去连字符)对得上某个已有 page(典型场景:
// route 输出 `zhong-yao-qi-qing-pei-wu`,DB 存 `zhongyao-qiqing-peiwu`)
existing = pageService.findByCanonicalSlug(kbId, slug);
if (existing != null && !existing.getSlug().equals(slug)) {
log.info("[Wiki] Phase B merge slug='{}' canonical-matches existing '{}', using canonical slug for LLM call",
slug, existing.getSlug());
slug = existing.getSlug();
} else {
log.warn("[Wiki] Phase B merge page slug='{}' planned for update but not found in DB (even by canonical), skipping", slug);
return false;
}
}
String configContent = kb.getConfigContent() != null ? kb.getConfigContent() : "";
String mergeSystem = PromptLoader.loadPrompt("wiki/merge-page-system");
// Trim existing content to prevent context overflow on small models (qwen-turbo: 4096 tokens).
// Merging a 3000-char page + 30K chunk blows past the limit → truncated JSON → parse failure.
// 1800 chars ≈ ~600 tokens, leaving ample room for the chunk and response.
final int MAX_EXISTING_CHARS = 1800;
String rawExisting = existing.getContent() != null ? existing.getContent() : "";
String trimmedExisting = rawExisting.length() > MAX_EXISTING_CHARS
? rawExisting.substring(0, MAX_EXISTING_CHARS) + "\n...(内容已截断,请基于以上内容合并新信息)"
: rawExisting;
String mergeUserTemplate = PromptLoader.loadPrompt("wiki/merge-page-user");
String mergeUser = mergeUserTemplate
.replace("{config}", configContent)
.replace("{page_slug}", existing.getSlug() != null ? existing.getSlug() : slug)
.replace("{page_title}", existing.getTitle() != null ? existing.getTitle() : "")
.replace("{page_last_updated_by}", existing.getLastUpdatedBy() != null ? existing.getLastUpdatedBy() : "ai")
.replace("{page_content}", trimmedExisting)
.replace("{raw_title}", raw.getTitle())
.replace("{raw_content}", chunkText);
Prompt prompt = new Prompt(List.of(
new SystemMessage(mergeSystem),
new UserMessage(mergeUser)
));
if (isAborted(rawId, "merge slug=" + slug)) return false;
String response = callLlmWithResilientRetry(prompt,
"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);
return false;
}
String content = mergeJson.path("content").asText("");
String summary = mergeJson.path("summary").asText("");
if (content.isBlank()) {
log.warn("[Wiki] Phase B merge page slug='{}' returned blank content, skipping", slug);
return false;
}
WikiPageEntity updated = pageService.updatePageByAi(kbId, slug, content, summary, rawId);
log.info("[Wiki] Phase B merge page slug='{}' done", slug);
// RFC-047 P2: merge paired source lineage on update
if (updated != null) {
pageService.mergeSourceLineage(updated.getId(), rawId, raw.getTitle());
}
// RFC-029: async citation rebuild
if (updated != null) {
citationService.buildCitationsAsync(updated.getId(), kbId);
}
return true;
}
/**
* 解析 LLM 响应并创建/更新 Wiki 页面
*
* @return 创建+更新的页面总数
*/
private int applyLlmResponse(Long kbId, Long rawId, String llmResponse) {
JsonNode root = parseJsonResponse(llmResponse);
if (root == null) {
log.warn("[Wiki] Failed to parse LLM response for kbId={}, rawId={}, responseLen={}, first200={}",
kbId, rawId, llmResponse != null ? llmResponse.length() : 0,
llmResponse != null ? llmResponse.substring(0, Math.min(200, llmResponse.length())) : "null");
return 0;
}
// 结构校验:必须有 pages 数组
if (!root.has("pages") || !root.get("pages").isArray()) {
log.warn("[Wiki] LLM response missing 'pages' array for kbId={}, rawId={}", kbId, rawId);
return 0;
}
String sourceRawIds = "[" + rawId + "]";
int created = 0;
int updated = 0;
// 新页面
JsonNode pagesNode = root.path("pages");
if (pagesNode.isArray()) {
for (JsonNode pageNode : pagesNode) {
String slug = pageNode.path("slug").asText("");
String title = pageNode.path("title").asText("");
String content = pageNode.path("content").asText("");
String summary = pageNode.path("summary").asText("");
if (slug.isBlank() || title.isBlank()) continue;
// 检查是否已存在(LLM 可能将已有页面误判为新页面)
WikiPageEntity existing = pageService.getBySlug(kbId, slug);
if (existing != null) {
pageService.updatePageByAi(kbId, slug, content, summary, rawId);
updated++;
} else {
pageService.createPage(kbId, slug, title, content, summary, sourceRawIds);
created++;
}
}
}
// 更新的页面
JsonNode updatedPagesNode = root.path("updated_pages");
if (updatedPagesNode.isArray()) {
for (JsonNode pageNode : updatedPagesNode) {
String slug = pageNode.path("slug").asText("");
String content = pageNode.path("content").asText("");
String summary = pageNode.path("summary").asText("");
if (slug.isBlank()) continue;
WikiPageEntity existing = pageService.getBySlug(kbId, slug);
if (existing != null) {
// 保护手动编辑的页面:仍然更新,但 LLM 已在 prompt 中被告知要保留手动内容
pageService.updatePageByAi(kbId, slug, content, summary, rawId);
updated++;
}
}
}
log.info("[Wiki] Applied LLM response: kbId={}, rawId={}, created={}, updated={}",
kbId, rawId, created, updated);
return created + updated;
}
/**
* RFC-047 follow-up: Document-level analysis pass.
* Single LLM call on a sample of the full document to produce a concept map
* (topics + key_concepts + structure_notes). The result is injected into
* every chunk's route prompt so the router has global document awareness,
* reducing concept omissions caused by chunk-local context blindness.
*
* @return pretty-printed JSON string of the concept map, or "" on failure
*/
private String analyzeDocument(WikiKnowledgeBaseEntity kb, WikiRawMaterialEntity raw, String textContent) {
int sampleChars = Math.max(1000, properties.getDocumentAnalysisSampleChars());
String sample = textContent.length() > sampleChars
? textContent.substring(0, sampleChars) + "\n...[文档较长,以上为节选]"
: textContent;
String system = PromptLoader.loadPrompt("wiki/analyze-system");
String userTemplate = PromptLoader.loadPrompt("wiki/analyze-user");
String user = userTemplate
.replace("{raw_title}", raw.getTitle())
.replace("{text_sample}", sample);
Prompt prompt = new Prompt(List.of(
new SystemMessage(system),
new UserMessage(user)
));
try {
if (isAborted(raw.getId(), "doc analysis")) return "";
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={}",
raw.getId(),
json.path("topics").size(),
json.path("key_concepts").size());
return json.toPrettyString();
}
} catch (Exception e) {
log.warn("[Wiki] Document analysis failed for raw={}, continuing without: {}", raw.getId(), e.getMessage());
}
return "";
}
/**
* 构建已有 Wiki 页面索引(供 LLM 参考)
*/
private String buildExistingPagesIndex(Long kbId) {
List
* 可重试(一直重试直到成功):网络抖动、5xx、429 限流、超时、连接中断、内容过滤偶发、JSON 空输出。
*
* 立即终止(模型不可用):401/403 认证失败、模型不存在、quota 用尽、非法 API key、
* InterruptedException(优雅关停)。
*
* 使用指数退避(1s → 2s → 4s → ... → 封顶 60s)。
*
* RFC-012 M1:加入 maxAttempts 与 maxTotalDurationMs 双重上限,避免 nginx 504 这种
* 反复瞬时错误把单 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();
ResolvedChatModel resolved = resolveChatModel(kbId, step);
ChatModel chatModel = resolved.chatModel;
Long currentModelId = resolved.modelId;
boolean alreadyFellBack = false;
// When we fail over after a primary fatal, hold onto the primary's
// exception so we can surface BOTH errors if the fallback also dies.
// Operators reading logs need to see "the GLM 1113 we tried first" —
// not just "DashScope timed out" with the original cause lost.
Throwable primaryFatalError = null;
String primaryRootInfo = null;
int attempt = 0;
while (true) {
attempt++;
try {
ChatResponse response = chatModel.call(prompt);
if (response == null || response.getResult() == null
|| response.getResult().getOutput() == null
|| response.getResult().getOutput().getText() == null
|| response.getResult().getOutput().getText().isBlank()) {
throw new TransientLlmException("Empty response from model");
}
if (attempt > 1) {
log.info("[Wiki] LLM call for {} succeeded on attempt {}", ctx, attempt);
}
// The provider that just answered is healthy: clear any pending
// cooldown so subsequent calls can pick it freely. No-op on the
// happy path (counter is already 0); cleanup after a fallback
// success.
recordProviderSuccess(currentModelId);
// P0 telemetry: log token usage per LLM call so we can baseline costs before optimizing
long durationMs = (System.nanoTime() - startNanos) / 1_000_000L;
try {
var usage = response.getMetadata() != null ? response.getMetadata().getUsage() : null;
long pt = (usage != null && usage.getPromptTokens() != null) ? usage.getPromptTokens() : -1L;
long ct = (usage != null && usage.getCompletionTokens() != null) ? usage.getCompletionTokens() : -1L;
log.info("[wiki-telemetry] ctx={} promptTokens={} completionTokens={} durationMs={}",
ctx, pt, ct, durationMs);
} catch (Exception ignored) {}
return response.getResult().getOutput().getText();
} catch (Throwable t) {
if (Thread.currentThread().isInterrupted()) {
throw new RuntimeException("LLM call interrupted for " + ctx, t);
}
String rootInfo = summarizeRoot(t);
if (isFatalModelError(t)) {
// Mark the wedged provider down so the chat-agent path,
// any future fallback walks, and the admin diagnostics
// know to skip it during cooldown.
recordProviderFailure(currentModelId);
// One-hop fallback: when the primary provider is wedged
// (auth / quota / 余额不足 / model-not-found), try the next
// configured chat model with a different provider before
// giving up. The hop is one-shot — if the fallback also
// fatals, we surface BOTH errors rather than walking the
// whole chain to avoid pathological loops.
if (!alreadyFellBack) {
ResolvedChatModel next = pickFallbackChatModel(currentModelId);
if (next != null) {
log.warn("[Wiki] LLM fatal on primary model={} for {}; failing over to model={} (rootCause={}): {}",
currentModelId, ctx, next.modelId, rootInfo, t.getMessage());
primaryFatalError = t;
primaryRootInfo = rootInfo;
chatModel = next.chatModel;
currentModelId = next.modelId;
alreadyFellBack = true;
// Reset attempt counter so the fallback gets a fresh budget; total time
// budget continues to count down.
attempt = 0;
backoffMs = 1000;
continue;
}
}
if (alreadyFellBack && primaryFatalError != null) {
log.error("[Wiki] LLM unavailable (fatal on BOTH primary + fallback) for {}: primary rootCause={}: {} | fallback rootCause={}: {}",
ctx, primaryRootInfo, primaryFatalError.getMessage(), rootInfo, t.getMessage());
RuntimeException re = new RuntimeException(
"LLM unavailable on both primary and fallback. Primary (rootCause="
+ primaryRootInfo + "): " + primaryFatalError.getMessage()
+ " | Fallback (rootCause=" + rootInfo + "): " + t.getMessage(),
t);
re.addSuppressed(primaryFatalError);
throw re;
}
log.error("[Wiki] LLM unavailable (fatal) for {} after {} attempts on model={} (rootCause={}): {}",
ctx, attempt, currentModelId, rootInfo, t.getMessage());
throw new RuntimeException("LLM unavailable (rootCause=" + rootInfo + "): " + t.getMessage(), t);
}
long elapsedMs = (System.nanoTime() - startNanos) / 1_000_000L;
if (attempt >= maxAttempts || elapsedMs >= maxTotalDurationMs) {
log.error("[Wiki] LLM exhausted for {} after {} attempts in {}ms (limits: maxAttempts={}, maxTotalDurationMs={}, rootCause={}): {}",
ctx, attempt, elapsedMs, maxAttempts, maxTotalDurationMs, rootInfo, t.getMessage());
throw new RuntimeException("LLM exhausted after " + attempt + " attempts in " + elapsedMs
+ "ms (rootCause=" + rootInfo + "): " + t.getMessage(), t);
}
long sleepMs = Math.min(backoffMs, Math.max(0L, maxTotalDurationMs - elapsedMs));
log.warn("[Wiki] LLM transient failure for {} attempt={}/{} elapsed={}ms, retrying in {}ms (rootCause={}): {}",
ctx, attempt, maxAttempts, elapsedMs, sleepMs, rootInfo, t.getMessage());
try {
Thread.sleep(sleepMs);
} catch (InterruptedException ie) {
Thread.currentThread().interrupt();
throw new RuntimeException("LLM retry interrupted for " + ctx, ie);
}
backoffMs = Math.min(maxBackoffMs, backoffMs * 2);
}
}
}
private void recordProviderFailure(Long modelId) {
if (providerHealthTracker == null || modelId == null) return;
try {
ModelConfigEntity m = modelConfigService.getModel(modelId);
if (m != null && m.getProvider() != null) {
providerHealthTracker.recordFailure(m.getProvider());
}
} catch (Exception ignored) {}
}
private void recordProviderSuccess(Long modelId) {
if (providerHealthTracker == null || modelId == null) return;
try {
ModelConfigEntity m = modelConfigService.getModel(modelId);
if (m != null && m.getProvider() != null) {
providerHealthTracker.recordSuccess(m.getProvider());
}
} catch (Exception ignored) {}
}
/** Pair of modelId + built ChatModel — null modelId means we used the system default. */
private record ResolvedChatModel(Long modelId, ChatModel chatModel) {}
private ResolvedChatModel resolveChatModel(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 new ResolvedChatModel(modelId,
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());
}
}
ModelConfigEntity defaultModel = modelConfigService.getDefaultModel();
Long id = defaultModel == null ? null : defaultModel.getId();
return new ResolvedChatModel(id, buildChatModel());
}
/**
* Pick the next enabled chat model whose provider differs from the failed
* one, so we cycle to a fresh credential / billing account rather than
* retrying a wedged one. Returns null when no alternative exists.
*
* Provider order is deterministic — driven by
* {@code mate_model_provider.fallback_priority} (asc, lowest first) when
* {@link #modelProviderService} is wired. Falls back to
* {@code listEnabledModels} in DB order only when the provider service
* isn't available (older test harnesses), so behavior is at-least
* stable, never random.
*
* Skips providers currently in cooldown ({@link
* vip.mate.llm.failover.ProviderHealthTracker}) so a flapping provider
* doesn't keep getting tried while we wait for it to recover.
*/
private ResolvedChatModel pickFallbackChatModel(Long failedModelId) {
try {
String failedProviderId = null;
if (failedModelId != null) {
ModelConfigEntity failed = modelConfigService.getModel(failedModelId);
if (failed != null) failedProviderId = failed.getProvider();
}
if (modelProviderService != null) {
for (var provider : modelProviderService.listFallbackChain()) {
String pid = provider.getProviderId();
if (pid == null) continue;
if (pid.equals(failedProviderId)) continue;
if (providerHealthTracker != null && providerHealthTracker.isInCooldown(pid)) continue;
ResolvedChatModel built = firstChatModelForProvider(pid, failedModelId);
if (built != null) return built;
}
}
// Defensive fallback path — only triggers in tests / minimal
// environments without ModelProviderService wired. Mirrors the
// pre-PR behavior so legacy tests don't regress.
for (ModelConfigEntity candidate : modelConfigService.listEnabledModels()) {
if (!isUsableChatFallback(candidate, failedModelId, failedProviderId)) continue;
ResolvedChatModel built = tryBuild(candidate);
if (built != null) return built;
}
} catch (Exception e) {
log.debug("[Wiki] fallback lookup failed: {}", e.getMessage());
}
return null;
}
private ResolvedChatModel firstChatModelForProvider(String providerId, Long failedModelId) {
List
* 三类视为 fatal:
*
* 说明:关键字启发式在极少数场景可能误判(例如瞬时错误的 message 恰好含 "authentication"),
* 但实际云厂商 SDK 的错误消息规范度较高,这个风险可接受。
*/
private boolean isFatalModelError(Throwable t) {
Throwable cur = t;
int depth = 0;
while (cur != null && depth < 8) {
// 按异常类型直接判 fatal —— DNS / 连接拒绝 / TLS 问题重试都是浪费
String className = cur.getClass().getSimpleName();
if ("UnknownHostException".equals(className)
|| "SSLHandshakeException".equals(className)
|| "CertificateException".equals(className)
|| "SSLPeerUnverifiedException".equals(className)) {
return true;
}
if ("ConnectException".equals(className) && cur.getMessage() != null
&& cur.getMessage().toLowerCase().contains("refused")) {
return true;
}
String msg = cur.getMessage();
if (msg != null) {
String m = msg.toLowerCase();
// 鉴权 / 配额 / 模型不存在(HTTP 401/403 + 提供方错误字段)
if (m.contains("401") || m.contains("unauthorized")
|| m.contains("403") || m.contains("forbidden")
|| m.contains("invalid api key") || m.contains("invalid_api_key")
|| m.contains("authentication") || m.contains("api key not valid")
|| m.contains("model not found") || m.contains("model_not_found")
|| m.contains("invalidapikey") || m.contains("invalid_request_error")
|| m.contains("quota") || m.contains("insufficient_quota")
|| m.contains("no default model") || m.contains("model configuration")) {
return true;
}
// 中文 provider 余额耗尽:Zhipu 1113 / DashScope "Throttling" 中文返回 / 通用 "余额不足"。
// 原本走 "transient retry" 5 次 × 1s 退避,3 分钟才放弃;归类为 fatal 后立即触发 fallback hop。
// Chinese patterns checked against original msg (case-insensitive in Chinese is moot);
// codes / English snippets against the already-lowercased m.
if (msg.contains("余额不足") || msg.contains("请充值")
|| m.contains("\"code\":\"1113\"") || m.contains("\"code\":1113")
|| m.contains("accountbalancenotenough")
|| m.contains("balance not enough")) {
return true;
}
// 【Review Bug 3】prompt 结构性错误:重试也得同样结果,立即终止
if (m.contains("context_length_exceeded")
|| m.contains("context length")
|| m.contains("maximum context")
|| m.contains("max_tokens")
|| m.contains("prompt too long")
|| m.contains("input is too long")
|| m.contains("token limit")) {
return true;
}
// 【Review Bug 2】内容审核过滤:被 safety 挡下的 prompt 重试也是同样结果
if (m.contains("content_filter")
|| m.contains("content filter")
|| m.contains("data_inspection_failed")
|| (m.contains("safety") && m.contains("block"))) {
return true;
}
// 基础设施类永久错误:关键字兜底(和类名判断互补,跨语言 SDK 也能抓到)
if (m.contains("unknown host") || m.contains("no such host")
|| m.contains("connection refused")
|| m.contains("pkix path building failed")
|| m.contains("certificate verify failed")
|| m.contains("certificate_unknown")
|| m.contains("ssl handshake")) {
return true;
}
}
cur = cur.getCause();
depth++;
}
return false;
}
/**
* 沿 getCause() 遍历到最深,返回根因异常的 "SimpleName: message" 形式。
*
* Spring RestClient 会把 HTTP 层异常包装成 ResourceAccessException,外层消息统一是
* "I/O error on POST request for ...:
* 拼进最终抛出的 RuntimeException 消息里,UI 就算截断也能在前几十字看见类名。
*/
private String summarizeRoot(Throwable t) {
Throwable cur = t;
int depth = 0;
while (cur != null && cur.getCause() != null && cur.getCause() != cur && depth < 8) {
cur = cur.getCause();
depth++;
}
String cls = cur != null ? cur.getClass().getSimpleName() : "Unknown";
String msg = cur != null ? cur.getMessage() : null;
if (msg == null) return cls;
// 截短 message 避免把整段 HTML 错误页塞进异常链
String trimmed = msg.replaceAll("\\s+", " ").trim();
if (trimmed.length() > 200) trimmed = trimmed.substring(0, 200) + "...";
return cls + ": " + trimmed;
}
/** Transient error marker to route empty responses through the retry path */
private static class TransientLlmException extends RuntimeException {
TransientLlmException(String msg) { super(msg); }
}
// ==================== RFC-030: Error classification ====================
/**
* Classify an exception into an error code aligned with RFC-009 ErrorType.
*
* @return error code string: AUTH_ERROR, BILLING, MODEL_NOT_FOUND,
* RATE_LIMIT, SERVER_ERROR, TIMEOUT, CONTENT_FILTER, UNKNOWN
*/
public String classifyErrorCode(Throwable t) {
Throwable cur = t;
int depth = 0;
while (cur != null && depth < 8) {
String className = cur.getClass().getSimpleName();
if ("UnknownHostException".equals(className)
|| "SSLHandshakeException".equals(className)) {
return "AUTH_ERROR";
}
String msg = cur.getMessage();
if (msg != null) {
String m = msg.toLowerCase();
if (m.contains("401") || m.contains("unauthorized") || m.contains("403")
|| m.contains("forbidden") || m.contains("invalid api key")
|| m.contains("invalid_api_key") || m.contains("authentication")) {
return "AUTH_ERROR";
}
if (m.contains("quota") || m.contains("insufficient_quota") || m.contains("billing")) {
return "BILLING";
}
if (m.contains("model not found") || m.contains("model_not_found")) {
return "MODEL_NOT_FOUND";
}
if (m.contains("429") || m.contains("rate_limit") || m.contains("too many requests")) {
return "RATE_LIMIT";
}
if (m.contains("content_filter") || m.contains("content filter")
|| m.contains("data_inspection_failed")) {
return "CONTENT_FILTER";
}
if (m.contains("timeout") || m.contains("timed out")) {
return "TIMEOUT";
}
if (m.contains("500") || m.contains("502") || m.contains("503") || m.contains("504")) {
return "SERVER_ERROR";
}
}
cur = cur.getCause();
depth++;
}
return "UNKNOWN";
}
// ==================== RFC-031: Methods for template delegation ====================
/**
* Repair a single page by regenerating its content.
* Used by LocalRepairTemplate.
*/
public void repairSinglePage(Long targetPageId, Long modelId) {
WikiPageEntity page = pageService.getById(targetPageId);
if (page == null) {
log.warn("[Wiki] repairSinglePage: page not found: {}", targetPageId);
return;
}
WikiKnowledgeBaseEntity kb = kbService.getById(page.getKbId());
if (kb == null) return;
// Find source raw material
List
* {@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) {
WikiRawMaterialEntity raw = rawService.getById(rawId);
if (raw == null) {
log.info("[Wiki] Aborting {} for raw={}: raw was deleted mid-processing", ctx, rawId);
return true;
}
if (Boolean.TRUE.equals(raw.getCancelRequested())) {
log.info("[Wiki] Aborting {} for raw={}: cancellation requested by user", ctx, rawId);
return true;
}
return false;
}
/**
* RFC-051 PR-1c: bridge from {@link DocumentPreprocessService.Chunker} to
* the existing sentence-boundary chunker. Returns {@code [start, end]}
* pairs over the supplied text.
*/
private List
* Intentionally minimal: reuses the legacy {@code persistChunks(List
*
* 其余(网络、超时、5xx、429 限流、偶发空响应)均视为瞬时,按指数退避持续重试。
* >() {});
} catch (Exception e) {
return List.of();
}
}
private JsonNode parseJsonResponse(String response) {
if (response == null || response.isBlank()) return null;
String cleaned = response.trim();
// 1. 剥离 markdown 代码块标记
if (cleaned.startsWith("```json")) {
cleaned = cleaned.substring(7);
} else if (cleaned.startsWith("```")) {
cleaned = cleaned.substring(3);
}
if (cleaned.endsWith("```")) {
cleaned = cleaned.substring(0, cleaned.length() - 3);
}
cleaned = cleaned.trim();
// 2. 清洗控制字符(保留 \n \r \t),防止 LLM 输出含不可见字符导致 JSON 解析失败
cleaned = cleaned.replaceAll("[\\x00-\\x08\\x0B\\x0C\\x0E-\\x1F]", "");
// 3. 第一次尝试直接解析
try {
return objectMapper.readTree(cleaned);
} catch (Exception e) {
// 4. 如果整体不是 JSON,尝试提取第一个 JSON 对象块(LLM 可能在 JSON 前后加了说明文字)
int jsonStart = cleaned.indexOf("{");
int jsonEnd = cleaned.lastIndexOf("}");
if (jsonStart >= 0 && jsonEnd > jsonStart) {
String extracted = cleaned.substring(jsonStart, jsonEnd + 1);
try {
return objectMapper.readTree(extracted);
} catch (Exception e2) {
log.warn("[Wiki] Failed to parse extracted JSON block: {}", e2.getMessage());
}
}
log.warn("[Wiki] Failed to parse JSON response: {}", e.getMessage());
return null;
}
}
/**
* RFC-051 PR-6b follow-up: KB-level override for structured route output,
* falling back to the global property when the KB hasn't set a preference.
* Parse failures fall back to global too — never block ingest on bad config.
*/
private boolean resolveStructuredRouteFlag(WikiKnowledgeBaseEntity kb) {
boolean fallback = properties.isUseStructuredRoute();
if (kb == null || kb.getConfigContent() == null) return fallback;
try {
WikiKbConfig config = objectMapper.readValue(kb.getConfigContent(), WikiKbConfig.class);
return config.getUseStructuredRoute() != null
? config.getUseStructuredRoute()
: fallback;
} catch (Exception e) {
return fallback;
}
}
/**
* RFC-051 PR-1b: read {@code ingestMode} from KB config JSON. Returns null
* on any parse error or missing field so the caller falls through to eager.
*/
private String resolveIngestMode(WikiKnowledgeBaseEntity kb) {
if (kb == null || kb.getConfigContent() == null) return null;
try {
WikiKbConfig config = objectMapper.readValue(kb.getConfigContent(), WikiKbConfig.class);
return config.getIngestMode();
} catch (Exception e) {
log.warn("[Wiki] Failed to parse KB config for ingest mode, falling back to eager: {}", e.getMessage());
return null;
}
}
/**
* 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.
*