diff --git a/resources/wcdb/win32/x64/wcdb_api.dll b/resources/wcdb/win32/x64/wcdb_api.dll index 530f6ed..f1f9815 100644 Binary files a/resources/wcdb/win32/x64/wcdb_api.dll and b/resources/wcdb/win32/x64/wcdb_api.dll differ diff --git a/src/main/index.ts b/src/main/index.ts index 8be05fb..9464d55 100644 --- a/src/main/index.ts +++ b/src/main/index.ts @@ -148,7 +148,10 @@ import type { AppLogEntry } from '../shared/app-log' import { appUpdateService } from './services/app-update-service' import { clearCache, getCacheSummary, openKnowledgeDirectory } from './services/cache-service' import { imageTextIndexService } from './services/image-text-index-service' -import type { ImageTextIndexStartOptions } from '../shared/image-text-index' +import { + IMAGE_TEXT_SEGMENT_MESSAGE_LIMIT, + type ImageTextIndexStartOptions +} from '../shared/image-text-index' import type { CacheClearScope } from './services/cache-service' import { configureRecallArchive, RecallArchiveMonitor } from './services/recall-archive-service' import { VideoAssetService } from './video-asset-service' @@ -703,12 +706,20 @@ app.whenReady().then(async () => { * 全量读取一个 20 万条消息的会话实测要 15s 以上,而其中 99% 以上的行 * 图片索引根本不看 —— 那是数据边界错了,不是 OCR 慢。 */ - listImageMessages: (conversationId) => - chat.listImageMessagesAsync(conversationId, undefined, 'image-text-index'), - countConversationImages: (conversationId, sinceMs) => - chat.countImageMessagesAsync(conversationId, sinceMs), - imageWatermark: (conversationId, sinceMs) => - chat.imageConversationWatermarkAsync(conversationId, sinceMs), + listImageMessages: (conversationId, window) => + chat.listImageMessagesAsync( + conversationId, + { + ...(window?.sinceMs !== undefined ? { sinceMs: window.sinceMs } : {}), + ...(window?.beforeMs !== undefined ? { beforeMs: window.beforeMs } : {}), + limit: IMAGE_TEXT_SEGMENT_MESSAGE_LIMIT + }, + 'image-text-index' + ), + countConversationImages: (conversationId, range) => + chat.countImageMessagesAsync(conversationId, range), + imageWatermark: (conversationId, range) => + chat.imageConversationWatermarkAsync(conversationId, range), decryptService: () => ensureImageDecryptService(), capability: () => systemOcrService.getCapability(), /** diff --git a/src/main/services/chat-service.ts b/src/main/services/chat-service.ts index 427d4ec..5fa81a5 100644 --- a/src/main/services/chat-service.ts +++ b/src/main/services/chat-service.ts @@ -17,6 +17,7 @@ import { import { mergeRecallArchiveMessages, recordRecallArchiveMessages } from './recall-archive-service' import type { ExportImageQuality } from '../../shared/image-quality' import type { ImageMessageCountProbe } from '../../shared/image-text-index' +import { imageTextWindowToSeconds } from '../../shared/image-text-index' import { wcdbDebugLog } from '../wcdb-debug' import { buildContactSearchIndex, @@ -852,23 +853,45 @@ export async function listMessagesAsync( */ export async function listImageMessagesAsync( userMd5: string, + window: { + /** 闭下界(epoch ms)。 */ + sinceMs?: number + /** 开上界(epoch ms)。 */ + beforeMs?: number + limit?: number + } = {}, requestId = '', caller: ListMessagesCaller = 'unknown' ): Promise { if (!dbRef) return [] const perf = emptyPerf(caller, requestId || nextListMessagesRequestId()) const totalStartedAt = Date.now() + // ms 半开区间 → 秒闭区间。换算只有共享契约里那一处实现。 + const { sinceSec, beforeSecInclusive } = imageTextWindowToSeconds(window) + const startTime = sinceSec ?? undefined + const endTime = beforeSecInclusive ?? undefined try { const rawReadStartedAt = Date.now() - const rawMessages = await dbRef - .getWcdb4Client() - .listImageMessagesAsync(userMd5, { requestId: perf.requestId }) + const rawMessages = await dbRef.getWcdb4Client().listImageMessagesAsync(userMd5, { + ...(window.sinceMs !== undefined ? { sinceMs: window.sinceMs } : {}), + ...(window.beforeMs !== undefined ? { beforeMs: window.beforeMs } : {}), + ...(window.limit !== undefined ? { limit: window.limit } : {}), + // recent-first:同一时间窗内**新的图片先处理**。 + order: 'desc', + requestId: perf.requestId + }) perf.rawReadMs += Date.now() - rawReadStartedAt + /** + * 时间边界必须同时交给格式化与召回归档合并。 + * + * 少了这一步,归档合并会把**窗口之外**的撤回图片补回来 —— 于是"最近 7 天" + * 这一段会混进十年前的消息,分段窗口形同虚设。 + */ const sourceMessages = listSourceMessages( userMd5, - undefined, - undefined, - undefined, + startTime, + endTime, + window.limit !== undefined ? { limit: window.limit } : undefined, rawMessages, perf.requestId, perf @@ -880,9 +903,9 @@ export async function listImageMessagesAsync( const result = mergeRecallArchiveMessages( userMd5, sourceMessages, - undefined, - undefined, - undefined + startTime, + endTime, + window.limit ) perf.sortMs += Date.now() - recallStartedAt return result @@ -933,10 +956,10 @@ export async function listMessagesForExport( */ export async function countImageMessagesAsync( userMd5: string, - sinceMs?: number + range?: number | { sinceMs?: number; beforeMs?: number } ): Promise { if (!dbRef) return { count: null, typeColumn: null, error: '微信数据库尚未就绪' } - return dbRef.getWcdb4Client().countImageMessagesAsync(userMd5, sinceMs) + return dbRef.getWcdb4Client().countImageMessagesAsync(userMd5, range) } /** @@ -947,10 +970,10 @@ export async function countImageMessagesAsync( */ export async function imageConversationWatermarkAsync( userMd5: string, - sinceMs?: number + range?: number | { sinceMs?: number; beforeMs?: number } ): Promise<{ count: number; maxLocalId: number } | null> { if (!dbRef) return null - return dbRef.getWcdb4Client().imageConversationWatermarkAsync(userMd5, sinceMs) + return dbRef.getWcdb4Client().imageConversationWatermarkAsync(userMd5, range) } export async function countVoiceMessagesAsync( diff --git a/src/main/services/image-text-index-service.ts b/src/main/services/image-text-index-service.ts index 3a4c696..9e2a8d4 100644 --- a/src/main/services/image-text-index-service.ts +++ b/src/main/services/image-text-index-service.ts @@ -17,12 +17,14 @@ import { createHash } from 'node:crypto' import { existsSync } from 'node:fs' import { DEFAULT_IMAGE_TEXT_OCR_CONCURRENCY, + IMAGE_TEXT_BACKFILL_TIER_ORDER, IMAGE_TEXT_INDEX_BATCH_SIZE, IMAGE_TEXT_INDEX_PROGRESS_INTERVAL_MS, IMAGE_TEXT_INDEX_RATE_MIN_SPAN_MS, IMAGE_TEXT_INDEX_RATE_WINDOW_MS, IMAGE_OCR_RETRIABLE_FAILURE_STATES, buildImageOcrArtifactKey, + buildImageTextBackfillSegments, imageTextProcessedPercent, isTerminalImageOcrState, resolveImageTextOcrConcurrency, @@ -30,8 +32,10 @@ import { type ImageMessageWatermark, type ImageOcrPersistedState, type ImageOcrProvenance, + type ImageTextBackfillSegment, type ImageTextIndexCountResult, type ImageTextIndexCoverage, + type ImageTextIndexPhase, type ImageTextIndexProgress, type ImageTextIndexRepairResult, type ImageTextIndexRunState, @@ -39,7 +43,8 @@ import { type ImageTextIndexStageTimings, type ImageTextIndexStartOptions, type ImageTextIndexStatus, - type ImageTextIndexStorageStats + type ImageTextIndexStorageStats, + type ImageTextTierCoverage } from '../../shared/image-text-index' import { detectSystemOcrImageFormat, @@ -86,7 +91,16 @@ export interface ImageTextIndexServiceDeps { * 由 WCDB 在 SQL 层过滤,而不是把整个会话读进来再筛。 * 缺省时回退到 `listMessages`(测试用),但生产必须接上 —— 否则大会话会拖垮一遍 pass。 */ - listImageMessages?: (conversationId: string) => Promise + listImageMessages?: ( + conversationId: string, + /** + * 时间窗(`[sinceMs, beforeMs)`,半开)。不传 = 整个会话。 + * + * recent-first 的分段计划靠它把"最近 7 天"和"更早"分开,而不是把整个会话 + * 读进来再在 JS 里筛 —— 那正是大会话跑不动的根因。 + */ + window?: { sinceMs?: number; beforeMs?: number } + ) => Promise /** * 单个会话的图片消息计数(SQL 统计,不解密)。 * @@ -94,7 +108,7 @@ export interface ImageTextIndexServiceDeps { */ countConversationImages?: ( conversationId: string, - sinceMs?: number + range?: number | { sinceMs?: number; beforeMs?: number } ) => Promise /** * 单个会话的图片消息增量水位(条数 + 最大插入序),SQL 聚合,不解密。 @@ -103,7 +117,7 @@ export interface ImageTextIndexServiceDeps { */ imageWatermark?: ( conversationId: string, - sinceMs?: number + range?: number | { sinceMs?: number; beforeMs?: number } ) => Promise decryptService?: () => ImageDecryptService | null /** 本地 OCR。 */ @@ -289,6 +303,14 @@ export class ImageTextIndexService { private counting = false private listeners = new Set<(status: ImageTextIndexStatus) => void>() private lastError: string | undefined + /** + * 当前阶段(recent-first 可见性)。 + * + * 它**不参与进度计算**:总进度永远是 `processed / totalImageMessages`。 + */ + private currentPhase: ImageTextIndexPhase = 'complete' + /** 本遍 pass 的起点,用于 `preLoop.startupMs`。 */ + private passStartedAt = 0 private startedAt: number | undefined /** 上一次清理实际重建(失效)了多少个会话的 Knowledge 索引;用于诊断与测试。 */ lastInvalidatedConversations = 0 @@ -602,6 +624,20 @@ export class ImageTextIndexService { * 它必须阻断 `complete` —— 否则 Query Agent 会拿着"覆盖完整"去回答"没有"。 */ const systemicFailure = processed > 0 && indexed === 0 && empty === 0 && missing === 0 + /** + * 覆盖完整性。 + * + * 分母取自流水线**真实走过**的集合时(`useScan`),不再要求 `counted.complete`: + * 那个标志表达的是"`countImageMessages()` 把每个会话都数上了",而进度现在已经 + * 不用那个分母了。继续要求它,会让一个**已经跑完**的索引因为"某个会话数不上" + * 而永远停在"部分完成"。 + */ + const complete = + total > 0 && + runtimeUnavailable === 0 && + !systemicFailure && + processed >= total && + (useScan || (counted !== null && counted.complete)) return { totalImageMessages: total, processed, @@ -613,22 +649,11 @@ export class ImageTextIndexService { pending: Math.max(0, total - processed - runtimeUnavailable), // 从未统计过总数 → 不算"已建立":不知道分母就不允许声称覆盖。 established: counted !== null && (processed > 0 || runtimeUnavailable > 0), - /** - * 覆盖完整性。 - * - * 分母取自流水线**真实走过**的集合时(`useScan`),不再要求 `counted.complete`: - * 那个标志表达的是"`countImageMessages()` 把每个会话都数上了",而进度现在已经 - * 不用那个分母了。继续要求它,会让一个**已经跑完**的索引因为"某个会话数不上" - * 而永远停在"部分完成"。 - */ - complete: - total > 0 && - runtimeUnavailable === 0 && - !systemicFailure && - processed >= total && - (useScan || (counted !== null && counted.complete)), + complete, systemicFailure, - countedAt: counted?.countedAt ?? null + countedAt: counted?.countedAt ?? null, + // recent-first 的时间维度:让调用方能回答"这一段时间能不能下确定性结论"。 + ...this.tierCoverageSnapshot(complete) } } @@ -672,7 +697,14 @@ export class ImageTextIndexService { cancellable: this.running, paused: this.runState === 'paused', ...this.rateSnapshot(coverage.processed, total), - ...(this.lastError ? { lastError: this.lastError } : {}) + ...(this.lastError ? { lastError: this.lastError } : {}), + // 阶段提示只在"真的在做这件事"时给:运行/暂停时报当前分段;整体完成时报完成。 + // 其余情况(idle / cancelled 且未完成)不给 —— 宁可不说,也不给一句过期的阶段。 + ...(this.running || this.runState === 'paused' + ? { currentPhase: this.currentPhase } + : coverage.complete + ? { currentPhase: 'complete' as ImageTextIndexPhase } + : {}) } } @@ -713,7 +745,9 @@ export class ImageTextIndexService { established: false, complete: false, systemicFailure: false, - countedAt: null + countedAt: null, + tiers: [], + coveredToMs: null }, storage: this.emptyStorage(), counting: this.counting @@ -1265,6 +1299,22 @@ export class ImageTextIndexService { return { started: true, state: this.runState } } + /** + * 一次索引 pass。 + * + * recent-first 的核心:**外层是时间分段,内层才是会话**。 + * + * 为什么不能只把会话列表按"最近活跃"排序就完事:那只保证"先把 A 群全部历史扫完", + * 而用户要的是"最近这几天的图片,不管在哪个群,都先能搜"。所以必须分段优先 —— + * 先把最近 7 天在所有会话上横着扫完,再退到下一个更老的分段。 + * + * 四条不可动摇的性质: + * 1. **锚点固定**:backfill 的 `anchorMs` 一旦落盘就不再变,分段边界因此稳定, + * 不会"跑几小时后 7 天窗口往前挪"。 + * 2. **新消息永远优先**:比锚点更新的图片由"增量补齐"负责,且它在每个调度点之前跑。 + * 3. **已完成的分段不重扫**;已 terminal 的 binding 永远跳过(不重复 OCR)。 + * 4. **进度不骗人**:总进度始终是 `processed / total`,分段只提供阶段文案。 + */ private async runPass(options: ImageTextIndexStartOptions): Promise { const store = this.ensureStore() if (!store) { @@ -1276,6 +1326,8 @@ export class ImageTextIndexService { // 进入图片流水线之前的一次性成本:单列出来,避免被摊进"每张图片"。 const preLoopStartedAt = this.now() + // `startupMs` 的参照点。逐会话逻辑被抽成独立方法之后,它必须放在实例上。 + this.passStartedAt = preLoopStartedAt const capability = (await this.deps.capability?.()) ?? null if (capability && !capability.available) { this.lastError = '当前系统不支持本地图片文字识别' @@ -1323,16 +1375,261 @@ export class ImageTextIndexService { * `startupMs` 的取样点必须是"第一张图片进入流水线的那一刻",不能在这里就记 —— * 否则会话级的准备成本会漏在外面,而那正是"单张很快、整遍很慢"的差额来源之一。 */ - let preLoopCaptured = false + const preLoopState = { captured: false } const scanState = store.readScanState() - let budget = options.messageLimit && options.messageLimit > 0 ? options.messageLimit : Infinity + const budget = { + remaining: options.messageLimit && options.messageLimit > 0 ? options.messageLimit : Infinity + } + /** + * 受控窗口(用于小样本验证):只跑这一个窗口,**不写任何分段状态**, + * 因此同一个窗口可以反复跑(checkpoint 是围绕全量集合建立的,混用会让"跳过" + * 变得不可解释)。 + */ + const windowed = Boolean(options.sinceMs && options.sinceMs > 0) + + const plan = this.resolveBackfillPlan(store, options) + + /** + * 升级场景:老版本可能已经把**全量**图片索引建完了。此时库里没有任何分段信息, + * 应当直接落成"全部完成",不重新回填 —— 用户已经拥有的东西不能被降级。 + * + * 但**不能只看库里自称的进度**:`readScanProgress()` 记的是"上一次跑完时留下了多少", + * 它不知道源里后来又新增了图片。只信它就会把"库里自称已完成、源侧其实有新增" + * 误判成"全部完成",于是那些新增的图片永远不会被索引。 + * + * 所以必须先向**源侧**核实:逐会话比对插入序水位,只有确实没有新内容时才算数。 + * 这次核实本身就是增量补齐(水位没涨的会话一条 SQL 就跳过),不会白跑。 + */ + if (!windowed && plan.created) { + const existing = this.coverageFromCounts(store.countByState()) + if (existing.complete) { + const verified = await this.processWindow({ + window: null, + incremental: true, + contacts, + provenance, + store, + budget, + scanState, + preLoopState, + marksConversationDone: true + }) + if (!verified.interrupted && !verified.truncated && verified.imageCount === 0) { + for (const tier of IMAGE_TEXT_BACKFILL_TIER_ORDER) { + store.writeBackfillTierState(tier, 'complete') + } + store.writeBackfillCoveredToMs(this.now()) + this.currentPhase = 'complete' + this.runState = 'completed' + this.running = false + await this.emit() + return + } + } + } + const lastSegment = plan.segments[plan.segments.length - 1] + /** + * 增量补齐跑过没有(本次 pass)。 + * + * 它必须**独立于"本分段要不要扫"**:所有分段都已完成时,仍然需要一次补齐, + * 否则"全部建完之后新到的图片"就再也没人接。 + */ + let swept = false + + for (const segment of plan.segments) { + if (this.cancelRequested || this.pauseRequested || budget.remaining <= 0) break + + const tierState = windowed ? null : store.readBackfillState().tierStates[segment.tier] + // 已完成的分段不再重扫 —— restart / resume 因此从**当前**分段继续, + // 而不是回到最近 7 天把已经做过的事再做一遍。 + const shouldScan = windowed || tierState !== 'complete' + + /** + * 增量补齐:**每个调度点先跑一次**,且本 pass 至少跑一次。 + * + * 排在历史分段之前,是为了"绝不会因为正在扫十年前的历史,让今天新收到的图片排队"; + * 即使分段全部完成、本 pass 没有任何分段要扫,也仍然要跑一次。 + */ + if (!windowed && (!swept || shouldScan)) { + const interrupted = await this.sweepIncremental({ + plan, + contacts, + provenance, + store, + budget, + scanState, + preLoopState + }) + swept = true + if (interrupted) break + } + + if (!shouldScan) continue + + this.currentPhase = windowed ? 'incremental' : segment.tier + if (!windowed) store.writeBackfillTierState(segment.tier, 'running') + await this.emit() + + const scanned = await this.processWindow({ + window: { sinceMs: segment.startMs, beforeMs: segment.endMs }, + incremental: false, + contacts, + provenance, + store, + budget, + scanState, + preLoopState, + marksConversationDone: !windowed && segment === lastSegment + }) + + // 被取消 / 暂停 / 预算截断 → 这一段**没有**跑完,绝不能标成 complete。 + if (scanned.interrupted || scanned.truncated) break + + if (!windowed) store.writeBackfillTierState(segment.tier, 'complete') + await this.emit() + } + + if (this.cancelRequested) this.runState = 'cancelled' + else if (this.pauseRequested) this.runState = 'paused' + else this.runState = 'completed' + // 只有真正跑完全部分段才把阶段切成"完成";中途停下时保留当前阶段(那才是实话)。 + if (this.runState === 'completed' && !windowed) this.currentPhase = 'complete' + this.running = false + await this.emit() + } + + /** + * 解析本次 pass 的分段计划。 + * + * 计划一律由**落盘的锚点**派生:锚点不随 pass 变化,所以"跑了几小时之后 + * 7 天窗口往前漂移、进而产生重复或遗漏"在结构上就不可能发生。 + */ + private resolveBackfillPlan( + store: ImageTextIndexStore, + options: ImageTextIndexStartOptions + ): { anchorMs: number; segments: ImageTextBackfillSegment[]; created: boolean } { + if (options.sinceMs && options.sinceMs > 0) { + // 受控窗口:单段、无上界、不落任何分段状态。 + return { + anchorMs: options.sinceMs, + segments: [ + { tier: 'recent_7d', startMs: options.sinceMs, endMs: Number.POSITIVE_INFINITY } + ], + created: false + } + } + const state = store.readBackfillState() + if (state.anchorMs !== null) { + return { + anchorMs: state.anchorMs, + segments: buildImageTextBackfillSegments(state.anchorMs), + created: false + } + } + const anchorMs = this.now() + store.writeBackfillAnchor(anchorMs) + for (const tier of IMAGE_TEXT_BACKFILL_TIER_ORDER) { + store.writeBackfillTierState(tier, 'pending') + } + return { anchorMs, segments: buildImageTextBackfillSegments(anchorMs), created: true } + } + + /** + * 增量补齐:把"锚点之后新到的东西"处理掉,永远排在历史分段之前。 + * + * 窗口 = `[max(锚点, 上次补齐水位), +∞)`,只覆盖新到的东西,所以刚建计划时 + * 它在时间上是空的、连枚举都省掉。水位可用时另有一层判据:插入序没涨的会话 + * 直接跳过,涨了的会话则**去掉时间窗**读(见 `processWindow`)。 + * + * 返回 true = 被取消 / 暂停 / 预算截断。 + */ + private async sweepIncremental(input: { + plan: { anchorMs: number; segments: ImageTextBackfillSegment[] } + contacts: Array<{ md5: string; m_nsUsrName: string; type: 'user' | 'group' }> + provenance: ImageOcrProvenance + store: ImageTextIndexStore + budget: { remaining: number } + scanState: Map< + string, + { state: string; imageTotal: number; processed: number; maxLocalId: number } + > + preLoopState: { captured: boolean } + }): Promise { + const incrementalSinceMs = Math.max( + input.plan.anchorMs, + input.store.readBackfillState().coveredToMs ?? 0 + ) + // 窗口在时间上必然为空 → 直接跳过,省掉一轮枚举。 + if (this.now() - incrementalSinceMs < 1_000) return false + + const swept = await this.processWindow({ + window: { sinceMs: incrementalSinceMs }, + incremental: true, + contacts: input.contacts, + provenance: input.provenance, + store: input.store, + budget: input.budget, + scanState: input.scanState, + preLoopState: input.preLoopState, + marksConversationDone: false + }) + // 只有真的扫完(没被取消 / 暂停 / 预算截断)才推进水位,否则会漏掉没扫到的部分。 + if (!swept.interrupted && !swept.truncated && input.budget.remaining > 0) { + input.store.writeBackfillCoveredToMs(this.now()) + } + return swept.interrupted || swept.truncated + } + + /** + * 处理一个窗口(`null` = 不看时间,用于受控小样本验证)。 + * + * 每个会话先做两条**纯 SQL 聚合**:源侧水位 + 窗口内图片条数。 + * 条数为 0 就整段跳过 —— 不读消息、不解密、不写任何"完成"标记。 + */ + private async processWindow(input: { + window: { sinceMs?: number; beforeMs?: number } | null + /** + * 增量补齐模式。 + * + * - 水位**可用**且水位涨了 → 去掉时间窗读这个会话(接得住晚到的旧时间消息); + * - 水位可用且没涨 → 跳过; + * - 水位**不可用** → 按时间窗兜底重扫:宁可慢,也不允许因为判据拿不到就漏。 + */ + incremental: boolean + contacts: Array<{ md5: string; m_nsUsrName: string; type: 'user' | 'group' }> + provenance: ImageOcrProvenance + store: ImageTextIndexStore + budget: { remaining: number } + scanState: Map< + string, + { state: string; imageTotal: number; processed: number; maxLocalId: number } + > + preLoopState: { captured: boolean } + /** 本窗口扫完是否意味着该会话**全部**图片都已定态(最后一个分段)。 */ + marksConversationDone: boolean + }): Promise<{ interrupted: boolean; imageCount: number; truncated: boolean }> { + const { incremental, contacts, provenance, store, budget, scanState } = input + let imageCount = 0 + /** + * 预算被截断。 + * + * 必须与"跑完了"区分开:被截断的分段**不能**被标记成 complete —— + * 否则下一次 pass 会以"这一段已完成"跳过,被截掉的那些图片就永远不会被处理。 + */ + let truncated = false for (const contact of contacts) { - if (this.cancelRequested || this.pauseRequested) break - if (budget <= 0) break + if (this.cancelRequested || this.pauseRequested) { + return { interrupted: true, imageCount, truncated } + } + if (budget.remaining <= 0) { + truncated = true + break + } const conversationId = contact.md5 + const previous = scanState.get(conversationId) /** * 会话级准备:水位 / 计数 / 让路 / 读消息。 * @@ -1341,45 +1638,66 @@ export class ImageTextIndexService { * 而 `perImageMs` 只覆盖 batch 循环,看不到它。 */ const setupStartedAt = this.now() + /** - * 是否只处理一个时间窗口(用于小样本验证)。 + * 源侧水位:`count` + `max(local_id)`,一条 SQL 聚合(**整会话**,不带窗口)。 * - * 带窗口时**不做增量跳过**:checkpoint 是围绕全量集合建立的, - * 窗口内的图片可能从未被处理过,继续按"该会话已完成"跳过会让窗口形同虚设。 + * 增量判据用 `maxLocalId` 而不是 create_time:`local_id` 是 WCDB 行内单调的 + * 插入序,因此"撤回一张旧图 + 新增一张新图"这种总数不变的变更也能被发现, + * 而 create_time 会被"晚到的旧时间消息"骗过。 */ - const windowed = Boolean(options.sinceMs && options.sinceMs > 0) - // 增量水位 = 条数 + 最大插入序。只比条数会漏掉「撤回一张旧图 + - // 新增一张新图」这种总数不变、集合却变了的会话。 - const watermark = await (this.deps.imageWatermark?.(conversationId, options.sinceMs) ?? - Promise.resolve(null)) - const imageTotal = - watermark?.count ?? - (await this.deps.countConversationImages?.(conversationId, options.sinceMs))?.count ?? - 0 - if (imageTotal === 0) { + const watermark = await (this.deps.imageWatermark?.(conversationId) ?? + Promise.resolve(null)) + + /** + * 这一段对这个会话实际要读的时间窗。 + * + * `null` = 不看时间(整会话,新 → 旧)。 + */ + let effectiveWindow = input.window + if (incremental && watermark !== null && previous !== undefined) { + // 水位可用:没涨就代表确实没有新内容,跳过(不读消息、不查窗口)。 + if (previous.maxLocalId > 0 && watermark.maxLocalId <= previous.maxLocalId) { + this.preLoop.conversationSetupMs += this.now() - setupStartedAt + continue + } + /** + * 判据可用且涨了 → **去掉时间窗**读这个会话。 + * + * 为什么不能只读"锚点之后":`local_id` 是插入序,而 `create_time` 是业务时间, + * 两者可以不一致 —— 网络补发、消息恢复、合并转发回填都会让一条**旧时间**的 + * 消息在今天才落库。只按时间窗读就永远接不到它,用户会"搜不到明明收到过的图"。 + * 代价只落在真的发生了插入的会话上,而每一轮之后水位即被推平。 + */ + effectiveWindow = null + } + const windowArg = effectiveWindow === null ? undefined : effectiveWindow + + /** + * 窗口内是否有图片:一条 SQL COUNT。 + * + * `count === null` 是**统计失败,不是 0 张** —— 必须跳过并且不写任何完成标记, + * 否则这段会被当成"已覆盖",把数不出来谎报成没有图片。 + */ + const probe = await (this.deps.countConversationImages?.( + conversationId, + windowArg + ) ?? Promise.resolve({ count: null, typeColumn: null })) + if (probe.count === null) { this.preLoop.conversationSetupMs += this.now() - setupStartedAt - store.writeScanState({ - conversationId, - state: 'done', - imageTotal: 0, - imageProcessed: 0, - maxLocalId: watermark?.maxLocalId ?? 0 - }) continue } - - // 增量:会话已完成且**水位完全未变** → 不读 WCDB、不 OCR。 - // 水位不可用时(数据库不支持该聚合)一律重扫:宁可慢,不可漏。 - const previous = scanState.get(conversationId) - if ( - !windowed && - watermark && - previous && - previous.state === 'done' && - previous.imageTotal === watermark.count && - previous.maxLocalId === watermark.maxLocalId - ) { + const imageTotal = probe.count + if (imageTotal === 0) { this.preLoop.conversationSetupMs += this.now() - setupStartedAt + this.rememberConversationWatermark({ + store, + scanState, + conversationId, + watermark, + processed: previous?.processed ?? 0, + marksConversationDone: input.marksConversationDone + }) continue } @@ -1401,34 +1719,47 @@ export class ImageTextIndexService { try { const source = await this.runStep('list-image-messages', () => this.deps.listImageMessages - ? this.deps.listImageMessages(conversationId) + ? this.deps.listImageMessages(conversationId, windowArg) : (this.deps.listMessages?.(conversationId) ?? Promise.resolve([])) ) imageMessages = source // 专用路径仍要过滤:召回归档合并可能补进非图片的撤回消息。 .filter(isImageMessage) - // 时间窗过滤:小样本验证时只看窗口内的图片,不然还是在跑全量。 - .filter((message) => - windowed ? (message.createTime || 0) * 1000 >= (options.sinceMs as number) : true - ) + // 时间窗过滤:兼容路径没有 SQL 层窗口,只能在这里筛;这也让小样本验证 + // 的语义与专用路径一致。 + .filter((message) => { + if (effectiveWindow === null) return true + const createTimeMs = (message.createTime || 0) * 1000 + if (effectiveWindow.sinceMs !== undefined && createTimeMs < effectiveWindow.sinceMs) { + return false + } + if (effectiveWindow.beforeMs !== undefined && createTimeMs >= effectiveWindow.beforeMs) { + return false + } + return true + }) } catch { imageMessages = [] } this.preLoop.listMessagesMs += this.now() - listStartedAt this.preLoop.conversationSetupMs += this.now() - setupStartedAt if (!imageMessages.length) { - store.writeScanState({ + this.rememberConversationWatermark({ + store, + scanState, conversationId, - state: 'done', - imageTotal: 0, - imageProcessed: 0, - maxLocalId: 0 + watermark, + processed: previous?.processed ?? 0, + marksConversationDone: input.marksConversationDone }) continue } - // 水位取**实际读到的**消息里最大的 local_id,而不是源侧水位: - // 万一在我们查水位之后、读消息之前又落了一条新图,用观测值会让下一轮 - // 发现"源水位更高"从而重扫(安全);用源侧水位则会把它永久跳过(漏索引)。 + imageCount += imageMessages.length + /** + * 水位取**实际读到的**消息里最大的 local_id,而不是源侧水位: + * 万一在我们查水位之后、读消息之前又落了一条新图,用观测值会让下一轮 + * 发现"源水位更高"从而重扫(安全);用源侧水位则会把它永久跳过(漏索引)。 + */ const observedMaxLocalId = imageMessages.reduce( (max, message) => Math.max(max, Number(message.localId) || 0), 0 @@ -1441,9 +1772,9 @@ export class ImageTextIndexService { this.knowledgeDirty = false for (let index = 0; index < imageMessages.length; index += IMAGE_TEXT_INDEX_BATCH_SIZE) { - if (!preLoopCaptured) { - preLoopCaptured = true - this.preLoop.startupMs = this.now() - preLoopStartedAt + if (!input.preLoopState.captured) { + input.preLoopState.captured = true + this.preLoop.startupMs = this.now() - this.passStartedAt } if (this.cancelRequested || this.pauseRequested) { interrupted = true @@ -1456,9 +1787,9 @@ export class ImageTextIndexService { conversationId, provenance, ocrByMessage, - () => budget > 0, + () => budget.remaining > 0, () => { - budget -= 1 + budget.remaining -= 1 } ) processedInConversation += batchResult.processed @@ -1481,23 +1812,39 @@ export class ImageTextIndexService { * 写 `done` 会让下一遍按水位错误跳过这些图片(**永久漏索引**), * 或者让用户以为这个会话已经处理完。预算耗尽只能记 `partial`。 */ - if (interrupted || budget <= 0) { + if (interrupted || budget.remaining <= 0) { store.writeScanState({ conversationId, state: 'partial', - imageTotal: imageMessages.length, + /** + * **整会话**的图片总数,不是本窗口的条数。 + * + * `image_ocr_scan_state.image_total` 是进度的分母(`readScanProgress()` 求和); + * 把某个分段的窗口条数写进去,分母就会缩到"这一段处理了多少", + * 于是总进度会突然跳到接近 100% —— 那是这个功能最不能犯的谎。 + */ + imageTotal: watermark?.count ?? previous?.imageTotal ?? imageMessages.length, imageProcessed: processedInConversation, maxLocalId: observedMaxLocalId }) - break + scanState.set(conversationId, { + state: 'partial', + imageTotal: imageMessages.length, + processed: previous?.processed ?? 0, + maxLocalId: observedMaxLocalId + }) + truncated = true + return { interrupted: true, imageCount, truncated } } - store.writeScanState({ + this.rememberConversationWatermark({ + store, + scanState, conversationId, - state: 'done', - imageTotal: imageMessages.length, - imageProcessed: processedInConversation, - maxLocalId: observedMaxLocalId + watermark, + processed: (previous?.processed ?? 0) + processedInConversation, + marksConversationDone: input.marksConversationDone, + observedMaxLocalId }) /** @@ -1530,13 +1877,77 @@ export class ImageTextIndexService { await this.emit() } - if (this.cancelRequested) this.runState = 'cancelled' - else if (this.pauseRequested) this.runState = 'paused' - else this.runState = 'completed' - this.running = false - await this.emit() + return { interrupted: false, imageCount, truncated } } + /** + * 写入会话 checkpoint。 + * + * 存的是**源侧水位**(整会话的 count / maxLocalId),不是窗口内的观测值: + * 增量补齐的判据必须是"源有没有变",用窗口观测值会让每次窗口扫描都改水位, + * 于是增量判断永远为真 —— 表现就是"每轮都重扫一遍"。 + */ + private rememberConversationWatermark(input: { + store: ImageTextIndexStore + scanState: Map< + string, + { state: string; imageTotal: number; processed: number; maxLocalId: number } + > + conversationId: string + watermark: ImageMessageWatermark | null + processed: number + marksConversationDone: boolean + observedMaxLocalId?: number + }): void { + const { store, scanState, conversationId, watermark, processed, marksConversationDone } = input + const previous = scanState.get(conversationId) + const state: 'done' | 'partial' = + marksConversationDone && previous?.state !== 'partial' ? 'done' : 'partial' + const imageTotal = watermark?.count ?? previous?.imageTotal ?? 0 + const maxLocalId = watermark?.maxLocalId ?? input.observedMaxLocalId ?? previous?.maxLocalId ?? 0 + store.writeScanState({ + conversationId, + state, + imageTotal, + imageProcessed: processed, + maxLocalId + }) + scanState.set(conversationId, { state, imageTotal, processed, maxLocalId }) + } + + /** + * 分段覆盖度(recent-first 的时间维度)。 + * + * 三种情形必须分清: + * - 从未规划过(新用户 / 旧版本)→ `tiers: []`,调用方只能按整体状态判断。 + * - 规划过 → 逐段给出真实运行态。 + * - **整体已经 complete** → 一律表达为"全部完成"。老版本已经把全量索引建完的账号, + * 升级后不能被重新拉回去做 backfill,也不能因为"没有分段信息"而让 Query Agent + * 以为历史还没扫完。 + */ + private tierCoverageSnapshot(complete: boolean): { + tiers: ImageTextTierCoverage[] + coveredToMs: number | null + } { + const store = this.store + if (!store) return { tiers: [], coveredToMs: null } + const state = store.readBackfillState() + const anchorMs = + state.anchorMs ?? (complete ? (store.readCountedTotal()?.countedAt ?? this.now()) : null) + if (anchorMs === null) return { tiers: [], coveredToMs: state.coveredToMs } + const tiers: ImageTextTierCoverage[] = buildImageTextBackfillSegments(anchorMs).map( + (segment) => ({ + tier: segment.tier, + state: complete ? 'complete' : (state.tierStates[segment.tier] ?? 'pending'), + startMs: segment.startMs, + endMs: segment.endMs + }) + ) + const coveredToMs = complete ? Math.max(state.coveredToMs ?? 0, anchorMs) : state.coveredToMs + return { tiers, coveredToMs } + } + + // --------------------------------------------------------------- 控制接口 pause(): { paused: boolean; state: ImageTextIndexRunState } { diff --git a/src/main/services/image-text-index-store.ts b/src/main/services/image-text-index-store.ts index 1dc2caf..446ad1d 100644 --- a/src/main/services/image-text-index-store.ts +++ b/src/main/services/image-text-index-store.ts @@ -13,11 +13,14 @@ import { mkdirSync, rmSync, statSync, existsSync } from 'node:fs' import { dirname, join, resolve } from 'node:path' import { DatabaseSync } from 'node:sqlite' import { + IMAGE_TEXT_BACKFILL_TIER_ORDER, IMAGE_TEXT_INDEX_SCHEMA_VERSION, type ImageOcrArtifact, type ImageOcrBinding, type ImageOcrPersistedState, - type ImageTextIndexStorageStats + type ImageTextBackfillTier, + type ImageTextIndexStorageStats, + type ImageTextTierRunState } from '../../shared/image-text-index' const MAX_SAFE_ACCOUNT_SEGMENT = /^[a-f0-9]{32}$/ @@ -374,6 +377,76 @@ export class ImageTextIndexStore { this.writeMeta('total_image_messages_complete', input.complete ? '1' : '0') } + // ---------------------------------------------------- recent-first 回填状态 + // + // 全部走 `image_ocr_meta`(key/value),**不加表、不加列** —— 这是 old checkpoint + // 兼容性的来源:旧库没有这些 key 时读出来就是"没有计划",于是下一次 pass + // 以当前时刻为锚点重新建计划。已经处理过的图片由 terminal binding 兜住, + // 不会因为"计划是新的"而重新 OCR。 + + /** + * 读取回填计划状态。 + * + * `anchorMs === null` = 从未规划过(新用户,或从旧版本升级且没有这些 key)。 + * 调用方此时必须**新建**计划,而不是假设"已完成"。 + */ + readBackfillState(): { + anchorMs: number | null + tierStates: Partial> + coveredToMs: number | null + } { + const anchorRaw = this.readMeta('backfill_anchor_ms') + const anchorMs = anchorRaw === null ? null : Number(anchorRaw) + let tierStates: Partial> = {} + const statesRaw = this.readMeta('backfill_tier_states') + if (statesRaw) { + try { + const parsed = JSON.parse(statesRaw) as Record + for (const tier of IMAGE_TEXT_BACKFILL_TIER_ORDER) { + const value = parsed[tier] + if (value === 'pending' || value === 'running' || value === 'complete') { + tierStates[tier] = value + } + } + } catch { + // 状态串损坏时按"没有分段状态"处理:最坏情况是重扫一遍, + // 而 terminal binding 保证不会重复 OCR。 + tierStates = {} + } + } + const coveredRaw = this.readMeta('backfill_covered_to_ms') + const coveredToMs = coveredRaw === null ? null : Number(coveredRaw) + return { + anchorMs: anchorMs !== null && Number.isFinite(anchorMs) ? anchorMs : null, + tierStates, + coveredToMs: coveredToMs !== null && Number.isFinite(coveredToMs) ? coveredToMs : null + } + } + + writeBackfillAnchor(anchorMs: number): void { + this.writeMeta('backfill_anchor_ms', String(Math.floor(anchorMs))) + } + + /** + * 写入单个分段的运行态。 + * + * 一个 pass 是单线程串行推进分段,所以"读-改-写"整体落在一条 meta 行里; + * 这里额外做的只有"保留其它分段"。 + */ + writeBackfillTierState(tier: ImageTextBackfillTier, state: ImageTextTierRunState): void { + const current = this.readBackfillState().tierStates + const next = { ...current, [tier]: state } + this.writeMeta('backfill_tier_states', JSON.stringify(next)) + } + + /** 增量补齐水位:处理到哪儿了(`create_time <= ms` 都已就绪)。只增不减。 */ + writeBackfillCoveredToMs(coveredToMs: number): void { + const previous = this.readBackfillState().coveredToMs + if (previous !== null && previous >= coveredToMs) return + this.writeMeta('backfill_covered_to_ms', String(Math.floor(coveredToMs))) + } + + readScanState(): Map< string, { state: string; imageTotal: number; processed: number; maxLocalId: number } @@ -563,7 +636,10 @@ export class ImageTextIndexStore { DELETE FROM image_ocr_meta WHERE key IN ( 'total_image_messages', 'total_image_counted_at', - 'total_image_messages_complete' + 'total_image_messages_complete', + 'backfill_anchor_ms', + 'backfill_tier_states', + 'backfill_covered_to_ms' ); `) } diff --git a/src/main/services/local-query-api-service.ts b/src/main/services/local-query-api-service.ts index 156e7a2..0033e7b 100644 --- a/src/main/services/local-query-api-service.ts +++ b/src/main/services/local-query-api-service.ts @@ -29,7 +29,9 @@ import { normalizeMessageIdentity } from '../../shared/local-query-api' import { + IMAGE_TEXT_BACKFILL_TIER_LABEL, describeImageTextCoverage, + describeImageTextCoveredRanges, imageTextCoverageState, type ImageTextIndexCoverage } from '../../shared/image-text-index' @@ -144,6 +146,17 @@ export function buildImageOcrCoverage( ? `(图片数量统计于 ${formatLocalMinute(coverage.countedAt)})` : '' const base = describeImageTextCoverage(coverage) + /** + * recent-first 之后必须把"哪段时间能下确定性结论"一起给出。 + * + * 否则模型只看一个总百分比:30% 时它会以为连最近一周都不可信(过度保守没坏处), + * 但 99% 时它会以为"去年也能放心下结论"(这就把索引缺口说成了事实空缺)。 + */ + const tiers = coverage.tiers ?? [] + const coveredRanges = tiers + .filter((entry) => entry.state === 'complete') + .map((entry) => IMAGE_TEXT_BACKFILL_TIER_LABEL[entry.tier]) + const rangeNote = tiers.length ? describeImageTextCoveredRanges(coverage) : '' return { state, totalImageMessages: coverage.totalImageMessages, @@ -154,10 +167,11 @@ export function buildImageOcrCoverage( failed: coverage.failed, pending: coverage.pending, ...(coverage.countedAt ? { countedAtLabel: formatLocalMinute(coverage.countedAt) } : {}), + ...(coveredRanges.length ? { coveredRanges } : {}), summary: state === 'complete' ? `${base}${countedNote}` - : `${base}${countedNote}${IMAGE_OCR_ZERO_RESULT_CAUTION}` + : `${base}${countedNote}${rangeNote}${IMAGE_OCR_ZERO_RESULT_CAUTION}` } } diff --git a/src/main/wcdb4-client.ts b/src/main/wcdb4-client.ts index ac7a364..d111151 100644 --- a/src/main/wcdb4-client.ts +++ b/src/main/wcdb4-client.ts @@ -7,6 +7,7 @@ import { createConnection, Socket } from 'net' import { getResourceRoots } from './resource-paths' import { wcdbDebugLog } from './wcdb-debug' import type { ImageMessageCountProbe } from '../shared/image-text-index' +import { imageTextWindowToSeconds } from '../shared/image-text-index' export interface Wcdb4Session { username: string @@ -1352,13 +1353,53 @@ export class Wcdb4Client { return value } - /** 图片消息的 WHERE 片段;`sinceMs` 用于只统计某个时间点之后的消息(测试小窗口)。 */ - private imageMessageWhere(column: string, sinceMs?: number): string { - const clauses = [`(${this.quoteSqlIdentifier(column)} & 65535) = 3`] - // 微信的 create_time 是**秒**,调用方给的是毫秒。 - if (sinceMs && Number.isFinite(sinceMs) && sinceMs > 0) { - clauses.push(`"create_time" >= ${Math.floor(sinceMs / 1000)}`) + /** + * 归一化图片消息的时间范围参数。 + * + * 同时接受旧的裸 `sinceMs` 与新的 `{ sinceMs, beforeMs }`: + * 前者有若干既有调用点(统计卡片、增量水位),不该为了新功能去改它们; + * 后者是 recent-first 分段计划要的半开区间。 + */ + private normalizeImageRange( + input?: number | { sinceMs?: number; beforeMs?: number } + ): { sinceMs?: number; beforeMs?: number } { + if (typeof input === 'number') { + return Number.isFinite(input) && input > 0 ? { sinceMs: input } : {} } + if (!input) return {} + const range: { sinceMs?: number; beforeMs?: number } = {} + if (typeof input.sinceMs === 'number' && Number.isFinite(input.sinceMs) && input.sinceMs > 0) { + range.sinceMs = input.sinceMs + } + if ( + typeof input.beforeMs === 'number' && + Number.isFinite(input.beforeMs) && + input.beforeMs > 0 + ) { + range.beforeMs = input.beforeMs + } + return range + } + + /** + * 图片消息的 WHERE 片段;`[sinceMs, beforeMs)` 半开区间。 + * + * 边界换算**只走** `imageTextWindowToSeconds`:COUNT 与列表两条路径必须用 + * 同一份换算,否则"统计说有 3 张、列表却返回 2 张",而调用方会据此把一段 + * 标成"已覆盖"。 + */ + private imageMessageWhere( + column: string, + input?: number | { sinceMs?: number; beforeMs?: number } + ): string { + const clauses = [`(${this.quoteSqlIdentifier(column)} & 65535) = 3`] + // 微信的 create_time 是**秒**,`[sinceMs, beforeMs)` 转成秒闭区间 + // `[sinceSec, beforeSecInclusive]`;两端同一套下取整,保证不重不漏。 + const { sinceSec, beforeSecInclusive } = imageTextWindowToSeconds( + this.normalizeImageRange(input) + ) + if (sinceSec !== null) clauses.push(`"create_time" >= ${sinceSec}`) + if (beforeSecInclusive !== null) clauses.push(`"create_time" <= ${beforeSecInclusive}`) return clauses.join(' AND ') } @@ -1373,7 +1414,7 @@ export class Wcdb4Client { */ async countImageMessagesAsync( md5OrUsername: string, - sinceMs?: number + input?: number | { sinceMs?: number; beforeMs?: number } ): Promise { if (!this.wcdbGetMessageTableStats || !this.wcdbExecQuery) { return { count: null, typeColumn: null, error: '当前数据服务不支持消息表统计' } @@ -1406,7 +1447,7 @@ export class Wcdb4Client { this.wcdbExecQuery as unknown as KoffiAsyncFunction, 'message', table.dbPath, - `SELECT COUNT(*) AS "image_count" FROM ${this.quoteSqlIdentifier(table.tableName)} WHERE ${this.imageMessageWhere(column, sinceMs)}` + `SELECT COUNT(*) AS "image_count" FROM ${this.quoteSqlIdentifier(table.tableName)} WHERE ${this.imageMessageWhere(column, input)}` ) const value = Number(this.pickValue(rows[0] || {}, ['image_count', 'count', 'COUNT(*)'])) if (Number.isFinite(value)) total += value @@ -1433,7 +1474,7 @@ export class Wcdb4Client { */ async imageConversationWatermarkAsync( md5OrUsername: string, - sinceMs?: number + input?: number | { sinceMs?: number; beforeMs?: number } ): Promise<{ count: number; maxLocalId: number } | null> { if (!this.wcdbGetMessageTableStats || !this.wcdbExecQuery) return null // 与计数同因:必须先把会话 md5 解析成原生接口要的 username,否则永远匹配不到消息表。 @@ -1458,7 +1499,7 @@ export class Wcdb4Client { this.wcdbExecQuery as unknown as KoffiAsyncFunction, 'message', table.dbPath, - `SELECT COUNT(*) AS "image_count", MAX("local_id") AS "image_max_local_id" FROM ${this.quoteSqlIdentifier(table.tableName)} WHERE ${this.imageMessageWhere(column, sinceMs)}` + `SELECT COUNT(*) AS "image_count", MAX("local_id") AS "image_max_local_id" FROM ${this.quoteSqlIdentifier(table.tableName)} WHERE ${this.imageMessageWhere(column, input)}` ) const row = rows[0] || {} const tableCount = Number(this.pickValue(row, ['image_count', 'count', 'COUNT(*)'])) @@ -1548,7 +1589,21 @@ export class Wcdb4Client { */ async listImageMessagesAsync( md5OrUsername: string, - options: { sinceMs?: number; limit?: number; requestId?: string } = {} + options: { + sinceMs?: number + /** + * 开区间上界。recent-first 的分段窗口靠它把"最近 30 天"和"更早"切开, + * 与 `sinceMs` 一起构成 `[sinceMs, beforeMs)`。 + */ + beforeMs?: number + limit?: number + /** + * 行序。默认 `asc` 保持既有行为不变;recent-first 的窗口用 `desc`, + * 让同一窗口内**新的图片先被处理**(用户先受益,且断点续跑更有意义)。 + */ + order?: 'asc' | 'desc' + requestId?: string + } = {} ): Promise { if (!this.wcdbExecQuery) return [] const requestId = options.requestId ?? 'NO-REQUEST' @@ -1573,10 +1628,11 @@ export class Wcdb4Client { const column = this.resolveMessageTypeColumn(table) if (!column) continue try { - const where = this.imageMessageWhere(column, options.sinceMs) + const where = this.imageMessageWhere(column, options) // `local_id` 参与排序:`create_time` 同秒的消息需要一个稳定次序, // 否则多次读取的行序可能不同,调用方无法做稳定游标。 - const sql = `SELECT * FROM ${this.quoteSqlIdentifier(table.tableName)} WHERE ${where} ORDER BY "create_time" ASC, "local_id" ASC LIMIT ${limit}` + const direction = options.order === 'desc' ? 'DESC' : 'ASC' + const sql = `SELECT * FROM ${this.quoteSqlIdentifier(table.tableName)} WHERE ${where} ORDER BY "create_time" ${direction}, "local_id" ${direction} LIMIT ${limit}` const queryStartedAt = Date.now() const rows = await this.callJsonAsync[]>( this.wcdbExecQuery as unknown as KoffiAsyncFunction, diff --git a/src/renderer/src/components/search/ImageTextIndexCard.tsx b/src/renderer/src/components/search/ImageTextIndexCard.tsx index 176d4e7..d2d2d3c 100644 --- a/src/renderer/src/components/search/ImageTextIndexCard.tsx +++ b/src/renderer/src/components/search/ImageTextIndexCard.tsx @@ -21,7 +21,10 @@ import { import { describeImageTextCoverage, imageTextCoverageState, - imageTextProcessedPercent + imageTextPhaseLabel, + imageTextProcessedPercent, + imageTextTierSearchableNotice, + type ImageTextBackfillTier } from '../../../../shared/image-text-index' import { useImageTextIndexStatus } from './hooks/useImageTextIndexStatus' @@ -89,6 +92,31 @@ export function ImageTextIndexCard({ dbReady, onNotice }: ImageTextIndexCardProp const percent = coverage ? imageTextProcessedPercent(coverage.processed, coverage.totalImageMessages) : 0 + /** + * 当前阶段文案(recent-first)。 + * + * 它**不参与进度计算**:进度永远是 `processed / totalImageMessages` + * (例如 12,800 / 349,838)。阶段只回答"现在在优先做什么", + * 绝不允许用"某个分段做完了"冒充整体完成。 + */ + const phaseLabel = progress?.currentPhase ? imageTextPhaseLabel(progress.currentPhase) : null + /** + * 已经**真正完整**的最新分段 —— 只有这时才敢告诉用户"这段时间已经能搜了"。 + * + * 判据是 `state === 'complete'`,不是"扫到过"。整体 complete 时不必重复宣告。 + */ + const searchableNotice = (() => { + if (!coverage || coverageState === 'complete') return null + const completed = new Set( + (coverage.tiers ?? []) + .filter((entry) => entry.state === 'complete') + .map((entry) => entry.tier) + ) + const newest = ( + ['recent_7d', 'recent_30d', 'recent_1y', 'archive'] as ImageTextBackfillTier[] + ).find((tier) => completed.has(tier)) + return newest ? imageTextTierSearchableNotice(newest) : null + })() const systemicFailure = coverage?.systemicFailure === true const visualState = progress?.state === 'error' || coverageState === 'failed' @@ -243,6 +271,15 @@ export function ImageTextIndexCard({ dbReady, onNotice }: ImageTextIndexCardProp }} /> + {/* + 阶段只是一句人话,没有百分比 —— 进度必须留在上面那行 + (processed / total 全量),否则用户会把"最近 7 天做完了"读成"整体做完了"。 + */} + {phaseLabel && ( +

+ {phaseLabel} +

+ )}

{`${progress.percent}% · 识别出文字 ${progress.indexed.toLocaleString()} · 没有文字 ${progress.empty.toLocaleString()} · 图片已清理 ${progress.missing.toLocaleString()} · 失败 ${progress.failed.toLocaleString()}`}

@@ -291,6 +328,18 @@ export function ImageTextIndexCard({ dbReady, onNotice }: ImageTextIndexCardProp )} + {/* + 只有**真正完整**的分段才敢说"已可搜索"。这是一句承诺,不是进度提示 —— + 所以判据是 coverage 里该分段的 state === 'complete',而不是"扫到过"。 + */} + {searchableNotice && ( +

+ {searchableNotice} +

+ )}

