From 11f452e899f03b4322a8e91af37acfc0706952a6 Mon Sep 17 00:00:00 2001 From: antra-tess Date: Tue, 22 Sep 2026 22:06:09 -0700 Subject: [PATCH 1/5] feat(history): semantic_search over messages + summaries via a shared embed-service MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Adds an optional fifth HistoryModule tool, `semantic_search`, for "I remember roughly what it was about" queries the substring/regex `search` can't answer. The vectors live server-side in a fleet-shared embed-service (one per deployment cluster; namespace = one store), so remote residents need no local model or index. - src/modules/history/semantic.ts: SemanticIndexClient (upsert/search/stats) and SemanticIndexer. Sync is incremental and idempotent: messages walk the native time index from the service's cursor watermark minus a 10-minute overlap (the service dedups by id + text hash, so re-sends are free); summaries use createdMs (a fresh L3 spans months of old timestamps — its creation time is the only monotonic signal). What gets embedded per message: text blocks plus the string args of think/journal/skip_reply/ private_note (the agent's own diary); tool results, other tool args, thinking blocks and media are skipped. Failures back off exponentially and never break the module. - HistoryModule({ semantic }): tool offered only when configured; bounded catch-up before each search (a backlog is reported as `index.behind`, not blocked on); background tick (unref'd) from bind(); channel labels resolve through the same ChannelRegistry path as the other tools; level implies summaries, channel implies messages. - Exports HistoryModuleOptions / SemanticIndexConfig. - test/history-module-semantic.test.ts: fake embed-service over node:http — tool gating, text extraction, incremental sync with cursors, filter mapping, clean failure + backoff, input validation. Co-Authored-By: Claude Fable 5.1 --- changelog.d/history-semantic-search.added.md | 9 + src/modules/history/index.ts | 149 +++++++++ src/modules/history/semantic.ts | 329 +++++++++++++++++++ src/modules/index.ts | 3 +- test/history-module-semantic.test.ts | 154 +++++++++ 5 files changed, 643 insertions(+), 1 deletion(-) create mode 100644 changelog.d/history-semantic-search.added.md create mode 100644 src/modules/history/semantic.ts create mode 100644 test/history-module-semantic.test.ts diff --git a/changelog.d/history-semantic-search.added.md b/changelog.d/history-semantic-search.added.md new file mode 100644 index 0000000..a83230b --- /dev/null +++ b/changelog.d/history-semantic-search.added.md @@ -0,0 +1,9 @@ +- `HistoryModule` gains an optional `semantic_search` tool: meaning-based search + over the agent's raw messages (text plus its own think/journal/skip_reply + notes) and every compression summary, backed by a shared remote + embed-service whose vector index lives server-side, one namespace per store + (`new HistoryModule({ semantic: { url, token, namespace } })`). The module + keeps the index in sync itself — a background tick every 60 s and a bounded + catch-up before each search, messages watermarked by timestamp with an + overlap re-scan and summaries by `createdMs` — and backs off cleanly when + the service is unreachable, so the other four history tools are unaffected. diff --git a/src/modules/history/index.ts b/src/modules/history/index.ts index 8ea3612..dca2812 100644 --- a/src/modules/history/index.ts +++ b/src/modules/history/index.ts @@ -49,6 +49,7 @@ import type { ToolDefinition, ToolCall, ToolResult, ProcessEvent } from '../../t import type { EventResponse, ProcessState } from '../../types/module.js'; import type { SearchWorkerMessage, SearchWorkerMatch } from './search-regex-worker.js'; import type { ChannelRegistry } from '../../mcpl/channel-registry.js'; +import { SemanticIndexClient, SemanticIndexer, type SemanticIndexConfig, type SyncReport } from './semantic.js'; // ============================================================================ // Tool input shapes @@ -88,6 +89,22 @@ interface OverviewInput { limit?: number; } +interface SemanticSearchInput { + query: string; + limit?: number; + from?: string; + to?: string; + channelId?: string; + kinds?: 'messages' | 'summaries' | 'both'; + level?: number; + minScore?: number; +} + +export interface HistoryModuleOptions { + /** Enable `semantic_search` against a shared embed-service (see ./semantic.ts). Absent = tool not offered. */ + semantic?: SemanticIndexConfig; +} + // ============================================================================ // Limits // ============================================================================ @@ -126,6 +143,10 @@ const SEARCH_MAX_MAX_SCAN = 50000; const OVERVIEW_DEFAULT_LIMIT = 50; const OVERVIEW_MAX_LIMIT = 200; +const SEMANTIC_DEFAULT_LIMIT = 10; +const SEMANTIC_MAX_LIMIT = 50; +const SEMANTIC_SNIPPET_CHARS = 400; + /** * A `<@id>`/`<@!id>` Discord-style mention — same shape `ChannelRegistry` * uses for its own DM mention-form matching. A string matching this can @@ -163,6 +184,15 @@ export class HistoryModule implements Module { private ctx: ModuleContext | null = null; private cm: ContextManager | null = null; private channelRegistry: ChannelRegistry | null = null; + private readonly semanticCfg: SemanticIndexConfig | null; + private semanticClient: SemanticIndexClient | null = null; + private indexer: SemanticIndexer | null = null; + private syncTimer: ReturnType | null = null; + + constructor(options: HistoryModuleOptions = {}) { + this.semanticCfg = options.semantic ?? null; + if (this.semanticCfg) this.semanticClient = new SemanticIndexClient(this.semanticCfg); + } /** * Wire the context-manager instance, and optionally the host's @@ -180,6 +210,10 @@ export class HistoryModule implements Module { bind(contextManager: ContextManager, channelRegistry?: ChannelRegistry): void { this.cm = contextManager; this.channelRegistry = channelRegistry ?? null; + if (this.semanticCfg && this.semanticClient) { + this.indexer = new SemanticIndexer(contextManager, this.semanticClient, this.semanticCfg, (m) => console.warn(m)); + this.startSyncTimer(); + } } /** @@ -260,10 +294,34 @@ export class HistoryModule implements Module { async start(ctx: ModuleContext): Promise { this.ctx = ctx; + this.startSyncTimer(); } async stop(): Promise { this.ctx = null; + if (this.syncTimer) { clearInterval(this.syncTimer); this.syncTimer = null; } + } + + /** + * Background sync into the semantic index: every `syncIntervalMs` (default + * 60 s) push up to `maxSyncPerTick` new items. Idempotent to call — starts + * once we have both a context-manager (bind) and a config; `unref`'d so it + * never keeps a shutting-down process alive. A first tick runs after 5 s so + * a fresh store starts backfilling immediately rather than a minute later. + */ + private startSyncTimer(): void { + if (this.syncTimer || !this.indexer || !this.semanticCfg) return; + const interval = this.semanticCfg.syncIntervalMs ?? 60_000; + if (interval <= 0) return; + const perTick = this.semanticCfg.maxSyncPerTick ?? 1024; + const tick = (): void => { void this.indexer?.catchUp(perTick); }; + const first = setTimeout(tick, 5_000); first.unref?.(); + this.syncTimer = setInterval(tick, interval); this.syncTimer.unref?.(); + } + + /** Exposed for hosts/tests: one sync pass now. */ + syncSemanticIndex(maxItems = 1024): Promise | null { + return this.indexer ? this.indexer.catchUp(maxItems) : null; } getTools(): ToolDefinition[] { @@ -373,9 +431,40 @@ export class HistoryModule implements Module { }, }, }, + ...(this.semanticCfg ? [this.semanticSearchTool()] : []), ]; } + private semanticSearchTool(): ToolDefinition { + return { + name: 'semantic_search', + description: + 'Search your own history by MEANING, not exact words: an embedding index over your raw messages ' + + '(text plus your own think/journal/skip_reply notes) and every compression summary. Use it when you ' + + 'remember roughly what something was about but not the words — "the night the fluid sim was read ' + + 'back to me as art" — then narrow with `from`/`to`/`channelId` and drill into the exact span with ' + + '`extract` or `overview`. Results are ranked by cosine similarity (score ~0.6+ is a strong match, ' + + '~0.3 is thematic, below ~0.2 is noise); each hit carries its id (`msg:` or `sum:`), ' + + 'timestamp, channel, kind/level and a snippet. The index catches up with recent messages before ' + + 'searching (bounded, so a huge backlog is reported as `index.behind` rather than blocking). ' + + 'Purely a read: nothing is written to your history.', + inputSchema: { + type: 'object' as const, + properties: { + query: { type: 'string', description: 'What you are looking for, in natural language. A sentence works better than keywords.' }, + limit: { type: 'number', description: `Max hits (default ${SEMANTIC_DEFAULT_LIMIT}, cap ${SEMANTIC_MAX_LIMIT}).` }, + from: { type: 'string', description: 'ISO 8601 inclusive lower bound on the message timestamp / summary span start. Omit for open-ended.' }, + to: { type: 'string', description: 'ISO 8601 inclusive upper bound. Omit for open-ended.' }, + channelId: { type: 'string', description: 'Only raw messages from this channel (label like "#general" or the raw internal id). Summaries are not per-channel and are excluded when this is set.' }, + kinds: { type: 'string', enum: ['messages', 'summaries', 'both'], description: 'What to search: raw messages, compression summaries, or both (default).' }, + level: { type: 'number', description: 'Only summaries of this exact level (implies kinds=summaries).' }, + minScore: { type: 'number', description: 'Drop hits below this cosine score (0..1). Default none.' }, + }, + required: ['query'], + }, + }; + } + async handleToolCall(call: ToolCall): Promise { try { if (!this.cm) { @@ -390,6 +479,8 @@ export class HistoryModule implements Module { return await this.handleSearch((call.input ?? {}) as SearchInput); case 'overview': return this.handleOverview((call.input ?? {}) as OverviewInput); + case 'semantic_search': + return await this.handleSemanticSearch((call.input ?? {}) as SemanticSearchInput); default: return { success: false, isError: true, error: `Unknown tool: ${call.name}` }; } @@ -649,6 +740,64 @@ export class HistoryModule implements Module { // overview // ========================================================================== + private async handleSemanticSearch(input: SemanticSearchInput): Promise { + if (!this.semanticCfg || !this.semanticClient) { + throw new Error('semantic_search is not configured for this resident (no embed-service in the recipe).'); + } + if (typeof input.query !== 'string' || !input.query.trim()) throw new Error('query must be a non-empty string'); + const limit = clampCount(input.limit, SEMANTIC_DEFAULT_LIMIT, SEMANTIC_MAX_LIMIT, 'limit'); + if (limit < 1) throw new Error('limit must be at least 1'); + const fromMs = parseIsoDate(input.from, 'from'); + const toMs = parseIsoDate(input.to, 'to'); + if (fromMs !== undefined && toMs !== undefined && fromMs > toMs) throw new Error('`from` must not be after `to`'); + const channelId = this.resolveChannel(input.channelId); + let kinds: string[] | undefined; + if (input.level !== undefined || input.kinds === 'summaries') kinds = ['summary']; + else if (input.kinds === 'messages' || channelId) kinds = ['message']; + if (input.level !== undefined && (!Number.isInteger(input.level) || input.level < 0)) throw new Error('level must be a non-negative integer'); + if (input.minScore !== undefined && (typeof input.minScore !== 'number' || input.minScore < -1 || input.minScore > 1)) { + throw new Error('minScore must be a number in -1..1'); + } + + // Bounded catch-up so the newest messages are searchable; never block on a backlog. + let sync: SyncReport | null = null; + if (this.indexer) { + try { sync = await this.indexer.catchUp(this.semanticCfg.maxSyncBeforeSearch ?? 256); } catch { /* reported via indexer.lastError */ } + } + + const res = await this.semanticClient.search({ + query: input.query, k: limit, + ts_from: fromMs === undefined ? undefined : fromMs / 1000, + ts_to: toMs === undefined ? undefined : toMs / 1000, + channel: channelId, kinds, level: input.level, min_score: input.minScore, snippet: SEMANTIC_SNIPPET_CHARS, + }); + const hits = res.hits.map((h) => ({ + id: h.id, + kind: h.kind, + level: h.level, + score: h.score, + timestamp: h.ts === null ? null : new Date(h.ts * 1000).toISOString(), + channelId: h.channel, + participant: (h.meta as { participant?: unknown }).participant ?? null, + author: (h.meta as { author?: unknown }).author ?? null, + snippet: h.text ?? '', + chars: h.chars, + })); + return { + success: true, + data: { + hits, + index: { + indexed: res.count_indexed, + behind: sync ? sync.more : this.indexer === null ? null : true, + syncedThisCall: sync ? sync.pushed : 0, + lastError: this.indexer?.lastError ?? null, + }, + timingMs: res.timing_ms, + }, + }; + } + private handleOverview(input: OverviewInput): ToolResult { const channelId = this.resolveChannel(input.channelId); const fromMs = parseIsoDate(input.from, 'from'); diff --git a/src/modules/history/semantic.ts b/src/modules/history/semantic.ts new file mode 100644 index 0000000..ae4f205 --- /dev/null +++ b/src/modules/history/semantic.ts @@ -0,0 +1,329 @@ +/** + * Semantic (embedding) search for HistoryModule, backed by a shared remote + * embed-service (one per fleet; the index lives server-side, keyed by a + * per-store namespace). + * + * Two pieces: + * - `SemanticIndexClient` — thin HTTP client for the service's index API + * (`/v1/index/{ns}/upsert|search|stats`). + * - `SemanticIndexer` — incremental sync from the resident's chronicle into + * that namespace. Messages are watermarked by TIMESTAMP (chronicle's native + * time index is the only cheap "what's new" query), with a re-scan overlap + * so a message stamped slightly out of order is still picked up; the + * service dedups by id + text hash, so re-sending is free. Summaries are + * watermarked by `createdMs` (a fresh L3 spans months of old timestamps — + * its creation time is the only monotonic signal). + * + * What gets embedded per message: text blocks verbatim, plus the string + * arguments of the agent's private-prose tools (think / journal / skip_reply + * / private_note) labelled `[think] …` — that is the resident's own diary and + * is exactly what "what did I think about X" should find. tool_result + * payloads, other tool arguments, thinking blocks and media are skipped. + * + * Failure posture: the service being down never breaks the module. Sync ticks + * log and back off; `semantic_search` returns a clean tool error. + */ + +import type { ContextManager, StoredMessage } from '@animalabs/context-manager'; +import type { ContentBlock } from '@animalabs/membrane'; + +export interface SemanticIndexConfig { + /** Base URL of the embed-service, e.g. `http://100.90.161.34:8804`. */ + url: string; + /** Bearer token (the service's EMBED_TOKEN). Optional if the service runs open. */ + token?: string; + /** Index namespace for this store — must be unique fleet-wide (e.g. `linn/9f9857cd`). */ + namespace: string; + /** Per-request timeout. Default 30 s (bulk upserts of long messages can take a while). */ + requestTimeoutMs?: number; + /** Background sync cadence. Default 60 s. 0 disables background sync (search still catches up). */ + syncIntervalMs?: number; + /** Items per upsert request. Default 128, max 256 (service limit). */ + syncBatch?: number; + /** Max items one background tick will push. Default 1024. */ + maxSyncPerTick?: number; + /** Max items a pre-search catch-up will push before searching anyway. Default 256. */ + maxSyncBeforeSearch?: number; + /** Re-scan window behind the message watermark, ms. Default 10 min. */ + overlapMs?: number; + /** Include private-prose tool arguments (think/journal/skip_reply/private_note). Default true. */ + includePrivateTools?: boolean; + /** Per-item text cap in chars (service caps at 200k; model truncates at its max_seq_length). Default 32k. */ + maxChars?: number; +} + +export interface IndexItem { + id: string; + text: string; + ts?: number; + channel?: string | null; + kind: 'message' | 'summary'; + level?: number; + cursor?: number; + meta?: Record; +} + +export interface IndexStats { + namespace: string; + model: string; + dim: number; + count: number; + by_kind: Record; +} + +export interface SearchHit { + id: string; + score: number; + ts: number | null; + channel: string | null; + kind: string; + level: number | null; + meta: Record; + text?: string; + chars: number; +} + +export interface SearchRequest { + query: string; + k?: number; + ts_from?: number; + ts_to?: number; + channel?: string; + kinds?: string[]; + level?: number; + min_score?: number; + snippet?: number; +} + +export interface SearchResponse { + namespace: string; + hits: SearchHit[]; + count_indexed: number; + timing_ms: { embed: number; search: number }; +} + +const PRIVATE_PROSE_TOOLS = new Set(['think', 'journal', 'skip_reply', 'private_note']); + +export class SemanticIndexClient { + private readonly base: string; + private readonly ns: string; + private readonly timeoutMs: number; + constructor(private readonly cfg: SemanticIndexConfig) { + this.base = cfg.url.replace(/\/+$/, ''); + this.ns = encodeURIComponent(cfg.namespace); + this.timeoutMs = cfg.requestTimeoutMs ?? 30_000; + } + + private async call(method: 'GET' | 'POST', path: string, body?: unknown): Promise { + const headers: Record = { 'content-type': 'application/json' }; + if (this.cfg.token) headers.authorization = `Bearer ${this.cfg.token}`; + const ctrl = new AbortController(); + const timer = setTimeout(() => ctrl.abort(), this.timeoutMs); + try { + const res = await fetch(`${this.base}${path}`, { + method, headers, body: body === undefined ? undefined : JSON.stringify(body), signal: ctrl.signal, + }); + const text = await res.text(); + let json: unknown; + try { json = JSON.parse(text); } catch { json = undefined; } + if (!res.ok) { + const msg = (json as { error?: { message?: string } } | undefined)?.error?.message ?? text.slice(0, 200); + throw new Error(`embed-service ${method} ${path} → ${res.status}: ${msg}`); + } + return json as T; + } finally { + clearTimeout(timer); + } + } + + /** Stats for this namespace; `null` when the namespace does not exist yet. */ + async stats(): Promise { + try { + return await this.call('GET', `/v1/index/${this.ns}/stats`); + } catch (e) { + if (e instanceof Error && /→ 404/.test(e.message)) return null; + throw e; + } + } + + async upsert(items: IndexItem[]): Promise<{ inserted: number; updated: number; unchanged: number; count: number }> { + return this.call('POST', `/v1/index/${this.ns}/upsert`, { items }); + } + + async search(req: SearchRequest): Promise { + return this.call('POST', `/v1/index/${this.ns}/search`, req); + } +} + +/** Text to embed for one message, or '' when there is nothing worth indexing. */ +export function messageIndexText(msg: StoredMessage, includePrivateTools = true): string { + const parts: string[] = []; + for (const block of msg.content as ContentBlock[]) { + if (block.type === 'text' && typeof block.text === 'string') { + parts.push(block.text); + } else if (includePrivateTools && block.type === 'tool_use' && PRIVATE_PROSE_TOOLS.has(block.name)) { + const input = (block as { input?: Record }).input ?? {}; + for (const v of Object.values(input)) { + if (typeof v === 'string' && v.trim().length > 20) parts.push(`[${block.name}] ${v}`); + } + } + } + return parts.join('\n').trim(); +} + +function channelOf(msg: StoredMessage): string | null { + const ext = (msg.metadata as { external?: { channelId?: unknown } } | undefined)?.external; + return typeof ext?.channelId === 'string' ? ext.channelId : null; +} + +function authorOf(msg: StoredMessage): string | null { + const md = msg.metadata as { external?: { authorName?: unknown }; authorName?: unknown } | undefined; + const a = md?.external?.authorName ?? md?.authorName; + return typeof a === 'string' ? a : null; +} + +export function messageToItem(msg: StoredMessage, cfg: { includePrivateTools?: boolean; maxChars?: number }): IndexItem | null { + const text = messageIndexText(msg, cfg.includePrivateTools ?? true); + if (!text) return null; + const tsMs = msg.timestamp.getTime(); + return { + id: `msg:${String(msg.id)}`, + text: text.slice(0, cfg.maxChars ?? 32_000), + ts: tsMs / 1000, + channel: channelOf(msg), + kind: 'message', + cursor: tsMs, + meta: { participant: msg.participant, author: authorOf(msg), seq: msg.sequence }, + }; +} + +export interface SummaryLike { + id: string; + level: number; + content: string; + tokens: number; + startMs: number; + endMs: number; + firstSequence: number; + lastSequence: number; + createdMs: number; + parentId?: string; +} + +export function summaryToItem(s: SummaryLike, cfg: { maxChars?: number }): IndexItem | null { + const text = s.content.trim(); + if (!text) return null; + return { + id: `sum:${s.id}`, + text: text.slice(0, cfg.maxChars ?? 32_000), + ts: s.startMs / 1000, + kind: 'summary', + level: s.level, + cursor: s.createdMs, + meta: { endTs: s.endMs / 1000, firstSequence: s.firstSequence, lastSequence: s.lastSequence, tokens: s.tokens, parentId: s.parentId ?? null }, + }; +} + +export interface SyncReport { + pushed: number; + inserted: number; + updated: number; + unchanged: number; + /** True when the tick hit its item cap before reaching the end of the store. */ + more: boolean; + messagesScanned: number; + summariesScanned: number; +} + +export class SemanticIndexer { + private inFlight: Promise | null = null; + private consecutiveFailures = 0; + private backoffUntil = 0; + lastError: string | null = null; + lastSyncAt: number | null = null; + + constructor( + private readonly cm: ContextManager, + private readonly client: SemanticIndexClient, + private readonly cfg: SemanticIndexConfig, + private readonly log: (msg: string) => void = () => {}, + ) {} + + get backingOff(): boolean { return Date.now() < this.backoffUntil; } + + /** + * Push what the index is missing, up to `maxItems`. Coalesces: a call while + * a sync is already running returns that run's promise instead of racing it. + */ + catchUp(maxItems: number): Promise { + if (this.inFlight) return this.inFlight; + const run = this.runCatchUp(maxItems).finally(() => { this.inFlight = null; }); + this.inFlight = run; + return run; + } + + private async runCatchUp(maxItems: number): Promise { + const report: SyncReport = { pushed: 0, inserted: 0, updated: 0, unchanged: 0, more: false, messagesScanned: 0, summariesScanned: 0 }; + if (this.backingOff) { report.more = true; return report; } + const batchSize = Math.min(256, Math.max(1, this.cfg.syncBatch ?? 128)); + try { + const stats = await this.client.stats(); + const msgWm = stats?.by_kind.message?.max_cursor ?? null; + const sumWm = stats?.by_kind.summary?.max_cursor ?? null; + let budget = maxItems; + let pending: IndexItem[] = []; + const flush = async (): Promise => { + if (pending.length === 0) return; + const r = await this.client.upsert(pending); + report.pushed += pending.length; report.inserted += r.inserted; report.updated += r.updated; report.unchanged += r.unchanged; + pending = []; + }; + + // Messages: walk forward from (watermark - overlap) via the time index. + let fromMs = msgWm === null ? undefined : Math.max(0, msgWm - (this.cfg.overlapMs ?? 600_000)); + const pageSize = 256; + for (;;) { + if (budget <= 0) { report.more = true; break; } + const page = this.cm.queryMessagesByTime({ fromMs, limit: pageSize }); + const msgs = page.messages; + report.messagesScanned += msgs.length; + for (const m of msgs) { + const item = messageToItem(m, this.cfg); + if (!item) continue; + pending.push(item); budget--; + if (pending.length >= batchSize) await flush(); + } + if (msgs.length < pageSize) break; + const lastMs = msgs[msgs.length - 1]!.timestamp.getTime(); + // Advance strictly: if a whole page shares one millisecond we would spin — step past it. + fromMs = lastMs === fromMs ? lastMs + 1 : lastMs; + } + await flush(); + + // Summaries: everything created after the summary watermark, any level. + if (!report.more) { + const all = this.cm.getSummariesInRange({ fromMs: 0, toMs: Number.MAX_SAFE_INTEGER }) as SummaryLike[]; + report.summariesScanned = all.length; + const fresh = all.filter((s) => sumWm === null || s.createdMs > sumWm).sort((a, b) => a.createdMs - b.createdMs); + for (const s of fresh) { + if (budget <= 0) { report.more = true; break; } + const item = summaryToItem(s, this.cfg); + if (!item) continue; + pending.push(item); budget--; + if (pending.length >= batchSize) await flush(); + } + await flush(); + } + this.consecutiveFailures = 0; this.lastError = null; this.lastSyncAt = Date.now(); + return report; + } catch (e) { + this.consecutiveFailures++; + const delay = Math.min(15 * 60_000, 5_000 * 2 ** Math.min(8, this.consecutiveFailures - 1)); + this.backoffUntil = Date.now() + delay; + this.lastError = e instanceof Error ? e.message : String(e); + this.log(`[history/semantic] sync failed (${this.consecutiveFailures}×, retry in ${Math.round(delay / 1000)}s): ${this.lastError}`); + report.more = true; + return report; + } + } +} diff --git a/src/modules/index.ts b/src/modules/index.ts index 1205967..8d89627 100644 --- a/src/modules/index.ts +++ b/src/modules/index.ts @@ -8,7 +8,8 @@ export type { ApiEvent } from './api/index.js'; export { HealthModule } from './health/index.js'; export type { HealthModuleConfig } from './health/index.js'; -export { HistoryModule } from './history/index.js'; +export { HistoryModule, type HistoryModuleOptions } from './history/index.js'; +export type { SemanticIndexConfig } from './history/semantic.js'; export { WorkspaceModule, WorkspaceReadError } from './workspace/index.js'; export type { WorkspaceReadErrorCode, WorkspaceReadStage, WorkspaceDiskReadResult, ReadFileFromDiskOptions } from './workspace/index.js'; diff --git a/test/history-module-semantic.test.ts b/test/history-module-semantic.test.ts new file mode 100644 index 0000000..b844a85 --- /dev/null +++ b/test/history-module-semantic.test.ts @@ -0,0 +1,154 @@ +import { describe, it, before, after } from 'node:test'; +import assert from 'node:assert/strict'; +import { createServer, type Server } from 'node:http'; +import { HistoryModule } from '../src/modules/history/index.js'; +import { messageIndexText } from '../src/modules/history/semantic.js'; +import type { ContextManager, StoredMessage } from '@animalabs/context-manager'; +import type { ContentBlock } from '@animalabs/membrane'; + +// ---- fake embed-service: in-memory namespace, cursor-aware stats, trivial "search" ------------- +interface Item { id: string; text: string; ts?: number; channel?: string | null; kind: string; level?: number; cursor?: number; meta?: Record } +class FakeService { + items = new Map(); + upserts: Item[][] = []; + searches: Record[] = []; + failNext = 0; + server!: Server; url = ''; + async start(): Promise { + this.server = createServer((req, res) => { + let body = ''; + req.on('data', (c) => { body += c; }); + req.on('end', () => { + if (this.failNext > 0) { this.failNext--; res.writeHead(500); res.end(JSON.stringify({ error: { message: 'boom' } })); return; } + if (req.headers.authorization !== 'Bearer tok') { res.writeHead(401); res.end(JSON.stringify({ error: { message: 'unauthorized' } })); return; } + const m = /^\/v1\/index\/([^/]+)\/(stats|upsert|search)$/.exec(req.url ?? ''); + if (!m) { res.writeHead(404); res.end('{}'); return; } + assert.equal(decodeURIComponent(m[1]!), 'test/ns'); + const json = (code: number, o: unknown): void => { res.writeHead(code, { 'content-type': 'application/json' }); res.end(JSON.stringify(o)); }; + if (m[2] === 'stats') { + if (this.items.size === 0) return json(404, { error: { message: 'no such namespace' } }); + const by: Record = {}; + for (const it of this.items.values()) { + const d = by[it.kind] ??= { count: 0, max_ts: null, min_ts: null, max_cursor: null }; + d.count++; if (it.cursor !== undefined) d.max_cursor = d.max_cursor === null ? it.cursor : Math.max(d.max_cursor, it.cursor); + } + return json(200, { namespace: 'test/ns', model: 'fake', dim: 4, count: this.items.size, by_kind: by }); + } + if (m[2] === 'upsert') { + const items = (JSON.parse(body) as { items: Item[] }).items; this.upserts.push(items); + let inserted = 0, unchanged = 0; + for (const it of items) { if (this.items.has(it.id) && this.items.get(it.id)!.text === it.text) unchanged++; else inserted++; this.items.set(it.id, it); } + return json(200, { namespace: 'test/ns', inserted, updated: 0, unchanged, count: this.items.size }); + } + const q = JSON.parse(body) as Record; this.searches.push(q); + const hits = [...this.items.values()].filter((it) => !q.kinds || (q.kinds as string[]).includes(it.kind)) + .filter((it) => !q.channel || it.channel === q.channel) + .map((it) => ({ id: it.id, score: it.text.includes(q.query as string) ? 0.9 : 0.1, ts: it.ts ?? null, channel: it.channel ?? null, kind: it.kind, level: it.level ?? null, meta: it.meta ?? {}, text: it.text.slice(0, 40), chars: it.text.length })) + .sort((a, b) => b.score - a.score).slice(0, (q.k as number) ?? 10); + return json(200, { namespace: 'test/ns', hits, count_indexed: this.items.size, timing_ms: { embed: 1, search: 1 } }); + }); + }); + await new Promise((r) => this.server.listen(0, '127.0.0.1', r)); + const a = this.server.address(); this.url = `http://127.0.0.1:${typeof a === 'object' && a ? a.port : 0}`; + } + stop(): Promise { return new Promise((r) => this.server.close(() => r())); } +} + +// ---- stub context-manager ---------------------------------------------------------------------- +function msg(id: string, ms: number, content: ContentBlock[], channelId?: string): StoredMessage { + return { id, sequence: Number(id.replace(/\D/g, '')), participant: 'Linn', content, timestamp: new Date(ms), + metadata: channelId ? { external: { source: 'discord', channelId, authorName: 'Linn' } } : undefined } as unknown as StoredMessage; +} +function stubCm(messages: StoredMessage[], summaries: Array> = []): ContextManager { + return { + queryMessagesByTime(o: { fromMs?: number; toMs?: number; limit?: number }) { + const all = messages.filter((m) => (o.fromMs === undefined || m.timestamp.getTime() >= o.fromMs) && (o.toMs === undefined || m.timestamp.getTime() <= o.toMs)) + .sort((a, b) => a.timestamp.getTime() - b.timestamp.getTime()); + return { messages: all.slice(0, o.limit ?? all.length), totalCount: all.length }; + }, + getSummariesInRange() { return summaries; }, + } as unknown as ContextManager; +} +const T0 = Date.UTC(2026, 8, 22); +const messages = [ + msg('m1', T0 + 1000, [{ type: 'text', text: 'the night the fluid sim was read back to me as art' }], 'chan-A'), + msg('m2', T0 + 2000, [{ type: 'tool_use', id: 'x', name: 'think', input: { content: 'a private note about vortices and the second small body' } }]), + msg('m3', T0 + 3000, [{ type: 'tool_result', toolUseId: 'x', content: 'ignored payload' }]), + msg('m4', T0 + 4000, [{ type: 'text', text: 'goodnight, crow' }], 'chan-B'), +]; +const summaries = [ + { id: 's1', level: 1, content: 'Summary of the evening', tokens: 10, startMs: T0, endMs: T0 + 5000, firstSequence: 1, lastSequence: 4, createdMs: T0 + 9000 }, +]; + +describe('HistoryModule semantic_search', () => { + const svc = new FakeService(); + before(() => svc.start()); + after(() => svc.stop()); + + it('offers the tool only when configured', () => { + assert.ok(!new HistoryModule().getTools().some((t) => t.name === 'semantic_search')); + const mod = new HistoryModule({ semantic: { url: 'http://x', namespace: 'n', syncIntervalMs: 0 } }); + assert.ok(mod.getTools().some((t) => t.name === 'semantic_search')); + }); + + it('messageIndexText keeps text + private-prose tool args, drops tool_result', () => { + assert.equal(messageIndexText(messages[0]!), 'the night the fluid sim was read back to me as art'); + assert.equal(messageIndexText(messages[1]!), '[think] a private note about vortices and the second small body'); + assert.equal(messageIndexText(messages[2]!), ''); + assert.equal(messageIndexText(messages[1]!, false), ''); + }); + + it('syncs messages and summaries incrementally with cursors, then searches with mapped filters', async () => { + const mod = new HistoryModule({ semantic: { url: svc.url, token: 'tok', namespace: 'test/ns', syncIntervalMs: 0, syncBatch: 2 } }); + mod.bind(stubCm(messages, summaries)); + const r1 = await mod.syncSemanticIndex()!; + assert.deepEqual({ pushed: r1.pushed, inserted: r1.inserted, more: r1.more }, { pushed: 4, inserted: 4, more: false }); + assert.deepEqual([...svc.items.keys()].sort(), ['msg:m1', 'msg:m2', 'msg:m4', 'sum:s1']); + assert.equal(svc.items.get('msg:m1')!.cursor, T0 + 1000); + assert.equal(svc.items.get('sum:s1')!.cursor, T0 + 9000); + assert.equal(svc.items.get('msg:m1')!.channel, 'chan-A'); + // second pass: overlap re-sends the recent messages but the service reports them unchanged; the summary is not re-sent + const r2 = await mod.syncSemanticIndex()!; + assert.equal(r2.inserted, 0); + assert.ok(!svc.upserts.at(-1)!.some((i) => i.id === 'sum:s1') || r2.pushed === 0); + + const res = await mod.handleToolCall({ id: 'c1', name: 'semantic_search', input: { query: 'fluid sim', from: '2026-09-22T00:00:00Z', kinds: 'messages', limit: 5 } }); + assert.equal(res.success, true, JSON.stringify(res)); + const data = res.data as { hits: Array<{ id: string; score: number; timestamp: string; channelId: string | null }>; index: { indexed: number; behind: boolean } }; + assert.equal(data.hits[0]!.id, 'msg:m1'); + assert.equal(data.hits[0]!.channelId, 'chan-A'); + assert.equal(data.hits[0]!.timestamp, new Date(T0 + 1000).toISOString()); + assert.equal(data.index.behind, false); + const q = svc.searches.at(-1)!; + assert.deepEqual(q.kinds, ['message']); + assert.equal(q.ts_from, Date.UTC(2026, 8, 22) / 1000); + assert.equal(q.k, 5); + // level implies summaries + await mod.handleToolCall({ id: 'c2', name: 'semantic_search', input: { query: 'evening', level: 1 } }); + assert.deepEqual(svc.searches.at(-1)!.kinds, ['summary']); + await mod.stop(); + }); + + it('service failure is a clean tool error and backs off sync', async () => { + const mod = new HistoryModule({ semantic: { url: svc.url, token: 'tok', namespace: 'test/ns', syncIntervalMs: 0 } }); + mod.bind(stubCm(messages, summaries)); + svc.failNext = 3; + const r = await mod.syncSemanticIndex()!; + assert.equal(r.more, true); + const res = await mod.handleToolCall({ id: 'c3', name: 'semantic_search', input: { query: 'anything' } }); + assert.equal(res.success, false); + assert.match(String(res.error), /embed-service|boom/); + await mod.stop(); + }); + + it('validates input', async () => { + const mod = new HistoryModule({ semantic: { url: svc.url, token: 'tok', namespace: 'test/ns', syncIntervalMs: 0 } }); + mod.bind(stubCm([])); + for (const input of [{}, { query: '' }, { query: 'x', from: 'nope' }, { query: 'x', from: '2026-02-01T00:00:00Z', to: '2026-01-01T00:00:00Z' }, { query: 'x', limit: 0 }, { query: 'x', minScore: 7 }]) { + const res = await mod.handleToolCall({ id: 'v', name: 'semantic_search', input }); + assert.equal(res.success, false, JSON.stringify(input)); + } + const unconfigured = await new HistoryModule().handleToolCall({ id: 'u', name: 'semantic_search', input: { query: 'x' } }); + assert.equal(unconfigured.success, false); + }); +}); From 9045646ef0605bcd21be8c851fff54b85677da05 Mon Sep 17 00:00:00 2001 From: slimepriestess Date: Wed, 23 Sep 2026 10:12:34 -0700 Subject: [PATCH 2/5] fix(history/semantic): overlap re-sends do not spend the sync budget; stop() cancels the first tick Each sync tick re-walks the overlap window behind the service's watermark so slightly out-of-order messages are still picked up; the service dedups the re-sends. Those re-sends also counted against the tick's item budget, so whenever the overlap window held more indexable messages than the budget (> 256 in 10 min for the pre-search catch-up, > 1024 for the background tick) every tick spent the whole budget re-sending indexed items, reported more: true, and never reached the first new message: the index froze at the watermark. Only items past the watermark spend budget now. The 5 s first-tick timer was not tracked, so a module stopped before it fired still called stats + upsert afterwards. stop() clears it. Regression: test/history-semantic-sync-converges.test.ts (both red on 11f452e). Co-Authored-By: Claude Fable 5.1 Claude-Session: https://claude.ai/code/session_01CoQK2cP55YhezE6ajSx58h --- src/modules/history/index.ts | 5 +- src/modules/history/semantic.ts | 9 +- test/history-semantic-sync-converges.test.ts | 125 +++++++++++++++++++ 3 files changed, 137 insertions(+), 2 deletions(-) create mode 100644 test/history-semantic-sync-converges.test.ts diff --git a/src/modules/history/index.ts b/src/modules/history/index.ts index dca2812..62104af 100644 --- a/src/modules/history/index.ts +++ b/src/modules/history/index.ts @@ -188,6 +188,7 @@ export class HistoryModule implements Module { private semanticClient: SemanticIndexClient | null = null; private indexer: SemanticIndexer | null = null; private syncTimer: ReturnType | null = null; + private firstSyncTimer: ReturnType | null = null; constructor(options: HistoryModuleOptions = {}) { this.semanticCfg = options.semantic ?? null; @@ -300,6 +301,7 @@ export class HistoryModule implements Module { async stop(): Promise { this.ctx = null; if (this.syncTimer) { clearInterval(this.syncTimer); this.syncTimer = null; } + if (this.firstSyncTimer) { clearTimeout(this.firstSyncTimer); this.firstSyncTimer = null; } } /** @@ -315,7 +317,8 @@ export class HistoryModule implements Module { if (interval <= 0) return; const perTick = this.semanticCfg.maxSyncPerTick ?? 1024; const tick = (): void => { void this.indexer?.catchUp(perTick); }; - const first = setTimeout(tick, 5_000); first.unref?.(); + this.firstSyncTimer = setTimeout(() => { this.firstSyncTimer = null; tick(); }, 5_000); + this.firstSyncTimer.unref?.(); this.syncTimer = setInterval(tick, interval); this.syncTimer.unref?.(); } diff --git a/src/modules/history/semantic.ts b/src/modules/history/semantic.ts index ae4f205..d7e1319 100644 --- a/src/modules/history/semantic.ts +++ b/src/modules/history/semantic.ts @@ -290,7 +290,14 @@ export class SemanticIndexer { for (const m of msgs) { const item = messageToItem(m, this.cfg); if (!item) continue; - pending.push(item); budget--; + pending.push(item); + // Only NEW items spend the budget. The overlap re-walk behind the + // watermark is a dedup no-op for the service, and it must be one + // for the budget too: when the overlap window held more messages + // than a tick's budget, every tick spent it all re-sending the same + // indexed items and never reached the first new one — the index + // froze at the watermark while reporting merely `more: true`. + if (msgWm === null || item.cursor === undefined || item.cursor > msgWm) budget--; if (pending.length >= batchSize) await flush(); } if (msgs.length < pageSize) break; diff --git a/test/history-semantic-sync-converges.test.ts b/test/history-semantic-sync-converges.test.ts new file mode 100644 index 0000000..ad480d9 --- /dev/null +++ b/test/history-semantic-sync-converges.test.ts @@ -0,0 +1,125 @@ +import { describe, it, before, after } from 'node:test'; +import assert from 'node:assert/strict'; +import { createServer, type Server } from 'node:http'; +import { HistoryModule } from '../src/modules/history/index.js'; +import type { ContextManager, StoredMessage } from '@animalabs/context-manager'; +import type { ContentBlock } from '@animalabs/membrane'; + +/** + * The incremental sync must CONVERGE: repeated bounded ticks over a store + * denser than the tick budget end with every message indexed. + * + * Each tick re-walks the overlap window behind the service's watermark + * (messages stamped slightly out of order are picked up that way; the + * service dedups the re-sends). Those re-sends are free for the service but + * they were not free for the tick's item budget: when the overlap window + * held more indexable messages than the budget, every tick spent the whole + * budget re-sending the same already-indexed items, reported `more: true`, + * and never reached the first new message — the index froze at the + * watermark while looking merely "behind". At 1 msg/s that is any 10-minute + * window with > 256 messages for the pre-search catch-up and > 1024 for the + * background tick. + */ + +interface Item { id: string; text: string; ts?: number; channel?: string | null; kind: string; level?: number; cursor?: number } + +/** Minimal embed-service: dedup by id, cursor-aware stats, counts requests. */ +class FakeService { + items = new Map(); + calls: string[] = []; + server!: Server; + url = ''; + async start(): Promise { + this.server = createServer((req, res) => { + let body = ''; + req.on('data', (c) => { body += c; }); + req.on('end', () => { + const m = /^\/v1\/index\/([^/]+)\/(stats|upsert|search)$/.exec(req.url ?? ''); + this.calls.push(m?.[2] ?? '?'); + const json = (code: number, o: unknown): void => { res.writeHead(code, { 'content-type': 'application/json' }); res.end(JSON.stringify(o)); }; + if (!m) return json(404, {}); + if (m[2] === 'stats') { + if (this.items.size === 0) return json(404, { error: { message: 'no such namespace' } }); + const by: Record = {}; + for (const it of this.items.values()) { + const d = by[it.kind] ??= { count: 0, max_ts: null, min_ts: null, max_cursor: null }; + d.count++; + if (it.cursor !== undefined) d.max_cursor = d.max_cursor === null ? it.cursor : Math.max(d.max_cursor, it.cursor); + } + return json(200, { namespace: 'n', model: 'fake', dim: 1, count: this.items.size, by_kind: by }); + } + if (m[2] === 'upsert') { + const items = (JSON.parse(body) as { items: Item[] }).items; + let inserted = 0, unchanged = 0; + for (const it of items) { if (this.items.has(it.id)) unchanged++; else inserted++; this.items.set(it.id, it); } + return json(200, { inserted, updated: 0, unchanged, count: this.items.size }); + } + return json(200, { namespace: 'n', hits: [], count_indexed: this.items.size, timing_ms: { embed: 0, search: 0 } }); + }); + }); + await new Promise((r) => this.server.listen(0, '127.0.0.1', r)); + const a = this.server.address(); + this.url = `http://127.0.0.1:${typeof a === 'object' && a ? a.port : 0}`; + } + stop(): Promise { + this.server.closeAllConnections(); + return new Promise((r) => this.server.close(() => r())); + } +} + +function msg(id: string, ms: number, text: string): StoredMessage { + return { + id, sequence: Number(id.replace(/\D/g, '')), participant: 'Linn', + content: [{ type: 'text', text }] as ContentBlock[], timestamp: new Date(ms), metadata: undefined, + } as unknown as StoredMessage; +} + +/** Stub CM with the real queryByTime contract: inclusive bounds, oldest first, `limit` = first N. */ +function stubCm(messages: StoredMessage[]): ContextManager { + return { + queryMessagesByTime(o: { fromMs?: number; toMs?: number; limit?: number }) { + const all = messages + .filter((m) => (o.fromMs === undefined || m.timestamp.getTime() >= o.fromMs) && (o.toMs === undefined || m.timestamp.getTime() <= o.toMs)) + .sort((a, b) => a.timestamp.getTime() - b.timestamp.getTime()); + return { messages: all.slice(0, o.limit ?? all.length), totalCount: all.length }; + }, + getSummariesInRange() { return []; }, + } as unknown as ContextManager; +} + +const T0 = Date.UTC(2026, 8, 22); + +describe('semantic sync converges', () => { + const svc = new FakeService(); + before(() => svc.start()); + after(() => svc.stop()); + + it('a store denser than the tick budget is fully indexed after enough ticks (overlap re-sends do not eat the budget)', async () => { + // 700 messages one second apart: the default 10-minute overlap window + // holds 600 of them, more than a 300-item tick. + const store = Array.from({ length: 700 }, (_, i) => msg(`m${i}`, T0 + i * 1000, `message number ${i} with enough text to index`)); + const mod = new HistoryModule({ semantic: { url: svc.url, namespace: 'n', syncIntervalMs: 0 } }); + mod.bind(stubCm(store)); + try { + const reports: Array<[number, boolean]> = []; + for (let tick = 0; tick < 12 && svc.items.size < 700; tick++) { + const r = await mod.syncSemanticIndex(300)!; + reports.push([r.pushed, r.more]); + } + assert.equal(svc.items.size, 700, `index froze at ${svc.items.size}; ticks (pushed, more): ${JSON.stringify(reports)}`); + const last = await mod.syncSemanticIndex(300)!; + assert.equal(last.more, false, 'a caught-up store reports more: false'); + } finally { + await mod.stop(); + } + }); + + it('stop() cancels the first-tick timer: nothing reaches the service after a stopped module', async () => { + const mod = new HistoryModule({ semantic: { url: svc.url, namespace: 'n', syncIntervalMs: 60_000 } }); + mod.bind(stubCm([msg('z1', T0, 'a message that would be pushed by the first tick')])); + await mod.stop(); + svc.calls.length = 0; + await new Promise((r) => setTimeout(r, 5_300)); + assert.deepEqual(svc.calls, [], `stopped module still called the service: ${JSON.stringify(svc.calls)}`); + }); +}); From 9a16605cdb828bb12a10bd999d9c399f65ece448 Mon Sep 17 00:00:00 2001 From: antra-tess Date: Fri, 25 Sep 2026 12:44:27 -0700 Subject: [PATCH 3/5] fix(history/semantic): reject channelId with a summaries-only search; note catch-up may wait on a running sync channelId + kinds="summaries" (or level) could only ever match nothing, since summaries carry no channel; it now throws a clear error instead of returning an empty result. level + kinds="messages" likewise. The tool description now says the bounded pre-search catch-up waits for an in-flight background sync when one is running. Co-Authored-By: Claude Opus 5.5 --- src/modules/history/index.ts | 16 +++++++++++++--- test/history-module-semantic.test.ts | 9 +++++++++ 2 files changed, 22 insertions(+), 3 deletions(-) diff --git a/src/modules/history/index.ts b/src/modules/history/index.ts index 62104af..c73b2b7 100644 --- a/src/modules/history/index.ts +++ b/src/modules/history/index.ts @@ -449,7 +449,8 @@ export class HistoryModule implements Module { '`extract` or `overview`. Results are ranked by cosine similarity (score ~0.6+ is a strong match, ' + '~0.3 is thematic, below ~0.2 is noise); each hit carries its id (`msg:` or `sum:`), ' + 'timestamp, channel, kind/level and a snippet. The index catches up with recent messages before ' + - 'searching (bounded, so a huge backlog is reported as `index.behind` rather than blocking). ' + + 'searching (bounded, so a huge backlog is reported as `index.behind` rather than blocking; if a ' + + 'background sync is already running, the search waits for that run to finish instead). ' + 'Purely a read: nothing is written to your history.', inputSchema: { type: 'object' as const, @@ -458,7 +459,7 @@ export class HistoryModule implements Module { limit: { type: 'number', description: `Max hits (default ${SEMANTIC_DEFAULT_LIMIT}, cap ${SEMANTIC_MAX_LIMIT}).` }, from: { type: 'string', description: 'ISO 8601 inclusive lower bound on the message timestamp / summary span start. Omit for open-ended.' }, to: { type: 'string', description: 'ISO 8601 inclusive upper bound. Omit for open-ended.' }, - channelId: { type: 'string', description: 'Only raw messages from this channel (label like "#general" or the raw internal id). Summaries are not per-channel and are excluded when this is set.' }, + channelId: { type: 'string', description: 'Only raw messages from this channel (label like "#general" or the raw internal id). Summaries are not per-channel and are excluded when this is set; combining it with kinds="summaries" or level is an error.' }, kinds: { type: 'string', enum: ['messages', 'summaries', 'both'], description: 'What to search: raw messages, compression summaries, or both (default).' }, level: { type: 'number', description: 'Only summaries of this exact level (implies kinds=summaries).' }, minScore: { type: 'number', description: 'Drop hits below this cosine score (0..1). Default none.' }, @@ -754,8 +755,17 @@ export class HistoryModule implements Module { const toMs = parseIsoDate(input.to, 'to'); if (fromMs !== undefined && toMs !== undefined && fromMs > toMs) throw new Error('`from` must not be after `to`'); const channelId = this.resolveChannel(input.channelId); + // Summaries carry no channel, so a channel filter on a summaries-only + // search can only ever match nothing. Say so instead of returning []. + const wantsSummariesOnly = input.level !== undefined || input.kinds === 'summaries'; + if (channelId && wantsSummariesOnly) { + throw new Error('channelId cannot be combined with kinds="summaries" or level: summaries are not per-channel. Drop channelId, or search kinds="messages".'); + } + if (input.level !== undefined && input.kinds === 'messages') { + throw new Error('level applies to summaries only and cannot be combined with kinds="messages".'); + } let kinds: string[] | undefined; - if (input.level !== undefined || input.kinds === 'summaries') kinds = ['summary']; + if (wantsSummariesOnly) kinds = ['summary']; else if (input.kinds === 'messages' || channelId) kinds = ['message']; if (input.level !== undefined && (!Number.isInteger(input.level) || input.level < 0)) throw new Error('level must be a non-negative integer'); if (input.minScore !== undefined && (typeof input.minScore !== 'number' || input.minScore < -1 || input.minScore > 1)) { diff --git a/test/history-module-semantic.test.ts b/test/history-module-semantic.test.ts index b844a85..0361565 100644 --- a/test/history-module-semantic.test.ts +++ b/test/history-module-semantic.test.ts @@ -148,6 +148,15 @@ describe('HistoryModule semantic_search', () => { const res = await mod.handleToolCall({ id: 'v', name: 'semantic_search', input }); assert.equal(res.success, false, JSON.stringify(input)); } + // contradictory filters are an error, not a silent empty result + const searchesBefore = svc.searches.length; + for (const input of [{ query: 'x', channelId: 'chan-A', kinds: 'summaries' }, { query: 'x', channelId: 'chan-A', level: 1 }, { query: 'x', kinds: 'messages', level: 1 }]) { + const res = await mod.handleToolCall({ id: 'v', name: 'semantic_search', input }); + assert.equal(res.success, false, JSON.stringify(input)); + assert.match(String(res.error), /summaries|level/); + } + assert.equal(svc.searches.length, searchesBefore); + await mod.stop(); const unconfigured = await new HistoryModule().handleToolCall({ id: 'u', name: 'semantic_search', input: { query: 'x' } }); assert.equal(unconfigured.success, false); }); From 8d8f38f306aed53a21b372623b0b2247aa9f9b78 Mon Sep 17 00:00:00 2001 From: slimepriestess Date: Fri, 25 Sep 2026 09:19:37 -0700 Subject: [PATCH 4/5] fix(history/semantic): drop hits that are not on the current branch MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The remote index is append-only and knows nothing about branches. /undo, /checkout, /restore, /branchto and /newtopic move the Chronicle head within one store; a message removed that way keeps its id (new messages get fresh ids) and stays indexed, so semantic_search kept returning it verbatim as a snippet. /undo is how an operator removes a poisoned or leaked turn, and with semantic search on it no longer did. Check every hit against the current branch before it goes back to the model: msg: through getMessage, sum: through getSummary. Both answer null for anything off the branch (MessageStore.lookupIndex rebuilds on branch switch). k is small, so this is at most `limit` local lookups. Dropped hits are counted in index.droppedOffBranch, and the tool description says so. Test: index m1..m4 + s1, undo m2, search again — msg:m2 no longer returned, droppedOffBranch 1; same for a summary. Red on 6b6c43c. Verified on a real Chronicle store that branchAt + switchBranch makes getMessage(undoneId) null while queryMessages still lists the survivor. Reported by Sol on connectome-host#144 (#3). Co-Authored-By: Claude Fable 5.1 Claude-Session: https://claude.ai/code/session_01CoQK2cP55YhezE6ajSx58h --- src/modules/history/index.ts | 26 +++++++++++++++++++-- test/history-module-semantic.test.ts | 34 +++++++++++++++++++++++++++- 2 files changed, 57 insertions(+), 3 deletions(-) diff --git a/src/modules/history/index.ts b/src/modules/history/index.ts index c73b2b7..1e4ed12 100644 --- a/src/modules/history/index.ts +++ b/src/modules/history/index.ts @@ -451,7 +451,9 @@ export class HistoryModule implements Module { 'timestamp, channel, kind/level and a snippet. The index catches up with recent messages before ' + 'searching (bounded, so a huge backlog is reported as `index.behind` rather than blocking; if a ' + 'background sync is already running, the search waits for that run to finish instead). ' + - 'Purely a read: nothing is written to your history.', + 'Purely a read: nothing is written to your history. Hits are checked against your current branch ' + + 'before they come back: a message you undid or a summary from a branch you left is dropped and ' + + 'counted in index.droppedOffBranch.', inputSchema: { type: 'object' as const, properties: { @@ -784,7 +786,25 @@ export class HistoryModule implements Module { ts_to: toMs === undefined ? undefined : toMs / 1000, channel: channelId, kinds, level: input.level, min_score: input.minScore, snippet: SEMANTIC_SNIPPET_CHARS, }); - const hits = res.hits.map((h) => ({ + // The remote index is append-only and knows nothing about branches: a + // message removed from the agent's reach by /undo, /checkout, /restore or + // /newtopic keeps its id (a new message gets a fresh one) and stays + // indexed. Check every hit against the CURRENT branch before it goes back + // to the model — getMessage/getSummary rebuild on branch switch, so an + // undone message or a summary minted on another branch answers null. + // k is small, so this is at most `limit` local lookups. + const cm = this.cm as ContextManager; + let droppedOffBranch = 0; + const onBranch = res.hits.filter((h) => { + const msgId = /^msg:(.+)$/.exec(h.id)?.[1]; + const sumId = /^sum:(.+)$/.exec(h.id)?.[1]; + const present = msgId !== undefined ? cm.getMessage(msgId) !== null + : sumId !== undefined ? cm.getSummary(sumId) !== null + : true; + if (!present) droppedOffBranch++; + return present; + }); + const hits = onBranch.map((h) => ({ id: h.id, kind: h.kind, level: h.level, @@ -805,6 +825,8 @@ export class HistoryModule implements Module { behind: sync ? sync.more : this.indexer === null ? null : true, syncedThisCall: sync ? sync.pushed : 0, lastError: this.indexer?.lastError ?? null, + /** Hits the index returned for messages/summaries not on the current branch (undone, checked out past). */ + droppedOffBranch, }, timingMs: res.timing_ms, }, diff --git a/test/history-module-semantic.test.ts b/test/history-module-semantic.test.ts index 0361565..a937b51 100644 --- a/test/history-module-semantic.test.ts +++ b/test/history-module-semantic.test.ts @@ -59,8 +59,12 @@ function msg(id: string, ms: number, content: ContentBlock[], channelId?: string return { id, sequence: Number(id.replace(/\D/g, '')), participant: 'Linn', content, timestamp: new Date(ms), metadata: channelId ? { external: { source: 'discord', channelId, authorName: 'Linn' } } : undefined } as unknown as StoredMessage; } -function stubCm(messages: StoredMessage[], summaries: Array> = []): ContextManager { +function stubCm(messages: StoredMessage[], summaries: Array> = [], offBranch: Set = new Set()): ContextManager { return { + // Branch-scoped like the real store: getMessage/getSummary answer null for + // anything not on the current branch (see MessageStore.lookupIndex). + getMessage(id: string) { return offBranch.has(id) ? null : (messages.find((m) => m.id === id) ?? null); }, + getSummary(id: string) { return offBranch.has(id) ? null : (summaries.find((x) => x.id === id) ?? null); }, queryMessagesByTime(o: { fromMs?: number; toMs?: number; limit?: number }) { const all = messages.filter((m) => (o.fromMs === undefined || m.timestamp.getTime() >= o.fromMs) && (o.toMs === undefined || m.timestamp.getTime() <= o.toMs)) .sort((a, b) => a.timestamp.getTime() - b.timestamp.getTime()); @@ -129,6 +133,34 @@ describe('HistoryModule semantic_search', () => { await mod.stop(); }); + it('drops hits for messages and summaries that are no longer on the current branch (/undo, /checkout)', async () => { + const svc2 = new FakeService(); await svc2.start(); + try { + const mod = new HistoryModule({ semantic: { url: svc2.url, token: 'tok', namespace: 'test/ns', syncIntervalMs: 0 } }); + const offBranch = new Set(); + mod.bind(stubCm(messages, summaries, offBranch)); + await mod.syncSemanticIndex()!; + assert.ok(svc2.items.has('msg:m2') && svc2.items.has('sum:s1')); + // Before the undo: the think note is a hit. + const before = await mod.handleToolCall({ id: 'b', name: 'semantic_search', input: { query: 'private note', kinds: 'messages' } }); + assert.equal(before.success, true, JSON.stringify(before)); + assert.ok((before.data as { hits: Array<{ id: string }> }).hits.some((h) => h.id === 'msg:m2')); + // /undo to m1: m2 leaves the branch but stays in the remote index (ids are never reused). + offBranch.add('m2'); + const after = await mod.handleToolCall({ id: 'a', name: 'semantic_search', input: { query: 'private note', kinds: 'messages' } }); + const data = after.data as { hits: Array<{ id: string }>; index: { droppedOffBranch: number } }; + assert.ok(!data.hits.some((h) => h.id === 'msg:m2'), JSON.stringify(data.hits)); + assert.equal(data.index.droppedOffBranch, 1); + // Same rule for summaries minted on a branch the agent has left. + offBranch.add('s1'); + const sum = await mod.handleToolCall({ id: 's', name: 'semantic_search', input: { query: 'evening', level: 1 } }); + const sdata = sum.data as { hits: Array<{ id: string }>; index: { droppedOffBranch: number } }; + assert.ok(!sdata.hits.some((h) => h.id === 'sum:s1'), JSON.stringify(sdata.hits)); + assert.equal(sdata.index.droppedOffBranch, 1); + await mod.stop(); + } finally { await svc2.stop(); } + }); + it('service failure is a clean tool error and backs off sync', async () => { const mod = new HistoryModule({ semantic: { url: svc.url, token: 'tok', namespace: 'test/ns', syncIntervalMs: 0 } }); mod.bind(stubCm(messages, summaries)); From 98ed6e1525faad250441bcbd612b6f4f60d8f29c Mon Sep 17 00:00:00 2001 From: antra-tess Date: Tue, 29 Sep 2026 13:03:32 -0700 Subject: [PATCH 5/5] fix(history/semantic): sync edits/removals, both channel shapes, same-ms pages, summaries first, short notes Addresses the Greptile review of #173: - channelOf reads metadata.channelId (MCPL ingestion) as well as metadata.external.channelId, matching getChannelId; channel-filtered searches no longer miss MCPL messages. - The message walk pages with offset within a millisecond instead of stepping past it, so a run of >256 messages sharing one timestamp is indexed rather than silently skipped. - Summaries sync before messages, so a message backlog that outlasts every tick's budget no longer keeps new summaries out of the index. - The indexer subscribes to CM onMessage: an edit re-upserts the message (the time walk never revisits a message older than the overlap), and a remove deletes it from the index via the service's delete endpoint. Events are held in memory only; removals are also hidden at search time by the current-branch check. - Private-prose tool args are indexed when non-empty (was > 20 chars), so a short think/journal/skip_reply note is searchable. Tests: test/history-semantic-review-findings.test.ts, six cases, all red before this change. Co-Authored-By: Claude Opus 5.5 --- src/modules/history/index.ts | 4 + src/modules/history/semantic.ts | 107 ++++++++--- test/history-semantic-review-findings.test.ts | 173 ++++++++++++++++++ 3 files changed, 263 insertions(+), 21 deletions(-) create mode 100644 test/history-semantic-review-findings.test.ts diff --git a/src/modules/history/index.ts b/src/modules/history/index.ts index 1e4ed12..d6daa2e 100644 --- a/src/modules/history/index.ts +++ b/src/modules/history/index.ts @@ -212,7 +212,9 @@ export class HistoryModule implements Module { this.cm = contextManager; this.channelRegistry = channelRegistry ?? null; if (this.semanticCfg && this.semanticClient) { + this.indexer?.dispose(); this.indexer = new SemanticIndexer(contextManager, this.semanticClient, this.semanticCfg, (m) => console.warn(m)); + this.indexer.attach(); this.startSyncTimer(); } } @@ -295,6 +297,7 @@ export class HistoryModule implements Module { async start(ctx: ModuleContext): Promise { this.ctx = ctx; + this.indexer?.attach(); this.startSyncTimer(); } @@ -302,6 +305,7 @@ export class HistoryModule implements Module { this.ctx = null; if (this.syncTimer) { clearInterval(this.syncTimer); this.syncTimer = null; } if (this.firstSyncTimer) { clearTimeout(this.firstSyncTimer); this.firstSyncTimer = null; } + this.indexer?.dispose(); } /** diff --git a/src/modules/history/semantic.ts b/src/modules/history/semantic.ts index d7e1319..d8829e9 100644 --- a/src/modules/history/semantic.ts +++ b/src/modules/history/semantic.ts @@ -150,6 +150,10 @@ export class SemanticIndexClient { return this.call('POST', `/v1/index/${this.ns}/upsert`, { items }); } + async delete(ids: string[]): Promise<{ deleted: number }> { + return this.call('POST', `/v1/index/${this.ns}/delete`, { ids }); + } + async search(req: SearchRequest): Promise { return this.call('POST', `/v1/index/${this.ns}/search`, req); } @@ -164,16 +168,18 @@ export function messageIndexText(msg: StoredMessage, includePrivateTools = true) } else if (includePrivateTools && block.type === 'tool_use' && PRIVATE_PROSE_TOOLS.has(block.name)) { const input = (block as { input?: Record }).input ?? {}; for (const v of Object.values(input)) { - if (typeof v === 'string' && v.trim().length > 20) parts.push(`[${block.name}] ${v}`); + if (typeof v === 'string' && v.trim().length > 0) parts.push(`[${block.name}] ${v}`); } } } return parts.join('\n').trim(); } +/** Same two metadata shapes HistoryModule's getChannelId reads: MCPL ingestion writes `metadata.channelId`, older paths `metadata.external.channelId`. */ function channelOf(msg: StoredMessage): string | null { - const ext = (msg.metadata as { external?: { channelId?: unknown } } | undefined)?.external; - return typeof ext?.channelId === 'string' ? ext.channelId : null; + const md = msg.metadata as { channelId?: unknown; external?: { channelId?: unknown } } | undefined; + if (typeof md?.channelId === 'string') return md.channelId; + return typeof md?.external?.channelId === 'string' ? md.external.channelId : null; } function authorOf(msg: StoredMessage): string | null { @@ -249,6 +255,32 @@ export class SemanticIndexer { private readonly log: (msg: string) => void = () => {}, ) {} + /** + * Edits and removals since the last successful sync. The time-cursor walk + * only revisits the overlap window, so an edit to an older message (its + * timestamp does not change) or a removal would otherwise never reach the + * index. In memory only: events that land while the module is detached or + * before a restart are not replayed (removals are still hidden at search + * time by the current-branch check in HistoryModule). + */ + private readonly editedIds = new Set(); + private readonly removedIds = new Set(); + private detach: (() => void) | null = null; + + /** Subscribe to message-store edits/removals. Idempotent; a CM without onMessage is tolerated. */ + attach(): void { + if (this.detach) return; + const on = (this.cm as { onMessage?: ContextManager['onMessage'] }).onMessage; + if (typeof on !== 'function') return; + this.detach = on.call(this.cm, (e) => { + if (e.type === 'edit') { this.editedIds.add(String(e.messageId)); } + else if (e.type === 'remove') { const id = String(e.messageId); this.editedIds.delete(id); this.removedIds.add(id); } + // removeRange carries only its endpoints; removed ids in it are hidden at search time instead. + }); + } + + dispose(): void { this.detach?.(); this.detach = null; } + get backingOff(): boolean { return Date.now() < this.backoffUntil; } /** @@ -279,12 +311,59 @@ export class SemanticIndexer { pending = []; }; + // Edits and removals first: they never come back round in the time walk. + if (this.removedIds.size > 0 || this.editedIds.size > 0) { + const removed = [...this.removedIds]; + const edited = [...this.editedIds]; + this.removedIds.clear(); this.editedIds.clear(); + try { + const toDelete = removed.map((id) => `msg:${id}`); + const toUpsert: IndexItem[] = []; + for (const id of edited) { + const m = this.cm.getMessage(id as never); + const item = m ? messageToItem(m, this.cfg) : null; + if (item) toUpsert.push(item); else toDelete.push(`msg:${id}`); + } + if (toDelete.length > 0) await this.client.delete(toDelete); + for (let i = 0; i < toUpsert.length; i += batchSize) { + const chunk = toUpsert.slice(i, i + batchSize); + const r = await this.client.upsert(chunk); + report.pushed += chunk.length; report.inserted += r.inserted; report.updated += r.updated; report.unchanged += r.unchanged; + } + } catch (e) { + for (const id of removed) this.removedIds.add(id); + for (const id of edited) if (!this.removedIds.has(id)) this.editedIds.add(id); + throw e; + } + } + + // Summaries before messages: they are few, and a message backlog that + // exceeds every tick's budget must not keep new summaries out forever. + { + const all = this.cm.getSummariesInRange({ fromMs: 0, toMs: Number.MAX_SAFE_INTEGER }) as SummaryLike[]; + report.summariesScanned = all.length; + const fresh = all.filter((s) => sumWm === null || s.createdMs > sumWm).sort((a, b) => a.createdMs - b.createdMs); + for (const s of fresh) { + if (budget <= 0) { report.more = true; break; } + const item = summaryToItem(s, this.cfg); + if (!item) continue; + pending.push(item); budget--; + if (pending.length >= batchSize) await flush(); + } + await flush(); + } + // Messages: walk forward from (watermark - overlap) via the time index. let fromMs = msgWm === null ? undefined : Math.max(0, msgWm - (this.cfg.overlapMs ?? 600_000)); + // Messages at exactly `fromMs` already seen this pass. The query bounds + // are inclusive, so the next page starts at the last page's final + // millisecond and skips the ones already walked with `offset` — a run + // of more than a page within one millisecond is walked, not jumped. + let offset = 0; const pageSize = 256; for (;;) { if (budget <= 0) { report.more = true; break; } - const page = this.cm.queryMessagesByTime({ fromMs, limit: pageSize }); + const page = this.cm.queryMessagesByTime({ fromMs, offset, limit: pageSize }); const msgs = page.messages; report.messagesScanned += msgs.length; for (const m of msgs) { @@ -302,25 +381,11 @@ export class SemanticIndexer { } if (msgs.length < pageSize) break; const lastMs = msgs[msgs.length - 1]!.timestamp.getTime(); - // Advance strictly: if a whole page shares one millisecond we would spin — step past it. - fromMs = lastMs === fromMs ? lastMs + 1 : lastMs; + const atLast = msgs.filter((m) => m.timestamp.getTime() === lastMs).length; + offset = lastMs === fromMs ? offset + atLast : atLast; + fromMs = lastMs; } await flush(); - - // Summaries: everything created after the summary watermark, any level. - if (!report.more) { - const all = this.cm.getSummariesInRange({ fromMs: 0, toMs: Number.MAX_SAFE_INTEGER }) as SummaryLike[]; - report.summariesScanned = all.length; - const fresh = all.filter((s) => sumWm === null || s.createdMs > sumWm).sort((a, b) => a.createdMs - b.createdMs); - for (const s of fresh) { - if (budget <= 0) { report.more = true; break; } - const item = summaryToItem(s, this.cfg); - if (!item) continue; - pending.push(item); budget--; - if (pending.length >= batchSize) await flush(); - } - await flush(); - } this.consecutiveFailures = 0; this.lastError = null; this.lastSyncAt = Date.now(); return report; } catch (e) { diff --git a/test/history-semantic-review-findings.test.ts b/test/history-semantic-review-findings.test.ts new file mode 100644 index 0000000..4229fc3 --- /dev/null +++ b/test/history-semantic-review-findings.test.ts @@ -0,0 +1,173 @@ +import { describe, it, before, after, beforeEach } from 'node:test'; +import assert from 'node:assert/strict'; +import { createServer, type Server } from 'node:http'; +import { HistoryModule } from '../src/modules/history/index.js'; +import type { ContextManager, StoredMessage } from '@animalabs/context-manager'; +import type { ContentBlock } from '@animalabs/membrane'; + +/** + * Regressions for the Greptile review of AF #173 (semantic sync edge cases). + * Each case was red before the fix it names. + */ + +interface Item { id: string; text: string; ts?: number; channel?: string | null; kind: string; level?: number; cursor?: number } + +/** Minimal embed-service: dedup by id, cursor-aware stats, counts requests. */ +class FakeService { + items = new Map(); + calls: string[] = []; + deleted: string[] = []; + server!: Server; + url = ''; + async start(): Promise { + this.server = createServer((req, res) => { + let body = ''; + req.on('data', (c) => { body += c; }); + req.on('end', () => { + const m = /^\/v1\/index\/([^/]+)\/(stats|upsert|search|delete)$/.exec(req.url ?? ''); + this.calls.push(m?.[2] ?? '?'); + const json = (code: number, o: unknown): void => { res.writeHead(code, { 'content-type': 'application/json' }); res.end(JSON.stringify(o)); }; + if (!m) return json(404, {}); + if (m[2] === 'stats') { + if (this.items.size === 0) return json(404, { error: { message: 'no such namespace' } }); + const by: Record = {}; + for (const it of this.items.values()) { + const d = by[it.kind] ??= { count: 0, max_ts: null, min_ts: null, max_cursor: null }; + d.count++; + if (it.cursor !== undefined) d.max_cursor = d.max_cursor === null ? it.cursor : Math.max(d.max_cursor, it.cursor); + } + return json(200, { namespace: 'n', model: 'fake', dim: 1, count: this.items.size, by_kind: by }); + } + if (m[2] === 'upsert') { + const items = (JSON.parse(body) as { items: Item[] }).items; + let inserted = 0, unchanged = 0; + for (const it of items) { if (this.items.get(it.id)?.text === it.text) unchanged++; else inserted++; this.items.set(it.id, it); } + return json(200, { inserted, updated: 0, unchanged, count: this.items.size }); + } + if (m[2] === 'delete') { + const ids = (JSON.parse(body) as { ids: string[] }).ids; + let n = 0; for (const id of ids) if (this.items.delete(id)) n++; + this.deleted.push(...ids); + return json(200, { deleted: n }); + } + return json(200, { namespace: 'n', hits: [], count_indexed: this.items.size, timing_ms: { embed: 0, search: 0 } }); + }); + }); + await new Promise((r) => this.server.listen(0, '127.0.0.1', r)); + const a = this.server.address(); + this.url = `http://127.0.0.1:${typeof a === 'object' && a ? a.port : 0}`; + } + stop(): Promise { + this.server.closeAllConnections(); + return new Promise((r) => this.server.close(() => r())); + } +} + +type Listener = (e: { type: string; messageId?: string }) => void; + +function msg(id: string, ms: number, content: ContentBlock[] | string, metadata?: Record): StoredMessage { + return { + id, sequence: Number(id.replace(/\D/g, '')), participant: 'Linn', + content: typeof content === 'string' ? [{ type: 'text', text: content }] : content, + timestamp: new Date(ms), metadata, + } as unknown as StoredMessage; +} + +/** Stub CM following the real contracts: inclusive bounds, oldest first, offset then limit, onMessage events. */ +class StubCm { + listeners = new Set(); + constructor(public messages: StoredMessage[], public summaries: Array> = []) {} + queryMessagesByTime(o: { fromMs?: number; toMs?: number; limit?: number; offset?: number }) { + const all = this.messages + .filter((m) => (o.fromMs === undefined || m.timestamp.getTime() >= o.fromMs) && (o.toMs === undefined || m.timestamp.getTime() <= o.toMs)) + .sort((a, b) => a.timestamp.getTime() - b.timestamp.getTime() || a.sequence - b.sequence); + const off = o.offset ?? 0; + return { messages: all.slice(off, o.limit === undefined ? undefined : off + o.limit), totalCount: all.length }; + } + getSummariesInRange() { return this.summaries; } + getMessage(id: string) { return this.messages.find((m) => m.id === id) ?? null; } + getSummary(id: string) { return this.summaries.find((s) => s.id === id) ?? null; } + onMessage(l: Listener) { this.listeners.add(l); return () => this.listeners.delete(l); } + edit(id: string, text: string) { + const m = this.getMessage(id)!; + (m as { content: ContentBlock[] }).content = [{ type: 'text', text }]; + for (const l of this.listeners) l({ type: 'edit', messageId: id }); + } + remove(id: string) { + this.messages = this.messages.filter((m) => m.id !== id); + for (const l of this.listeners) l({ type: 'remove', messageId: id }); + } +} + +const T0 = Date.UTC(2026, 8, 22); + +describe('semantic sync: Greptile #173 findings', () => { + const svc = new FakeService(); + before(() => svc.start()); + after(() => svc.stop()); + beforeEach(() => { svc.items.clear(); svc.calls.length = 0; svc.deleted.length = 0; }); + + function mod(cm: StubCm): HistoryModule { + const m = new HistoryModule({ semantic: { url: svc.url, namespace: 'n', syncIntervalMs: 0 } }); + m.bind(cm as unknown as ContextManager); + return m; + } + + it('#2 indexes the channel from metadata.channelId (MCPL ingestion shape), not only metadata.external.channelId', async () => { + const m = mod(new StubCm([ + msg('m1', T0, 'hello from the mcpl side of things', { channelId: 'chan-mcpl' }), + msg('m2', T0 + 1, 'hello from the legacy discord shape', { external: { channelId: 'chan-ext' } }), + ])); + await m.syncSemanticIndex(); + assert.equal(svc.items.get('msg:m1')?.channel, 'chan-mcpl'); + assert.equal(svc.items.get('msg:m2')?.channel, 'chan-ext'); + await m.stop(); + }); + + it('#3 more than one page of messages sharing one millisecond are all indexed', async () => { + const store = Array.from({ length: 600 }, (_, i) => msg(`m${i}`, T0, `same-millisecond message ${i}`)); + const m = mod(new StubCm(store)); + for (let i = 0; i < 5; i++) await m.syncSemanticIndex(10_000); + assert.equal(svc.items.size, 600); + await m.stop(); + }); + + it('#5 a message backlog larger than the tick budget does not starve new summaries', async () => { + const store = Array.from({ length: 2000 }, (_, i) => msg(`m${i}`, T0 + i * 1000, `backlog message ${i}`)); + const summaries = [{ id: 's1', level: 1, content: 'a fresh summary', tokens: 3, startMs: T0, endMs: T0 + 1, firstSequence: 0, lastSequence: 1, createdMs: T0 + 5 }]; + const m = mod(new StubCm(store, summaries)); + const r = await m.syncSemanticIndex(300)!; + assert.equal(r.more, true); + assert.ok(svc.items.has('sum:s1'), 'summary indexed while the message backlog is still draining'); + await m.stop(); + }); + + it('#6 an edit to a message older than the overlap window reaches the index on the next sync', async () => { + const cm = new StubCm([msg('old', T0, 'the original wording'), msg('new', T0 + 3_600_000, 'an hour later')]); + const m = mod(cm); + await m.syncSemanticIndex(); + cm.edit('old', 'the corrected wording'); + await m.syncSemanticIndex(); + assert.equal(svc.items.get('msg:old')?.text, 'the corrected wording'); + await m.stop(); + }); + + it('#1 a removed message is deleted from the index on the next sync', async () => { + const cm = new StubCm([msg('gone', T0, 'something that gets deleted'), msg('kept', T0 + 1, 'something that stays')]); + const m = mod(cm); + await m.syncSemanticIndex(); + cm.remove('gone'); + await m.syncSemanticIndex(); + assert.deepEqual(svc.deleted, ['msg:gone']); + assert.ok(!svc.items.has('msg:gone')); + assert.ok(svc.items.has('msg:kept')); + await m.stop(); + }); + + it('#8 a short private note is still indexed', async () => { + const m = mod(new StubCm([msg('t', T0, [{ type: 'tool_use', id: 'x', name: 'think', input: { content: 'wait for Ben' } }] as ContentBlock[])])); + await m.syncSemanticIndex(); + assert.equal(svc.items.get('msg:t')?.text, '[think] wait for Ben'); + await m.stop(); + }); +});