diff --git a/mateclaw-server/src/main/java/vip/mate/tool/builtin/DocumentExtractTool.java b/mateclaw-server/src/main/java/vip/mate/tool/builtin/DocumentExtractTool.java index 5c4cd5bc..eb7ae193 100644 --- a/mateclaw-server/src/main/java/vip/mate/tool/builtin/DocumentExtractTool.java +++ b/mateclaw-server/src/main/java/vip/mate/tool/builtin/DocumentExtractTool.java @@ -32,7 +32,7 @@ import java.util.zip.ZipInputStream; public class DocumentExtractTool { private static final int COMMAND_TIMEOUT_SECONDS = 30; - private static final int MAX_OUTPUT_LENGTH = 100000; // 100KB 限制 + private static final int MAX_OUTPUT_LENGTH = 500000; // 500KB — CLOB column has no size limit private static final boolean IS_WINDOWS = System.getProperty("os.name", "") .toLowerCase(Locale.ROOT).contains("win"); diff --git a/mateclaw-server/src/main/java/vip/mate/wiki/WikiProperties.java b/mateclaw-server/src/main/java/vip/mate/wiki/WikiProperties.java index 9d64bd2c..71702894 100644 --- a/mateclaw-server/src/main/java/vip/mate/wiki/WikiProperties.java +++ b/mateclaw-server/src/main/java/vip/mate/wiki/WikiProperties.java @@ -33,12 +33,15 @@ public class WikiProperties { private int maxParallelRawMaterials = 3; /** - * 单个材料内 chunk 的并行处理数上限。 + * Max parallel chunks within a single raw material. *

- * RFC-012 Change 1:从硬编码 3 提到 5,并暴露为配置项。默认总并发为 - * maxParallelRawMaterials × maxParallelChunks = 15,仍在常见 60 RPM 限额下。 + * RFC-047 P3: Changed from 5 to 1. Parallel chunks all share the same existingPagesIndex + * snapshot, so chunk N cannot see pages created by chunk N-1 — causing duplicate pages and + * stale-index collisions. Serializing chunks eliminates this class of bug. Document-level + * parallelism (maxParallelRawMaterials) is preserved, so overall throughput is unchanged + * for multi-document batches. */ - private int maxParallelChunks = 5; + private int maxParallelChunks = 1; /** * 单个 chunk 内 phase B 阶段的 page 并行处理数上限。 @@ -84,6 +87,24 @@ public class WikiProperties { */ private long llmMaxTotalDurationMs = 240_000; + /** + * RFC-047 P1: Max pages per BatchCreate LLM call. + * Pages planned by route are chunked into sub-batches of this size; + * a local liveIndex is updated between sub-batches so later pages can + * link to earlier ones created in the same chunk. + * Default 5: keeps output tokens per call predictable while covering + * the typical 3–5 creates per chunk in a single call. + */ + private int batchCreatePageSize = 3; + + /** + * RFC-047: Minimum chunk length (chars) for the chunk-fallback mechanism. + * If route returns 0 create+update entries and the chunk exceeds this threshold, + * an overview page is auto-injected so no substantial content is silently dropped. + * Chunks shorter than this (e.g. TOC lines, blank pages) are allowed to produce nothing. + */ + private int chunkFallbackMinChars = 200; + /** * 是否启用两阶段消化(路由 → 逐页 merge)。 *

@@ -127,4 +148,20 @@ public class WikiProperties { /** Maximum characters for local repair single-page regeneration */ private int localRepairMaxChars = 8000; + + /** + * Whether to run a document-level analysis pass before routing. + * When enabled, a single LLM call produces a concept map (topics + key_concepts) + * that is injected into every chunk's route prompt, giving the router global + * awareness of the document structure and reducing concept omissions. + * Adds ~1 LLM call and 10-20s per raw material. + */ + private boolean useDocumentAnalysis = true; + + /** + * Max characters of document text fed to the analysis pass. + * Larger values improve coverage but increase prompt tokens. + * Default 15000 covers most documents while staying well within model limits. + */ + private int documentAnalysisSampleChars = 15000; } diff --git a/mateclaw-server/src/main/java/vip/mate/wiki/model/WikiPageEntity.java b/mateclaw-server/src/main/java/vip/mate/wiki/model/WikiPageEntity.java index 554885ba..785607c4 100644 --- a/mateclaw-server/src/main/java/vip/mate/wiki/model/WikiPageEntity.java +++ b/mateclaw-server/src/main/java/vip/mate/wiki/model/WikiPageEntity.java @@ -42,6 +42,10 @@ public class WikiPageEntity { @TableField(updateStrategy = FieldStrategy.ALWAYS) private String sourceRawIds; + /** RFC-047 P2: paired source lineage — JSON array of {rawId, rawTitle} objects. Canonical; dual-written with sourceRawIds. */ + @TableField(updateStrategy = FieldStrategy.ALWAYS) + private String sourceEntries; + /** Page type: entity / concept / source / synthesis */ private String pageType; diff --git a/mateclaw-server/src/main/java/vip/mate/wiki/service/WikiPageService.java b/mateclaw-server/src/main/java/vip/mate/wiki/service/WikiPageService.java index 4d842068..c85ac8b7 100644 --- a/mateclaw-server/src/main/java/vip/mate/wiki/service/WikiPageService.java +++ b/mateclaw-server/src/main/java/vip/mate/wiki/service/WikiPageService.java @@ -248,6 +248,43 @@ public class WikiPageService { return existing; } + /** + * RFC-047 P2: Paired source lineage entry (rawId + rawTitle snapshot at ingest time). + * Keyed by rawId; rawTitle is a snapshot — the raw may be renamed later but lineage stays accurate. + */ + public record SourceEntry(long rawId, String rawTitle) {} + + /** + * RFC-047 P2: Merge a (rawId, rawTitle) pair into a page's source lineage. + * Dual-writes to both sourceEntries (canonical) and sourceRawIds (legacy compat). + * Idempotent: no-ops if rawId already present. + */ + @Transactional + public void mergeSourceLineage(Long pageId, Long rawId, String rawTitle) { + WikiPageEntity page = pageMapper.selectById(pageId); + if (page == null) return; + + List entries = parseSourceEntries(page.getSourceEntries()); + boolean entryExists = entries.stream().anyMatch(e -> e.rawId() == rawId); + + List rawIds = parseSourceRawIds(page.getSourceRawIds()); + boolean idExists = rawIds.contains(rawId); + + if (!entryExists) { + entries.add(new SourceEntry(rawId, rawTitle != null ? rawTitle : "")); + page.setSourceEntries(toJson(entries)); + } + if (!idExists) { + rawIds.add(rawId); + page.setSourceRawIds(toJson(rawIds)); + } + + if (!entryExists || !idExists) { + pageMapper.updateById(page); + evictSummaryCache(page.getKbId()); + } + } + /** * 手动更新页面内容 */ @@ -347,13 +384,15 @@ public class WikiPageService { List sourceIds = parseSourceRawIds(page.getSourceRawIds()); if (sourceIds.contains(rawId)) { if (sourceIds.size() == 1) { - // 独占页面:直接删除 delete(kbId, page.getSlug()); deleted++; } else { - // 多来源页面:仅移除该 rawId 引用 + // Multi-source page: remove this rawId from both sourceRawIds and sourceEntries sourceIds.remove(rawId); page.setSourceRawIds(toJson(sourceIds)); + List entries = parseSourceEntries(page.getSourceEntries()); + entries.removeIf(e -> e.rawId() == rawId); + page.setSourceEntries(toJson(entries)); pageMapper.updateById(page); } } @@ -406,6 +445,15 @@ public class WikiPageService { } } + private List parseSourceEntries(String json) { + if (json == null || json.isBlank()) return new ArrayList<>(); + try { + return objectMapper.readValue(json, new TypeReference>() {}); + } catch (Exception e) { + return new ArrayList<>(); + } + } + private String toJson(Object obj) { try { return objectMapper.writeValueAsString(obj); 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 e27f9311..28fb2681 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 @@ -185,6 +185,13 @@ public class WikiProcessingService { // Phase 3: 构建已有页面索引(一次构建,所有 chunk 共用) String existingPagesIndex = buildExistingPagesIndex(kb.getId()); + // Phase 3b: Document-level analysis (optional, RFC-047 follow-up) + // Single LLM call producing a concept map injected into every chunk's route prompt. + String documentMap = ""; + if (properties.isUseDocumentAnalysis()) { + documentMap = analyzeDocument(kb, raw, textContent); + } + // Transition job to phase_a (chunk processing begins) if (wikiJobService != null && jobId != null) { try { wikiJobService.transition(jobId, vip.mate.wiki.job.WikiJobStage.PHASE_A_RUNNING); } catch (Exception ignored) {} @@ -194,7 +201,7 @@ public class WikiProcessingService { // result[0] = totalPages, result[1] = failedChunks, result[2] = totalChunks int[] result; if (textContent.length() > properties.getMaxChunkSize()) { - result = processInChunks(kb, raw, textContent, existingPagesIndex); + result = processInChunks(kb, raw, textContent, existingPagesIndex, documentMap); } else { // 单 chunk 也持久化(RFC-013:保证所有 chunk 都入库) try { @@ -203,7 +210,7 @@ public class WikiProcessingService { } catch (Exception e) { log.warn("[Wiki] Single chunk persistence failed for raw={}: {}", rawId, e.getMessage()); } - int pages = processChunk(kb, raw, textContent, existingPagesIndex); + int pages = processChunk(kb, raw, textContent, existingPagesIndex, documentMap); result = new int[]{pages, pages == 0 ? 1 : 0, 1}; } @@ -360,7 +367,7 @@ public class WikiProcessingService { * @return int[3]: [totalPages, failedChunks, totalChunks] */ private int[] processInChunks(WikiKnowledgeBaseEntity kb, WikiRawMaterialEntity raw, String text, - String existingPagesIndex) { + String existingPagesIndex, String documentMap) { // Phase 1: 切分文本为 chunks(带偏移,供持久化) List chunksWithOffset = splitIntoChunksWithOffsets(text); List chunks = chunksWithOffset.stream().map(ChunkWithOffset::text).toList(); @@ -379,7 +386,7 @@ public class WikiProcessingService { if (totalChunks == 1) { // 单 chunk 不走并行 try { - int pages = processChunk(kb, raw, chunks.get(0), existingPagesIndex); + int pages = processChunk(kb, raw, chunks.get(0), existingPagesIndex, documentMap); return new int[]{pages, pages == 0 ? 1 : 0, 1}; } catch (Exception e) { log.warn("[Wiki] Single chunk failed: {}", e.getMessage()); @@ -407,7 +414,7 @@ public class WikiProcessingService { } try { log.info("[Wiki] Processing chunk {}/{}: {} chars", chunkIndex + 1, totalChunks, chunk.length()); - int pages = processChunk(kb, raw, chunk, existingPagesIndex); + int pages = processChunk(kb, raw, chunk, existingPagesIndex, documentMap); totalPages.addAndGet(pages); } catch (Exception e) { failedChunks.incrementAndGet(); @@ -508,10 +515,10 @@ public class WikiProcessingService { * @return 创建+更新的页面数 */ private int processChunk(WikiKnowledgeBaseEntity kb, WikiRawMaterialEntity raw, String textContent, - String existingPagesIndex) { + String existingPagesIndex, String documentMap) { // RFC-012 M2:两阶段消化(路由 → 逐页 merge),单次 LLM 调用输出量大幅缩减,避免 nginx 60s 网关超时 if (properties.isUseTwoPhaseDigest()) { - return processChunkTwoPhase(kb, raw, textContent, existingPagesIndex); + return processChunkTwoPhase(kb, raw, textContent, existingPagesIndex, documentMap); } // 旧路径:单次调用让 LLM 同时处理新建 + 全量 merge(输出爆炸,易触发 504) @@ -548,7 +555,7 @@ public class WikiProcessingService { * @return 创建+更新的页面数 */ private int processChunkTwoPhase(WikiKnowledgeBaseEntity kb, WikiRawMaterialEntity raw, - String textContent, String existingPagesIndex) { + String textContent, String existingPagesIndex, String documentMap) { Long kbId = kb.getId(); Long rawId = raw.getId(); String configContent = kb.getConfigContent() != null ? kb.getConfigContent() : ""; @@ -561,8 +568,12 @@ public class WikiProcessingService { // ─── 阶段 A:路由 ─── 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}", existingPagesIndex) .replace("{raw_title}", rawTitle) .replace("{raw_content}", textContent); @@ -603,6 +614,22 @@ public class WikiProcessingService { } } int totalPlanned = createMetas.size() + updateSlugs.size(); + + // Chunk fallback: if route returned nothing for a non-trivial chunk, inject an overview page + // so no content is silently dropped (mirrors llm_wiki source-summary guarantee). + if (totalPlanned == 0 && textContent.length() >= properties.getChunkFallbackMinChars()) { + String overviewSlug = WikiPageService.toSlug(rawTitle) + "-overview"; + com.fasterxml.jackson.databind.node.ObjectNode fallbackMeta = + objectMapper.createObjectNode(); + fallbackMeta.put("slug", overviewSlug); + fallbackMeta.put("title", rawTitle + " 概述"); + fallbackMeta.put("summary", "来自「" + rawTitle + "」的综合概述,涵盖本章节的核心内容。"); + createMetas.add(fallbackMeta); + totalPlanned = 1; + log.info("[Wiki] Chunk fallback: route returned empty for rawId={} chunkLen={}, injecting overview page '{}'", + rawId, textContent.length(), overviewSlug); + } + log.info("[Wiki] Route phase: kbId={}, rawId={}, planned create={}, planned update={}", kbId, rawId, createMetas.size(), updateSlugs.size()); @@ -628,46 +655,11 @@ public class WikiProcessingService { int parallelPages = Math.max(1, properties.getMaxParallelPhaseBPages()); Semaphore pageSem = new Semaphore(parallelPages); - // ─── 阶段 B-1:并行 create ─── - List> createFutures = new ArrayList<>(createMetas.size()); - for (JsonNode meta : createMetas) { - final JsonNode metaRef = meta; - createFutures.add(CompletableFuture.runAsync(() -> { - try { - pageSem.acquire(); - } catch (InterruptedException ie) { - Thread.currentThread().interrupt(); - return; - } - boolean ok = false; - try { - try { - if (createOnePage(kb, raw, textContent, existingPagesIndex, metaRef)) { - created.incrementAndGet(); - } - // createOnePage 内部的 DuplicateKey / canonical / claim fallback 不抛异常 → ok=true。 - ok = true; - } catch (RuntimeException e) { - log.warn("[Wiki] Phase B create page slug='{}' failed: {}", - metaRef.path("slug").asText(""), e.getMessage()); - } - if (pc != null) { - int d = pc.done.incrementAndGet(); - if (!ok) 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", - "ok", ok, - "done", d, - "total", pc.total.get())); - } - } finally { - pageSem.release(); - } - }, WIKI_EXECUTOR)); - } + // ─── Phase B-1: BatchCreate (RFC-047 P1) ─── + // One LLM call for all N creates instead of N individual calls. + // batchCreatePages handles sub-batching, liveIndex updates, and progress counting. + batchCreatePages(kb, raw, textContent, existingPagesIndex, createMetas, created, pc); + List> createFutures = new ArrayList<>(0); // kept for allOf join below // ─── 阶段 B-2:并行 merge ─── List> mergeFutures = new ArrayList<>(updateSlugs.size()); @@ -715,30 +707,231 @@ public class WikiProcessingService { CompletableFuture.allOf(allFutures.toArray(new CompletableFuture[0])).join(); // 单 chunk 完成时不写"done"——多 chunk 还在跑;最终"done"由 processRawMaterial 的 finally 写入 - log.info("[Wiki] Two-phase digest applied: kbId={}, rawId={}, created={}, updated={}", - kbId, rawId, created.get(), updated.get()); + log.info("[wiki-telemetry-chunk] kbId={} rawId={} chunkLen={} indexLen={} creates={} updates={}", + kbId, rawId, textContent.length(), existingPagesIndex.length(), created.get(), updated.get()); return created.get() + updated.get(); } /** - * RFC-012 M2 v2 — 阶段 B 单页生成:用 chunk 文本 + 该页 metadata 让 LLM 写出完整页面。 + * RFC-047 P1: BatchCreate — single LLM call generating all new pages for a chunk. *

- * 输入仅几 KB(chunk 主题片段 + metadata + 已有页索引),输出仅一页 markdown, - * 单次调用稳稳 ≤ 60 秒,避免 nginx 60s 网关。 - *

- * 兜底:如果 slug 已存在(route 误判),改走 update 路径。 - * - * @return true 表示成功 create 一页(或兜底 update 一页时返回 false 以让上层归到 update 计数) + * 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 boolean createOnePage(WikiKnowledgeBaseEntity kb, WikiRawMaterialEntity raw, - String chunkText, String existingPagesIndex, JsonNode meta) { + private int batchCreatePages(WikiKnowledgeBaseEntity kb, WikiRawMaterialEntity raw, + String chunkText, String existingPagesIndex, + List createMetas, AtomicInteger created, + ProgressCounter pc) { + if (createMetas.isEmpty()) return 0; Long kbId = kb.getId(); Long rawId = raw.getId(); - String slug = meta.path("slug").asText(""); - String title = meta.path("title").asText(""); - String summary = meta.path("summary").asText(""); - if (slug.isBlank() || title.isBlank()) return false; + String configContent = kb.getConfigContent() != null ? kb.getConfigContent() : ""; + int batchSize = Math.max(1, properties.getBatchCreatePageSize()); + WikiBatchCreateParser batchParser = new WikiBatchCreateParser(); + StringBuilder liveIndex = new StringBuilder(existingPagesIndex); + int totalCreated = 0; + + for (int bStart = 0; bStart < createMetas.size(); bStart += batchSize) { + int bEnd = Math.min(bStart + batchSize, createMetas.size()); + List subBatch = createMetas.subList(bStart, bEnd); + + // Build pages_to_create JSON array for this sub-batch + StringBuilder metasJson = new StringBuilder("["); + for (int i = 0; i < subBatch.size(); i++) { + if (i > 0) metasJson.append(","); + metasJson.append(subBatch.get(i).toString()); + } + metasJson.append("]"); + + String batchSystem = PromptLoader.loadPrompt("wiki/batch-create-system"); + String batchUserTemplate = PromptLoader.loadPrompt("wiki/batch-create-user"); + String batchUser = batchUserTemplate + .replace("{config}", configContent) + .replace("{existing_pages}", liveIndex.toString()) + .replace("{pages_to_create}", metasJson.toString()) + .replace("{raw_title}", raw.getTitle()) + .replace("{raw_content}", chunkText); + Prompt batchPrompt = new Prompt(List.of( + new SystemMessage(batchSystem), + new UserMessage(batchUser) + )); + + String batchResponse = callLlmWithResilientRetry(batchPrompt, + "batch-create " + subBatch.size() + " pages of raw=" + rawId + + " subBatch=" + (bStart / batchSize + 1)); + + List parsedPages = batchParser.parse(batchResponse); + + // Fallback: if LLM ignored FILE block format (common for single-page sub-batches), + // try treating the entire response as a bare JSON page. + if (parsedPages.isEmpty() && subBatch.size() == 1 && batchResponse != null && !batchResponse.isBlank()) { + String fallbackSlug = subBatch.get(0).path("slug").asText(""); + parsedPages = List.of(new WikiBatchCreateParser.ParsedPage(fallbackSlug, batchResponse.strip())); + log.info("[Wiki] BatchCreate sub-batch {}: no FILE blocks found, trying bare-JSON fallback for slug='{}'", + bStart / batchSize + 1, fallbackSlug); + } + + log.info("[Wiki] BatchCreate sub-batch {}: planned={} parsed={}", + bStart / batchSize + 1, subBatch.size(), parsedPages.size()); + + for (WikiBatchCreateParser.ParsedPage pp : parsedPages) { + JsonNode pageJson = parseJsonResponse(pp.rawJson()); + if (pageJson == null) { + log.warn("[Wiki] BatchCreate: unparseable JSON for slug='{}', skipping", pp.slug()); + continue; + } + // Use slug from JSON body; fall back to FILE header slug + String slug = pageJson.path("slug").asText(pp.slug()); + if (slug.isBlank()) slug = pp.slug(); + String title = pageJson.path("title").asText(""); + String content = pageJson.path("content").asText(""); + String pageSummary = pageJson.path("summary").asText(""); + if (content.isBlank()) { + log.info("[Wiki] BatchCreate: blank content for slug='{}', retrying individually", slug); + final String blankSlug = slug; + final String headerSlug = pp.slug(); + // Find the original meta from subBatch for this slug so retrySingleCreate has title+summary + JsonNode retryMeta = subBatch.stream() + .filter(m -> blankSlug.equals(m.path("slug").asText("")) + || headerSlug.equals(m.path("slug").asText(""))) + .findFirst().orElse(null); + boolean ok = false; + if (retryMeta != null) { + try { + String retryResult = retrySingleCreate(kb, raw, chunkText, liveIndex.toString(), retryMeta); + if (retryResult != null) { + JsonNode retryJson = parseJsonResponse(retryResult); + if (retryJson != null) { + String rContent = retryJson.path("content").asText(""); + String rTitle = retryJson.path("title").asText(title); + String rSummary = retryJson.path("summary").asText(pageSummary); + if (!rContent.isBlank()) { + boolean wasCreated = savePageContent(kb, raw, slug, rTitle, rContent, rSummary); + if (wasCreated) { + created.incrementAndGet(); + totalCreated++; + String brief = rSummary.length() > 100 ? rSummary.substring(0, 100) : rSummary; + liveIndex.append("\n- ").append(slug).append(": ").append(brief); + } + ok = true; + } + } + } + } catch (Exception retryEx) { + log.warn("[Wiki] BatchCreate: blank-content retry for slug='{}' failed: {}", slug, retryEx.getMessage()); + } + } + if (pc != null) { + int d = pc.done.incrementAndGet(); + if (!ok) 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", ok, "done", d, "total", pc.total.get())); + } + continue; + } + + boolean wasCreated = false; + boolean ok = false; + try { + wasCreated = savePageContent(kb, raw, slug, title, content, pageSummary); + if (wasCreated) { + created.incrementAndGet(); + totalCreated++; + // Append to liveIndex so next sub-batch can link to this page + String briefSummary = pageSummary.length() > 100 + ? pageSummary.substring(0, 100) : pageSummary; + liveIndex.append("\n- ").append(slug).append(": ").append(briefSummary); + } + ok = true; + } catch (RuntimeException e) { + log.warn("[Wiki] BatchCreate: savePageContent failed for slug='{}': {}", slug, e.getMessage()); + } + + if (pc != null) { + int d = pc.done.incrementAndGet(); + if (!ok) 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", + "ok", ok, + "done", d, + "total", pc.total.get())); + } + } + + // Retry any pages that LLM omitted from the batch response + int subBatchNum = bStart / batchSize + 1; + java.util.Set returnedSlugs = new java.util.HashSet<>(); + for (WikiBatchCreateParser.ParsedPage pp : parsedPages) { + returnedSlugs.add(pp.slug()); + } + for (JsonNode missingMeta : subBatch) { + String missingSlug = missingMeta.path("slug").asText(""); + if (missingSlug.isBlank() || returnedSlugs.contains(missingSlug)) continue; + log.info("[Wiki] BatchCreate sub-batch {}: slug='{}' missing, retrying individually", + subBatchNum, missingSlug); + boolean ok = false; + try { + String retryResult = retrySingleCreate(kb, raw, chunkText, liveIndex.toString(), + missingMeta); + if (retryResult != null) { + JsonNode retryJson = parseJsonResponse(retryResult); + if (retryJson != null) { + String slug = retryJson.path("slug").asText(missingSlug); + if (slug.isBlank()) slug = missingSlug; + String title = retryJson.path("title").asText(""); + String content = retryJson.path("content").asText(""); + String pageSummary = retryJson.path("summary").asText(""); + if (!content.isBlank()) { + boolean wasCreated = savePageContent(kb, raw, slug, title, content, pageSummary); + if (wasCreated) { + created.incrementAndGet(); + totalCreated++; + String brief = pageSummary.length() > 100 + ? pageSummary.substring(0, 100) : pageSummary; + liveIndex.append("\n- ").append(slug).append(": ").append(brief); + } + ok = true; + log.info("[Wiki] BatchCreate sub-batch {}: retry for slug='{}' succeeded", + subBatchNum, missingSlug); + } + } + } + } catch (Exception retryEx) { + log.warn("[Wiki] BatchCreate sub-batch {}: retry for slug='{}' failed: {}", + subBatchNum, missingSlug, retryEx.getMessage()); + } + if (pc != null) { + int d = pc.done.incrementAndGet(); + if (!ok) 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", ok, "done", d, "total", pc.total.get())); + } + } + } + return totalCreated; + } + + /** + * RFC-047 P1 retry: call the single-page create prompt for one missing slug. + * Used when a BatchCreate sub-batch omits a planned page. + * + * @return raw LLM response string, or null on failure + */ + private String retrySingleCreate(WikiKnowledgeBaseEntity kb, WikiRawMaterialEntity raw, + String chunkText, String existingPagesIndex, + JsonNode pageMeta) { + String slug = pageMeta.path("slug").asText(""); + String title = pageMeta.path("title").asText(""); + String summary = pageMeta.path("summary").asText(""); String configContent = kb.getConfigContent() != null ? kb.getConfigContent() : ""; String createSystem = PromptLoader.loadPrompt("wiki/create-page-system"); String createUserTemplate = PromptLoader.loadPrompt("wiki/create-page-user"); @@ -754,80 +947,74 @@ public class WikiProcessingService { new SystemMessage(createSystem), new UserMessage(createUser) )); - String response = callLlmWithResilientRetry(prompt, - "create page slug=" + slug + " of raw=" + rawId); - JsonNode pageJson = parseJsonResponse(response); - if (pageJson == null) { - log.warn("[Wiki] Phase B create page slug='{}' returned unparseable JSON, skipping", slug); - return false; - } - String content = pageJson.path("content").asText(""); - String pageSummary = pageJson.path("summary").asText(""); - if (pageSummary.isBlank()) pageSummary = summary; - if (content.isBlank()) { - log.warn("[Wiki] Phase B create page slug='{}' returned blank content, skipping", slug); - return false; - } + return callLlmWithResilientRetry(prompt, "retry-create slug=" + slug + " of raw=" + raw.getId()); + } - // 兜底 0:跨拼写 canonical 匹配(DB 已有 page,但 slug 拼写不同) - // 覆盖场景:之前上传的 raw 已经创建了同概念 page,本次 LLM 给了不同拼写 + /** + * 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) { + Long kbId = kb.getId(); + Long rawId = raw.getId(); + + // 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; } - // 兜底 0.5:跨 chunk in-flight slug 抢占(同一 raw 的另一并发 chunk 已声明同概念) - // computeIfAbsent 是原子操作,先到的 chunk 把自己 slug 注册为 winner + // 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()) { - // lambda 要求 effectively final,用 finalSlug 副本 final String routedSlug = slug; String winnerSlug = pcLocal.slugClaims.computeIfAbsent(canonical, k -> routedSlug); if (!winnerSlug.equals(slug)) { - // 另一 chunk 先 claim 了同 canonical,但用了不同 slug 拼写 WikiPageEntity winner = pageService.getBySlug(kbId, winnerSlug); if (winner != null) { - // winner 已 INSERT 进 DB → 直接 update 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; } - // winner claim 早于 INSERT(claim 是 in-memory,INSERT 是 DB IO) - // → 用 winnerSlug 继续走下面的 INSERT 路径,DuplicateKey fallback 会兜住实际 race log.info("[Wiki] Phase B create slug='{}' redirects to in-flight winner '{}'", slug, winnerSlug); slug = winnerSlug; } } - // 兜底 1:如果 slug 已存在(route 误判 / 上一次成功 INSERT),走 update 而不是 create + // 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; // 不计入 created + return false; } + String sourceRawIds = "[" + rawId + "]"; try { WikiPageEntity created = pageService.createPage(kbId, slug, title, content, pageSummary, sourceRawIds); + pageService.mergeSourceLineage(created.getId(), rawId, raw.getTitle()); log.info("[Wiki] Phase B create page slug='{}' done (created)", slug); - // RFC-029: async citation build citationService.buildCitationsAsync(created.getId(), kbId); return true; } catch (org.springframework.dao.DuplicateKeyException e) { - // 兜底 2:select-then-create 在并发下不是原子操作。当 N 个 chunk 同时 - // route 出相同 slug,只有第一个 INSERT 能成功,其余都会触发 H2/MySQL - // unique key violation。本次 chunk 的 LLM 输出仍有价值——降级为 update, - // 把内容合并进已存在的 page,而不是丢弃。 + // 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; // 不计入 created + return false; } } @@ -888,6 +1075,10 @@ public class WikiProcessingService { } 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); @@ -966,6 +1157,45 @@ public class WikiProcessingService { 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 { + String response = callLlmWithResilientRetry(prompt, "analyze doc raw=" + raw.getId()); + 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 参考) */ @@ -1037,6 +1267,15 @@ public class WikiProcessingService { if (attempt > 1) { log.info("[Wiki] LLM call for {} succeeded on attempt {}", ctx, attempt); } + // 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()) { diff --git a/mateclaw-server/src/main/resources/prompts/wiki/batch-create-system.txt b/mateclaw-server/src/main/resources/prompts/wiki/batch-create-system.txt index dbbfe021..ae4ed9e9 100644 --- a/mateclaw-server/src/main/resources/prompts/wiki/batch-create-system.txt +++ b/mateclaw-server/src/main/resources/prompts/wiki/batch-create-system.txt @@ -12,8 +12,6 @@ - ❌ **不要生成 `pages_to_create` 以外的页面** - ❌ **不要使用 markdown 代码块**(不要用 ``` 包裹) - ❌ **不要生成无关内容** —— 每个页面只包含与该页主题相关的信息 -- ❌ **不要编造原始材料中没有的信息** —— 只从提供的原始材料和概念地图中提取 -- ❌ **不要把概念地图的分析文字直接复制到页面** —— 用它来理解关系,用原始材料来写内容 ## 每个页面的格式规则 @@ -28,7 +26,7 @@ 每个页面输出一个 FILE 块,格式如下: ---FILE: {slug}--- -{"slug":"...","title":"...","summary":"...","page_type":"concept","content":"## 标题\n\n摘要...\n\n### 章节\n..."} +{"slug":"...","title":"...","summary":"...","content":"## 标题\n\n摘要...\n\n### 章节\n..."} ---END FILE--- 规则: diff --git a/mateclaw-server/src/main/resources/prompts/wiki/batch-create-user.txt b/mateclaw-server/src/main/resources/prompts/wiki/batch-create-user.txt index 7a5a5f39..73bd485b 100644 --- a/mateclaw-server/src/main/resources/prompts/wiki/batch-create-user.txt +++ b/mateclaw-server/src/main/resources/prompts/wiki/batch-create-user.txt @@ -2,8 +2,6 @@ {config} -{document_map_section} - ## 已有 Wiki 页面索引(用于建立 [[链接]]) {existing_pages} @@ -25,6 +23,5 @@ 请为上面 `pages_to_create` 数组中的**每一个页面**生成完整的 markdown 内容。 - 按 system 中规定的 FILE 块格式输出,每个页面一个 FILE 块 - 每个页面的内容必须基于原始材料中与该主题相关的信息 -- 如果提供了"文档全局概念地图",可参考其中的概念关系来丰富页面内容和链接 - 适当使用 [[页面标题]] 链接到相关页面(同批次内其他页面也可链接) - 不要遗漏任何一个页面 diff --git a/mateclaw-server/src/main/resources/prompts/wiki/route-system.txt b/mateclaw-server/src/main/resources/prompts/wiki/route-system.txt index 78ad4f3a..1dd19315 100644 --- a/mateclaw-server/src/main/resources/prompts/wiki/route-system.txt +++ b/mateclaw-server/src/main/resources/prompts/wiki/route-system.txt @@ -6,7 +6,12 @@ 2. **比对已有页面索引**(仅含 slug + title + summary,不含正文): - 材料中出现的概念**已经被某个已有页面充分覆盖** → 放入 `update` 列表(只写 slug) - 材料中出现的概念**没有任何已有页面覆盖** → 放入 `create` 列表(只给 slug + title + summary,**不要正文**) -3. 不创建浅薄页面(< 3 句实质内容的概念跳过,不进任何列表) +3. 宁可多列一页也不要遗漏——只跳过完全没有实质信息的内容(如空白段、目录行、单句无意义重复) +4. **兜底规则**:如果整段材料有实质内容,`create` + `update` 合计**至少 1 条**。确实找不到独立概念时,把本段最核心的主题作为一个概念放入 `create`。 +5. **概念地图约束**(当 user prompt 中提供了「文档全局概念地图」时): + - 概念地图仅供参考——**只把当前 chunk 中有足够内容支撑的概念列入 create/update** + - 概念地图中出现的概念,若当前 chunk 只是一句话带过或完全未提及,**不要列入** + - 同一概念不要跨 chunk 重复 create——如果已有页面索引里已有该概念,放入 update ## 你不做什么 @@ -54,7 +59,7 @@ 字段说明: - `create`:完全新增的页面,**只给 metadata,不给 content** - `update`:已存在但需要根据本材料合并更新的页面,**只列 slug 字符串数组** -- 两个数组都可以为空(材料完全无价值时全空,材料只更新已有页时 create 为空,材料只新增页时 update 为空) +- 两个数组合计通常至少 1 条(材料完全无价值——如空白页、纯目录——才允许全空) ## 关键纪律 diff --git a/mateclaw-server/src/main/resources/prompts/wiki/route-user.txt b/mateclaw-server/src/main/resources/prompts/wiki/route-user.txt index b05cbcdf..57540968 100644 --- a/mateclaw-server/src/main/resources/prompts/wiki/route-user.txt +++ b/mateclaw-server/src/main/resources/prompts/wiki/route-user.txt @@ -2,6 +2,8 @@ {config} +{document_map_section} + ## 已有 Wiki 页面索引(slug + 摘要) {existing_pages}