feat: 图片文字索引按时间分段优先处理最近图片

- 首次索引先处理最近 7 天,再依次回溯 30 天 / 近一年 / 更早历史
 - 完成后新到的图片单独补齐,不受历史回填影响
 - 覆盖度增加时间维度,可区分「最近已完整」与「更早仍在补齐」
This commit is contained in:
Wxw-Gu
2026-09-18 14:39:05 +08:00
parent 12b8fb34c6
commit 0f270366aa
15 changed files with 1917 additions and 183 deletions
Binary file not shown.
+18 -7
View File
@@ -148,7 +148,10 @@ import type { AppLogEntry } from '../shared/app-log'
import { appUpdateService } from './services/app-update-service' import { appUpdateService } from './services/app-update-service'
import { clearCache, getCacheSummary, openKnowledgeDirectory } from './services/cache-service' import { clearCache, getCacheSummary, openKnowledgeDirectory } from './services/cache-service'
import { imageTextIndexService } from './services/image-text-index-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 type { CacheClearScope } from './services/cache-service'
import { configureRecallArchive, RecallArchiveMonitor } from './services/recall-archive-service' import { configureRecallArchive, RecallArchiveMonitor } from './services/recall-archive-service'
import { VideoAssetService } from './video-asset-service' import { VideoAssetService } from './video-asset-service'
@@ -703,12 +706,20 @@ app.whenReady().then(async () => {
* 全量读取一个 20 万条消息的会话实测要 15s 以上,而其中 99% 以上的行 * 全量读取一个 20 万条消息的会话实测要 15s 以上,而其中 99% 以上的行
* 图片索引根本不看 —— 那是数据边界错了,不是 OCR 慢。 * 图片索引根本不看 —— 那是数据边界错了,不是 OCR 慢。
*/ */
listImageMessages: (conversationId) => listImageMessages: (conversationId, window) =>
chat.listImageMessagesAsync(conversationId, undefined, 'image-text-index'), chat.listImageMessagesAsync(
countConversationImages: (conversationId, sinceMs) => conversationId,
chat.countImageMessagesAsync(conversationId, sinceMs), {
imageWatermark: (conversationId, sinceMs) => ...(window?.sinceMs !== undefined ? { sinceMs: window.sinceMs } : {}),
chat.imageConversationWatermarkAsync(conversationId, 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(), decryptService: () => ensureImageDecryptService(),
capability: () => systemOcrService.getCapability(), capability: () => systemOcrService.getCapability(),
/** /**
+36 -13
View File
@@ -17,6 +17,7 @@ import {
import { mergeRecallArchiveMessages, recordRecallArchiveMessages } from './recall-archive-service' import { mergeRecallArchiveMessages, recordRecallArchiveMessages } from './recall-archive-service'
import type { ExportImageQuality } from '../../shared/image-quality' import type { ExportImageQuality } from '../../shared/image-quality'
import type { ImageMessageCountProbe } from '../../shared/image-text-index' import type { ImageMessageCountProbe } from '../../shared/image-text-index'
import { imageTextWindowToSeconds } from '../../shared/image-text-index'
import { wcdbDebugLog } from '../wcdb-debug' import { wcdbDebugLog } from '../wcdb-debug'
import { import {
buildContactSearchIndex, buildContactSearchIndex,
@@ -852,23 +853,45 @@ export async function listMessagesAsync(
*/ */
export async function listImageMessagesAsync( export async function listImageMessagesAsync(
userMd5: string, userMd5: string,
window: {
/** 闭下界(epoch ms)。 */
sinceMs?: number
/** 开上界(epoch ms)。 */
beforeMs?: number
limit?: number
} = {},
requestId = '', requestId = '',
caller: ListMessagesCaller = 'unknown' caller: ListMessagesCaller = 'unknown'
): Promise<FormattedMessage[]> { ): Promise<FormattedMessage[]> {
if (!dbRef) return [] if (!dbRef) return []
const perf = emptyPerf(caller, requestId || nextListMessagesRequestId()) const perf = emptyPerf(caller, requestId || nextListMessagesRequestId())
const totalStartedAt = Date.now() const totalStartedAt = Date.now()
// ms 半开区间 → 秒闭区间。换算只有共享契约里那一处实现。
const { sinceSec, beforeSecInclusive } = imageTextWindowToSeconds(window)
const startTime = sinceSec ?? undefined
const endTime = beforeSecInclusive ?? undefined
try { try {
const rawReadStartedAt = Date.now() const rawReadStartedAt = Date.now()
const rawMessages = await dbRef const rawMessages = await dbRef.getWcdb4Client().listImageMessagesAsync(userMd5, {
.getWcdb4Client() ...(window.sinceMs !== undefined ? { sinceMs: window.sinceMs } : {}),
.listImageMessagesAsync(userMd5, { requestId: perf.requestId }) ...(window.beforeMs !== undefined ? { beforeMs: window.beforeMs } : {}),
...(window.limit !== undefined ? { limit: window.limit } : {}),
// recent-first:同一时间窗内**新的图片先处理**。
order: 'desc',
requestId: perf.requestId
})
perf.rawReadMs += Date.now() - rawReadStartedAt perf.rawReadMs += Date.now() - rawReadStartedAt
/**
* 时间边界必须同时交给格式化与召回归档合并。
*
* 少了这一步,归档合并会把**窗口之外**的撤回图片补回来 —— 于是"最近 7 天"
* 这一段会混进十年前的消息,分段窗口形同虚设。
*/
const sourceMessages = listSourceMessages( const sourceMessages = listSourceMessages(
userMd5, userMd5,
undefined, startTime,
undefined, endTime,
undefined, window.limit !== undefined ? { limit: window.limit } : undefined,
rawMessages, rawMessages,
perf.requestId, perf.requestId,
perf perf
@@ -880,9 +903,9 @@ export async function listImageMessagesAsync(
const result = mergeRecallArchiveMessages( const result = mergeRecallArchiveMessages(
userMd5, userMd5,
sourceMessages, sourceMessages,
undefined, startTime,
undefined, endTime,
undefined window.limit
) )
perf.sortMs += Date.now() - recallStartedAt perf.sortMs += Date.now() - recallStartedAt
return result return result
@@ -933,10 +956,10 @@ export async function listMessagesForExport(
*/ */
export async function countImageMessagesAsync( export async function countImageMessagesAsync(
userMd5: string, userMd5: string,
sinceMs?: number range?: number | { sinceMs?: number; beforeMs?: number }
): Promise<ImageMessageCountProbe> { ): Promise<ImageMessageCountProbe> {
if (!dbRef) return { count: null, typeColumn: null, error: '微信数据库尚未就绪' } 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( export async function imageConversationWatermarkAsync(
userMd5: string, userMd5: string,
sinceMs?: number range?: number | { sinceMs?: number; beforeMs?: number }
): Promise<{ count: number; maxLocalId: number } | null> { ): Promise<{ count: number; maxLocalId: number } | null> {
if (!dbRef) return null if (!dbRef) return null
return dbRef.getWcdb4Client().imageConversationWatermarkAsync(userMd5, sinceMs) return dbRef.getWcdb4Client().imageConversationWatermarkAsync(userMd5, range)
} }
export async function countVoiceMessagesAsync( export async function countVoiceMessagesAsync(
+499 -88
View File
@@ -17,12 +17,14 @@ import { createHash } from 'node:crypto'
import { existsSync } from 'node:fs' import { existsSync } from 'node:fs'
import { import {
DEFAULT_IMAGE_TEXT_OCR_CONCURRENCY, DEFAULT_IMAGE_TEXT_OCR_CONCURRENCY,
IMAGE_TEXT_BACKFILL_TIER_ORDER,
IMAGE_TEXT_INDEX_BATCH_SIZE, IMAGE_TEXT_INDEX_BATCH_SIZE,
IMAGE_TEXT_INDEX_PROGRESS_INTERVAL_MS, IMAGE_TEXT_INDEX_PROGRESS_INTERVAL_MS,
IMAGE_TEXT_INDEX_RATE_MIN_SPAN_MS, IMAGE_TEXT_INDEX_RATE_MIN_SPAN_MS,
IMAGE_TEXT_INDEX_RATE_WINDOW_MS, IMAGE_TEXT_INDEX_RATE_WINDOW_MS,
IMAGE_OCR_RETRIABLE_FAILURE_STATES, IMAGE_OCR_RETRIABLE_FAILURE_STATES,
buildImageOcrArtifactKey, buildImageOcrArtifactKey,
buildImageTextBackfillSegments,
imageTextProcessedPercent, imageTextProcessedPercent,
isTerminalImageOcrState, isTerminalImageOcrState,
resolveImageTextOcrConcurrency, resolveImageTextOcrConcurrency,
@@ -30,8 +32,10 @@ import {
type ImageMessageWatermark, type ImageMessageWatermark,
type ImageOcrPersistedState, type ImageOcrPersistedState,
type ImageOcrProvenance, type ImageOcrProvenance,
type ImageTextBackfillSegment,
type ImageTextIndexCountResult, type ImageTextIndexCountResult,
type ImageTextIndexCoverage, type ImageTextIndexCoverage,
type ImageTextIndexPhase,
type ImageTextIndexProgress, type ImageTextIndexProgress,
type ImageTextIndexRepairResult, type ImageTextIndexRepairResult,
type ImageTextIndexRunState, type ImageTextIndexRunState,
@@ -39,7 +43,8 @@ import {
type ImageTextIndexStageTimings, type ImageTextIndexStageTimings,
type ImageTextIndexStartOptions, type ImageTextIndexStartOptions,
type ImageTextIndexStatus, type ImageTextIndexStatus,
type ImageTextIndexStorageStats type ImageTextIndexStorageStats,
type ImageTextTierCoverage
} from '../../shared/image-text-index' } from '../../shared/image-text-index'
import { import {
detectSystemOcrImageFormat, detectSystemOcrImageFormat,
@@ -86,7 +91,16 @@ export interface ImageTextIndexServiceDeps {
* 由 WCDB 在 SQL 层过滤,而不是把整个会话读进来再筛。 * 由 WCDB 在 SQL 层过滤,而不是把整个会话读进来再筛。
* 缺省时回退到 `listMessages`(测试用),但生产必须接上 —— 否则大会话会拖垮一遍 pass。 * 缺省时回退到 `listMessages`(测试用),但生产必须接上 —— 否则大会话会拖垮一遍 pass。
*/ */
listImageMessages?: (conversationId: string) => Promise<chat.FormattedMessage[]> listImageMessages?: (
conversationId: string,
/**
* 时间窗(`[sinceMs, beforeMs)`,半开)。不传 = 整个会话。
*
* recent-first 的分段计划靠它把"最近 7 天"和"更早"分开,而不是把整个会话
* 读进来再在 JS 里筛 —— 那正是大会话跑不动的根因。
*/
window?: { sinceMs?: number; beforeMs?: number }
) => Promise<chat.FormattedMessage[]>
/** /**
* 单个会话的图片消息计数(SQL 统计,不解密)。 * 单个会话的图片消息计数(SQL 统计,不解密)。
* *
@@ -94,7 +108,7 @@ export interface ImageTextIndexServiceDeps {
*/ */
countConversationImages?: ( countConversationImages?: (
conversationId: string, conversationId: string,
sinceMs?: number range?: number | { sinceMs?: number; beforeMs?: number }
) => Promise<ImageMessageCountProbe> ) => Promise<ImageMessageCountProbe>
/** /**
* 单个会话的图片消息增量水位(条数 + 最大插入序),SQL 聚合,不解密。 * 单个会话的图片消息增量水位(条数 + 最大插入序),SQL 聚合,不解密。
@@ -103,7 +117,7 @@ export interface ImageTextIndexServiceDeps {
*/ */
imageWatermark?: ( imageWatermark?: (
conversationId: string, conversationId: string,
sinceMs?: number range?: number | { sinceMs?: number; beforeMs?: number }
) => Promise<ImageMessageWatermark | null> ) => Promise<ImageMessageWatermark | null>
decryptService?: () => ImageDecryptService | null decryptService?: () => ImageDecryptService | null
/** 本地 OCR。 */ /** 本地 OCR。 */
@@ -289,6 +303,14 @@ export class ImageTextIndexService {
private counting = false private counting = false
private listeners = new Set<(status: ImageTextIndexStatus) => void>() private listeners = new Set<(status: ImageTextIndexStatus) => void>()
private lastError: string | undefined private lastError: string | undefined
/**
* 当前阶段(recent-first 可见性)。
*
* 它**不参与进度计算**:总进度永远是 `processed / totalImageMessages`。
*/
private currentPhase: ImageTextIndexPhase = 'complete'
/** 本遍 pass 的起点,用于 `preLoop.startupMs`。 */
private passStartedAt = 0
private startedAt: number | undefined private startedAt: number | undefined
/** 上一次清理实际重建(失效)了多少个会话的 Knowledge 索引;用于诊断与测试。 */ /** 上一次清理实际重建(失效)了多少个会话的 Knowledge 索引;用于诊断与测试。 */
lastInvalidatedConversations = 0 lastInvalidatedConversations = 0
@@ -602,6 +624,20 @@ export class ImageTextIndexService {
* 它必须阻断 `complete` —— 否则 Query Agent 会拿着"覆盖完整"去回答"没有"。 * 它必须阻断 `complete` —— 否则 Query Agent 会拿着"覆盖完整"去回答"没有"。
*/ */
const systemicFailure = processed > 0 && indexed === 0 && empty === 0 && missing === 0 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 { return {
totalImageMessages: total, totalImageMessages: total,
processed, processed,
@@ -613,22 +649,11 @@ export class ImageTextIndexService {
pending: Math.max(0, total - processed - runtimeUnavailable), pending: Math.max(0, total - processed - runtimeUnavailable),
// 从未统计过总数 → 不算"已建立":不知道分母就不允许声称覆盖。 // 从未统计过总数 → 不算"已建立":不知道分母就不允许声称覆盖。
established: counted !== null && (processed > 0 || runtimeUnavailable > 0), established: counted !== null && (processed > 0 || runtimeUnavailable > 0),
/** complete,
* 覆盖完整性。
*
* 分母取自流水线**真实走过**的集合时(`useScan`),不再要求 `counted.complete`:
* 那个标志表达的是"`countImageMessages()` 把每个会话都数上了",而进度现在已经
* 不用那个分母了。继续要求它,会让一个**已经跑完**的索引因为"某个会话数不上"
* 而永远停在"部分完成"。
*/
complete:
total > 0 &&
runtimeUnavailable === 0 &&
!systemicFailure &&
processed >= total &&
(useScan || (counted !== null && counted.complete)),
systemicFailure, systemicFailure,
countedAt: counted?.countedAt ?? null countedAt: counted?.countedAt ?? null,
// recent-first 的时间维度:让调用方能回答"这一段时间能不能下确定性结论"。
...this.tierCoverageSnapshot(complete)
} }
} }
@@ -672,7 +697,14 @@ export class ImageTextIndexService {
cancellable: this.running, cancellable: this.running,
paused: this.runState === 'paused', paused: this.runState === 'paused',
...this.rateSnapshot(coverage.processed, total), ...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, established: false,
complete: false, complete: false,
systemicFailure: false, systemicFailure: false,
countedAt: null countedAt: null,
tiers: [],
coveredToMs: null
}, },
storage: this.emptyStorage(), storage: this.emptyStorage(),
counting: this.counting counting: this.counting
@@ -1265,6 +1299,22 @@ export class ImageTextIndexService {
return { started: true, state: this.runState } 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<void> { private async runPass(options: ImageTextIndexStartOptions): Promise<void> {
const store = this.ensureStore() const store = this.ensureStore()
if (!store) { if (!store) {
@@ -1276,6 +1326,8 @@ export class ImageTextIndexService {
// 进入图片流水线之前的一次性成本:单列出来,避免被摊进"每张图片"。 // 进入图片流水线之前的一次性成本:单列出来,避免被摊进"每张图片"。
const preLoopStartedAt = this.now() const preLoopStartedAt = this.now()
// `startupMs` 的参照点。逐会话逻辑被抽成独立方法之后,它必须放在实例上。
this.passStartedAt = preLoopStartedAt
const capability = (await this.deps.capability?.()) ?? null const capability = (await this.deps.capability?.()) ?? null
if (capability && !capability.available) { if (capability && !capability.available) {
this.lastError = '当前系统不支持本地图片文字识别' this.lastError = '当前系统不支持本地图片文字识别'
@@ -1323,16 +1375,261 @@ export class ImageTextIndexService {
* `startupMs` 的取样点必须是"第一张图片进入流水线的那一刻",不能在这里就记 —— * `startupMs` 的取样点必须是"第一张图片进入流水线的那一刻",不能在这里就记 ——
* 否则会话级的准备成本会漏在外面,而那正是"单张很快、整遍很慢"的差额来源之一。 * 否则会话级的准备成本会漏在外面,而那正是"单张很快、整遍很慢"的差额来源之一。
*/ */
let preLoopCaptured = false const preLoopState = { captured: false }
const scanState = store.readScanState() 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<boolean> {
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) { for (const contact of contacts) {
if (this.cancelRequested || this.pauseRequested) break if (this.cancelRequested || this.pauseRequested) {
if (budget <= 0) break return { interrupted: true, imageCount, truncated }
}
if (budget.remaining <= 0) {
truncated = true
break
}
const conversationId = contact.md5 const conversationId = contact.md5
const previous = scanState.get(conversationId)
/** /**
* 会话级准备:水位 / 计数 / 让路 / 读消息。 * 会话级准备:水位 / 计数 / 让路 / 读消息。
* *
@@ -1341,45 +1638,66 @@ export class ImageTextIndexService {
* 而 `perImageMs` 只覆盖 batch 循环,看不到它。 * 而 `perImageMs` 只覆盖 batch 循环,看不到它。
*/ */
const setupStartedAt = this.now() 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) ??
// 增量水位 = 条数 + 最大插入序。只比条数会漏掉「撤回一张旧图 + Promise.resolve<ImageMessageWatermark | null>(null))
// 新增一张新图」这种总数不变、集合却变了的会话。
const watermark = await (this.deps.imageWatermark?.(conversationId, options.sinceMs) ?? /**
Promise.resolve(null)) * 这一段对这个会话实际要读的时间窗。
const imageTotal = *
watermark?.count ?? * `null` = 不看时间(整会话,新 → 旧)。
(await this.deps.countConversationImages?.(conversationId, options.sinceMs))?.count ?? */
0 let effectiveWindow = input.window
if (imageTotal === 0) { 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<ImageMessageCountProbe>({ count: null, typeColumn: null }))
if (probe.count === null) {
this.preLoop.conversationSetupMs += this.now() - setupStartedAt this.preLoop.conversationSetupMs += this.now() - setupStartedAt
store.writeScanState({
conversationId,
state: 'done',
imageTotal: 0,
imageProcessed: 0,
maxLocalId: watermark?.maxLocalId ?? 0
})
continue continue
} }
const imageTotal = probe.count
// 增量:会话已完成且**水位完全未变** → 不读 WCDB、不 OCR。 if (imageTotal === 0) {
// 水位不可用时(数据库不支持该聚合)一律重扫:宁可慢,不可漏。
const previous = scanState.get(conversationId)
if (
!windowed &&
watermark &&
previous &&
previous.state === 'done' &&
previous.imageTotal === watermark.count &&
previous.maxLocalId === watermark.maxLocalId
) {
this.preLoop.conversationSetupMs += this.now() - setupStartedAt this.preLoop.conversationSetupMs += this.now() - setupStartedAt
this.rememberConversationWatermark({
store,
scanState,
conversationId,
watermark,
processed: previous?.processed ?? 0,
marksConversationDone: input.marksConversationDone
})
continue continue
} }
@@ -1401,34 +1719,47 @@ export class ImageTextIndexService {
try { try {
const source = await this.runStep('list-image-messages', () => const source = await this.runStep('list-image-messages', () =>
this.deps.listImageMessages this.deps.listImageMessages
? this.deps.listImageMessages(conversationId) ? this.deps.listImageMessages(conversationId, windowArg)
: (this.deps.listMessages?.(conversationId) ?? Promise.resolve([])) : (this.deps.listMessages?.(conversationId) ?? Promise.resolve([]))
) )
imageMessages = source imageMessages = source
// 专用路径仍要过滤:召回归档合并可能补进非图片的撤回消息。 // 专用路径仍要过滤:召回归档合并可能补进非图片的撤回消息。
.filter(isImageMessage) .filter(isImageMessage)
// 时间窗过滤:小样本验证时只看窗口内的图片,不然还是在跑全量。 // 时间窗过滤:兼容路径没有 SQL 层窗口,只能在这里筛;这也让小样本验证
.filter((message) => // 的语义与专用路径一致。
windowed ? (message.createTime || 0) * 1000 >= (options.sinceMs as number) : true .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 { } catch {
imageMessages = [] imageMessages = []
} }
this.preLoop.listMessagesMs += this.now() - listStartedAt this.preLoop.listMessagesMs += this.now() - listStartedAt
this.preLoop.conversationSetupMs += this.now() - setupStartedAt this.preLoop.conversationSetupMs += this.now() - setupStartedAt
if (!imageMessages.length) { if (!imageMessages.length) {
store.writeScanState({ this.rememberConversationWatermark({
store,
scanState,
conversationId, conversationId,
state: 'done', watermark,
imageTotal: 0, processed: previous?.processed ?? 0,
imageProcessed: 0, marksConversationDone: input.marksConversationDone
maxLocalId: 0
}) })
continue continue
} }
// 水位取**实际读到的**消息里最大的 local_id,而不是源侧水位: imageCount += imageMessages.length
// 万一在我们查水位之后、读消息之前又落了一条新图,用观测值会让下一轮 /**
// 发现"源水位更高"从而重扫(安全);用源侧水位则会把它永久跳过(漏索引)。 * 水位取**实际读到的**消息里最大的 local_id,而不是源侧水位:
* 万一在我们查水位之后、读消息之前又落了一条新图,用观测值会让下一轮
* 发现"源水位更高"从而重扫(安全);用源侧水位则会把它永久跳过(漏索引)。
*/
const observedMaxLocalId = imageMessages.reduce( const observedMaxLocalId = imageMessages.reduce(
(max, message) => Math.max(max, Number(message.localId) || 0), (max, message) => Math.max(max, Number(message.localId) || 0),
0 0
@@ -1441,9 +1772,9 @@ export class ImageTextIndexService {
this.knowledgeDirty = false this.knowledgeDirty = false
for (let index = 0; index < imageMessages.length; index += IMAGE_TEXT_INDEX_BATCH_SIZE) { for (let index = 0; index < imageMessages.length; index += IMAGE_TEXT_INDEX_BATCH_SIZE) {
if (!preLoopCaptured) { if (!input.preLoopState.captured) {
preLoopCaptured = true input.preLoopState.captured = true
this.preLoop.startupMs = this.now() - preLoopStartedAt this.preLoop.startupMs = this.now() - this.passStartedAt
} }
if (this.cancelRequested || this.pauseRequested) { if (this.cancelRequested || this.pauseRequested) {
interrupted = true interrupted = true
@@ -1456,9 +1787,9 @@ export class ImageTextIndexService {
conversationId, conversationId,
provenance, provenance,
ocrByMessage, ocrByMessage,
() => budget > 0, () => budget.remaining > 0,
() => { () => {
budget -= 1 budget.remaining -= 1
} }
) )
processedInConversation += batchResult.processed processedInConversation += batchResult.processed
@@ -1481,23 +1812,39 @@ export class ImageTextIndexService {
* 写 `done` 会让下一遍按水位错误跳过这些图片(**永久漏索引**), * 写 `done` 会让下一遍按水位错误跳过这些图片(**永久漏索引**),
* 或者让用户以为这个会话已经处理完。预算耗尽只能记 `partial`。 * 或者让用户以为这个会话已经处理完。预算耗尽只能记 `partial`。
*/ */
if (interrupted || budget <= 0) { if (interrupted || budget.remaining <= 0) {
store.writeScanState({ store.writeScanState({
conversationId, conversationId,
state: 'partial', state: 'partial',
imageTotal: imageMessages.length, /**
* **整会话**的图片总数,不是本窗口的条数。
*
* `image_ocr_scan_state.image_total` 是进度的分母(`readScanProgress()` 求和);
* 把某个分段的窗口条数写进去,分母就会缩到"这一段处理了多少",
* 于是总进度会突然跳到接近 100% —— 那是这个功能最不能犯的谎。
*/
imageTotal: watermark?.count ?? previous?.imageTotal ?? imageMessages.length,
imageProcessed: processedInConversation, imageProcessed: processedInConversation,
maxLocalId: observedMaxLocalId 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, conversationId,
state: 'done', watermark,
imageTotal: imageMessages.length, processed: (previous?.processed ?? 0) + processedInConversation,
imageProcessed: processedInConversation, marksConversationDone: input.marksConversationDone,
maxLocalId: observedMaxLocalId observedMaxLocalId
}) })
/** /**
@@ -1530,13 +1877,77 @@ export class ImageTextIndexService {
await this.emit() await this.emit()
} }
if (this.cancelRequested) this.runState = 'cancelled' return { interrupted: false, imageCount, truncated }
else if (this.pauseRequested) this.runState = 'paused'
else this.runState = 'completed'
this.running = false
await this.emit()
} }
/**
* 写入会话 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 } { pause(): { paused: boolean; state: ImageTextIndexRunState } {
+78 -2
View File
@@ -13,11 +13,14 @@ import { mkdirSync, rmSync, statSync, existsSync } from 'node:fs'
import { dirname, join, resolve } from 'node:path' import { dirname, join, resolve } from 'node:path'
import { DatabaseSync } from 'node:sqlite' import { DatabaseSync } from 'node:sqlite'
import { import {
IMAGE_TEXT_BACKFILL_TIER_ORDER,
IMAGE_TEXT_INDEX_SCHEMA_VERSION, IMAGE_TEXT_INDEX_SCHEMA_VERSION,
type ImageOcrArtifact, type ImageOcrArtifact,
type ImageOcrBinding, type ImageOcrBinding,
type ImageOcrPersistedState, type ImageOcrPersistedState,
type ImageTextIndexStorageStats type ImageTextBackfillTier,
type ImageTextIndexStorageStats,
type ImageTextTierRunState
} from '../../shared/image-text-index' } from '../../shared/image-text-index'
const MAX_SAFE_ACCOUNT_SEGMENT = /^[a-f0-9]{32}$/ 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') 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<Record<ImageTextBackfillTier, ImageTextTierRunState>>
coveredToMs: number | null
} {
const anchorRaw = this.readMeta('backfill_anchor_ms')
const anchorMs = anchorRaw === null ? null : Number(anchorRaw)
let tierStates: Partial<Record<ImageTextBackfillTier, ImageTextTierRunState>> = {}
const statesRaw = this.readMeta('backfill_tier_states')
if (statesRaw) {
try {
const parsed = JSON.parse(statesRaw) as Record<string, unknown>
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< readScanState(): Map<
string, string,
{ state: string; imageTotal: number; processed: number; maxLocalId: number } { state: string; imageTotal: number; processed: number; maxLocalId: number }
@@ -563,7 +636,10 @@ export class ImageTextIndexStore {
DELETE FROM image_ocr_meta WHERE key IN ( DELETE FROM image_ocr_meta WHERE key IN (
'total_image_messages', 'total_image_messages',
'total_image_counted_at', 'total_image_counted_at',
'total_image_messages_complete' 'total_image_messages_complete',
'backfill_anchor_ms',
'backfill_tier_states',
'backfill_covered_to_ms'
); );
`) `)
} }
+15 -1
View File
@@ -29,7 +29,9 @@ import {
normalizeMessageIdentity normalizeMessageIdentity
} from '../../shared/local-query-api' } from '../../shared/local-query-api'
import { import {
IMAGE_TEXT_BACKFILL_TIER_LABEL,
describeImageTextCoverage, describeImageTextCoverage,
describeImageTextCoveredRanges,
imageTextCoverageState, imageTextCoverageState,
type ImageTextIndexCoverage type ImageTextIndexCoverage
} from '../../shared/image-text-index' } from '../../shared/image-text-index'
@@ -144,6 +146,17 @@ export function buildImageOcrCoverage(
? `(图片数量统计于 ${formatLocalMinute(coverage.countedAt)})` ? `(图片数量统计于 ${formatLocalMinute(coverage.countedAt)})`
: '' : ''
const base = describeImageTextCoverage(coverage) 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 { return {
state, state,
totalImageMessages: coverage.totalImageMessages, totalImageMessages: coverage.totalImageMessages,
@@ -154,10 +167,11 @@ export function buildImageOcrCoverage(
failed: coverage.failed, failed: coverage.failed,
pending: coverage.pending, pending: coverage.pending,
...(coverage.countedAt ? { countedAtLabel: formatLocalMinute(coverage.countedAt) } : {}), ...(coverage.countedAt ? { countedAtLabel: formatLocalMinute(coverage.countedAt) } : {}),
...(coveredRanges.length ? { coveredRanges } : {}),
summary: summary:
state === 'complete' state === 'complete'
? `${base}${countedNote}` ? `${base}${countedNote}`
: `${base}${countedNote}${IMAGE_OCR_ZERO_RESULT_CAUTION}` : `${base}${countedNote}${rangeNote}${IMAGE_OCR_ZERO_RESULT_CAUTION}`
} }
} }
+69 -13
View File
@@ -7,6 +7,7 @@ import { createConnection, Socket } from 'net'
import { getResourceRoots } from './resource-paths' import { getResourceRoots } from './resource-paths'
import { wcdbDebugLog } from './wcdb-debug' import { wcdbDebugLog } from './wcdb-debug'
import type { ImageMessageCountProbe } from '../shared/image-text-index' import type { ImageMessageCountProbe } from '../shared/image-text-index'
import { imageTextWindowToSeconds } from '../shared/image-text-index'
export interface Wcdb4Session { export interface Wcdb4Session {
username: string username: string
@@ -1352,13 +1353,53 @@ export class Wcdb4Client {
return value return value
} }
/** 图片消息的 WHERE 片段;`sinceMs` 用于只统计某个时间点之后的消息(测试小窗口)。 */ /**
private imageMessageWhere(column: string, sinceMs?: number): string { * 归一化图片消息的时间范围参数。
const clauses = [`(${this.quoteSqlIdentifier(column)} & 65535) = 3`] *
// 微信的 create_time 是**秒**,调用方给的是毫秒。 * 同时接受旧的裸 `sinceMs` 与新的 `{ sinceMs, beforeMs }`:
if (sinceMs && Number.isFinite(sinceMs) && sinceMs > 0) { * 前者有若干既有调用点(统计卡片、增量水位),不该为了新功能去改它们;
clauses.push(`"create_time" >= ${Math.floor(sinceMs / 1000)}`) * 后者是 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 ') return clauses.join(' AND ')
} }
@@ -1373,7 +1414,7 @@ export class Wcdb4Client {
*/ */
async countImageMessagesAsync( async countImageMessagesAsync(
md5OrUsername: string, md5OrUsername: string,
sinceMs?: number input?: number | { sinceMs?: number; beforeMs?: number }
): Promise<ImageMessageCountProbe> { ): Promise<ImageMessageCountProbe> {
if (!this.wcdbGetMessageTableStats || !this.wcdbExecQuery) { if (!this.wcdbGetMessageTableStats || !this.wcdbExecQuery) {
return { count: null, typeColumn: null, error: '当前数据服务不支持消息表统计' } return { count: null, typeColumn: null, error: '当前数据服务不支持消息表统计' }
@@ -1406,7 +1447,7 @@ export class Wcdb4Client {
this.wcdbExecQuery as unknown as KoffiAsyncFunction, this.wcdbExecQuery as unknown as KoffiAsyncFunction,
'message', 'message',
table.dbPath, 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(*)'])) const value = Number(this.pickValue(rows[0] || {}, ['image_count', 'count', 'COUNT(*)']))
if (Number.isFinite(value)) total += value if (Number.isFinite(value)) total += value
@@ -1433,7 +1474,7 @@ export class Wcdb4Client {
*/ */
async imageConversationWatermarkAsync( async imageConversationWatermarkAsync(
md5OrUsername: string, md5OrUsername: string,
sinceMs?: number input?: number | { sinceMs?: number; beforeMs?: number }
): Promise<{ count: number; maxLocalId: number } | null> { ): Promise<{ count: number; maxLocalId: number } | null> {
if (!this.wcdbGetMessageTableStats || !this.wcdbExecQuery) return null if (!this.wcdbGetMessageTableStats || !this.wcdbExecQuery) return null
// 与计数同因:必须先把会话 md5 解析成原生接口要的 username,否则永远匹配不到消息表。 // 与计数同因:必须先把会话 md5 解析成原生接口要的 username,否则永远匹配不到消息表。
@@ -1458,7 +1499,7 @@ export class Wcdb4Client {
this.wcdbExecQuery as unknown as KoffiAsyncFunction, this.wcdbExecQuery as unknown as KoffiAsyncFunction,
'message', 'message',
table.dbPath, 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 row = rows[0] || {}
const tableCount = Number(this.pickValue(row, ['image_count', 'count', 'COUNT(*)'])) const tableCount = Number(this.pickValue(row, ['image_count', 'count', 'COUNT(*)']))
@@ -1548,7 +1589,21 @@ export class Wcdb4Client {
*/ */
async listImageMessagesAsync( async listImageMessagesAsync(
md5OrUsername: string, 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<Wcdb4Message[]> { ): Promise<Wcdb4Message[]> {
if (!this.wcdbExecQuery) return [] if (!this.wcdbExecQuery) return []
const requestId = options.requestId ?? 'NO-REQUEST' const requestId = options.requestId ?? 'NO-REQUEST'
@@ -1573,10 +1628,11 @@ export class Wcdb4Client {
const column = this.resolveMessageTypeColumn(table) const column = this.resolveMessageTypeColumn(table)
if (!column) continue if (!column) continue
try { try {
const where = this.imageMessageWhere(column, options.sinceMs) const where = this.imageMessageWhere(column, options)
// `local_id` 参与排序:`create_time` 同秒的消息需要一个稳定次序, // `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 queryStartedAt = Date.now()
const rows = await this.callJsonAsync<Record<string, unknown>[]>( const rows = await this.callJsonAsync<Record<string, unknown>[]>(
this.wcdbExecQuery as unknown as KoffiAsyncFunction, this.wcdbExecQuery as unknown as KoffiAsyncFunction,
@@ -21,7 +21,10 @@ import {
import { import {
describeImageTextCoverage, describeImageTextCoverage,
imageTextCoverageState, imageTextCoverageState,
imageTextProcessedPercent imageTextPhaseLabel,
imageTextProcessedPercent,
imageTextTierSearchableNotice,
type ImageTextBackfillTier
} from '../../../../shared/image-text-index' } from '../../../../shared/image-text-index'
import { useImageTextIndexStatus } from './hooks/useImageTextIndexStatus' import { useImageTextIndexStatus } from './hooks/useImageTextIndexStatus'
@@ -89,6 +92,31 @@ export function ImageTextIndexCard({ dbReady, onNotice }: ImageTextIndexCardProp
const percent = coverage const percent = coverage
? imageTextProcessedPercent(coverage.processed, coverage.totalImageMessages) ? imageTextProcessedPercent(coverage.processed, coverage.totalImageMessages)
: 0 : 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 systemicFailure = coverage?.systemicFailure === true
const visualState = const visualState =
progress?.state === 'error' || coverageState === 'failed' progress?.state === 'error' || coverageState === 'failed'
@@ -243,6 +271,15 @@ export function ImageTextIndexCard({ dbReady, onNotice }: ImageTextIndexCardProp
}} }}
/> />
</div> </div>
{/*
阶段只是一句人话,没有百分比 —— 进度必须留在上面那行
(processed / total 全量),否则用户会把"最近 7 天做完了"读成"整体做完了"。
*/}
{phaseLabel && (
<p className="ai-search-knowledge-pass-line" data-testid="image-text-index-phase">
{phaseLabel}
</p>
)}
<p className="ai-search-knowledge-pass-line"> <p className="ai-search-knowledge-pass-line">
{`${progress.percent}% · 识别出文字 ${progress.indexed.toLocaleString()} · 没有文字 ${progress.empty.toLocaleString()} · 图片已清理 ${progress.missing.toLocaleString()} · 失败 ${progress.failed.toLocaleString()}`} {`${progress.percent}% · 识别出文字 ${progress.indexed.toLocaleString()} · 没有文字 ${progress.empty.toLocaleString()} · 图片已清理 ${progress.missing.toLocaleString()} · 失败 ${progress.failed.toLocaleString()}`}
</p> </p>
@@ -291,6 +328,18 @@ export function ImageTextIndexCard({ dbReady, onNotice }: ImageTextIndexCardProp
</strong> </strong>
</div> </div>
)} )}
{/*
只有**真正完整**的分段才敢说"已可搜索"。这是一句承诺,不是进度提示 ——
所以判据是 coverage 里该分段的 state === 'complete',而不是"扫到过"。
*/}
{searchableNotice && (
<p
className="ai-search-knowledge-pass-line"
data-testid="image-text-index-searchable-notice"
>
{searchableNotice}
</p>
)}
<p className="ai-search-knowledge-pass-line">{describeImageTextCoverage(coverage)}</p> <p className="ai-search-knowledge-pass-line">{describeImageTextCoverage(coverage)}</p>
</div> </div>
)} )}
+322
View File
@@ -70,6 +70,308 @@ export function resolveImageTextOcrConcurrency(raw?: string | number | null): nu
/** 已完成一批之后、回到会话循环前的让出时间。 */ /** 已完成一批之后、回到会话循环前的让出时间。 */
export const IMAGE_TEXT_INDEX_YIELD_MS = 0 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<ImageTextBackfillTier, number> = {
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<ImageTextBackfillTier, string> = {
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 结果状态。 */ /** 单张图片的 OCR 结果状态。 */
export type ImageOcrState = export type ImageOcrState =
/** 尚未处理 */ /** 尚未处理 */
@@ -269,6 +571,13 @@ export interface ImageTextIndexProgress {
speedPerSec?: number | null speedPerSec?: number | null
/** 按当前窗口速度估算的剩余时间(毫秒);速度不可用或分母不可信时为 null。 */ /** 按当前窗口速度估算的剩余时间(毫秒);速度不可用或分母不可信时为 null。 */
etaMs?: number | null etaMs?: number | null
/**
* 当前正在处理的阶段(recent-first 的可见性)。
*
* **它只是阶段提示,不是进度**:总进度仍然必须是 `processed / totalImageMessages`
* (全量图片消息),绝不允许用"某个分段处理完了"冒充整体完成。
*/
currentPhase?: ImageTextIndexPhase
} }
/** /**
@@ -311,6 +620,19 @@ export interface ImageTextIndexCoverage {
* 冒充成「现在完整」。 * 冒充成「现在完整」。
*/ */
countedAt: number | null countedAt: number | null
/**
* 各历史分段的完成状态与边界(新 → 旧)。未建立 / 旧快照时为 `[]`。
*
* 边界来自**当时固定的锚点**,因此可以和用户查询的时间范围直接求交。
*/
tiers: ImageTextTierCoverage[]
/**
* 增量补齐水位:`create_time <= coveredToMs` 的新图片都已处理完;null = 尚无。
*
* 它单独存在的原因:锚点之后新到的图片不属于任何历史分段,必须有独立水位
* 才能回答"最近这几个小时是否已经可搜"。
*/
coveredToMs: number | null
} }
/** 覆盖度状态(外加"未建立")。UI 与 Query Agent 共用同一判据,避免两处各推一套口径漂移。 */ /** 覆盖度状态(外加"未建立")。UI 与 Query Agent 共用同一判据,避免两处各推一套口径漂移。 */
+8
View File
@@ -369,6 +369,14 @@ export interface QueryImageTextCoverage {
pending: number pending: number
/** 图片数量统计时刻(本地时间 `MM-DD HH:mm`);从未统计时为 undefined。 */ /** 图片数量统计时刻(本地时间 `MM-DD HH:mm`);从未统计时为 undefined。 */
countedAtLabel?: string countedAtLabel?: string
/**
* 已经**真正完整**的时间段描述(新 → 旧),例如"最近 7 天"。
*
* recent-first 之后必须有这一维:总进度 30% 不代表"最近一周不可信",
* 反过来总进度 99% 也不代表"去年可以下确定性结论"。模型只能引用这里列出的
* 时间段去下"没有"的结论,其余范围一律只能说"仍在补齐"。
*/
coveredRanges?: string[]
/** 可直接引用的结论句;模型只引用,不要自己换算或推断。 */ /** 可直接引用的结论句;模型只引用,不要自己换算或推断。 */
summary: string summary: string
} }
@@ -445,3 +445,105 @@ describe('图片文字索引卡片 — 修复图片搜索索引', () => {
expect(String(onNotice.mock.calls.at(-1)?.[0])).toContain('正在进行中') expect(String(onNotice.mock.calls.at(-1)?.[0])).toContain('正在进行中')
}) })
}) })
/**
* recent-first 的阶段可见性。
*
* 这里要守住的是两件**不能混**的事:
* - **总进度**永远是 `processed / totalImageMessages`(全量图片消息);
* - **阶段文案**只回答"现在在优先做什么"。
*
* 用"某个分段做完了"去冒充整体完成,是这次改动最容易撒的谎,所以两条都断言。
*/
describe('recent-first 阶段显示', () => {
/** 已建立、且不在运行/暂停 —— 卡片明细块的渲染条件。 */
const settled = (
coverageOverride: Partial<ImageTextIndexStatus['coverage']>
): 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()
})
})
@@ -19,6 +19,7 @@ import {
} from '../../src/main/services/image-text-index-store' } from '../../src/main/services/image-text-index-store'
import { buildImageOcrCoverage } from '../../src/main/services/local-query-api-service' import { buildImageOcrCoverage } from '../../src/main/services/local-query-api-service'
import type { ImageTextIndexCoverage } from '../../src/shared/image-text-index' import type { ImageTextIndexCoverage } from '../../src/shared/image-text-index'
import { sourceMessageId } from '../../src/main/knowledge/message-identity'
const ACCOUNT = 'wxid_fixture_account' const ACCOUNT = 'wxid_fixture_account'
const CONVERSATION = 'conversation-md5-fixture' 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 * recent-first 之后,"什么时候读 WCDB"由**分段完成状态 + 增量水位**共同决定,
harness.watermark.maxLocalId = 20 * 但两条硬约束一个字都没变:
* 1. 没有新内容时**不得**重复读会话消息、更不得重复 OCR;
await harness.service.startPass() * 2. 有新内容(含"总数不变但集合变了")时**必须**处理到 —— 宁可慢,也不漏。
// 走到完成态需要等内部 promise 收敛。 *
await vi.waitFor(() => expect(harness.service.isRunning()).toBe(false)) * 这里的夹具刻意做成**窗口感知**的:新架构的"这段没有图片就整段跳过"完全依赖
expect(harness.listMessages).toHaveBeenCalledTimes(1) * 计数说实话;忽略窗口的夹具测出来的只是"夹具不过滤",不是调度器的行为。
*/
// 第二遍:水位完全一致 → 跳过,不再读会话消息。 describe('增量判据:水位不得漏掉新图片', () => {
await harness.service.startPass() function makeIncrementalHarness(options: {
await vi.waitFor(() => expect(harness.service.isRunning()).toBe(false)) messages: chat.FormattedMessage[]
expect(harness.listMessages).toHaveBeenCalledTimes(1) withWatermark?: boolean
}) }): {
service: ImageTextIndexService
it('总数相同但最大插入序前进 → 必须重扫(撤回一张旧图 + 新增一张新图)', async () => { databasePath: string
const harness = makeHarness({ messages: [imageMessage(10, 1000), imageMessage(20, 2000)] }) state: { messages: chat.FormattedMessage[]; count: number; maxLocalId: number; now: number }
harness.watermark.count = 2 listMessages: ReturnType<typeof vi.fn>
harness.watermark.maxLocalId = 20 listImageMessages: ReturnType<typeof vi.fn>
} {
await harness.service.startPass() const databaseRoot = makeDatabaseRoot()
await vi.waitFor(() => expect(harness.service.isRunning()).toBe(false)) const databasePath = getImageTextIndexDatabasePath(databaseRoot, ACCOUNT)
expect(harness.listMessages).toHaveBeenCalledTimes(1) const state = {
messages: [...options.messages],
// 集合变了、条数没变:localId 10 被撤回,新增 localId 30。 count: options.messages.length,
harness.listMessages.mockImplementation(async () => [ maxLocalId: options.messages.reduce((max, m) => Math.max(max, Number(m.localId) || 0), 0),
imageMessage(20, 2000), now: 1_800_000_000_000
imageMessage(30, 3000) }
]) const inWindow = (
harness.watermark.maxLocalId = 30 message: chat.FormattedMessage,
window?: { sinceMs?: number; beforeMs?: number }
await harness.service.startPass() ): boolean => {
await vi.waitFor(() => expect(harness.service.isRunning()).toBe(false)) const createTimeMs = (message.createTime || 0) * 1000
// 只看 count 的实现会在这里静默跳过 —— 那正是会漏掉新图片的洞。 if (window?.sinceMs !== undefined && createTimeMs < window.sinceMs) return false
expect(harness.listMessages).toHaveBeenCalledTimes(2) if (window?.beforeMs !== undefined && createTimeMs >= window.beforeMs) return false
}) return true
}
it('水位不可用(数据库不支持该聚合)时一律重扫,宁可慢也不漏', async () => { const listMessages = vi.fn(async () => state.messages)
const harness = makeHarness({ messages: [imageMessage(10, 1000)] }) const listImageMessages = vi.fn(
harness.watermark.count = 1 async (_conversationId: string, window?: { sinceMs?: number; beforeMs?: number }) =>
state.messages.filter((message) => inWindow(message, window))
)
const service = new ImageTextIndexService() const service = new ImageTextIndexService()
const listMessages = vi.fn(async () => [imageMessage(10, 1000)])
service.bind({ service.bind({
databaseRoot: harness.databaseRoot, databaseRoot,
resolveAccountId: () => ACCOUNT, 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, listMessages,
countConversationImages: async () => ({ count: 1, typeColumn: 'local_type' }), listImageMessages,
// 关键:不提供 imageWatermark 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, decryptService: () => ({ findImageFile: () => null, decryptImage: () => null }) as never,
capability: async () => ({ capability: async () => ({
available: true, available: true,
engine: 'windows-system-ocr', engine: 'windows-system-ocr',
platform: 'win32', platform: 'win32',
runtimeVersion: null, runtimeVersion: '1.2.0',
language: null 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 harness.service.startPass()
await vi.waitFor(() => expect(service.isRunning()).toBe(false)) await vi.waitFor(() => expect(harness.service.isRunning()).toBe(false))
await service.startPass() // 两条图片的 create_time 都很老 → 只落在归档段,只被那一段读到。
await vi.waitFor(() => expect(service.isRunning()).toBe(false)) expect(harness.listImageMessages).toHaveBeenCalledTimes(1)
expect(listMessages).toHaveBeenCalledTimes(2) 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)))
}) })
}) })
@@ -86,8 +86,24 @@ function createHarness(): Harness {
const all = mixedMessages() const all = mixedMessages()
const imagesOnly = all.filter((message) => message.contentData?.type === 'image') 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 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() const service = new ImageTextIndexService()
service.bind({ service.bind({
@@ -98,7 +114,13 @@ function createHarness(): Harness {
listContacts: async () => [ listContacts: async () => [
{ md5: CONVERSATION, m_nsUsrName: 'boundary', type: 'user' as const } { 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 }), imageWatermark: async () => ({ count: IMAGE_COUNT, maxLocalId: TEXT_COUNT + IMAGE_COUNT }),
listMessages, listMessages,
listImageMessages, listImageMessages,
@@ -117,8 +117,15 @@ function createService(options: { count: number; notifyIntervalMs: number; perIm
return { service, notifications } return { service, notifications }
} }
/**
* 等 pass 收尾。
*
* 显式给足超时:这些用例故意让 240 张图各睡几毫秒来制造可观测的持续时间,
* 而 `vi.waitFor` 的默认超时是 1000ms —— 机器稍慢就会以**断言失败**而不是
* "超时"的形态报出来。这里等的是"跑完",不是"跑得快"。
*/
const finish = async (service: ImageTextIndexService): Promise<void> => { const finish = async (service: ImageTextIndexService): Promise<void> => {
await vi.waitFor(() => expect(service.isRunning()).toBe(false)) await vi.waitFor(() => expect(service.isRunning()).toBe(false), { timeout: 30_000 })
} }
describe('进度通知节流', () => { describe('进度通知节流', () => {
@@ -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<typeof vi.fn>
/** 真正"过了一遍处理"的图片数 —— 用它证明断点续跑没有重做。 */
findImageFile: ReturnType<typeof vi.fn>
}
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<Record<ImageTextBackfillTier, ImageTextTierRunState>> => {
const store = new ImageTextIndexStore(databasePath, ACCOUNT)
const result = store.readBackfillState()
store.close()
return result.tierStates
}
const runPass = async (
harness: HarnessResult,
options?: Parameters<ImageTextIndexService['startPass']>[0]
): Promise<void> => {
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<string, number[]>
} {
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<string, number[]>()
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>): 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'
])
})
})