import { Disposable } from '#/_base/di/lifecycle'; import { LifecycleScope } from '#/app/scopes'; import { ScopeActivation, registerScopedService } from '#/_base/di/scope'; import { ILogService } from '#/_base/log/log'; import { IntervalTimer } from '#/_base/utils/timer'; import { IBootstrapService } from '#/app/bootstrap/bootstrap'; import { IConfigService } from '#/app/config/config'; import { ITelemetryService } from '#/app/telemetry/telemetry'; import type { SessionIndexDegradedEvent } from '#/app/telemetry/events'; import { isError2 } from '#/errors'; import { databaseBaseEnabled } from '#/persistence/configSection'; import { IAtomicDocumentStore } from '#/persistence/interface/atomicDocumentStore'; import { IQueryStore, type Checkpoint, type ColumnBounds, type Page, type QueryFilter, } from '#/persistence/interface/queryStore'; import { IFileSystemStorageService } from '#/persistence/interface/storage'; import { CHILD_SESSION_KIND, CHILD_SESSION_KIND_KEY, ISessionIndex, ISessionIndexMirror, PARENT_SESSION_ID_KEY, type SessionCountQuery, type SessionIndexStatus, type SessionIndexState, type SessionListQuery, type SessionSummary, } from './sessionIndex'; import { markSessionDirty } from './sessionIndexDirtyJournal'; import { PARENT_INDEX_NAME, SESSION_INDEX_MANIFEST, SESSION_INDEX_SCHEMA_VERSION, recencyColumn, sessionCollection, sessionCountersCollection, stripRecencyField, type SessionWorkspaceCounts, } from './sessionIndexModel'; import { SessionIndexProjector } from './sessionIndexProjector'; import { listSessionIds, listWorkspaceIds, readSessionSummary, scanSessionsFreshness, summaryMatchesChildOf, } from './sessionIndexSource'; const RECONCILE_INTERVAL_MS = 60_000; const DEGRADED_RETRY_MS = 5_000; const TIE_REPAIR_LIMIT = 1_000; const UNBOUNDED = Number.MAX_SAFE_INTEGER; function canonicalOrder(a: SessionSummary, b: SessionSummary): number { if (a.updatedAt !== b.updatedAt) return b.updatedAt - a.updatedAt; return a.id < b.id ? 1 : a.id > b.id ? -1 : 0; } function isSessionSummaryShape(value: unknown): value is SessionSummary { if (value === null || typeof value !== 'object') return false; const summary = value as Record; return ( typeof summary['id'] === 'string' && typeof summary['workspaceId'] === 'string' && typeof summary['createdAt'] === 'number' && typeof summary['updatedAt'] === 'number' && typeof summary['archived'] === 'boolean' ); } export class FileSessionIndex extends Disposable implements ISessionIndex { declare readonly _serviceBrand: undefined; private state: SessionIndexState = 'uninitialized'; private generation: number | undefined; private statusReason: string | undefined; private degradedCount = 0; private nextPrepareRetryAt = 0; private lastDegradedKey: string | undefined; private prepareFlight: Promise | undefined; private projectFlight: Promise | undefined; private readonly reconcileTimer = this._register(new IntervalTimer({ unref: true })); private readonly projector: SessionIndexProjector; constructor( @IBootstrapService private readonly bootstrap: IBootstrapService, @IFileSystemStorageService private readonly storage: IFileSystemStorageService, @IAtomicDocumentStore private readonly docs: IAtomicDocumentStore, @IQueryStore private readonly queryStore: IQueryStore, @IConfigService private readonly config: IConfigService, @ISessionIndexMirror private readonly mirror: ISessionIndexMirror, @ITelemetryService private readonly telemetry: ITelemetryService, @ILogService private readonly log: ILogService, ) { super(); this.projector = new SessionIndexProjector({ storage, docs, queryStore, log, sessionsScope: bootstrap.scope('sessions'), }); } private ensureReconcileTimer(): void { if (!this.reconcileTimer.isSet()) { this.reconcileTimer.cancelAndSet(() => void this.tick(), RECONCILE_INTERVAL_MS); } } async prepare(options?: { deadlineMs?: number }): Promise { if (!this.readModelEnabled()) return this.status(); this.prepareFlight ??= this.doPrepare(options?.deadlineMs).finally(() => { this.prepareFlight = undefined; }); return this.prepareFlight; } status(): SessionIndexStatus { return { state: this.readModelEnabled() ? this.state : 'uninitialized', generation: this.generation, reason: this.statusReason, degradedCount: this.degradedCount, }; } private async doPrepare(deadlineMs?: number): Promise { if (this.state === 'ready') return this.status(); this.state = 'preparing'; try { const manifest = await this.queryStore.getCheckpoint(SESSION_INDEX_MANIFEST); if (manifest === undefined || manifest.schemaVersion !== SESSION_INDEX_SCHEMA_VERSION) { const projection = this.ensureProjection(); if (deadlineMs === undefined) { await projection; } else { await Promise.race([ projection, new Promise((resolve) => { setTimeout(resolve, deadlineMs); }), ]); } } else { this.generation = manifest.seq; await this.ensureSchema(manifest.seq); if (!(await this.manifestFresh(manifest))) { try { const reconciliation = this.projector.reconcile(manifest.seq); if (deadlineMs === undefined) { await reconciliation; } else { await Promise.race([ reconciliation, new Promise((resolve) => { setTimeout(resolve, deadlineMs); }), ]); } } catch (error) { const published = await this.queryStore .getCheckpoint(SESSION_INDEX_MANIFEST) .catch(() => undefined); if (published === undefined) throw error; this.log.warn('session index startup reconciliation failed; serving the published generation', { error: String(error), }); } } } const published = await this.queryStore.getCheckpoint(SESSION_INDEX_MANIFEST); if (published !== undefined) { this.generation = published.seq; this.markReady(); } } catch (error) { this.markDegraded('prepare failed', error); } return this.status(); } private async manifestFresh(manifest: Checkpoint): Promise { if (manifest.sourceSessionCount === undefined) return false; try { const scan = await scanSessionsFreshness(this.storage, this.sessionsScope); return scan.dirtyMarkCount === 0 && scan.sessionCount === manifest.sourceSessionCount; } catch (error) { this.log.warn('session index freshness check failed; treating the index as stale', { error: String(error), }); return false; } } private ensureProjection(): Promise { this.projectFlight ??= this.runProjection().finally(() => { this.projectFlight = undefined; }); return this.projectFlight; } private async runProjection(): Promise { const startedAt = Date.now(); try { const manifest = await this.queryStore.getCheckpoint(SESSION_INDEX_MANIFEST); const next = (manifest?.seq ?? 0) + 1; const result = await this.projector.project(next); this.generation = result.generation; this.markReady(); this.telemetry.track2('session_index_projected', { duration_ms: Date.now() - startedAt, session_count: result.sessions, generation: result.generation, }); } catch (error) { const published = await this.queryStore .getCheckpoint(SESSION_INDEX_MANIFEST) .catch(() => undefined); if (published !== undefined) { this.generation = published.seq; this.markReady(); this.log.warn('session index re-projection failed; staying on the previous generation', { generation: published.seq, error: String(error), }); } else { this.markDegraded('projection failed', error); } } } async reconcileNow(): Promise { if (!this.readModelEnabled()) return; const manifest = await this.queryStore.getCheckpoint(SESSION_INDEX_MANIFEST); if (manifest === undefined) return; this.generation = manifest.seq; await this.projector.reconcile(manifest.seq); } async reprojectNow(): Promise { if (!this.readModelEnabled()) return; await this.ensureProjection(); } stopReconcileLoop(): void { this.reconcileTimer.cancel(); } private async tick(): Promise { if (!this.readModelEnabled()) return; if (this.state === 'degraded') { void this.prepare(); return; } if (this.state !== 'ready') return; try { const manifest = await this.queryStore.getCheckpoint(SESSION_INDEX_MANIFEST); if (manifest === undefined) { this.markDegraded('published generation lost'); void this.prepare(); return; } this.generation = manifest.seq; if (await this.manifestFresh(manifest)) return; await this.projector.reconcile(manifest.seq); } catch (error) { this.log.warn('session index reconciliation failed', { error: String(error) }); } } private markReady(): void { if (this.state === 'degraded') { this.log.info('session index read model recovered', { degradedCount: this.degradedCount }); } this.state = 'ready'; this.statusReason = undefined; this.lastDegradedKey = undefined; this.ensureReconcileTimer(); } private markDegraded(reason: string, error?: unknown): void { this.state = 'degraded'; this.statusReason = reason; this.degradedCount += 1; this.nextPrepareRetryAt = Date.now() + DEGRADED_RETRY_MS; this.ensureReconcileTimer(); const detail = error instanceof Error ? error.message : typeof error === 'string' ? error : undefined; const episodeKey = `${reason}:${detail ?? ''}`; if (episodeKey === this.lastDegradedKey) return; this.lastDegradedKey = episodeKey; this.log.warn('session index read model degraded; serving authoritative reads', { reason, ...(detail !== undefined ? { error: detail } : {}), degradedCount: this.degradedCount, }); const properties: SessionIndexDegradedEvent = { reason, degraded_count: this.degradedCount, }; if (error !== undefined) { properties.error_type = isError2(error) ? error.code : error instanceof Error ? error.name : 'Unknown'; } this.telemetry.track2('session_index_degraded', properties); } private async ensureSchema(generation: number): Promise { await this.queryStore.ensureIndex(sessionCollection(generation), { kind: 'value', name: PARENT_INDEX_NAME, field: `custom.${PARENT_SESSION_ID_KEY}`, }); } async get(id: string): Promise { return this.withReadModel( (generation) => this.getFromReadModel(generation, id), () => this.getLegacy(id), ); } async listRecent(query: SessionListQuery): Promise> { return this.withReadModel( (generation) => this.listRecentFromReadModel(generation, query), () => this.listLegacy(query), ); } async count(query: SessionCountQuery): Promise { return this.withReadModel( (generation) => this.countFromReadModel(generation, query), () => this.countLegacy(query), ); } async remove(id: string): Promise { await this.mirror.evict(id); await this.withReadModel( async (generation) => { await this.queryStore.delete(sessionCollection(generation), id); }, () => Promise.resolve(), ); try { await markSessionDirty(this.storage, this.sessionsScope, id); } catch (error) { this.log.warn('session index dirty mark failed', { error: String(error) }); } } private async withReadModel( op: (generation: number) => Promise, legacy: () => Promise, ): Promise { if (!this.readModelEnabled()) return legacy(); if (this.state === 'uninitialized') { void this.prepare(); return legacy(); } if (this.state === 'preparing') return legacy(); if (this.state === 'degraded') { if (Date.now() >= this.nextPrepareRetryAt) void this.prepare(); return legacy(); } let manifest: Checkpoint | undefined; try { manifest = await this.queryStore.getCheckpoint(SESSION_INDEX_MANIFEST); } catch (error) { this.markDegraded('read model read failed', error); return legacy(); } if (manifest === undefined) { this.markDegraded('published generation lost'); void this.prepare(); return legacy(); } this.generation = manifest.seq; try { return await op(manifest.seq); } catch (error) { this.markDegraded('read model read failed', error); return legacy(); } } private async getFromReadModel( generation: number, id: string, ): Promise { const queued = this.mirror.pending().find((summary) => summary.id === id); if (queued !== undefined) return queued; const cached: unknown = await this.queryStore.get(sessionCollection(generation), id); if (isSessionSummaryShape(cached)) return stripRecencyField(generation, cached); const summary = await this.getLegacy(id); if (summary !== undefined) this.mirror.record(summary); return summary; } private async listRecentFromReadModel( generation: number, query: SessionListQuery, ): Promise> { const collection = sessionCollection(generation); if (query.sessionId !== undefined) { const summary = await this.getFromReadModel(generation, query.sessionId); const items = summary !== undefined && (!summary.archived || query.includeArchived === true) ? [summary] : []; return { items: query.limit !== undefined ? items.slice(0, query.limit) : items }; } const cursor = await this.resolveCursor(generation, query); if (cursor === undefined) return { items: [] }; const limit = query.limit ?? UNBOUNDED; const filter = { ...this.baseFilter(query), ...cursor.filter }; const column = recencyColumn(generation); const strip = (records: SessionSummary[]): SessionSummary[] => records.map((record) => stripRecencyField(generation, record)); const page = query.childOf !== undefined ? await this.windowedPage( (bounds, fetchLimit) => { const base = this.queryStore .query(collection) .where(filter) .orderBy('updatedAt', 'desc') .limit(fetchLimit); const q = Object.keys(bounds).length > 0 ? base.whereColumn(column, bounds) : base; return q.execute().then((p) => strip([...p.items])); }, cursor.bounds, limit, ) : await this.windowedPage( (bounds, fetchLimit) => this.queryStore .pageByColumn(collection, { column, dir: 'desc', filter, bounds, limit: fetchLimit, }) .then((p) => strip([...p.items])), cursor.bounds, limit, ); return this.mergePending(page, query, cursor.position); } private async countFromReadModel( generation: number, query: SessionCountQuery, ): Promise { const counters = sessionCountersCollection(generation); const restricted = query.workspaceIds; const workspaceIds = restricted ?? (await this.queryStore.listKeys(counters)); const counts = await this.queryStore.getMany(counters, workspaceIds); let total = 0; for (const entry of counts.values()) { total += query.includeArchived === true ? entry.active + entry.archived : entry.active; } const pending = this.mirror .pending() .filter((summary) => restricted === undefined || restricted.includes(summary.workspaceId)); if (pending.length === 0) return total; const stored = await this.queryStore.getMany( sessionCollection(generation), pending.map((summary) => summary.id), ); const weight = (archived: boolean): number => query.includeArchived === true || !archived ? 1 : 0; for (const summary of pending) { const old = stored.get(summary.id); total += weight(summary.archived) - (old === undefined ? 0 : weight(old.archived)); } return total; } private async windowedPage( fetch: (bounds: ColumnBounds, limit: number) => Promise, bounds: ColumnBounds, limit: number, ): Promise> { const raw = await fetch(bounds, limit + 1); if (raw.length <= limit) { return { items: raw.toSorted(canonicalOrder) }; } const minUpdatedAt = Math.min(...raw.map((summary) => summary.updatedAt)); const tie = await fetch({ gte: minUpdatedAt, lte: minUpdatedAt }, TIE_REPAIR_LIMIT); const merged = new Map(); for (const summary of tie) merged.set(summary.id, summary); for (const summary of raw) merged.set(summary.id, summary); const items = [...merged.values()].toSorted(canonicalOrder); const kept = items.slice(0, limit); const hasMore = items.length > limit || tie.length >= TIE_REPAIR_LIMIT; return { items: kept, nextCursor: hasMore ? kept.at(-1)!.id : undefined }; } private mergePending( page: Page, query: SessionListQuery, position?: { u: number; id: string; before: boolean }, ): Page { const pending = this.mirror .pending() .filter( (summary) => (query.workspaceIds === undefined || query.workspaceIds.includes(summary.workspaceId)) && (query.includeArchived === true || !summary.archived) && summaryMatchesChildOf(summary, query.childOf) && (position === undefined || (position.before ? summary.updatedAt < position.u || (summary.updatedAt === position.u && summary.id < position.id) : summary.updatedAt > position.u || (summary.updatedAt === position.u && summary.id > position.id))), ); if (pending.length === 0) return page; const merged = new Map(); for (const summary of pending) merged.set(summary.id, summary); for (const summary of page.items) { if (!merged.has(summary.id)) merged.set(summary.id, summary); } const items = [...merged.values()].toSorted(canonicalOrder); if (query.limit === undefined) return { items }; const kept = items.slice(0, query.limit); const hasMore = page.nextCursor !== undefined || items.length > query.limit; return { items: kept, nextCursor: hasMore ? kept.at(-1)!.id : undefined }; } private async resolveCursor( generation: number, query: SessionListQuery, ): Promise< | { filter: QueryFilter; bounds: ColumnBounds; position?: { u: number; id: string; before: boolean } } | undefined > { const id = query.before ?? query.after; if (id === undefined) return { filter: {}, bounds: {} }; const storedValue: unknown = await this.queryStore.get(sessionCollection(generation), id); const stored = isSessionSummaryShape(storedValue) ? storedValue : undefined; const cursor = stored ?? this.mirror.pending().find((summary) => summary.id === id); if (cursor === undefined) return undefined; const u = cursor.updatedAt; if (query.before !== undefined) { return { bounds: { lte: u }, filter: { $or: [ { updatedAt: { $lt: u } }, { updatedAt: u, id: { $lt: id } }, ], }, position: { u, id, before: true }, }; } return { bounds: { gte: u }, filter: { $or: [ { updatedAt: { $gt: u } }, { updatedAt: u, id: { $gt: id } }, ], }, position: { u, id, before: false }, }; } private baseFilter(query: SessionListQuery): QueryFilter { const filter: Record = {}; if (query.workspaceIds !== undefined) { filter['workspaceId'] = query.workspaceIds.length === 1 ? query.workspaceIds[0] : { $in: [...query.workspaceIds] }; } if (query.childOf !== undefined) { filter[`custom.${PARENT_SESSION_ID_KEY}`] = query.childOf; filter[`custom.${CHILD_SESSION_KIND_KEY}`] = CHILD_SESSION_KIND; } if (query.includeArchived !== true) filter['archived'] = { $ne: true }; return filter; } private get sessionsScope(): string { return this.bootstrap.scope('sessions'); } private async listLegacy(query: SessionListQuery): Promise> { if (query.sessionId !== undefined) { const summary = await this.getLegacy(query.sessionId); const items = summary !== undefined && (!summary.archived || query.includeArchived === true) ? [summary] : []; return { items: query.limit !== undefined ? items.slice(0, query.limit) : items }; } const collected = (await this.collectAuthoritative(query.workspaceIds)).filter( (summary) => (query.includeArchived === true || !summary.archived) && summaryMatchesChildOf(summary, query.childOf), ); const items = collected.toSorted(canonicalOrder); let start = 0; let end = items.length; const cursorId = query.before ?? query.after; if (cursorId !== undefined) { const index = items.findIndex((summary) => summary.id === cursorId); if (index === -1) return { items: [] }; if (query.before !== undefined) start = index + 1; else end = index; } const window = items.slice(start, end); if (query.limit === undefined) return { items: window }; const kept = window.slice(0, query.limit); return { items: kept, nextCursor: window.length > query.limit ? kept.at(-1)!.id : undefined, }; } private async getLegacy(id: string): Promise { for (const workspaceId of await listWorkspaceIds(this.storage, this.sessionsScope)) { const sessionIds = await listSessionIds(this.storage, this.sessionsScope, workspaceId); if (!sessionIds.includes(id)) continue; const summary = await readSessionSummary(this.docs, this.sessionsScope, workspaceId, id); if (summary !== undefined) return summary; } return undefined; } private async countLegacy(query: SessionCountQuery): Promise { let count = 0; for (const summary of await this.collectAuthoritative(query.workspaceIds)) { if (query.includeArchived === true || !summary.archived) count += 1; } return count; } private async collectAuthoritative( workspaceIds: readonly string[] | undefined, ): Promise { let collected: SessionSummary[]; if ( this.readModelEnabled() && (this.state === 'uninitialized' || this.state === 'preparing') ) { const { summaries } = await this.projector.sharedScanForRead(); collected = workspaceIds === undefined ? summaries : summaries.filter((summary) => workspaceIds.includes(summary.workspaceId)); } else { const ids = workspaceIds ?? (await listWorkspaceIds(this.storage, this.sessionsScope)); collected = []; for (const workspaceId of ids) { for (const sessionId of await listSessionIds(this.storage, this.sessionsScope, workspaceId)) { const summary = await readSessionSummary(this.docs, this.sessionsScope, workspaceId, sessionId); if (summary !== undefined) collected.push(summary); } } } const pending = this.mirror.pending(); if (pending.length === 0) return collected; const byId = new Map(collected.map((summary) => [summary.id, summary])); for (const summary of pending) { if (workspaceIds !== undefined && !workspaceIds.includes(summary.workspaceId)) continue; byId.set(summary.id, summary); } return [...byId.values()]; } private readModelEnabled(): boolean { return databaseBaseEnabled(this.config); } } registerScopedService( LifecycleScope.App, ISessionIndex, FileSessionIndex, ScopeActivation.OnScopeCreated, 'sessionIndex', );