{describeImageTextCoverage(coverage)}

)} diff --git a/src/shared/image-text-index.ts b/src/shared/image-text-index.ts index 2738863..f1ab0c7 100644 --- a/src/shared/image-text-index.ts +++ b/src/shared/image-text-index.ts @@ -70,6 +70,308 @@ export function resolveImageTextOcrConcurrency(raw?: string | number | null): nu /** 已完成一批之后、回到会话循环前的让出时间。 */ export const IMAGE_TEXT_INDEX_YIELD_MS = 0 +/** + * 单段时间窗口的图片消息读取上限。 + * + * 必须显式给 limit:底层在不给 limit 时按 `create_time ASC` 排序并截断 —— + * 那是"从最老开始读",正好与 recent-first 相反,而且会把窗口内较新的图片悄悄丢掉。 + */ +export const IMAGE_TEXT_SEGMENT_MESSAGE_LIMIT = 200_000 + +// ---------------------------------------------------------------- recent-first +// +// 首次建立索引时,用户要的不是"从十年前开始扫",而是"最近聊天的图片先能搜"。 +// 这里把"最近优先"固化成**显式时间分段计划**,而不是靠反转全库排序碰运气。 +// +// 三条硬约束(任何实现都必须同时满足): +// 1. 分段唯一来源是本文件:`[startInclusive, endExclusive)`,不允许别的模块自己推边界。 +// 2. 一次 backfill session 的 `anchorMs` **固定不变**,否则跑几小时后窗口会漂移, +// 产生重复 / 遗漏 / checkpoint 不稳定。 +// 3. 比锚点更新的消息由**增量补齐**单独负责,永远排在所有历史分段之前。 + +/** 历史回填分段。刻意用时间语义而不是 Tier 编号(内部也不要出现 Tier 1/2)。 */ +export type ImageTextBackfillTier = 'recent_7d' | 'recent_30d' | 'recent_1y' | 'archive' + +/** 处理顺序 = 优先级顺序:越靠前越新。 */ +export const IMAGE_TEXT_BACKFILL_TIER_ORDER: readonly ImageTextBackfillTier[] = [ + 'recent_7d', + 'recent_30d', + 'recent_1y', + 'archive' +] + +/** + * 「最近」窗口的天数。 + * + * 取 7 天而不是 3 天:用户对"最近"的直觉通常是"这一周",而 7 天窗口在真实库里 + * 通常只有几百到几千张图,几分钟内就能把"最近聊天里的图"变成可搜索。 + */ +export const IMAGE_TEXT_RECENT_WINDOW_DAYS = 7 + +/** 各段下界的回看天数;`archive` 无下界。 */ +export const IMAGE_TEXT_BACKFILL_WINDOW_DAYS: Record = { + recent_7d: IMAGE_TEXT_RECENT_WINDOW_DAYS, + recent_30d: 30, + recent_1y: 365, + archive: Number.POSITIVE_INFINITY +} + +/** 各段的用户可见名称。 */ +export const IMAGE_TEXT_BACKFILL_TIER_LABEL: Record = { + recent_7d: '最近图片', + recent_30d: '最近 30 天', + recent_1y: '近一年图片', + archive: '更早图片' +} + +const DAY_MS = 24 * 60 * 60 * 1000 + +/** 一段回填窗口。`endMs` 为开区间上界;`archive` 的 `startMs` 是 -∞。 */ +export interface ImageTextBackfillSegment { + tier: ImageTextBackfillTier + /** epoch ms,**闭**下界。 */ + startMs: number + /** epoch ms,**开**上界。 */ + endMs: number +} + +/** + * 由固定锚点推出全部分段。 + * + * 结果按优先级从新到旧排列,且**互不重叠、无空隙**: + * `[now-7d, now)` → `[now-30d, now-7d)` → `[now-365d, now-30d)` → `(-∞, now-365d)`。 + * + * 相邻段的边界由同一个锚点派生,所以不会 off-by-one:上一段的 `endMs` 恒等于下一段的 + * `startMs`,而区间语义是"下界闭、上界开",落在边界上的消息只会被**一段**认领。 + */ +export function buildImageTextBackfillSegments(anchorMs: number): ImageTextBackfillSegment[] { + const segments: ImageTextBackfillSegment[] = [] + let upper = anchorMs + for (const tier of IMAGE_TEXT_BACKFILL_TIER_ORDER) { + const days = IMAGE_TEXT_BACKFILL_WINDOW_DAYS[tier] + const lower = Number.isFinite(days) ? anchorMs - days * DAY_MS : Number.NEGATIVE_INFINITY + segments.push({ tier, startMs: lower, endMs: upper }) + upper = lower + } + return segments +} + +/** + * 把 ms 半开区间转换成底层查询要的**秒**闭区间。 + * + * **唯一的换算实现**:`sinceMs` 下取整做闭下界,`beforeMs` 下取整减一放开区间上界。 + * 两端必须用同一种取整方式 —— 一边 ceil 一边 floor 的话,边界那一秒会被 + * 相邻两段同时认领(重复 OCR)或被同时跳过(漏索引)。 + * + * 微信的 `create_time` 只有秒级精度,所以 1 秒以内的边界歧义无法在源头消除; + * 能做的是让它**确定且不重叠**。 + */ +export function imageTextWindowToSeconds(range: { sinceMs?: number; beforeMs?: number }): { + sinceSec: number | null + beforeSecInclusive: number | null +} { + const { sinceMs, beforeMs } = range + const sinceSec = + typeof sinceMs === 'number' && Number.isFinite(sinceMs) && sinceMs > 0 + ? Math.floor(sinceMs / 1000) + : null + const beforeSecInclusive = + typeof beforeMs === 'number' && Number.isFinite(beforeMs) && beforeMs > 0 + ? Math.floor(beforeMs / 1000) - 1 + : null + return { sinceSec, beforeSecInclusive } +} + +/** 分段 → 秒闭区间;语义同 `imageTextWindowToSeconds`。 */ +export function imageTextSegmentToSecondRange(segment: ImageTextBackfillSegment): { + sinceSec: number | null + beforeSecInclusive: number | null +} { + return imageTextWindowToSeconds({ + ...(Number.isFinite(segment.startMs) ? { sinceMs: segment.startMs } : {}), + ...(Number.isFinite(segment.endMs) ? { beforeMs: segment.endMs } : {}) + }) +} + +// ---------------------------------------------------------------- 进度阶段 + +/** 索引运行阶段:增量补齐 / 某个历史分段 / 全部完成。 */ +export type ImageTextIndexPhase = 'incremental' | ImageTextBackfillTier | 'complete' + +export function imageTextPhaseLabel(phase: ImageTextIndexPhase): string { + switch (phase) { + case 'incremental': + case 'recent_7d': + return '正在优先索引最近图片' + case 'recent_30d': + return '正在补齐最近 30 天' + case 'recent_1y': + return '正在补齐近一年图片' + case 'archive': + return '正在补齐更早图片' + case 'complete': + return '图片文字索引已完成' + } +} + +/** + * 「某段已可搜索」的宣告文案。 + * + * 只允许在**该段真的 complete** 时使用 —— 它是一句承诺,不是进度提示。 + */ +export function imageTextTierSearchableNotice(tier: ImageTextBackfillTier): string { + switch (tier) { + case 'recent_7d': + return '最近图片已可搜索' + case 'recent_30d': + return '最近 30 天图片已可搜索' + case 'recent_1y': + return '近一年图片已可搜索' + case 'archive': + return '更早图片已可搜索' + } +} + +/** 分段的运行态。只有 `complete` 才允许被当作"这段已经完整可搜"。 */ +export type ImageTextTierRunState = 'pending' | 'running' | 'complete' + +/** 单个历史分段的完成状态与边界。 */ +export interface ImageTextTierCoverage { + tier: ImageTextBackfillTier + state: ImageTextTierRunState + /** epoch ms,闭下界;`archive` 为 -∞。 */ + startMs: number + /** epoch ms,开上界。 */ + endMs: number +} + +/** 时间范围询问的结论。 */ +export interface ImageTextRangeCoverage { + state: 'not_built' | 'partial' | 'complete' + /** 请求范围内**已真正完整**的子区间;无交集时为 null。 */ + coveredFromMs: number | null + coveredToMs: number | null + /** 请求范围是否有一部分落在"尚未覆盖"的区域。 */ + hasUncovered: boolean +} + +/** + * 给定查询时间范围,回答"这一段的图片文字覆盖是否完整"。 + * + * 判据刻意严格:只有当**整个请求区间**都落在已完成的覆盖并集里才回 `complete`。 + * 请求范围开放(不传 sinceMs / beforeMs)时视为"全部历史",只要有任何一段没完成 + * 就是 `partial` —— 这正是零结果诚实性要的:0 条证据不能回答"没有"。 + */ +export function imageTextRangeCoverage( + coverage: ImageTextIndexCoverage, + range: { sinceMs?: number; beforeMs?: number } = {} +): ImageTextRangeCoverage { + if (!coverage.established) { + return { state: 'not_built', coveredFromMs: null, coveredToMs: null, hasUncovered: true } + } + // 整体 complete = 全部图片消息都已定态 → 任意时间范围都完整,开放区间也一样。 + // 少了这一条,"全历史"这种没有上界的查询会永远因为"上界之后还没覆盖"被判成 partial。 + if (coverage.complete) { + return { state: 'complete', coveredFromMs: null, coveredToMs: null, hasUncovered: false } + } + // 已完成分段 + 增量水位合并成"已覆盖并集"。 + // + // `?? []` 不是多余的防御:覆盖度快照可能来自**旧版本落盘的库**(那时还没有分段概念), + // 也可能来自手写的夹具。缺字段只应该让"时间维度"暂时不可用,绝不能让查询直接崩。 + const tiers = coverage.tiers ?? [] + const covered: { startMs: number; endMs: number }[] = tiers + .filter((entry) => entry.state === 'complete') + .map((entry) => ({ startMs: entry.startMs, endMs: entry.endMs })) + if (coverage.coveredToMs !== null && coverage.coveredToMs !== undefined && tiers.length) { + const anchorMs = tiers[0]?.endMs ?? coverage.coveredToMs + covered.push({ startMs: anchorMs, endMs: coverage.coveredToMs }) + } + const merged = mergeImageTextRanges(covered) + + // 查询范围:不给边界 = 全历史(-∞, +∞)。 + const queryStart = range.sinceMs ?? Number.NEGATIVE_INFINITY + const queryEnd = range.beforeMs ?? Number.POSITIVE_INFINITY + if (!(queryEnd > queryStart)) { + return { state: 'complete', coveredFromMs: null, coveredToMs: null, hasUncovered: false } + } + + let cursor = queryStart + let coveredFromMs: number | null = null + let coveredToMs: number | null = null + for (const span of merged) { + if (span.endMs <= cursor) continue + if (span.startMs > cursor) break + if (coveredFromMs === null) coveredFromMs = Math.max(span.startMs, queryStart) + coveredToMs = Math.min(span.endMs, queryEnd) + cursor = Math.min(span.endMs, queryEnd) + if (cursor >= queryEnd) break + } + const hasUncovered = cursor < queryEnd + return { + state: hasUncovered ? 'partial' : 'complete', + coveredFromMs: coveredFromMs === null ? null : coveredFromMs, + coveredToMs, + hasUncovered + } +} + +function mergeImageTextRanges( + ranges: { startMs: number; endMs: number }[] +): { startMs: number; endMs: number }[] { + const sorted = [...ranges].sort((left, right) => left.startMs - right.startMs) + const merged: { startMs: number; endMs: number }[] = [] + for (const current of sorted) { + const last = merged[merged.length - 1] + if (last && current.startMs <= last.endMs) { + if (current.endMs > last.endMs) last.endMs = current.endMs + continue + } + merged.push({ ...current }) + } + return merged +} + +/** + * 「哪些时间段已经**真正**可以放心搜」的人话描述。 + * + * recent-first 之后,"索引建了多少"和"哪段时间能下确定结论"是两件事: + * 最近 7 天可能已经 100% 可用,而更早的历史还在补齐。只给一个总百分比, + * 模型会把"最近能搜"误当成"全历史都搜过了",于是对"去年有没有发过 XXX" + * 给出"没有"这种不该下的结论。 + */ +export function describeImageTextCoveredRanges(coverage: ImageTextIndexCoverage): string { + if (!coverage.established) return '还没有任何一个时间段完成图片文字索引。' + const completed = (coverage.tiers ?? []).filter((entry) => entry.state === 'complete') + if (coverage.complete) return '全部历史时间的图片文字都已可搜索。' + if (!completed.length) return '还没有任何一个时间段完成图片文字索引。' + const labels = completed.map((entry) => IMAGE_TEXT_BACKFILL_TIER_LABEL[entry.tier]) + // 只有归档段也完成才可能覆盖到最早;否则一定还有更老的历史没扫。 + const hasArchive = completed.some((entry) => entry.tier === 'archive') + return hasArchive + ? `已完整覆盖:${labels.join('、')}。` + : `已完整覆盖:${labels.join('、')};更早的图片仍在补齐。` +} + +/** + * 时间范围覆盖度的人话结论(Query Agent 只引用,不自己换算)。 + * + * 与 `describeImageTextCoverage` 的分工:那个回答"整体建了多少", + * 这个回答"**这次查的这个时间段**能不能下确定性结论"。 + */ +export function describeImageTextRangeCoverage( + coverage: ImageTextIndexCoverage, + range: { sinceMs?: number; beforeMs?: number } = {} +): string { + const result = imageTextRangeCoverage(coverage, range) + if (result.state === 'not_built') { + return '图片文字索引尚未建立:当前范围内图片里的文字还搜不到,不能据此回答"没有"。' + } + if (result.state === 'complete') { + return '图片文字索引已覆盖该时间范围:范围内没有匹配的图片文字,可以据此回答。' + } + return '图片文字历史仍在补齐,这次查询的时间范围尚未完整索引:当前结果不能排除尚未索引的图片。' +} + /** 单张图片的 OCR 结果状态。 */ export type ImageOcrState = /** 尚未处理 */ @@ -269,6 +571,13 @@ export interface ImageTextIndexProgress { speedPerSec?: number | null /** 按当前窗口速度估算的剩余时间(毫秒);速度不可用或分母不可信时为 null。 */ etaMs?: number | null + /** + * 当前正在处理的阶段(recent-first 的可见性)。 + * + * **它只是阶段提示,不是进度**:总进度仍然必须是 `processed / totalImageMessages` + * (全量图片消息),绝不允许用"某个分段处理完了"冒充整体完成。 + */ + currentPhase?: ImageTextIndexPhase } /** @@ -311,6 +620,19 @@ export interface ImageTextIndexCoverage { * 冒充成「现在完整」。 */ countedAt: number | null + /** + * 各历史分段的完成状态与边界(新 → 旧)。未建立 / 旧快照时为 `[]`。 + * + * 边界来自**当时固定的锚点**,因此可以和用户查询的时间范围直接求交。 + */ + tiers: ImageTextTierCoverage[] + /** + * 增量补齐水位:`create_time <= coveredToMs` 的新图片都已处理完;null = 尚无。 + * + * 它单独存在的原因:锚点之后新到的图片不属于任何历史分段,必须有独立水位 + * 才能回答"最近这几个小时是否已经可搜"。 + */ + coveredToMs: number | null } /** 覆盖度状态(外加"未建立")。UI 与 Query Agent 共用同一判据,避免两处各推一套口径漂移。 */ diff --git a/src/shared/local-query-api.ts b/src/shared/local-query-api.ts index 3ae1784..2be8bb7 100644 --- a/src/shared/local-query-api.ts +++ b/src/shared/local-query-api.ts @@ -369,6 +369,14 @@ export interface QueryImageTextCoverage { pending: number /** 图片数量统计时刻(本地时间 `MM-DD HH:mm`);从未统计时为 undefined。 */ countedAtLabel?: string + /** + * 已经**真正完整**的时间段描述(新 → 旧),例如"最近 7 天"。 + * + * recent-first 之后必须有这一维:总进度 30% 不代表"最近一周不可信", + * 反过来总进度 99% 也不代表"去年可以下确定性结论"。模型只能引用这里列出的 + * 时间段去下"没有"的结论,其余范围一律只能说"仍在补齐"。 + */ + coveredRanges?: string[] /** 可直接引用的结论句;模型只引用,不要自己换算或推断。 */ summary: string } diff --git a/tests/component/image-text-index-card.test.tsx b/tests/component/image-text-index-card.test.tsx index 03d6679..decc234 100644 --- a/tests/component/image-text-index-card.test.tsx +++ b/tests/component/image-text-index-card.test.tsx @@ -445,3 +445,105 @@ describe('图片文字索引卡片 — 修复图片搜索索引', () => { expect(String(onNotice.mock.calls.at(-1)?.[0])).toContain('正在进行中') }) }) + +/** + * recent-first 的阶段可见性。 + * + * 这里要守住的是两件**不能混**的事: + * - **总进度**永远是 `processed / totalImageMessages`(全量图片消息); + * - **阶段文案**只回答"现在在优先做什么"。 + * + * 用"某个分段做完了"去冒充整体完成,是这次改动最容易撒的谎,所以两条都断言。 + */ +describe('recent-first 阶段显示', () => { + /** 已建立、且不在运行/暂停 —— 卡片明细块的渲染条件。 */ + const settled = ( + coverageOverride: Partial + ): ImageTextIndexStatus => + status({ + ...running, + progress: { ...running.progress, state: 'idle', cancellable: false, paused: false }, + coverage: { ...running.coverage, ...coverageOverride } + }) + + const tiers = ( + completeTier: 'recent_7d' | 'recent_30d' | 'recent_1y' | 'archive' | null + ): ImageTextIndexStatus['coverage']['tiers'] => [ + { + tier: 'recent_7d', + state: completeTier === 'recent_7d' ? 'complete' : 'running', + startMs: 7, + endMs: 8 + }, + { + tier: 'recent_30d', + state: completeTier === 'recent_30d' ? 'complete' : 'pending', + startMs: 6, + endMs: 7 + }, + { + tier: 'recent_1y', + state: completeTier === 'recent_1y' ? 'complete' : 'pending', + startMs: 5, + endMs: 6 + }, + { + tier: 'archive', + state: completeTier === 'archive' ? 'complete' : 'pending', + startMs: 4, + endMs: 5 + } + ] + + it('运行中显示阶段文案,同时总进度仍以全量为分母', async () => { + api.getImageTextIndexStatus.mockResolvedValue( + status({ + ...running, + progress: { ...running.progress, currentPhase: 'recent_30d' } + }) + ) + + await renderCard() + + expect(screen.getByTestId('image-text-index-phase').textContent).toBe('正在补齐最近 30 天') + // 阶段 ≠ 进度:分子分母仍然是全量数字,不是"这一段处理了多少张"。 + expect(screen.getByTestId('image-text-index-progress').textContent).toBe( + `${running.progress.processed.toLocaleString()} / ${running.progress.totalImageMessages.toLocaleString()}` + ) + }) + + it('阶段文案是用户语言,不出现工程术语', async () => { + api.getImageTextIndexStatus.mockResolvedValue( + status({ ...running, progress: { ...running.progress, currentPhase: 'archive' } }) + ) + + await renderCard() + + const text = screen.getByTestId('image-text-index-phase').textContent ?? '' + expect(text).toBe('正在补齐更早图片') + expect(text).not.toMatch(/Tier/i) + }) + + it('分段真的完成时才宣告"已可搜索"', async () => { + api.getImageTextIndexStatus.mockResolvedValue( + settled({ tiers: tiers('recent_7d'), coveredToMs: 8 }) + ) + + await renderCard() + + expect(screen.getByTestId('image-text-index-searchable-notice').textContent).toBe( + '最近图片已可搜索' + ) + }) + + it('分段还没完成时不得宣告"已可搜索"', async () => { + api.getImageTextIndexStatus.mockResolvedValue( + settled({ tiers: tiers(null), coveredToMs: 8 }) + ) + + await renderCard() + + // 这是一句承诺,不是进度提示:没有真正 complete 就不许说。 + expect(screen.queryByTestId('image-text-index-searchable-notice')).toBeNull() + }) +}) diff --git a/tests/integration/image-text-index-checkpoint-and-coverage.test.ts b/tests/integration/image-text-index-checkpoint-and-coverage.test.ts index f8fd2d4..aca6c8b 100644 --- a/tests/integration/image-text-index-checkpoint-and-coverage.test.ts +++ b/tests/integration/image-text-index-checkpoint-and-coverage.test.ts @@ -19,6 +19,7 @@ import { } from '../../src/main/services/image-text-index-store' import { buildImageOcrCoverage } from '../../src/main/services/local-query-api-service' import type { ImageTextIndexCoverage } from '../../src/shared/image-text-index' +import { sourceMessageId } from '../../src/main/knowledge/message-identity' const ACCOUNT = 'wxid_fixture_account' const CONVERSATION = 'conversation-md5-fixture' @@ -93,72 +94,161 @@ function makeHarness(options: { messages?: chat.FormattedMessage[] } = {}): Harn } } -describe('增量水位:只比条数会漏掉「等量替换」', () => { - it('水位(条数 + 最大插入序)都没变时才跳过,不读 WCDB', async () => { - const harness = makeHarness({ messages: [imageMessage(10, 1000), imageMessage(20, 2000)] }) - harness.watermark.count = 2 - harness.watermark.maxLocalId = 20 - - await harness.service.startPass() - // 走到完成态需要等内部 promise 收敛。 - await vi.waitFor(() => expect(harness.service.isRunning()).toBe(false)) - expect(harness.listMessages).toHaveBeenCalledTimes(1) - - // 第二遍:水位完全一致 → 跳过,不再读会话消息。 - await harness.service.startPass() - await vi.waitFor(() => expect(harness.service.isRunning()).toBe(false)) - expect(harness.listMessages).toHaveBeenCalledTimes(1) - }) - - it('总数相同但最大插入序前进 → 必须重扫(撤回一张旧图 + 新增一张新图)', async () => { - const harness = makeHarness({ messages: [imageMessage(10, 1000), imageMessage(20, 2000)] }) - harness.watermark.count = 2 - harness.watermark.maxLocalId = 20 - - await harness.service.startPass() - await vi.waitFor(() => expect(harness.service.isRunning()).toBe(false)) - expect(harness.listMessages).toHaveBeenCalledTimes(1) - - // 集合变了、条数没变:localId 10 被撤回,新增 localId 30。 - harness.listMessages.mockImplementation(async () => [ - imageMessage(20, 2000), - imageMessage(30, 3000) - ]) - harness.watermark.maxLocalId = 30 - - await harness.service.startPass() - await vi.waitFor(() => expect(harness.service.isRunning()).toBe(false)) - // 只看 count 的实现会在这里静默跳过 —— 那正是会漏掉新图片的洞。 - expect(harness.listMessages).toHaveBeenCalledTimes(2) - }) - - it('水位不可用(数据库不支持该聚合)时一律重扫,宁可慢也不漏', async () => { - const harness = makeHarness({ messages: [imageMessage(10, 1000)] }) - harness.watermark.count = 1 +/** + * 增量判据。 + * + * recent-first 之后,"什么时候读 WCDB"由**分段完成状态 + 增量水位**共同决定, + * 但两条硬约束一个字都没变: + * 1. 没有新内容时**不得**重复读会话消息、更不得重复 OCR; + * 2. 有新内容(含"总数不变但集合变了")时**必须**处理到 —— 宁可慢,也不漏。 + * + * 这里的夹具刻意做成**窗口感知**的:新架构的"这段没有图片就整段跳过"完全依赖 + * 计数说实话;忽略窗口的夹具测出来的只是"夹具不过滤",不是调度器的行为。 + */ +describe('增量判据:水位不得漏掉新图片', () => { + function makeIncrementalHarness(options: { + messages: chat.FormattedMessage[] + withWatermark?: boolean + }): { + service: ImageTextIndexService + databasePath: string + state: { messages: chat.FormattedMessage[]; count: number; maxLocalId: number; now: number } + listMessages: ReturnType + listImageMessages: ReturnType + } { + const databaseRoot = makeDatabaseRoot() + const databasePath = getImageTextIndexDatabasePath(databaseRoot, ACCOUNT) + const state = { + messages: [...options.messages], + count: options.messages.length, + maxLocalId: options.messages.reduce((max, m) => Math.max(max, Number(m.localId) || 0), 0), + now: 1_800_000_000_000 + } + const inWindow = ( + message: chat.FormattedMessage, + window?: { sinceMs?: number; beforeMs?: number } + ): boolean => { + const createTimeMs = (message.createTime || 0) * 1000 + if (window?.sinceMs !== undefined && createTimeMs < window.sinceMs) return false + if (window?.beforeMs !== undefined && createTimeMs >= window.beforeMs) return false + return true + } + const listMessages = vi.fn(async () => state.messages) + const listImageMessages = vi.fn( + async (_conversationId: string, window?: { sinceMs?: number; beforeMs?: number }) => + state.messages.filter((message) => inWindow(message, window)) + ) const service = new ImageTextIndexService() - const listMessages = vi.fn(async () => [imageMessage(10, 1000)]) service.bind({ - databaseRoot: harness.databaseRoot, + databaseRoot, resolveAccountId: () => ACCOUNT, - listContacts: async () => [{ md5: CONVERSATION, m_nsUsrName: 'fixture', type: 'group' }], + resolveAccountRoot: () => 'C:/fixture/account', + now: () => state.now, + listContacts: async () => [ + { md5: CONVERSATION, m_nsUsrName: 'fixture', type: 'group' as const } + ], listMessages, - countConversationImages: async () => ({ count: 1, typeColumn: 'local_type' }), - // 关键:不提供 imageWatermark + listImageMessages, + countConversationImages: async ( + _conversationId: string, + window?: { sinceMs?: number; beforeMs?: number } + ) => ({ + count: state.messages.filter((message) => inWindow(message, window)).length, + typeColumn: 'local_type' + }), + ...(options.withWatermark === false + ? {} + : { imageWatermark: async () => ({ count: state.count, maxLocalId: state.maxLocalId }) }), + // 没有解密服务 → 每张图片都会被判成 image_missing。这样测试完全不碰真实图片。 decryptService: () => ({ findImageFile: () => null, decryptImage: () => null }) as never, capability: async () => ({ available: true, engine: 'windows-system-ocr', platform: 'win32', - runtimeVersion: null, - language: null - }) + runtimeVersion: '1.2.0', + language: 'zh-Hans-CN' + }), + recognize: async () => ({ success: true, text: '', language: 'zh-Hans-CN' }) + }) + return { service, databasePath, state, listMessages, listImageMessages } + } + + /** 直接读派生库的绑定,回答"到底处理了哪几条",而不是只看调用次数。 */ + const bindingIds = (databasePath: string): string[] => { + const store = new ImageTextIndexStore(databasePath, ACCOUNT) + const ids = [...store.getConversationOcr(CONVERSATION).keys()] + store.close() + return ids + } + + it('全部完成之后:第二遍不再读会话消息,也不重复 OCR', async () => { + const harness = makeIncrementalHarness({ + messages: [imageMessage(10, 1000), imageMessage(20, 2000)] }) - await service.startPass() - await vi.waitFor(() => expect(service.isRunning()).toBe(false)) - await service.startPass() - await vi.waitFor(() => expect(service.isRunning()).toBe(false)) - expect(listMessages).toHaveBeenCalledTimes(2) + await harness.service.startPass() + await vi.waitFor(() => expect(harness.service.isRunning()).toBe(false)) + // 两条图片的 create_time 都很老 → 只落在归档段,只被那一段读到。 + expect(harness.listImageMessages).toHaveBeenCalledTimes(1) + const firstIds = bindingIds(harness.databasePath) + expect(firstIds).toHaveLength(2) + + // 第二遍:水位完全一致 + 各分段已完成 → 一个字都不再读。 + await harness.service.startPass() + await vi.waitFor(() => expect(harness.service.isRunning()).toBe(false)) + expect(harness.listImageMessages).toHaveBeenCalledTimes(1) + expect(bindingIds(harness.databasePath)).toEqual(firstIds) + }) + + it('总数相同但最大插入序前进 → 必须处理到新图片(即使是旧时间)', async () => { + const harness = makeIncrementalHarness({ + messages: [imageMessage(10, 1000), imageMessage(20, 2000)] + }) + + await harness.service.startPass() + await vi.waitFor(() => expect(harness.service.isRunning()).toBe(false)) + expect(bindingIds(harness.databasePath)).toHaveLength(2) + + // 集合变了、条数没变:localId 10 被撤回,新增 localId 30 —— 而且它带着**旧时间**, + // 只按时间窗读的实现在这里就会漏掉它。 + harness.state.messages = [imageMessage(20, 2000), imageMessage(30, 2000)] + harness.state.count = 2 + harness.state.maxLocalId = 30 + harness.state.now += 5_000 + + await harness.service.startPass() + await vi.waitFor(() => expect(harness.service.isRunning()).toBe(false)) + + // 只看 count 的实现会在这里静默跳过 —— 那正是会漏掉新图片的洞。 + const ids = bindingIds(harness.databasePath) + expect(ids).toHaveLength(3) + expect(ids).toContain(sourceMessageId(imageMessage(30, 2000))) + }) + + it('水位不可用(数据库不支持该聚合)时按时间窗兜底重扫,宁可慢也不漏', async () => { + const harness = makeIncrementalHarness({ + messages: [imageMessage(10, 1000)], + withWatermark: false + }) + const anchorSeconds = Math.floor(harness.state.now / 1000) + + await harness.service.startPass() + await vi.waitFor(() => expect(harness.service.isRunning()).toBe(false)) + expect(bindingIds(harness.databasePath)).toHaveLength(1) + + // 锚点之后新到一张图;并且时间确实往前走了一段(否则"窗口必然为空"的短路会生效, + // 那不是漏,而是正确地判定"还没有新东西")。 + harness.state.messages = [imageMessage(10, 1000), imageMessage(40, anchorSeconds + 10)] + harness.state.count = 2 + harness.state.maxLocalId = 40 + harness.state.now += 20_000 + + await harness.service.startPass() + await vi.waitFor(() => expect(harness.service.isRunning()).toBe(false)) + + const ids = bindingIds(harness.databasePath) + expect(ids).toHaveLength(2) + expect(ids).toContain(sourceMessageId(imageMessage(40, anchorSeconds + 10))) }) }) diff --git a/tests/integration/image-text-index-image-query.test.ts b/tests/integration/image-text-index-image-query.test.ts index fe75500..8d4d948 100644 --- a/tests/integration/image-text-index-image-query.test.ts +++ b/tests/integration/image-text-index-image-query.test.ts @@ -86,8 +86,24 @@ function createHarness(): Harness { const all = mixedMessages() const imagesOnly = all.filter((message) => message.contentData?.type === 'image') + /** + * recent-first 之后,"读哪些行"由**时间分段**决定,所以夹具必须像真 WCDB 一样 + * 按窗口说话。忽略窗口的夹具测不出调度行为,只会让"读了几次"变成 4 次。 + */ + const inWindow = ( + message: chat.FormattedMessage, + window?: { sinceMs?: number; beforeMs?: number } + ): boolean => { + const createTimeMs = (message.createTime || 0) * 1000 + if (window?.sinceMs !== undefined && createTimeMs < window.sinceMs) return false + if (window?.beforeMs !== undefined && createTimeMs >= window.beforeMs) return false + return true + } const listMessages = vi.fn(async () => all) - const listImageMessages = vi.fn(async () => imagesOnly) + const listImageMessages = vi.fn( + async (_conversationId: string, window?: { sinceMs?: number; beforeMs?: number }) => + imagesOnly.filter((message) => inWindow(message, window)) + ) const service = new ImageTextIndexService() service.bind({ @@ -98,7 +114,13 @@ function createHarness(): Harness { listContacts: async () => [ { md5: CONVERSATION, m_nsUsrName: 'boundary', type: 'user' as const } ], - countConversationImages: async () => ({ count: IMAGE_COUNT, typeColumn: 'local_type' }), + countConversationImages: async ( + _conversationId: string, + window?: { sinceMs?: number; beforeMs?: number } + ) => ({ + count: imagesOnly.filter((message) => inWindow(message, window)).length, + typeColumn: 'local_type' + }), imageWatermark: async () => ({ count: IMAGE_COUNT, maxLocalId: TEXT_COUNT + IMAGE_COUNT }), listMessages, listImageMessages, diff --git a/tests/integration/image-text-index-progress-notify.test.ts b/tests/integration/image-text-index-progress-notify.test.ts index 008be9a..eb88757 100644 --- a/tests/integration/image-text-index-progress-notify.test.ts +++ b/tests/integration/image-text-index-progress-notify.test.ts @@ -117,8 +117,15 @@ function createService(options: { count: number; notifyIntervalMs: number; perIm return { service, notifications } } +/** + * 等 pass 收尾。 + * + * 显式给足超时:这些用例故意让 240 张图各睡几毫秒来制造可观测的持续时间, + * 而 `vi.waitFor` 的默认超时是 1000ms —— 机器稍慢就会以**断言失败**而不是 + * "超时"的形态报出来。这里等的是"跑完",不是"跑得快"。 + */ const finish = async (service: ImageTextIndexService): Promise => { - await vi.waitFor(() => expect(service.isRunning()).toBe(false)) + await vi.waitFor(() => expect(service.isRunning()).toBe(false), { timeout: 30_000 }) } describe('进度通知节流', () => { diff --git a/tests/integration/image-text-index-recent-first.test.ts b/tests/integration/image-text-index-recent-first.test.ts new file mode 100644 index 0000000..92f71a6 --- /dev/null +++ b/tests/integration/image-text-index-recent-first.test.ts @@ -0,0 +1,543 @@ +/** + * 「recent-first 图片文字索引」的行为测试。 + * + * 这一轮改的是**调度**:不再从最老的历史往下扫,而是先让"最近聊天的图片"可搜索。 + * 调度改动最容易骗人的地方是"看起来更快了,其实漏了东西",所以这里的断言都落在 + * **可观察的结果**上(读了哪些窗口、处理了哪些消息、覆盖度怎么说),而不是内部变量。 + * + * 被测的性质: + * - 分段计划:互不重叠、无空隙,边界唯一; + * - 处理顺序按时间分段从新到旧; + * - 断点续跑不重复 OCR; + * - 历史分段跑着的时候,新图片仍然优先; + * - 覆盖度能按时间范围回答"这段能不能下确定性结论"; + * - 老版本既有索引 / 旧 checkpoint 不被降级,源侧有新增时必须补上; + * - UI 文案不暴露工程术语。 + */ +import { mkdtempSync, rmSync } from 'node:fs' +import { tmpdir } from 'node:os' +import { join } from 'node:path' +import { afterEach, describe, expect, it, vi } from 'vitest' +import type * as chat from '../../src/main/services/chat-service' +import { ImageTextIndexService } from '../../src/main/services/image-text-index-service' +import { + ImageTextIndexStore, + getImageTextIndexDatabasePath +} from '../../src/main/services/image-text-index-store' +import { sourceMessageId } from '../../src/main/knowledge/message-identity' +import { + IMAGE_TEXT_BACKFILL_WINDOW_DAYS, + buildImageTextBackfillSegments, + describeImageTextRangeCoverage, + imageTextPhaseLabel, + imageTextRangeCoverage, + imageTextWindowToSeconds, + type ImageTextBackfillTier, + type ImageTextIndexCoverage, + type ImageTextTierRunState +} from '../../src/shared/image-text-index' + +const ACCOUNT = 'wxid_recent_first_fixture' +const CONVERSATION = 'conversation-fixture-md5' +const DAY_MS = 24 * 60 * 60 * 1000 +const ANCHOR_MS = 1_800_000_000_000 +const roots: string[] = [] + +afterEach(() => { + for (const root of roots.splice(0)) { + try { + rmSync(root, { recursive: true, force: true }) + } catch { + // 测试收尾尽力而为。 + } + } +}) + +function makeDatabaseRoot(): string { + const root = mkdtempSync(join(tmpdir(), 'tm-recent-first-')) + roots.push(root) + return root +} + +function imageMessage(localId: number, createTimeSeconds: number): chat.FormattedMessage { + return { + localId: String(localId), + createTime: createTimeSeconds, + content: '[图片]', + contentData: { type: 'image', md5: `md5-${localId}`, datName: `dat-${localId}` } + } as unknown as chat.FormattedMessage +} + +type Window = { sinceMs?: number; beforeMs?: number } + +interface HarnessState { + messages: chat.FormattedMessage[] + count: number + maxLocalId: number + now: number +} + +interface HarnessResult { + service: ImageTextIndexService + databasePath: string + state: HarnessState + /** 每次向底层请求图片消息时记录的**时间边界**(不含分页上限)。 */ + windows: Window[] + listImageMessages: ReturnType + /** 真正"过了一遍处理"的图片数 —— 用它证明断点续跑没有重做。 */ + findImageFile: ReturnType +} + +function inWindow(message: chat.FormattedMessage, window?: Window): boolean { + const createTimeMs = (message.createTime || 0) * 1000 + if (window?.sinceMs !== undefined && createTimeMs < window.sinceMs) return false + if (window?.beforeMs !== undefined && createTimeMs >= window.beforeMs) return false + return true +} + +function makeHarness(options: { + messages: chat.FormattedMessage[] + onWindow?: (window: Window | undefined, callIndex: number) => void +}): HarnessResult { + const databaseRoot = makeDatabaseRoot() + const databasePath = getImageTextIndexDatabasePath(databaseRoot, ACCOUNT) + const state = { + messages: [...options.messages], + count: options.messages.length, + maxLocalId: options.messages.reduce((max, m) => Math.max(max, Number(m.localId) || 0), 0), + now: ANCHOR_MS + } + const windows: Window[] = [] + const findImageFile = vi.fn(() => null) + + const listImageMessages = vi.fn(async (_conversationId: string, window?: Window) => { + // 只记时间边界:调用方还会带一个分页上限,那与"读了哪个时间窗"无关。 + const bounds: Window = {} + if (window?.sinceMs !== undefined) bounds.sinceMs = window.sinceMs + if (window?.beforeMs !== undefined) bounds.beforeMs = window.beforeMs + windows.push(bounds) + options.onWindow?.(window, windows.length - 1) + return state.messages.filter((message) => inWindow(message, window)) + }) + + const service = new ImageTextIndexService() + service.bind({ + databaseRoot, + resolveAccountId: () => ACCOUNT, + resolveAccountRoot: () => 'C:/fixture/account', + now: () => state.now, + listContacts: async () => [ + { md5: CONVERSATION, m_nsUsrName: 'fixture', type: 'group' as const } + ], + listImageMessages, + // 窗口感知的计数:新架构"这段没图片就整段跳过"完全依赖它说实话。 + countConversationImages: async (_conversationId: string, window?: Window) => ({ + count: state.messages.filter((message) => inWindow(message, window)).length, + typeColumn: 'local_type' + }), + imageWatermark: async () => ({ count: state.count, maxLocalId: state.maxLocalId }), + // 没有解密能力 → 每张图片都会被判成 image_missing,测试完全不碰真实图片。 + decryptService: () => ({ findImageFile, decryptImage: () => null }) as never, + capability: async () => ({ + available: true, + engine: 'windows-system-ocr', + platform: 'win32', + runtimeVersion: '1.2.0', + language: 'zh-Hans-CN' + }), + recognize: async () => ({ success: true, text: '', language: 'zh-Hans-CN' }) + }) + + return { service, databasePath, state, windows, listImageMessages, findImageFile } +} + +const bindingIds = (databasePath: string): string[] => { + const store = new ImageTextIndexStore(databasePath, ACCOUNT) + const ids = [...store.getConversationOcr(CONVERSATION).keys()] + store.close() + return ids +} + +const tierStates = ( + databasePath: string +): Partial> => { + const store = new ImageTextIndexStore(databasePath, ACCOUNT) + const result = store.readBackfillState() + store.close() + return result.tierStates +} + +const runPass = async ( + harness: HarnessResult, + options?: Parameters[0] +): Promise => { + harness.service.startPass(options) + await vi.waitFor(() => expect(harness.service.isRunning()).toBe(false), { timeout: 15_000 }) +} + +/** 合成数据集:今天 / 3 天前 / 15 天前 / 6 个月前 / 3 年前,各 10 张。 */ +function syntheticDataset(anchorMs: number): { + messages: chat.FormattedMessage[] + byBucket: Map +} { + const buckets: Array<{ key: string; offsetDays: number }> = [ + { key: 'today', offsetDays: 0 }, + { key: '3d', offsetDays: 3 }, + { key: '15d', offsetDays: 15 }, + { key: '6mo', offsetDays: 183 }, + { key: '3y', offsetDays: 1095 } + ] + const messages: chat.FormattedMessage[] = [] + const byBucket = new Map() + let localId = 1 + for (const bucket of buckets) { + const ids: number[] = [] + const baseSeconds = Math.floor((anchorMs - bucket.offsetDays * DAY_MS) / 1000) - 60 + for (let index = 0; index < 10; index += 1) { + messages.push(imageMessage(localId, baseSeconds + index)) + ids.push(localId) + localId += 1 + } + byBucket.set(bucket.key, ids) + } + return { messages, byBucket } +} + +describe('分段计划:唯一边界、不重不漏', () => { + it('四段从新到旧、相邻且互不重叠', () => { + const segments = buildImageTextBackfillSegments(ANCHOR_MS) + expect(segments.map((segment) => segment.tier)).toEqual([ + 'recent_7d', + 'recent_30d', + 'recent_1y', + 'archive' + ]) + expect(segments[0].endMs).toBe(ANCHOR_MS) + for (let index = 1; index < segments.length; index += 1) { + // 上一段的上界必须等于下一段的下界:半开区间才既不重叠也不留缝。 + expect(segments[index].endMs).toBe(segments[index - 1].startMs) + } + expect(segments[2].startMs).toBe( + ANCHOR_MS - IMAGE_TEXT_BACKFILL_WINDOW_DAYS.recent_1y * DAY_MS + ) + // 归档段没有下界。 + expect(segments[3].startMs).toBe(Number.NEGATIVE_INFINITY) + }) + + it('边界那一秒只属于一段(两端同一套取整)', () => { + const boundaryMs = ANCHOR_MS - 7 * DAY_MS + const { beforeSecInclusive } = imageTextWindowToSeconds({ beforeMs: boundaryMs }) + const { sinceSec } = imageTextWindowToSeconds({ sinceMs: boundaryMs }) + // 上界段的"含"到 beforeSecInclusive,下界段从 sinceSec 起 —— 必须正好衔接。 + expect(sinceSec! - beforeSecInclusive!).toBe(1) + }) +}) + +describe('recent-first:合成数据集的真实处理顺序', () => { + it('今天 / 3 天 → 15 天 → 6 个月前 → 3 年前,先新后旧', async () => { + const harness = makeHarness({ messages: [] }) + const { messages, byBucket } = syntheticDataset(ANCHOR_MS) + harness.state.messages = messages + harness.state.count = messages.length + harness.state.maxLocalId = Math.max(...messages.map((m) => Number(m.localId))) + + await runPass(harness) + + // 只有"有图片的分段"才会去读消息:今天+3 天 / 15 天 / 6 个月 / 3 年,共 4 次。 + expect(harness.windows).toHaveLength(4) + const [first, second, third, fourth] = harness.windows + expect(first).toEqual({ sinceMs: ANCHOR_MS - 7 * DAY_MS, beforeMs: ANCHOR_MS }) + expect(second).toEqual({ sinceMs: ANCHOR_MS - 30 * DAY_MS, beforeMs: ANCHOR_MS - 7 * DAY_MS }) + expect(third).toEqual({ sinceMs: ANCHOR_MS - 365 * DAY_MS, beforeMs: ANCHOR_MS - 30 * DAY_MS }) + expect(fourth).toEqual({ sinceMs: Number.NEGATIVE_INFINITY, beforeMs: ANCHOR_MS - 365 * DAY_MS }) + + // 第一次读就拿到了"最近的图"—— 这正是用户要的价值。 + const ids = bindingIds(harness.databasePath) + expect(ids).toHaveLength(50) + const recentIds = [...byBucket.get('today')!, ...byBucket.get('3d')!] + for (const id of recentIds) expect(ids).toContain(sourceMessageId(imageMessage(id, 0))) + expect(Math.max(...byBucket.get('today')!)).toBeLessThan(Math.min(...byBucket.get('3y')!)) + for (const id of byBucket.get('3y')!) { + expect(ids).toContain(sourceMessageId(imageMessage(id, 0))) + } + }) +}) + +describe('断点续跑:不重复 OCR', () => { + it('第一批被预算截断后,重启只补剩下的,已完成的不再走一遍', async () => { + const harness = makeHarness({ messages: [] }) + const { messages } = syntheticDataset(ANCHOR_MS) + harness.state.messages = messages + harness.state.count = messages.length + harness.state.maxLocalId = Math.max(...messages.map((m) => Number(m.localId))) + + // 第一轮:只允许处理 6 张 → 必然停在"最近图片"这一段中间。 + await runPass(harness, { messageLimit: 6 }) + expect(harness.findImageFile).toHaveBeenCalledTimes(6) + expect(harness.windows).toHaveLength(1) + let states = tierStates(harness.databasePath) + // 该分段没跑完 → 不允许被标成 complete(否则剩下的图永远不会被处理)。 + expect(states.recent_7d).not.toBe('complete') + + // 第二轮:预算放开 → 先把最近这段的剩余 4 张补完,再继续往下。 + await runPass(harness) + expect(harness.findImageFile).toHaveBeenCalledTimes(50) + states = tierStates(harness.databasePath) + expect(states.recent_7d).toBe('complete') + expect(states.archive).toBe('complete') + expect(bindingIds(harness.databasePath)).toHaveLength(50) + }) +}) + +describe('增量优先:历史分段跑着的时候新图片先处理', () => { + it('新图片在下一次调度点被处理,而不是排到历史之后', async () => { + let injected = false + const harness = makeHarness({ + messages: [], + onWindow: (window) => { + // 第二轮窗口请求(recent_30d)之后,模拟"用户刚收到一张新图"。 + if (injected || window?.sinceMs === undefined) return + if (window.beforeMs === ANCHOR_MS - 7 * DAY_MS) { + injected = true + harness.state.messages = [ + ...harness.state.messages, + imageMessage(999, Math.floor((ANCHOR_MS + 5_000) / 1000)) + ] + harness.state.count += 1 + harness.state.maxLocalId = 999 + harness.state.now = ANCHOR_MS + 5_000 + } + } + }) + const { messages } = syntheticDataset(ANCHOR_MS) + harness.state.messages = messages + harness.state.count = messages.length + harness.state.maxLocalId = Math.max(...messages.map((m) => Number(m.localId))) + + await runPass(harness) + + const newMessageId = sourceMessageId(imageMessage(999, 0)) + expect(bindingIds(harness.databasePath)).toContain(newMessageId) + /** + * 新图片必须在"更老的分段"之前被处理。 + * + * 增量补齐对"插入序涨了"的会话是**不设时间窗**读取的(这样才接得住晚到的旧时间消息), + * 所以这里找的是"没有任何时间边界的那个窗口",它必须早于 1 年 / 归档分段。 + */ + const sweepWindowIndex = harness.windows.findIndex( + (window) => window.sinceMs === undefined && window.beforeMs === undefined + ) + const yearWindowIndex = harness.windows.findIndex( + (window) => window?.beforeMs === ANCHOR_MS - 30 * DAY_MS + ) + expect(sweepWindowIndex).toBeGreaterThanOrEqual(0) + expect(yearWindowIndex).toBeGreaterThan(sweepWindowIndex) + }) +}) + +describe('兼容性:老索引不被降级,源侧有新增必须补上', () => { + it('老版本已经全量建完且源侧无变化 → 保持完成,一个字都不读', async () => { + const harness = makeHarness({ messages: [] }) + // 造一份"旧版本建完"的库:总数 3、三条 binding 全部 terminal、没有任何分段元数据。 + const store = new ImageTextIndexStore(harness.databasePath, ACCOUNT) + store.writeCountedTotal({ total: 3, countedAt: ANCHOR_MS - 1000, complete: true }) + // 旧版本每跑完一个会话都会留下 scan_state —— 升级后正是靠它证明源侧没动过。 + store.writeScanState({ + conversationId: CONVERSATION, + state: 'done', + imageTotal: 3, + imageProcessed: 3, + maxLocalId: 3 + }) + for (let index = 1; index <= 3; index += 1) { + store.putBinding({ + accountId: ACCOUNT, + conversationId: CONVERSATION, + messageId: `local:${index}`, + createTime: ANCHOR_MS - 1000, + imageIdentity: `sha256:${index}`, + artifactKey: `sha256:${index}|legacy`, + state: 'indexed', + updatedAt: 1 + }) + } + store.close() + + harness.state.messages = [imageMessage(1, 1000), imageMessage(2, 2000), imageMessage(3, 3000)] + harness.state.count = 3 + harness.state.maxLocalId = 3 + + await runPass(harness) + + // 源侧水位一致 → 一个字都不该读:用户已经拥有的东西不能被降级成"重新建立"。 + expect(harness.listImageMessages).not.toHaveBeenCalled() + expect(harness.findImageFile).not.toHaveBeenCalled() + const states = tierStates(harness.databasePath) + expect(states.recent_7d).toBe('complete') + expect(states.archive).toBe('complete') + }) + + it('老版本自称"已建完"但源侧已新增 → 必须补上新增,不能信库里的进度', async () => { + const harness = makeHarness({ messages: [] }) + const store = new ImageTextIndexStore(harness.databasePath, ACCOUNT) + // 库里自称 2/2 已完成(旧版本跑完时的状态)。 + store.writeCountedTotal({ total: 2, countedAt: ANCHOR_MS - 1000, complete: true }) + for (const localId of [1, 2]) { + store.writeScanState({ + conversationId: CONVERSATION, + state: 'done', + imageTotal: 2, + imageProcessed: 2, + maxLocalId: 2 + }) + store.putBinding({ + accountId: ACCOUNT, + conversationId: CONVERSATION, + messageId: `local:${localId}`, + createTime: ANCHOR_MS - 1000, + imageIdentity: `sha256:old-${localId}`, + artifactKey: `sha256:old-${localId}|legacy`, + state: 'indexed', + updatedAt: 1 + }) + } + store.close() + + // 源侧现在是 3 张(插入序 1..3)—— 只看库里自称的进度会把它误判成"已完成"。 + harness.state.messages = [imageMessage(1, 1000), imageMessage(2, 2000), imageMessage(3, 3000)] + harness.state.count = 3 + harness.state.maxLocalId = 3 + + await runPass(harness) + + const ids = bindingIds(harness.databasePath) + expect(ids).toHaveLength(3) + expect(ids).toContain(sourceMessageId(imageMessage(3, 3000))) + // 老的 terminal 结果必须被复用,不能重新 OCR。 + expect(harness.findImageFile).toHaveBeenCalledTimes(1) + }) +}) + +describe('时间范围覆盖度', () => { + const coverageOf = (input: Partial): ImageTextIndexCoverage => ({ + totalImageMessages: 100, + processed: 20, + indexed: 18, + empty: 2, + missing: 0, + failed: 0, + runtimeUnavailable: 0, + pending: 80, + established: true, + complete: false, + systemicFailure: false, + countedAt: ANCHOR_MS, + tiers: [ + { tier: 'recent_7d', state: 'complete', startMs: ANCHOR_MS - 7 * DAY_MS, endMs: ANCHOR_MS }, + { + tier: 'recent_30d', + state: 'pending', + startMs: ANCHOR_MS - 30 * DAY_MS, + endMs: ANCHOR_MS - 7 * DAY_MS + }, + { + tier: 'recent_1y', + state: 'pending', + startMs: ANCHOR_MS - 365 * DAY_MS, + endMs: ANCHOR_MS - 30 * DAY_MS + }, + { + tier: 'archive', + state: 'pending', + startMs: Number.NEGATIVE_INFINITY, + endMs: ANCHOR_MS - 365 * DAY_MS + } + ], + coveredToMs: ANCHOR_MS, + ...input + }) + + it('A. 最近 7 天已完成 → 查询最近 24h 可以下确定性结论', () => { + const coverage = coverageOf({}) + const range = { sinceMs: ANCHOR_MS - DAY_MS, beforeMs: ANCHOR_MS } + const result = imageTextRangeCoverage(coverage, range) + expect(result.state).toBe('complete') + expect(result.hasUncovered).toBe(false) + expect(describeImageTextRangeCoverage(coverage, range)).toContain('已覆盖该时间范围') + }) + + it('B. 更早历史未完成 → 全历史查询只能是 partial,并明确"不能排除尚未索引的图片"', () => { + const coverage = coverageOf({}) + expect(imageTextRangeCoverage(coverage, {}).state).toBe('partial') + expect(describeImageTextRangeCoverage(coverage, {})).toContain('不能排除尚未索引的图片') + }) + + it('C. 全部完成 → 任意范围都是 complete', () => { + const coverage = coverageOf({ + processed: 100, + indexed: 98, + empty: 2, + pending: 0, + complete: true, + tiers: coverageOf({}).tiers.map((entry) => ({ ...entry, state: 'complete' as const })) + }) + expect(imageTextRangeCoverage(coverage, {}).state).toBe('complete') + expect( + imageTextRangeCoverage(coverage, { sinceMs: ANCHOR_MS - 900 * DAY_MS }).state + ).toBe('complete') + }) + + it('D. 在更老的历史里被取消 → 最近一段仍可完整,全历史仍是 partial', () => { + const coverage = coverageOf({}) + expect( + imageTextRangeCoverage(coverage, { + sinceMs: ANCHOR_MS - 3 * DAY_MS, + beforeMs: ANCHOR_MS + }).state + ).toBe('complete') + expect(imageTextRangeCoverage(coverage, {}).state).toBe('partial') + }) + + it('未建立时不得声称任何范围完整', () => { + const coverage = coverageOf({ established: false, tiers: [], coveredToMs: null }) + expect(imageTextRangeCoverage(coverage, { sinceMs: ANCHOR_MS - DAY_MS }).state).toBe( + 'not_built' + ) + }) +}) + +describe('UI 文案与进度', () => { + it('阶段文案是用户语言,不出现工程术语', () => { + expect(imageTextPhaseLabel('recent_7d')).toBe('正在优先索引最近图片') + expect(imageTextPhaseLabel('incremental')).toBe('正在优先索引最近图片') + expect(imageTextPhaseLabel('recent_30d')).toBe('正在补齐最近 30 天') + expect(imageTextPhaseLabel('recent_1y')).toBe('正在补齐近一年图片') + expect(imageTextPhaseLabel('archive')).toBe('正在补齐更早图片') + expect(imageTextPhaseLabel('complete')).toBe('图片文字索引已完成') + const phases = ['incremental', 'recent_7d', 'recent_30d', 'recent_1y', 'archive', 'complete'] as const + for (const phase of phases) { + expect(imageTextPhaseLabel(phase)).not.toMatch(/Tier/i) + } + }) + + it('阶段字段在运行中出现,且总进度仍以全量图片数为分母', async () => { + const harness = makeHarness({ messages: [] }) + const { messages } = syntheticDataset(ANCHOR_MS) + harness.state.messages = messages + harness.state.count = messages.length + harness.state.maxLocalId = Math.max(...messages.map((m) => Number(m.localId))) + + await runPass(harness, { messageLimit: 3 }) + const status = await harness.service.getStatus() + // 总分母必须是全部图片消息,而不是"这一段处理了多少"。 + expect(status.progress.totalImageMessages).toBe(50) + expect(status.progress.processed).toBeGreaterThan(0) + expect(status.progress.percent).toBeLessThan(100) + expect(status.coverage.tiers.map((entry) => entry.tier)).toEqual([ + 'recent_7d', + 'recent_30d', + 'recent_1y', + 'archive' + ]) + }) +})