From 6ca10f0bf1e2250a638ca3930443144a2a35b2ec Mon Sep 17 00:00:00 2001 From: antra-tess Date: Tue, 22 Sep 2026 18:05:47 -0700 Subject: [PATCH 1/6] feat(history): author filter, extract around id, wholeWord, newest-first search Resident (Fable) diary-work feedback on the history tools: - author / excludeAuthor on search + extract. Incoming MCPL messages are all participant "user"; the author lives in metadata.author, so the filter matches author name/id, or participant for the agent's own turns. Exact match, never substring. Results carry `author`. - extract({aroundId, before, after, allChannels}): the conversation around a message id from search, so a search hit no longer needs its timestamp copied by hand into extract. - search wholeWord (Unicode-aware; "mission" no longer hits "uncommissioned"), order:"newest", and scannedThrough + hint whenever the scan stops early. Found on a real store (not by the stub): a time-only native query's matchedCount is the PAGE size, and channel queries come back in append order, not timestamp order. New edgeWindow() picks the n messages nearest a range edge in true timestamp order: native reverse when there is no channel, a widening time window when there is one. Plain extract's matchedCount had the same page-size problem: it now reports hasMore, and an exact matchedCount only when a channel is given. The test stubs now mimic both real contracts. Co-Authored-By: Claude Opus 5.5 (1M context) --- changelog.d/history-author-around.added.md | 10 + changelog.d/history-author-around.fixed.md | 4 + src/modules/history/index.ts | 470 +++++++++++++++++++-- test/history-module-ux.test.ts | 351 +++++++++++++++ test/history-module.test.ts | 22 +- 5 files changed, 832 insertions(+), 25 deletions(-) create mode 100644 changelog.d/history-author-around.added.md create mode 100644 changelog.d/history-author-around.fixed.md create mode 100644 test/history-module-ux.test.ts diff --git a/changelog.d/history-author-around.added.md b/changelog.d/history-author-around.added.md new file mode 100644 index 00000000..e0fb4d45 --- /dev/null +++ b/changelog.d/history-author-around.added.md @@ -0,0 +1,10 @@ +- History tools, from a resident's diary-work feedback: `search` and `extract` + take `author` / `excludeAuthor` (exact, case-insensitive match on + `metadata.author` name or id, or the stored participant for the agent's own + turns; MCPL-ingested messages are all participant `user`, so participant + alone could not tell authors apart). Results now carry `author`. + `extract({aroundId, before, after})` returns the conversation around a + message id from `search` (its own channel by default, `allChannels` to + interleave). `search` adds `wholeWord` (Unicode-aware) and + `order: "newest"`, and when it stops early (at `limit` or `maxScan`) it + reports `scannedThrough` plus a hint naming the `from`/`to` to continue with. diff --git a/changelog.d/history-author-around.fixed.md b/changelog.d/history-author-around.fixed.md new file mode 100644 index 00000000..6aa09a31 --- /dev/null +++ b/changelog.d/history-author-around.fixed.md @@ -0,0 +1,4 @@ +- History `extract` without a channel no longer reports a `matchedCount` that + was only the page size (a time-only native query returns a page-sized count, + not a total). It now reports `hasMore`, and keeps the exact `matchedCount` + only for channel-scoped queries. diff --git a/src/modules/history/index.ts b/src/modules/history/index.ts index 8ea36121..b97a3dc3 100644 --- a/src/modules/history/index.ts +++ b/src/modules/history/index.ts @@ -60,6 +60,9 @@ interface StatsInput { channelId?: string; } +/** One author spec or several — see `matchesAuthorFilter`. */ +type AuthorSpec = string | string[]; + interface ExtractInput { from?: string; to?: string; @@ -67,15 +70,26 @@ interface ExtractInput { limit?: number; offset?: number; format?: 'text' | 'raw'; + author?: AuthorSpec; + excludeAuthor?: AuthorSpec; + maxScan?: number; + aroundId?: string; + before?: number; + after?: number; + allChannels?: boolean; } interface SearchInput { query: string; regex?: boolean; caseSensitive?: boolean; + wholeWord?: boolean; from?: string; to?: string; channelId?: string; + author?: AuthorSpec; + excludeAuthor?: AuthorSpec; + order?: 'oldest' | 'newest'; limit?: number; maxScan?: number; } @@ -109,6 +123,25 @@ const EXTRACT_MAX_LIMIT = 200; */ const NATIVE_OFFSET_MAX = 0xffffffff; // 4294967295 +/** + * `extract` with an author filter can't lean on a native index (chronicle + * indexes time and channel, not author), so it pages through the time/channel + * window and filters in-process. Bounded like `search`'s maxScan, with the + * same explicit truncated/scannedThrough signal when the window is larger. + */ +const EXTRACT_FILTER_DEFAULT_MAX_SCAN = 10000; +const EXTRACT_FILTER_MAX_MAX_SCAN = 50000; +/** Native page size for the in-process filtered scans above. */ +const FILTER_SCAN_PAGE = 1000; + +/** `extract({aroundId})` window sizes, per side. */ +const AROUND_DEFAULT = 10; +const AROUND_MAX = 100; +/** Upper bound on same-millisecond neighbours fetched around an anchor — + * real stores see a handful at most; this only keeps a pathological burst + * from turning one call into an unbounded fetch. */ +const AROUND_TIE_CAP = 500; + const SEARCH_DEFAULT_LIMIT = 20; const SEARCH_MAX_LIMIT = 100; const SEARCH_DEFAULT_MAX_SCAN = 5000; @@ -157,6 +190,21 @@ const SEARCH_REGEX_TIMEOUT_MS = 2000; * directory this module's own compiled output lives in. */ const SEARCH_WORKER_PATH = join(dirname(fileURLToPath(import.meta.url)), 'search-regex-worker.js'); +// Declared as an array (tool-schema dialects vary on anyOf); a bare string is +// also accepted at runtime — see toAuthorSet. +const AUTHOR_SPEC_JSON_SCHEMA = { type: 'array', items: { type: 'string' } }; +const AUTHOR_SCHEMA = { + ...AUTHOR_SPEC_JSON_SCHEMA, + description: + 'Only messages written by one of these authors (a single string is fine too). Matches, case-insensitively and exactly, ' + + 'the author\'s display name or user id (a leading "@" or a <@id> mention is accepted), or the ' + + 'stored participant for messages without author metadata — e.g. your own turns, under your name.', +}; +const EXCLUDE_AUTHOR_SCHEMA = { + ...AUTHOR_SPEC_JSON_SCHEMA, + description: 'Drop messages written by any of these authors (a single string is fine too). Same matching as `author`.', +}; + export class HistoryModule implements Module { readonly name = 'history'; @@ -258,6 +306,83 @@ export class HistoryModule implements Module { return cm.queryMessagesByTime({ fromMs: ms, toMs: ms }).messages; } + /** + * The `n` messages of a time/channel range nearest one of its edges, in + * TIMESTAMP order walking away from that edge (newest-first for + * side:"newest", oldest-first for side:"oldest"), plus whether the range + * holds more than `n`. + * + * Why not just `queryMessagesByTimeAndChannel` + offset: (a) a time-only + * query's `matchedCount` is the returned PAGE size, not a total (see + * context-manager MessageStore.queryByTime), so it can't locate a tail; + * (b) channel-scoped results come back in APPEND (ordinal) order, and + * Discord catch-up appends old-timestamped messages late — an ordinal + * tail is not the timestamp tail. + * + * No channel: chronicle's timestamp index answers directly (`reverse` + * for the newest edge). With a channel: channel queries report an exact + * matchedCount, so grow a time window out from the edge (×4 per step) + * until it holds ≥ n channel messages or covers the whole range, then + * fetch just that window and sort. A window that overshoots wildly (a + * burst) falls back to its ordinal tail rather than decoding thousands + * of messages — correct for everything but same-burst backfill. + */ + private edgeWindow(opts: { + fromMs?: number; + toMs?: number; + channelId?: string; + n: number; + side: 'newest' | 'oldest'; + }): { messages: StoredMessage[]; more: boolean } { + const cm = this.cm as ContextManager; + const { fromMs, toMs, channelId, n, side } = opts; + const cmp = (a: StoredMessage, b: StoredMessage) => + a.timestamp.getTime() - b.timestamp.getTime() || a.sequence - b.sequence; + const dirSort = (ms: StoredMessage[]) => ms.sort(side === 'newest' ? (a, b) => cmp(b, a) : cmp); + + if (channelId === undefined) { + const r = cm.queryMessagesByTime({ fromMs, toMs, limit: n + 1, reverse: side === 'newest' }); + return { messages: r.messages.slice(0, n), more: r.messages.length > n }; + } + + const count = (f?: number, t?: number) => + cm.queryMessagesByTimeAndChannel({ fromMs: f, toMs: t, channelId, limit: 0 }).matchedCount; + const total = count(fromMs, toMs); + if (n === 0) return { messages: [], more: total > 0 }; + let wFrom = fromMs; + let wTo = toMs; + if (total > n) { + // Concrete bounds of the range (open ends resolved to the store's + // first message / now), so growth always terminates. + const lo = fromMs ?? this.earliestMessageMs(toMs) ?? 0; + const hi = toMs ?? Date.now(); + for (let width = 10 * 60_000; ; width *= 4) { + const f = side === 'newest' ? Math.max(hi - width, lo) : lo; + const t = side === 'newest' ? hi : Math.min(lo + width, hi); + const covered = side === 'newest' ? f <= lo : t >= hi; + if (covered) break; // whole range: keep wFrom/wTo = the range itself + if (count(f, t) >= n) { + wFrom = f; + wTo = t; + break; + } + } + } + const inWindow = count(wFrom, wTo); + const fetchCap = Math.max(4 * n, n + 1000); + const page = + inWindow <= fetchCap + ? cm.queryMessagesByTimeAndChannel({ fromMs: wFrom, toMs: wTo, channelId, limit: inWindow }).messages + : cm.queryMessagesByTimeAndChannel({ + fromMs: wFrom, + toMs: wTo, + channelId, + limit: n, + offset: side === 'newest' ? inWindow - n : 0, + }).messages; + return { messages: dirSort(page).slice(0, n), more: total > n }; + } + async start(ctx: ModuleContext): Promise { this.ctx = ctx; } @@ -291,15 +416,28 @@ export class HistoryModule implements Module { 'Fetch raw, uncompressed messages for a time range and/or channel, paginated oldest-first. ' + 'Use after `stats` to pull the actual content. format:"text" flattens each message to a short ' + 'readable string (tool_use/tool_result/thinking/images rendered as bracketed labels); ' + - 'format:"raw" returns the content blocks unmodified.', + 'format:"raw" returns the content blocks unmodified. ' + + 'Two extra modes: (1) aroundId — pass a message id (e.g. from `search`) to get the conversation ' + + 'around it: `before` messages before it, the message itself (anchor:true), and `after` messages ' + + 'after it, in the anchor\'s own channel unless channelId/allChannels says otherwise; from/to/offset ' + + 'do not apply in this mode. (2) author/excludeAuthor — keep only (or drop) messages by these ' + + 'authors; this is filtered in-process over at most maxScan messages of the window, so a response ' + + 'may report truncated:true + scannedThrough (continue with from:scannedThrough).', inputSchema: { type: 'object' as const, properties: { from: { type: 'string', description: 'ISO 8601 inclusive lower bound. Omit for open-ended.' }, to: { type: 'string', description: 'ISO 8601 inclusive upper bound. Omit for open-ended.' }, channelId: { type: 'string', description: 'Restrict to one channel. Accepts a channel label (e.g. "#general") or the raw internal channel id.' }, + author: AUTHOR_SCHEMA, + excludeAuthor: EXCLUDE_AUTHOR_SCHEMA, + aroundId: { type: 'string', description: 'Message id to center on (as returned by search/extract). Returns the surrounding conversation instead of a range.' }, + before: { type: 'number', description: `With aroundId: messages to include before the anchor (default ${AROUND_DEFAULT}, cap ${AROUND_MAX}).` }, + after: { type: 'number', description: `With aroundId: messages to include after the anchor (default ${AROUND_DEFAULT}, cap ${AROUND_MAX}).` }, + allChannels: { type: 'boolean', description: 'With aroundId: interleave every channel around the anchor\'s time instead of staying in its channel (default false).' }, limit: { type: 'number', description: `Max messages to return (default ${EXTRACT_DEFAULT_LIMIT}, hard cap ${EXTRACT_MAX_LIMIT}). Must be a non-negative integer.` }, offset: { type: 'number', description: `Number of matching messages to skip (default 0, capped to ${NATIVE_OFFSET_MAX}). Must be a non-negative integer.` }, + maxScan: { type: 'number', description: `With author/excludeAuthor: max window messages to examine (default ${EXTRACT_FILTER_DEFAULT_MAX_SCAN}, hard cap ${EXTRACT_FILTER_MAX_MAX_SCAN}).` }, format: { type: 'string', enum: ['text', 'raw'], description: 'Content rendering (default "text").' }, }, }, @@ -311,7 +449,13 @@ export class HistoryModule implements Module { 'range/channel window. Narrows candidates via the same query as `extract` before matching, up ' + 'to maxScan candidates — if the narrowed window is larger than maxScan, the response reports ' + 'truncated:true up front rather than silently missing later matches; narrow the filter or raise ' + - 'maxScan and retry. regex:true matching runs under a wall-clock deadline and is cleanly failed ' + + 'maxScan and retry. When the scan stops early (truncated, or `limit` matches reached) the response ' + + 'carries scannedThrough: pass it as `from` (order:"oldest") or `to` (order:"newest") to continue. ' + + 'order:"newest" scans the most recent maxScan messages of the window first — usually what you want ' + + 'for "when did this last come up". author/excludeAuthor narrow by who wrote the message; ' + + 'wholeWord:true stops "mission" from matching inside "uncommissioned". Each match carries an `id` ' + + 'you can hand to extract({aroundId}) to read the conversation around it. ' + + 'regex:true matching runs under a wall-clock deadline and is cleanly failed ' + `(not silently empty) if a pattern is too slow — avoid nested-quantifier patterns like (a+)+.`, inputSchema: { type: 'object' as const, @@ -319,9 +463,13 @@ export class HistoryModule implements Module { query: { type: 'string', description: 'Substring (or regex source, when regex:true) to search for.' }, regex: { type: 'boolean', description: 'Treat query as a regular expression (default false).' }, caseSensitive: { type: 'boolean', description: 'Case-sensitive match (default false).' }, + wholeWord: { type: 'boolean', description: 'Substring mode only: match whole words, not inside longer words (default false). In regex mode use \\b yourself.' }, from: { type: 'string', description: 'ISO 8601 inclusive lower bound. Omit for open-ended.' }, to: { type: 'string', description: 'ISO 8601 inclusive upper bound. Omit for open-ended.' }, channelId: { type: 'string', description: 'Restrict to one channel. Accepts a channel label (e.g. "#general") or the raw internal channel id.' }, + author: AUTHOR_SCHEMA, + excludeAuthor: EXCLUDE_AUTHOR_SCHEMA, + order: { type: 'string', enum: ['oldest', 'newest'], description: 'Scan and return oldest-first (default) or newest-first.' }, limit: { type: 'number', description: `Max matches to return (default ${SEARCH_DEFAULT_LIMIT}, hard cap ${SEARCH_MAX_LIMIT}). Must be a non-negative integer.` }, maxScan: { type: 'number', description: `Max candidate messages to scan (default ${SEARCH_DEFAULT_MAX_SCAN}, hard cap ${SEARCH_MAX_MAX_SCAN}). Must be a non-negative integer.` }, }, @@ -459,28 +607,172 @@ export class HistoryModule implements Module { // ========================================================================== private handleExtract(input: ExtractInput): ToolResult { + if (input.aroundId !== undefined) return this.handleExtractAround(input); + if (input.before !== undefined || input.after !== undefined || input.allChannels !== undefined) { + throw new Error('"before"/"after"/"allChannels" only apply together with "aroundId".'); + } const channelId = this.resolveChannel(input.channelId); const fromMs = parseIsoDate(input.from, 'from'); const toMs = parseIsoDate(input.to, 'to'); const limit = clampCount(input.limit, EXTRACT_DEFAULT_LIMIT, EXTRACT_MAX_LIMIT, 'limit'); const offset = clampCount(input.offset, 0, NATIVE_OFFSET_MAX, 'offset'); const format = input.format ?? 'text'; + const authorFilter = buildAuthorFilter(input.author, input.excludeAuthor); const cm = this.cm as ContextManager; + + if (authorFilter) { + // No native author index: page through the time/channel window and + // filter here. offset/limit apply to the FILTERED sequence. Stop as + // soon as the requested page is full; otherwise scan to maxScan. Only + // an exhausted window yields an exact matchedCount — anything else is + // reported as truncated with a resume point, never as a total. + const maxScan = clampCount(input.maxScan, EXTRACT_FILTER_DEFAULT_MAX_SCAN, EXTRACT_FILTER_MAX_MAX_SCAN, 'maxScan'); + const page: StoredMessage[] = []; + let kept = 0; + let scanned = 0; + let lastScanned: StoredMessage | undefined; + let exhausted = false; + let pageFull = false; + outer: for (;;) { + const want = Math.min(FILTER_SCAN_PAGE, maxScan - scanned); + if (want <= 0) break; + const chunk = cm.queryMessagesByTimeAndChannel({ fromMs, toMs, channelId, limit: want, offset: scanned }).messages; + for (const m of chunk) { + scanned++; + lastScanned = m; + if (!authorFilter(m)) continue; + if (kept >= offset && page.length < limit) page.push(m); + kept++; + if (page.length >= limit && kept > offset + limit) { + // One filtered match beyond the page proves there is more; + // no need to keep scanning. + pageFull = true; + break outer; + } + } + if (chunk.length < want) { + exhausted = true; + break; + } + } + if (!exhausted && !pageFull && scanned >= maxScan) { + // Probe one past maxScan: a window of exactly maxScan is complete. + const next = cm.queryMessagesByTimeAndChannel({ fromMs, toMs, channelId, limit: 1, offset: scanned }).messages; + if (next.length === 0) exhausted = true; + } + return { + success: true, + data: { + ...(exhausted ? { matchedCount: kept } : { matchedCountAtLeast: kept }), + returned: page.length, + scanned, + truncated: !exhausted, + ...(!exhausted && lastScanned + ? { + scannedThrough: lastScanned.timestamp.toISOString(), + hint: pageFull + ? 'More matching messages exist; raise offset (or continue with from:scannedThrough and offset:0).' + : `Stopped after scanning ${scanned} messages of the window without reaching its end; continue with ` + + 'from:scannedThrough (messages at exactly that instant may repeat), or narrow with channelId/dates.', + } + : {}), + messages: page.map((msg) => projectMessage(msg, format)), + }, + }; + } + + // One extra row answers "is there another page" in every case. The + // native matchedCount is a true total only for channel-scoped queries; + // for time-only/unfiltered ones it is just the page size (see + // edgeWindow), so it is reported only where it means what it says. const result = cm.queryMessagesByTimeAndChannel({ fromMs, toMs, channelId, - limit, + limit: limit + 1, offset, }); + const page = result.messages.slice(0, limit); + + return { + success: true, + data: { + ...(channelId !== undefined ? { matchedCount: result.matchedCount } : {}), + returned: page.length, + hasMore: result.messages.length > limit, + messages: page.map((msg) => projectMessage(msg, format)), + }, + }; + } + + /** + * `extract({aroundId})` — the conversation around one message, like + * Discord's fetch_around. Built from two index-backed time queries meeting + * at the anchor's millisecond: the newest of [.., t] and the oldest of + * [t, ..] (both via edgeWindow). Both include every message sharing t, so each side is + * over-fetched by that tie count; the union is deduped by id, ordered by + * (timestamp, sequence), and cut to `before`/`after` around the anchor. + */ + private handleExtractAround(input: ExtractInput): ToolResult { + if (input.from !== undefined || input.to !== undefined || input.offset !== undefined) { + throw new Error('"aroundId" cannot be combined with from/to/offset — it picks its own window.'); + } + if (input.author !== undefined || input.excludeAuthor !== undefined) { + throw new Error('"aroundId" returns the whole conversation around a message; author filters do not apply. Use search/extract with author instead.'); + } + if (input.allChannels && input.channelId !== undefined) { + throw new Error('Pass either "channelId" or "allChannels", not both.'); + } + const before = clampCount(input.before, AROUND_DEFAULT, AROUND_MAX, 'before'); + const after = clampCount(input.after, AROUND_DEFAULT, AROUND_MAX, 'after'); + const format = input.format ?? 'text'; + const cm = this.cm as ContextManager; + const anchor = cm.getMessage(String(input.aroundId)); + if (!anchor) { + throw new Error(`No message with id ${JSON.stringify(input.aroundId)} in this history.`); + } + const anchorId = String(anchor.id); + const anchorChannel = getChannelId(anchor); + const channelId = input.allChannels ? undefined : (this.resolveChannel(input.channelId) ?? anchorChannel); + const t = anchor.timestamp.getTime(); + + // Every message sharing the anchor's millisecond lands on BOTH sides + // below, so over-fetch each side by that count; the dedupe + sort then + // orders them by sequence around the anchor. + const ties = Math.min( + cm.queryMessagesByTimeAndChannel({ fromMs: t, toMs: t, channelId, limit: AROUND_TIE_CAP }).messages.length, + AROUND_TIE_CAP, + ); + const beforeSide = this.edgeWindow({ toMs: t, channelId, n: before + ties, side: 'newest' }).messages; + const afterSide = this.edgeWindow({ fromMs: t, channelId, n: after + ties, side: 'oldest' }).messages; + + const byId = new Map(); + for (const m of [...beforeSide, ...afterSide]) byId.set(String(m.id), m); + const ordered = [...byId.values()].sort( + (a, b) => a.timestamp.getTime() - b.timestamp.getTime() || a.sequence - b.sequence, + ); + const idx = ordered.findIndex((m) => String(m.id) === anchorId); + if (idx === -1) { + // Only reachable when an explicit channelId excludes the anchor. + throw new Error( + `Message ${anchorId} is not in channel ${JSON.stringify(input.channelId)}` + + (anchorChannel ? ` (it is in ${anchorChannel}).` : '.') + + ' Omit channelId to use the anchor\'s own channel.', + ); + } + const window = ordered.slice(Math.max(0, idx - before), idx + after + 1); return { success: true, data: { - matchedCount: result.matchedCount, - returned: result.messages.length, - messages: result.messages.map((msg) => projectMessage(msg, format)), + anchorId, + channelId: channelId ?? null, + returned: window.length, + messages: window.map((m) => ({ + ...projectMessage(m, format), + ...(String(m.id) === anchorId ? { anchor: true } : {}), + })), }, }; } @@ -500,6 +792,14 @@ export class HistoryModule implements Module { const maxScan = clampCount(input.maxScan, SEARCH_DEFAULT_MAX_SCAN, SEARCH_MAX_MAX_SCAN, 'maxScan'); const caseSensitive = input.caseSensitive ?? false; const flags = caseSensitive ? '' : 'i'; + const order = input.order ?? 'oldest'; + if (order !== 'oldest' && order !== 'newest') { + throw new Error(`"order" must be "oldest" or "newest", got ${JSON.stringify(input.order)}.`); + } + if (input.wholeWord && input.regex) { + throw new Error('"wholeWord" applies to substring search only; in regex mode write \\b around the pattern instead.'); + } + const authorFilter = buildAuthorFilter(input.author, input.excludeAuthor); // Validate regex SYNTAX up front — an invalid pattern is a clean tool // error, not a crash mid-scan. This does NOT bound match TIME (a @@ -516,19 +816,49 @@ export class HistoryModule implements Module { const needle = caseSensitive ? input.query : input.query.toLowerCase(); const cm = this.cm as ContextManager; - // Fetch one more than maxScan so an oversized candidate window is - // detected up front (before any matching is attempted), rather than + // The candidate window is the oldest (order:"oldest") or newest + // (order:"newest", via edgeWindow — see there for why not an offset + // tail) maxScan messages of the time/channel range. Oversize is + // detected up front rather than // silently scanning maxScan candidates and returning as if that were the // whole window. See the tool description's `truncated` contract. - const probe = cm.queryMessagesByTimeAndChannel({ - fromMs, - toMs, - channelId, - limit: maxScan + 1, - offset: 0, - }); - const truncated = probe.messages.length > maxScan; - const candidates = truncated ? probe.messages.slice(0, maxScan) : probe.messages; + let windowMessages: StoredMessage[]; + let truncated: boolean; + { + const w = this.edgeWindow({ fromMs, toMs, channelId, n: maxScan, side: order }); + truncated = w.more; + windowMessages = w.messages; + } + const candidatePoolSize = windowMessages.length; + const candidates = authorFilter ? windowMessages.filter(authorFilter) : windowMessages; + const poolInfo = { + candidatePoolSize, + ...(authorFilter ? { afterAuthorFilter: candidates.length } : {}), + order, + }; + // Where the scan stopped, when it stopped before the end of the window + // the caller asked about: either `limit` matches came first (the rest of + // the pool was never looked at) or the pool itself was truncated. The + // timestamp is the continuation point — from: (oldest) / to: (newest). + // Boundary is inclusive, so a message at exactly that instant may repeat. + const resumeInfo = (scanned: number): Record => { + const stoppedEarly = scanned < candidates.length; + if (!stoppedEarly && !truncated) return {}; + // Stopped at limit: the last candidate actually looked at. Scanned the + // whole (truncated) pool: the pool's own far edge — which may lie past + // the last author-filtered candidate. + const last = stoppedEarly ? candidates[scanned - 1] : windowMessages[windowMessages.length - 1]; + if (!last) return {}; + const at = last.timestamp.toISOString(); + const bound = order === 'oldest' ? 'from' : 'to'; + return { + scannedThrough: at, + hint: stoppedEarly + ? `Stopped at limit; more candidates remain. Continue with ${bound}:"${at}" (or raise limit).` + : `Only the ${order} ${candidatePoolSize} messages of this range were scanned (maxScan); continue with ` + + `${bound}:"${at}", raise maxScan, or narrow with channelId/author/dates.`, + }; + }; if (input.regex) { // Regex matching against caller-supplied patterns is ReDoS-shaped: a @@ -541,7 +871,7 @@ export class HistoryModule implements Module { // on a deadline. See search-regex-worker.ts's header for the full // rationale. const { matches, scanned } = await this.searchWithRegexWorker(candidates, input.query, flags, limit); - return { success: true, data: { scanned, candidatePoolSize: candidates.length, truncated, matches } }; + return { success: true, data: { scanned, ...poolInfo, truncated, ...resumeInfo(scanned), matches } }; } // Plain substring search (String.prototype.indexOf) is inherently @@ -559,12 +889,13 @@ export class HistoryModule implements Module { if (matches.length >= limit) break; scanned++; const text = flattenContent(msg.content); - const hit = matchSubstring(text, needle, caseSensitive); + const hit = matchSubstring(text, needle, caseSensitive, input.wholeWord ?? false); if (!hit) continue; matches.push({ id: String(msg.id), timestamp: msg.timestamp.toISOString(), participant: msg.participant, + author: authorName(msg), channelId: getChannelId(msg) ?? null, snippet: snippetAround(text, hit.index, hit.length), }); @@ -572,7 +903,7 @@ export class HistoryModule implements Module { return { success: true, - data: { scanned, candidatePoolSize: candidates.length, truncated, matches }, + data: { scanned, ...poolInfo, truncated, ...resumeInfo(scanned), matches }, }; } @@ -632,6 +963,7 @@ export class HistoryModule implements Module { id: String(msg.id), timestamp: msg.timestamp.toISOString(), participant: msg.participant, + author: authorName(msg), channelId: getChannelId(msg) ?? null, snippet: snippetAround(text, matchIndex, matchLength), }; @@ -1200,6 +1532,7 @@ function projectMessage(msg: StoredMessage, format: 'text' | 'raw'): Record/<@!id> mention to its id. */ +function normalizeAuthorSpec(spec: string): string { + const s = spec.trim(); + const mention = /^<@!?(\d+)>$/.exec(s); + if (mention) return mention[1]!; + return s.replace(/^@/, '').toLowerCase(); +} + +function toAuthorSet(spec: AuthorSpec | undefined, field: string): Set | null { + if (spec === undefined) return null; + const list = Array.isArray(spec) ? spec : [spec]; + if (list.some((s) => typeof s !== 'string')) { + throw new Error(`"${field}" must be a string or an array of strings.`); + } + const set = new Set(list.map(normalizeAuthorSpec).filter((s) => s.length > 0)); + if (set.size === 0) throw new Error(`"${field}" must name at least one author.`); + return set; +} + +/** + * Predicate for `author`/`excludeAuthor`, or null when neither is given. + * A message's identities are its metadata.author name and id plus its + * stored participant — the last is what identifies the agent's own turns + * (and anything else ingested without author metadata). Exact match after + * normalization, never substring: "ann" must not pull in "joanne". + */ +function buildAuthorFilter(include: AuthorSpec | undefined, exclude: AuthorSpec | undefined): ((m: StoredMessage) => boolean) | null { + const inc = toAuthorSet(include, 'author'); + const exc = toAuthorSet(exclude, 'excludeAuthor'); + if (!inc && !exc) return null; + return (m) => { + const a = authorOf(m); + const keys = [a?.name?.toLowerCase(), a?.id, m.participant?.toLowerCase()].filter((k): k is string => !!k); + if (inc && !keys.some((k) => inc.has(k))) return false; + if (exc && keys.some((k) => exc.has(k))) return false; + return true; + }; } /** ~SNIPPET_CONTEXT_CHARS of surrounding context on each side of a match, or diff --git a/test/history-module-ux.test.ts b/test/history-module-ux.test.ts new file mode 100644 index 00000000..f05ebc41 --- /dev/null +++ b/test/history-module-ux.test.ts @@ -0,0 +1,351 @@ +import { describe, it } from 'node:test'; +import assert from 'node:assert/strict'; +import { HistoryModule } from '../src/modules/history/index.js'; +import type { ContextManager, StoredMessage, IndexedMessageQueryResult } from '@animalabs/context-manager'; +import type { ToolCall } from '../src/types/events.js'; + +/** + * Author filter, extract-around-id, wholeWord, order:"newest" and the + * scannedThrough resume point (Fable's history-tool UX feedback, 09-22). + * Same stub posture as history-module.test.ts: in-memory filtering that + * mirrors the native query contract (time-ordered, both-ends-inclusive, + * offset/limit applied after filtering, matchedCount = pre-pagination). + */ + +interface Spec { + id: string; + ms: number; + text: string; + channel?: string; + /** metadata.author.name — MCPL-ingested messages are participant "user". */ + author?: string; + authorId?: string; + participant?: string; +} + +function mk(s: Spec): StoredMessage { + const metadata: Record = {}; + if (s.channel) metadata.channelId = s.channel; + if (s.author || s.authorId) metadata.author = { id: s.authorId ?? `id-${s.author}`, name: s.author }; + return { + id: s.id, + sequence: Number(s.id.replace(/\D/g, '')), + participant: s.participant ?? (s.author ? 'user' : 'Fable'), + content: [{ type: 'text', text: s.text }], + metadata, + timestamp: new Date(s.ms), + } as unknown as StoredMessage; +} + +/** + * Mirrors the REAL MessageStore contracts, including the unkind ones (the + * first version of this feature passed a kinder stub and failed on a real + * store): a time-only/unfiltered query is timestamp-ordered and its + * matchedCount is just the PAGE size; a channel-scoped query is APPEND + * (sequence) ordered with an exact matchedCount; queryMessagesByTime + * supports reverse. + */ +function stub(messages: StoredMessage[]) { + const calls: Array> = []; + const byTs = [...messages].sort((a, b) => a.timestamp.getTime() - b.timestamp.getTime() || a.sequence - b.sequence); + const bySeq = [...messages].sort((a, b) => a.sequence - b.sequence); + const inRange = (m: StoredMessage, f?: number, t?: number) => { + const ts = m.timestamp.getTime(); + return (f === undefined || ts >= f) && (t === undefined || ts <= t); + }; + const chan = (m: StoredMessage) => (m.metadata as { channelId?: string }).channelId; + const queryMessagesByTime = (args: { fromMs?: number; toMs?: number; limit?: number; offset?: number; reverse?: boolean }): IndexedMessageQueryResult => { + calls.push({ method: 'time', ...args }); + let list = byTs.filter((m) => inRange(m, args.fromMs, args.toMs)); + if (args.reverse) list = list.reverse(); + const offset = args.offset ?? 0; + const page = list.slice(offset, args.limit === undefined ? undefined : offset + args.limit); + return { messages: page, matchedCount: page.length }; + }; + const cm = { + queryMessagesByTime, + queryMessagesByTimeAndChannel(args: { fromMs?: number; toMs?: number; channelId?: string; limit?: number; offset?: number }): IndexedMessageQueryResult { + if (args.channelId === undefined) return queryMessagesByTime(args); + calls.push({ method: 'timeAndChannel', ...args }); + const filtered = bySeq.filter((m) => inRange(m, args.fromMs, args.toMs) && chan(m) === args.channelId); + const offset = args.offset ?? 0; + const limit = args.limit ?? filtered.length; + return { messages: filtered.slice(offset, offset + limit), matchedCount: filtered.length }; + }, + getMessage(id: string): StoredMessage | null { + return messages.find((m) => m.id === id) ?? null; + }, + } as unknown as ContextManager; + const mod = new HistoryModule(); + mod.bind(cm); + return { mod, calls }; +} + +async function call(mod: HistoryModule, name: string, input: Record) { + return mod.handleToolCall({ id: 't', name, input } as ToolCall); +} + +function data(r: { success: boolean; data?: unknown; error?: string }): any { + assert.equal(r.success, true, r.error); + return r.data; +} + +const T0 = Date.parse('2026-09-20T00:00:00Z'); +const min = (n: number) => T0 + n * 60_000; + +// Diary-shaped fixture: antra and Fable talking in #diary, plus a heartbeat +// and Fable's own journal checklist that mention "mission" too. +const fixture = [ + mk({ id: 'm1', ms: min(1), text: 'the mission today is the museum', channel: 'diary', author: 'antra' }), + mk({ id: 'm2', ms: min(2), text: 'museum of the uncommissioned — checklist', channel: 'diary' }), // Fable's own + mk({ id: 'm3', ms: min(3), text: 'heartbeat: mission status nominal', channel: 'diary', author: 'heartbeat' }), + mk({ id: 'm4', ms: min(4), text: 'Mission accomplished, I think', channel: 'diary', author: 'antra' }), + mk({ id: 'm5', ms: min(5), text: 'meanwhile elsewhere', channel: 'general', author: 'antra' }), + mk({ id: 'm6', ms: min(6), text: 'what was the mission again?', channel: 'diary', author: 'Joanne' }), + mk({ id: 'm7', ms: min(7), text: 'миссия выполнена', channel: 'diary', author: 'antra' }), +]; + +describe('HistoryModule UX: author filter', () => { + it('search author keeps only that author (case-insensitive, @ tolerated) and reports author on matches', async () => { + const { mod } = stub(fixture); + const d = data(await call(mod, 'search', { query: 'mission', author: '@Antra' })); + assert.deepEqual(d.matches.map((m: any) => m.id), ['m1', 'm4']); + assert.equal(d.matches[0].author, 'antra'); + assert.equal(d.matches[0].participant, 'user'); + assert.equal(d.candidatePoolSize, 7); + assert.equal(d.afterAuthorFilter, 4); + }); + + it('author is exact, never substring ("ann" does not pull in Joanne/antra)', async () => { + const { mod } = stub(fixture); + const d = data(await call(mod, 'search', { query: 'mission', author: 'ann' })); + assert.deepEqual(d.matches, []); + }); + + it('author matches author id, a <@id> mention, and participant for own turns', async () => { + const { mod } = stub([ + mk({ id: 'a1', ms: min(1), text: 'x', author: 'antra', authorId: '12345' }), + mk({ id: 'a2', ms: min(2), text: 'x' }), // participant Fable + mk({ id: 'a3', ms: min(3), text: 'x', author: 'bob' }), + ]); + assert.deepEqual(data(await call(mod, 'search', { query: 'x', author: '<@12345>' })).matches.map((m: any) => m.id), ['a1']); + assert.deepEqual(data(await call(mod, 'search', { query: 'x', author: 'fable' })).matches.map((m: any) => m.id), ['a2']); + assert.deepEqual(data(await call(mod, 'search', { query: 'x', author: ['bob', 'Fable'] })).matches.map((m: any) => m.id), ['a2', 'a3']); + }); + + it('excludeAuthor drops heartbeats and own turns', async () => { + const { mod } = stub(fixture); + const d = data(await call(mod, 'search', { query: 'mission', excludeAuthor: ['heartbeat', 'Fable'] })); + assert.deepEqual(d.matches.map((m: any) => m.id), ['m1', 'm4', 'm6']); + }); + + it('rejects an empty or non-string author spec', async () => { + const { mod } = stub(fixture); + assert.equal((await call(mod, 'search', { query: 'x', author: [] })).success, false); + assert.equal((await call(mod, 'search', { query: 'x', author: [5] })).success, false); + }); + + it('extract author: offset/limit over the filtered sequence, exact matchedCount when exhausted', async () => { + const { mod } = stub(fixture); + const d = data(await call(mod, 'extract', { author: 'antra', limit: 2, offset: 1 })); + assert.deepEqual(d.messages.map((m: any) => m.id), ['m4', 'm5']); + assert.equal(d.messages[0].author, 'antra'); + // page full + one more antra message (m7) exists → not a total + assert.equal(d.truncated, true); + assert.equal(d.matchedCountAtLeast, 4); + const all = data(await call(mod, 'extract', { author: 'antra' })); + assert.equal(all.truncated, false); + assert.equal(all.matchedCount, 4); + assert.equal(all.scanned, 7); + }); + + it('extract author: maxScan bound reports truncated + scannedThrough, and a window of exactly maxScan is complete', async () => { + const { mod } = stub(fixture); + const d = data(await call(mod, 'extract', { author: 'antra', maxScan: 3 })); + assert.equal(d.truncated, true); + assert.equal(d.scanned, 3); + assert.equal(d.scannedThrough, new Date(min(3)).toISOString()); + assert.deepEqual(d.messages.map((m: any) => m.id), ['m1']); + const exact = data(await call(mod, 'extract', { author: 'antra', maxScan: 7 })); + assert.equal(exact.truncated, false); + assert.equal(exact.matchedCount, 4); + }); +}); + +describe('HistoryModule UX: wholeWord', () => { + it('"mission" no longer matches inside "uncommissioned"', async () => { + const { mod } = stub([ + mk({ id: 'w1', ms: min(1), text: 'museum of the uncommissioned' }), + mk({ id: 'w2', ms: min(2), text: 'the Mission, finally' }), + ]); + assert.deepEqual(data(await call(mod, 'search', { query: 'mission' })).matches.map((m: any) => m.id), ['w1', 'w2']); + assert.deepEqual(data(await call(mod, 'search', { query: 'mission', wholeWord: true })).matches.map((m: any) => m.id), ['w2']); + }); + + it('is Unicode-aware and finds a later whole-word occurrence after an embedded one', async () => { + const { mod } = stub([ + mk({ id: 'u1', ms: min(1), text: 'суперкот и кот' }), + mk({ id: 'u2', ms: min(2), text: 'котёнок' }), + ]); + const d = data(await call(mod, 'search', { query: 'кот', wholeWord: true })); + assert.deepEqual(d.matches.map((m: any) => m.id), ['u1']); + assert.match(d.matches[0].snippet, /и кот/); + }); + + it('does not demand word boundaries next to a needle\'s own punctuation', async () => { + const { mod } = stub([mk({ id: 'p1', ms: min(1), text: 'see #general2 and #general' })]); + assert.equal(data(await call(mod, 'search', { query: '#general', wholeWord: true })).matches.length, 1); + }); + + it('is rejected together with regex', async () => { + const { mod } = stub(fixture); + const r = await call(mod, 'search', { query: 'a', regex: true, wholeWord: true }); + assert.equal(r.success, false); + assert.match(r.error!, /\\b/); + }); +}); + +describe('HistoryModule UX: order + scannedThrough', () => { + it('order:newest scans the most recent maxScan and returns newest-first, with a to: resume point', async () => { + const { mod } = stub(fixture); + const d = data(await call(mod, 'search', { query: 'mission', order: 'newest', maxScan: 4 })); + // window = m4..m7, newest first; m7 is Cyrillic so no match + assert.deepEqual(d.matches.map((m: any) => m.id), ['m6', 'm4']); + assert.equal(d.truncated, true); + assert.equal(d.scannedThrough, new Date(min(4)).toISOString()); + assert.match(d.hint, /to:/); + }); + + it('order:oldest truncated → from: resume point at the pool edge', async () => { + const { mod } = stub(fixture); + const d = data(await call(mod, 'search', { query: 'zzz', maxScan: 2 })); + assert.equal(d.truncated, true); + assert.equal(d.scannedThrough, new Date(min(2)).toISOString()); + assert.match(d.hint, /from:/); + }); + + it('stopping at limit also yields scannedThrough at the last scanned candidate', async () => { + const { mod } = stub(fixture); + const d = data(await call(mod, 'search', { query: 'mission', limit: 1 })); + assert.equal(d.truncated, false); + assert.equal(d.scannedThrough, new Date(min(1)).toISOString()); + }); + + it('a complete, un-limited scan carries no resume point', async () => { + const { mod } = stub(fixture); + const d = data(await call(mod, 'search', { query: 'mission' })); + assert.equal(d.scannedThrough, undefined); + assert.equal(d.order, 'oldest'); + }); + + it('rejects an unknown order', async () => { + const { mod } = stub(fixture); + assert.equal((await call(mod, 'search', { query: 'x', order: 'sideways' })).success, false); + }); +}); + +describe('HistoryModule UX: extract aroundId', () => { + const conv = Array.from({ length: 30 }, (_, i) => + mk({ id: `c${i + 1}`, ms: min(i + 1), text: `line ${i + 1}`, channel: i % 3 === 0 ? 'other' : 'diary', author: 'antra' }), + ); + + it('defaults to the anchor\'s own channel, before/after around it, anchor flagged', async () => { + const { mod } = stub(conv); + // c11 is index 10 → 10%3=1 → diary + const d = data(await call(mod, 'extract', { aroundId: 'c11', before: 2, after: 3 })); + assert.equal(d.channelId, 'diary'); + assert.deepEqual(d.messages.map((m: any) => m.id), ['c8', 'c9', 'c11', 'c12', 'c14', 'c15']); + assert.equal(d.messages.find((m: any) => m.anchor)?.id, 'c11'); + }); + + it('allChannels interleaves every channel', async () => { + const { mod } = stub(conv); + const d = data(await call(mod, 'extract', { aroundId: 'c11', before: 2, after: 2, allChannels: true })); + assert.deepEqual(d.messages.map((m: any) => m.id), ['c9', 'c10', 'c11', 'c12', 'c13']); + assert.equal(d.channelId, null); + }); + + it('clips at the ends of history', async () => { + const { mod } = stub(conv); + const d = data(await call(mod, 'extract', { aroundId: 'c2', before: 5, after: 1 })); + assert.deepEqual(d.messages.map((m: any) => m.id), ['c2', 'c3']); + }); + + it('handles same-millisecond neighbours in sequence order', async () => { + const t = min(100); + const { mod } = stub([ + mk({ id: 's1', ms: t - 1, text: 'a', channel: 'x' }), + mk({ id: 's2', ms: t, text: 'b', channel: 'x' }), + mk({ id: 's3', ms: t, text: 'c', channel: 'x' }), + mk({ id: 's4', ms: t, text: 'd', channel: 'x' }), + mk({ id: 's5', ms: t + 1, text: 'e', channel: 'x' }), + ]); + const d = data(await call(mod, 'extract', { aroundId: 's3', before: 2, after: 1 })); + assert.deepEqual(d.messages.map((m: any) => m.id), ['s1', 's2', 's3', 's4']); + }); + + it('errors cleanly: unknown id, anchor outside an explicit channel, mixing with from/to or author', async () => { + const { mod } = stub(conv); + assert.match((await call(mod, 'extract', { aroundId: 'nope' })).error!, /No message/); + assert.match((await call(mod, 'extract', { aroundId: 'c11', channelId: 'other' })).error!, /not in channel/); + assert.equal((await call(mod, 'extract', { aroundId: 'c11', from: '2026-09-20T00:00:00Z' })).success, false); + assert.equal((await call(mod, 'extract', { aroundId: 'c11', author: 'antra' })).success, false); + assert.equal((await call(mod, 'extract', { before: 3 })).success, false); + }); +}); + +describe('HistoryModule UX: real-store ordering (catch-up backfill)', () => { + // Discord catch-up appends messages late with OLDER timestamps: b1/b2 + // have high sequence numbers but timestamps inside the conversation. + const now = Date.now(); + const ago = (m: number) => now - m * 60_000; + const store = [ + mk({ id: 'r1', ms: ago(60), text: 'one', channel: 'diary', author: 'antra' }), + mk({ id: 'r2', ms: ago(50), text: 'two', channel: 'diary', author: 'antra' }), + mk({ id: 'r3', ms: ago(40), text: 'three', channel: 'diary', author: 'antra' }), + mk({ id: 'r4', ms: ago(10), text: 'four', channel: 'diary', author: 'antra' }), + mk({ id: 'r5', ms: ago(5), text: 'five', channel: 'diary', author: 'antra' }), + mk({ id: 'r90', ms: ago(45), text: 'backfilled 45', channel: 'diary', author: 'antra' }), + mk({ id: 'r91', ms: ago(20), text: 'backfilled 20', channel: 'diary', author: 'antra' }), + mk({ id: 'r6', ms: ago(30), text: 'elsewhere', channel: 'general', author: 'antra' }), + ]; + + it('aroundId in a channel uses timestamp neighbours, not append order', async () => { + const { mod } = stub(store); + const d = data(await call(mod, 'extract', { aroundId: 'r3', before: 2, after: 2 })); + assert.deepEqual(d.messages.map((m: any) => m.id), ['r2', 'r90', 'r3', 'r91', 'r4']); + }); + + it('aroundId with allChannels uses the timestamp index', async () => { + const { mod } = stub(store); + const d = data(await call(mod, 'extract', { aroundId: 'r3', before: 1, after: 2, allChannels: true })); + assert.deepEqual(d.messages.map((m: any) => m.id), ['r90', 'r3', 'r6', 'r91']); + }); + + it('order:newest without a channel does not trust a page-size matchedCount', async () => { + const { mod } = stub(store); + const d = data(await call(mod, 'search', { query: 'e', order: 'newest', maxScan: 3 })); + // newest three by timestamp: r5 (five), r4 (four), r91 (backfilled 20) + assert.equal(d.candidatePoolSize, 3); + assert.equal(d.truncated, true); + assert.deepEqual(d.matches.map((m: any) => m.id), ['r5', 'r91']); + }); + + it('order:newest in a channel picks the timestamp tail, not the append tail', async () => { + const { mod } = stub(store); + const d = data(await call(mod, 'search', { query: 'o', channelId: 'diary', order: 'newest', maxScan: 2 })); + // timestamp tail of #diary = r5, r4 (the append tail would be r91, r90) + assert.deepEqual(d.matches.map((m: any) => m.id), ['r4']); + assert.equal(d.scannedThrough, new Date(ago(10)).toISOString()); + }); + + it('plain extract without a channel reports hasMore instead of a page-size matchedCount', async () => { + const { mod } = stub(store); + const d = data(await call(mod, 'extract', { limit: 3 })); + assert.equal(d.matchedCount, undefined); + assert.equal(d.hasMore, true); + const c = data(await call(mod, 'extract', { channelId: 'diary', limit: 3 })); + assert.equal(c.matchedCount, 7); + assert.equal(c.hasMore, true); + }); +}); diff --git a/test/history-module.test.ts b/test/history-module.test.ts index 64709ca6..74fe9b45 100644 --- a/test/history-module.test.ts +++ b/test/history-module.test.ts @@ -77,6 +77,25 @@ function buildStub(messages: StoredMessage[], opts: StubOptions = {}): { cm: Con const limit = args.limit ?? filtered.length; return { messages: filtered.slice(offset, offset + limit), matchedCount }; }, + // Time-only native query (search's candidate window and extract's + // aroundId use it for timestamp-ordered edges). Real contract: + // timestamp-ordered, `reverse` supported, matchedCount = PAGE size. + queryMessagesByTime(args: { fromMs?: number; toMs?: number; limit?: number; offset?: number; reverse?: boolean }): IndexedMessageQueryResult { + calls.push({ method: 'queryMessagesByTime', args }); + if (opts.throwUnsupported) return unsupported(); + let filtered = messages + .filter((m) => { + const ts = m.timestamp.getTime(); + if (args.fromMs !== undefined && ts < args.fromMs) return false; + if (args.toMs !== undefined && ts > args.toMs) return false; + return true; + }) + .sort((a, b) => a.timestamp.getTime() - b.timestamp.getTime()); + if (args.reverse) filtered = filtered.reverse(); + const offset = args.offset ?? 0; + const page = filtered.slice(offset, args.limit === undefined ? undefined : offset + args.limit); + return { messages: page, matchedCount: page.length }; + }, getChannelMessageCounts(): ChannelCount[] { calls.push({ method: 'getChannelMessageCounts', args: undefined }); if (opts.throwUnsupported) return unsupported(); @@ -238,7 +257,8 @@ describe('HistoryModule', () => { await h.handleToolCall(call('extract', { limit: 999999 })); const queryCall = calls.find((c) => c.method === 'queryMessagesByTimeAndChannel'); - assert.equal((queryCall?.args as { limit?: number }).limit, 200); + // 200 = the hard cap; +1 is extract's own "is there another page" probe row. + assert.equal((queryCall?.args as { limit?: number }).limit, 201); }); it('clamps an out-of-u32-range offset to the native ceiling instead of passing it through to wrap (reviewer repro)', async () => { From 62b57f0641ae8ca578b53dc701d1499b8ebe1061 Mon Sep 17 00:00:00 2001 From: antra-tess Date: Fri, 25 Sep 2026 12:53:07 -0700 Subject: [PATCH 2/6] =?UTF-8?q?fix(history):=20address=20PR=20#174=20revie?= =?UTF-8?q?w=20=E2=80=94=20positional=20resume,=20exact=20edge=20windows?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - extract author/excludeAuthor: resume a truncated filtered scan by window position (`resume: {windowOffset, offset}`), not by timestamp — channel windows page in append order, so a timestamp resume skipped late-appended backfill. Fixes the pageFull hint too. scannedThrough dropped from extract. - edgeWindow: on overshoot of fetchCap, bisect the far bound (exact counts) instead of taking the ordinal tail; a single over-full millisecond is resolved by sequence, which is exact. Newest side stays open when `to` is. - aroundId: honest error for the >AROUND_TIE_CAP tie case. - author ids folded to lower case; participant is a key only for messages without author metadata; wholeWord treats \p{M} as word characters. - Changelog: matchedCount removal moved to a breaking fragment. - New real-store tests (JsStore + MessageStore, timestamps edited after append) reproducing both review findings and notes 3–6. Co-Authored-By: Claude Opus 5.5 --- changelog.d/history-author-around.added.md | 8 +- changelog.d/history-author-around.breaking.md | 4 + changelog.d/history-author-around.fixed.md | 4 - src/modules/history/index.ts | 192 +++++++++++++----- test/history-module-realstore.test.ts | 173 ++++++++++++++++ test/history-module-ux.test.ts | 6 +- 6 files changed, 324 insertions(+), 63 deletions(-) create mode 100644 changelog.d/history-author-around.breaking.md delete mode 100644 changelog.d/history-author-around.fixed.md create mode 100644 test/history-module-realstore.test.ts diff --git a/changelog.d/history-author-around.added.md b/changelog.d/history-author-around.added.md index e0fb4d45..6f9352ef 100644 --- a/changelog.d/history-author-around.added.md +++ b/changelog.d/history-author-around.added.md @@ -1,10 +1,12 @@ - History tools, from a resident's diary-work feedback: `search` and `extract` take `author` / `excludeAuthor` (exact, case-insensitive match on - `metadata.author` name or id, or the stored participant for the agent's own - turns; MCPL-ingested messages are all participant `user`, so participant - alone could not tell authors apart). Results now carry `author`. + `metadata.author` name or id; messages without author metadata, such as + the agent's own turns, match on their stored participant). Results now carry `author`. `extract({aroundId, before, after})` returns the conversation around a message id from `search` (its own channel by default, `allChannels` to interleave). `search` adds `wholeWord` (Unicode-aware) and `order: "newest"`, and when it stops early (at `limit` or `maxScan`) it reports `scannedThrough` plus a hint naming the `from`/`to` to continue with. + An author-filtered `extract` that stops early returns + `resume: {windowOffset, offset}` to repeat the call with (position, not + timestamp, so late-appended backfill in a channel is not skipped). diff --git a/changelog.d/history-author-around.breaking.md b/changelog.d/history-author-around.breaking.md new file mode 100644 index 00000000..30f288f9 --- /dev/null +++ b/changelog.d/history-author-around.breaking.md @@ -0,0 +1,4 @@ +- **Resident scripts parsing history `extract`:** without a `channelId`, + `extract` no longer returns `matchedCount` (a time-only native query only + ever reported the page size there, not a total). It returns `hasMore` + instead; channel-scoped `extract` keeps its exact `matchedCount`. diff --git a/changelog.d/history-author-around.fixed.md b/changelog.d/history-author-around.fixed.md deleted file mode 100644 index 6aa09a31..00000000 --- a/changelog.d/history-author-around.fixed.md +++ /dev/null @@ -1,4 +0,0 @@ -- History `extract` without a channel no longer reports a `matchedCount` that - was only the page size (a time-only native query returns a page-sized count, - not a total). It now reports `hasMore`, and keeps the exact `matchedCount` - only for channel-scoped queries. diff --git a/src/modules/history/index.ts b/src/modules/history/index.ts index b97a3dc3..d40f1c76 100644 --- a/src/modules/history/index.ts +++ b/src/modules/history/index.ts @@ -73,6 +73,8 @@ interface ExtractInput { author?: AuthorSpec; excludeAuthor?: AuthorSpec; maxScan?: number; + /** Resume point of an author-filtered scan (from a previous `resume`). */ + windowOffset?: number; aroundId?: string; before?: number; after?: number; @@ -142,6 +144,9 @@ const AROUND_MAX = 100; * from turning one call into an unbounded fetch. */ const AROUND_TIE_CAP = 500; +/** Largest valid Date value (ms) — an open newest edge. */ +const MAX_DATE_MS = 8.64e15; + const SEARCH_DEFAULT_LIMIT = 20; const SEARCH_MAX_LIMIT = 100; const SEARCH_DEFAULT_MAX_SCAN = 5000; @@ -322,10 +327,15 @@ export class HistoryModule implements Module { * No channel: chronicle's timestamp index answers directly (`reverse` * for the newest edge). With a channel: channel queries report an exact * matchedCount, so grow a time window out from the edge (×4 per step) - * until it holds ≥ n channel messages or covers the whole range, then - * fetch just that window and sort. A window that overshoots wildly (a - * burst) falls back to its ordinal tail rather than decoding thousands - * of messages — correct for everything but same-burst backfill. + * until it holds ≥ n channel messages or covers the whole range. If that + * window holds more than `fetchCap` (a density jump across one step, or a + * burst), bisect its far bound between the last step that held < n and + * this one — counts are exact, so it lands on a window of [n, fetchCap]. + * Only when a single millisecond holds the overflow does bisection bottom + * out; then the window is everything strictly nearer the edge plus that + * millisecond's nearest messages by sequence, which within one millisecond + * IS the (timestamp, sequence) order — still exact. The newest side is + * left open when `toMs` is (clock-skewed future stamps stay visible). */ private edgeWindow(opts: { fromMs?: number; @@ -336,51 +346,76 @@ export class HistoryModule implements Module { }): { messages: StoredMessage[]; more: boolean } { const cm = this.cm as ContextManager; const { fromMs, toMs, channelId, n, side } = opts; + const newest = side === 'newest'; const cmp = (a: StoredMessage, b: StoredMessage) => a.timestamp.getTime() - b.timestamp.getTime() || a.sequence - b.sequence; - const dirSort = (ms: StoredMessage[]) => ms.sort(side === 'newest' ? (a, b) => cmp(b, a) : cmp); + const dirSort = (ms: StoredMessage[]) => ms.sort(newest ? (a, b) => cmp(b, a) : cmp); if (channelId === undefined) { - const r = cm.queryMessagesByTime({ fromMs, toMs, limit: n + 1, reverse: side === 'newest' }); + const r = cm.queryMessagesByTime({ fromMs, toMs, limit: n + 1, reverse: newest }); return { messages: r.messages.slice(0, n), more: r.messages.length > n }; } const count = (f?: number, t?: number) => cm.queryMessagesByTimeAndChannel({ fromMs: f, toMs: t, channelId, limit: 0 }).matchedCount; + const fetch = (f: number | undefined, t: number | undefined, limit: number, offset = 0) => + limit <= 0 ? [] : cm.queryMessagesByTimeAndChannel({ fromMs: f, toMs: t, channelId, limit, offset }).messages; const total = count(fromMs, toMs); if (n === 0) return { messages: [], more: total > 0 }; - let wFrom = fromMs; - let wTo = toMs; - if (total > n) { - // Concrete bounds of the range (open ends resolved to the store's - // first message / now), so growth always terminates. - const lo = fromMs ?? this.earliestMessageMs(toMs) ?? 0; - const hi = toMs ?? Date.now(); - for (let width = 10 * 60_000; ; width *= 4) { - const f = side === 'newest' ? Math.max(hi - width, lo) : lo; - const t = side === 'newest' ? hi : Math.min(lo + width, hi); - const covered = side === 'newest' ? f <= lo : t >= hi; - if (covered) break; // whole range: keep wFrom/wTo = the range itself - if (count(f, t) >= n) { - wFrom = f; - wTo = t; - break; - } + if (total <= n) return { messages: dirSort(fetch(fromMs, toMs, total)), more: false }; + + // The window is parameterised by its FAR bound `e`: [e, toMs] for the + // newest side, [fromMs, e] for the oldest; its count is monotone in e. + const win = (e: number): [number | undefined, number | undefined] => (newest ? [e, toMs] : [fromMs, e]); + const cnt = (e: number) => count(...win(e)); + const lo = fromMs ?? this.earliestMessageMs(toMs) ?? 0; + const hi = toMs ?? Date.now(); // growth anchor only — never a bound + // `outside`: a far bound whose window holds < n (initially empty). + // `inside`: one whose window holds ≥ n (initially the whole range). + let outside: number | null = null; + let inside = newest ? lo : (toMs ?? MAX_DATE_MS); + let inCount = total; + for (let width = 10 * 60_000; ; width *= 4) { + const e = newest ? hi - width : lo + width; + if (newest ? e <= lo : e >= hi) break; // covered: keep the whole range + const c = cnt(e); + if (c >= n) { + inside = e; + inCount = c; + break; } + outside = e; } - const inWindow = count(wFrom, wTo); const fetchCap = Math.max(4 * n, n + 1000); - const page = - inWindow <= fetchCap - ? cm.queryMessagesByTimeAndChannel({ fromMs: wFrom, toMs: wTo, channelId, limit: inWindow }).messages - : cm.queryMessagesByTimeAndChannel({ - fromMs: wFrom, - toMs: wTo, - channelId, - limit: n, - offset: side === 'newest' ? inWindow - n : 0, - }).messages; - return { messages: dirSort(page).slice(0, n), more: total > n }; + if (inCount <= fetchCap) { + return { messages: dirSort(fetch(...win(inside), inCount)).slice(0, n), more: true }; + } + // Overshoot: bisect the far bound. The empty side starts one past the + // range's near end (newest: past toMs / any date; oldest: before lo). + const emptyOut = newest ? (toMs ?? MAX_DATE_MS) + 1 : lo - 1; + let out = outside ?? emptyOut; + while (Math.abs(inside - out) > 1) { + const mid = Math.floor((inside + out) / 2); + const c = cnt(mid); + if (c < n) { + out = mid; + } else { + inside = mid; + inCount = c; + if (c <= fetchCap) { + return { messages: dirSort(fetch(...win(mid), c)).slice(0, n), more: true }; + } + } + } + // Millisecond `inside` alone carries the overflow. Everything nearer the + // edge (the `out` window, < n messages) plus the nearest-by-sequence + // messages of that millisecond (channel pages are append = sequence + // ordered, which is the tie order within one millisecond). + const near = out === emptyOut ? [] : fetch(...win(out), cnt(out)); + const need = n - near.length; + const ties = count(inside, inside); + const edgeMs = fetch(inside, inside, need, newest ? Math.max(0, ties - need) : 0); + return { messages: dirSort([...near, ...edgeMs]).slice(0, n), more: true }; } async start(ctx: ModuleContext): Promise { @@ -422,7 +457,8 @@ export class HistoryModule implements Module { 'after it, in the anchor\'s own channel unless channelId/allChannels says otherwise; from/to/offset ' + 'do not apply in this mode. (2) author/excludeAuthor — keep only (or drop) messages by these ' + 'authors; this is filtered in-process over at most maxScan messages of the window, so a response ' + - 'may report truncated:true + scannedThrough (continue with from:scannedThrough).', + 'may report truncated:true with a `resume` object ({windowOffset, offset}) — repeat the call with those ' + + 'fields added to continue exactly where it stopped.', inputSchema: { type: 'object' as const, properties: { @@ -438,6 +474,7 @@ export class HistoryModule implements Module { limit: { type: 'number', description: `Max messages to return (default ${EXTRACT_DEFAULT_LIMIT}, hard cap ${EXTRACT_MAX_LIMIT}). Must be a non-negative integer.` }, offset: { type: 'number', description: `Number of matching messages to skip (default 0, capped to ${NATIVE_OFFSET_MAX}). Must be a non-negative integer.` }, maxScan: { type: 'number', description: `With author/excludeAuthor: max window messages to examine (default ${EXTRACT_FILTER_DEFAULT_MAX_SCAN}, hard cap ${EXTRACT_FILTER_MAX_MAX_SCAN}).` }, + windowOffset: { type: 'number', description: 'With author/excludeAuthor: resume a truncated scan at this position of the window — copy it from the previous response\'s `resume`.' }, format: { type: 'string', enum: ['text', 'raw'], description: 'Content rendering (default "text").' }, }, }, @@ -620,6 +657,9 @@ export class HistoryModule implements Module { const authorFilter = buildAuthorFilter(input.author, input.excludeAuthor); const cm = this.cm as ContextManager; + if (!authorFilter && input.windowOffset !== undefined) { + throw new Error('"windowOffset" only applies together with author/excludeAuthor (it resumes a filtered scan). Use offset.'); + } if (authorFilter) { // No native author index: page through the time/channel window and @@ -627,22 +667,37 @@ export class HistoryModule implements Module { // soon as the requested page is full; otherwise scan to maxScan. Only // an exhausted window yields an exact matchedCount — anything else is // reported as truncated with a resume point, never as a total. + // + // Resume is by WINDOW POSITION (`windowOffset`), never by timestamp: a + // channel-scoped window pages in append order, and Discord catch-up + // appends old-timestamped messages late, so "continue from the last + // timestamp seen" would skip them. const maxScan = clampCount(input.maxScan, EXTRACT_FILTER_DEFAULT_MAX_SCAN, EXTRACT_FILTER_MAX_MAX_SCAN, 'maxScan'); + const windowOffset = clampCount(input.windowOffset, 0, NATIVE_OFFSET_MAX, 'windowOffset'); const page: StoredMessage[] = []; let kept = 0; let scanned = 0; - let lastScanned: StoredMessage | undefined; + // Window position just past the last message put on the page. + let afterPage = windowOffset; let exhausted = false; let pageFull = false; outer: for (;;) { const want = Math.min(FILTER_SCAN_PAGE, maxScan - scanned); if (want <= 0) break; - const chunk = cm.queryMessagesByTimeAndChannel({ fromMs, toMs, channelId, limit: want, offset: scanned }).messages; + const chunk = cm.queryMessagesByTimeAndChannel({ + fromMs, + toMs, + channelId, + limit: want, + offset: Math.min(windowOffset + scanned, NATIVE_OFFSET_MAX), + }).messages; for (const m of chunk) { scanned++; - lastScanned = m; if (!authorFilter(m)) continue; - if (kept >= offset && page.length < limit) page.push(m); + if (kept >= offset && page.length < limit) { + page.push(m); + afterPage = windowOffset + scanned; + } kept++; if (page.length >= limit && kept > offset + limit) { // One filtered match beyond the page proves there is more; @@ -658,7 +713,13 @@ export class HistoryModule implements Module { } if (!exhausted && !pageFull && scanned >= maxScan) { // Probe one past maxScan: a window of exactly maxScan is complete. - const next = cm.queryMessagesByTimeAndChannel({ fromMs, toMs, channelId, limit: 1, offset: scanned }).messages; + const next = cm.queryMessagesByTimeAndChannel({ + fromMs, + toMs, + channelId, + limit: 1, + offset: Math.min(windowOffset + scanned, NATIVE_OFFSET_MAX), + }).messages; if (next.length === 0) exhausted = true; } return { @@ -668,14 +729,24 @@ export class HistoryModule implements Module { returned: page.length, scanned, truncated: !exhausted, - ...(!exhausted && lastScanned - ? { - scannedThrough: lastScanned.timestamp.toISOString(), - hint: pageFull - ? 'More matching messages exist; raise offset (or continue with from:scannedThrough and offset:0).' - : `Stopped after scanning ${scanned} messages of the window without reaching its end; continue with ` + - 'from:scannedThrough (messages at exactly that instant may repeat), or narrow with channelId/dates.', - } + ...(!exhausted + ? (() => { + // pageFull: resume just past the last returned message, no + // skip. maxScan: resume past everything scanned, still owing + // whatever part of `offset` was not yet consumed. + const resume = pageFull + ? { windowOffset: afterPage, offset: 0 } + : { windowOffset: windowOffset + scanned, offset: Math.max(0, offset - kept) }; + return { + resume, + hint: + (pageFull + ? 'More matching messages exist. ' + : `Stopped after scanning ${scanned} messages of the window without reaching its end. `) + + `Continue by repeating this call with the same from/to/channelId/author plus ` + + `windowOffset:${resume.windowOffset} and offset:${resume.offset} (the fields of \`resume\`).`, + }; + })() : {}), messages: page.map((msg) => projectMessage(msg, format)), }, @@ -715,8 +786,8 @@ export class HistoryModule implements Module { * (timestamp, sequence), and cut to `before`/`after` around the anchor. */ private handleExtractAround(input: ExtractInput): ToolResult { - if (input.from !== undefined || input.to !== undefined || input.offset !== undefined) { - throw new Error('"aroundId" cannot be combined with from/to/offset — it picks its own window.'); + if (input.from !== undefined || input.to !== undefined || input.offset !== undefined || input.windowOffset !== undefined) { + throw new Error('"aroundId" cannot be combined with from/to/offset/windowOffset — it picks its own window.'); } if (input.author !== undefined || input.excludeAuthor !== undefined) { throw new Error('"aroundId" returns the whole conversation around a message; author filters do not apply. Use search/extract with author instead.'); @@ -755,7 +826,14 @@ export class HistoryModule implements Module { ); const idx = ordered.findIndex((m) => String(m.id) === anchorId); if (idx === -1) { - // Only reachable when an explicit channelId excludes the anchor. + if (input.channelId === undefined || channelId === anchorChannel) { + // edgeWindow is exact, so the anchor can only be missed when more + // than AROUND_TIE_CAP messages share its millisecond. + throw new Error( + `Could not place message ${anchorId} among its neighbours: more than ${AROUND_TIE_CAP} messages share ` + + 'its timestamp. Use extract with from/to set to that instant instead.', + ); + } throw new Error( `Message ${anchorId} is not in channel ${JSON.stringify(input.channelId)}` + (anchorChannel ? ` (it is in ${anchorChannel}).` : '.') + @@ -1591,7 +1669,8 @@ interface SearchMatch { snippet: string; } -const WORD_CHAR_RE = /[\p{L}\p{N}_]/u; +/** Letters, combining marks (NFD accents, Indic vowel signs), digits, underscore. */ +const WORD_CHAR_RE = /[\p{L}\p{M}\p{N}_]/u; function isWordChar(ch: string | undefined): boolean { return ch !== undefined && WORD_CHAR_RE.test(ch); @@ -1678,7 +1757,12 @@ function buildAuthorFilter(include: AuthorSpec | undefined, exclude: AuthorSpec if (!inc && !exc) return null; return (m) => { const a = authorOf(m); - const keys = [a?.name?.toLowerCase(), a?.id, m.participant?.toLowerCase()].filter((k): k is string => !!k); + // participant only for messages WITHOUT author metadata (the agent's own + // turns): every MCPL-ingested message has participant "user", which + // would otherwise make author:"user" match everyone. + const keys = (a ? [a.name?.toLowerCase(), a.id?.toLowerCase()] : [m.participant?.toLowerCase()]).filter( + (k): k is string => !!k, + ); if (inc && !keys.some((k) => inc.has(k))) return false; if (exc && keys.some((k) => exc.has(k))) return false; return true; diff --git a/test/history-module-realstore.test.ts b/test/history-module-realstore.test.ts new file mode 100644 index 00000000..40b7943c --- /dev/null +++ b/test/history-module-realstore.test.ts @@ -0,0 +1,173 @@ +import { describe, it, after } from 'node:test'; +import assert from 'node:assert/strict'; +import { mkdtempSync, rmSync } from 'node:fs'; +import { tmpdir } from 'node:os'; +import { join } from 'node:path'; +import { JsStore } from '@animalabs/chronicle'; +import { MessageStore } from '@animalabs/context-manager'; +import type { ContextManager } from '@animalabs/context-manager'; +import { HistoryModule } from '../src/modules/history/index.js'; +import type { ToolCall } from '../src/types/events.js'; + +/** + * PR #174 review repros on a REAL chronicle-backed MessageStore (not a stub): + * messages are appended in one order and their timestamps rewritten after + * append (the way context-manager's own history-index tests do it), so + * append order and timestamp order genuinely disagree — Discord catch-up + * backfill. Channel-scoped native queries come back in APPEND order here. + */ + +const dirs: string[] = []; +after(() => { + for (const d of dirs) rmSync(d, { recursive: true, force: true }); +}); + +interface Row { + text: string; + channel: string; + ms: number; + author?: string; + authorId?: string; +} + +function build(rows: Row[]) { + const dir = mkdtempSync(join(tmpdir(), 'af-history-real-')); + dirs.push(dir); + const store = JsStore.openOrCreate({ path: join(dir, 'store') }); + try { + MessageStore.register(store); + } catch {} + const ms = new MessageStore(store); + const ids = new Map(); + rows.forEach((r, i) => { + const meta: Record = { channelId: r.channel }; + if (r.author) meta.author = { id: r.authorId ?? `id-${r.author}`, name: r.author }; + const m = ms.append('user', [{ type: 'text', text: r.text }], meta as never); + ids.set(r.text, String(m.id)); + const item = store.getStateItemJson('messages', i) as Record; + store.editStateItem('messages', i, Buffer.from(JSON.stringify({ ...item, timestamp: r.ms }))); + }); + const cm = { + queryMessagesByTime: (o: never) => ms.queryByTime(o), + queryMessagesByTimeAndChannel: (o: never) => ms.queryByTimeAndChannel(o), + getMessage: (id: string) => ms.get(id as never), + } as unknown as ContextManager; + const mod = new HistoryModule(); + mod.bind(cm); + return { mod, ids }; +} + +async function call(mod: HistoryModule, name: string, input: Record): Promise { + const r = await mod.handleToolCall({ id: 't', name, input } as ToolCall); + assert.equal(r.success, true, r.error); + return r.data; +} +const texts = (ms: Array<{ text?: string; content?: string }>) => ms.map((m: any) => String(m.text ?? m.content).replace(/^.*?:\s*/, '')); +const MIN = 60_000; + +describe('HistoryModule on a real store (PR #174 review)', () => { + it('finding 1: extract author+channel continuation reaches late-appended older messages', async () => { + const now = Date.now(); + const rows: Row[] = [60, 50, 40, 30, 20].map((m, i) => ({ text: `a${i + 1}`, channel: 'diary', author: 'antra', ms: now - m * MIN })); + rows.push({ text: 'b1', channel: 'diary', author: 'antra', ms: now - 55 * MIN }); + rows.push({ text: 'b2', channel: 'diary', author: 'antra', ms: now - 25 * MIN }); + const { mod } = build(rows); + const seen: string[] = []; + let d = await call(mod, 'extract', { author: 'antra', channelId: 'diary', maxScan: 5 }); + seen.push(...d.messages.map((m: any) => m.id)); + for (let guard = 0; d.truncated && guard < 10; guard++) { + assert.ok(d.resume, `truncated response must carry a resume object, got ${JSON.stringify(d)}`); + d = await call(mod, 'extract', { author: 'antra', channelId: 'diary', maxScan: 5, ...d.resume }); + seen.push(...d.messages.map((m: any) => m.id)); + } + assert.equal(new Set(seen).size, 7, `continuation must reach all 7 messages, saw ${seen.length} (${new Set(seen).size} distinct)`); + }); + + it('finding 1: pageFull continuation also resumes by position', async () => { + const now = Date.now(); + const rows: Row[] = [60, 50, 40, 30, 20].map((m, i) => ({ text: `a${i + 1}`, channel: 'diary', author: 'antra', ms: now - m * MIN })); + rows.push({ text: 'b1', channel: 'diary', author: 'antra', ms: now - 55 * MIN }); + const { mod } = build(rows); + const seen: string[] = []; + let d = await call(mod, 'extract', { author: 'antra', channelId: 'diary', limit: 2 }); + seen.push(...d.messages.map((m: any) => m.id)); + for (let guard = 0; d.truncated && guard < 10; guard++) { + d = await call(mod, 'extract', { author: 'antra', channelId: 'diary', limit: 2, ...d.resume }); + seen.push(...d.messages.map((m: any) => m.id)); + } + assert.equal(new Set(seen).size, 6); + assert.equal(seen.length, 6, 'no repeats'); + }); + + it('finding 2: newest edge window is the timestamp tail even when it overshoots fetchCap', async () => { + const now = Date.now(); + const rows: Row[] = [{ text: 'elsewhere', channel: 'x', ms: now - 30 * 24 * 60 * MIN }]; + const dense0 = now - 3 * 24 * 60 * MIN; + for (let i = 0; i < 2000; i++) rows.push({ text: `dense${i}`, channel: 'd', ms: dense0 + i * 1000 }); + [55, 45, 35, 25, 15].forEach((m, i) => rows.push({ text: `q${i + 1}`, channel: 'd', ms: now - m * MIN })); + rows.push({ text: 'OLD-backfill', channel: 'd', ms: now - 4 * 24 * 60 * MIN }); + const { mod, ids } = build(rows); + + const s = await call(mod, 'search', { query: 'e', channelId: 'd', order: 'newest', maxScan: 10 }); + // every dense* contains "e"; q* don't — check the pool edge via scannedThrough + assert.equal(s.truncated, true); + assert.equal(s.scannedThrough, new Date(dense0 + 1995 * 1000).toISOString()); + assert.ok(!s.matches.some((m: any) => m.id === ids.get('OLD-backfill'))); + + const a = await call(mod, 'extract', { aroundId: ids.get('q1'), before: 5, after: 2 }); + assert.deepEqual( + a.messages.map((m: any) => m.id), + ['dense1995', 'dense1996', 'dense1997', 'dense1998', 'dense1999', 'q1', 'q2', 'q3'].map((t) => ids.get(t)), + ); + }); + + it('finding 2: a single millisecond larger than fetchCap still yields the exact (timestamp, sequence) tail', async () => { + const now = Date.now(); + const T = now - 60 * MIN; + const rows: Row[] = [{ text: 'm-old', channel: 'x', ms: now - 30 * 24 * 60 * MIN }]; + for (let i = 0; i < 1100; i++) rows.push({ text: `m-tie${i}`, channel: 'd', ms: T }); + [30, 20, 10].forEach((m, i) => rows.push({ text: `m-late${i}`, channel: 'd', ms: now - m * MIN })); + rows.push({ text: 'm-backfill', channel: 'd', ms: T - 1000 }); + const { mod, ids } = build(rows); + const s = await call(mod, 'search', { query: 'm-', channelId: 'd', order: 'newest', maxScan: 10, limit: 50 }); + assert.deepEqual( + s.matches.map((m: any) => m.id), + ['m-late2', 'm-late1', 'm-late0', ...[1099, 1098, 1097, 1096, 1095, 1094, 1093].map((i) => `m-tie${i}`)].map((t) => ids.get(t)), + ); + }); + + it('note 6: author "user" does not match every MCPL-ingested message', async () => { + const { mod } = build([{ text: 'hello', channel: 'z', ms: Date.now() - MIN, author: 'antra' }]); + assert.equal((await call(mod, 'search', { query: 'hello', author: 'user' })).matches.length, 0); + }); + + it('note 5: channel newest window is not capped at Date.now()', async () => { + const now = Date.now(); + const rows: Row[] = [{ text: 'old', channel: 'x', ms: now - 60 * 24 * 60 * MIN }]; + for (let i = 0; i < 20; i++) rows.push({ text: `c${i}`, channel: 'c', ms: now - 5 * MIN + i * 10_000 }); + rows.push({ text: 'future', channel: 'c', ms: now + 60_000 }); + const { mod, ids } = build(rows); + const s = await call(mod, 'search', { query: 'c', channelId: 'c', order: 'newest', maxScan: 5, limit: 50 }); + const noCh = await call(mod, 'search', { query: 'u', order: 'newest', maxScan: 5, limit: 50 }); + assert.equal(noCh.matches[0].id, ids.get('future')); + assert.equal(s.scannedThrough, new Date(now - 5 * MIN + 16 * 10_000).toISOString()); + // pool = future c19 c18 c17 c16 — "future" contains no "c", so first match is c19 + assert.equal(s.matches[0].id, ids.get('c19')); + assert.equal(s.matches.length, 4); + }); + + it('note 3: author id with uppercase letters matches', async () => { + const { mod, ids } = build([{ text: 'hi', channel: 'z', ms: Date.now() - MIN, author: 'Ann', authorId: 'U0ABC' }]); + const d = await call(mod, 'search', { query: 'hi', author: 'U0ABC' }); + assert.deepEqual(d.matches.map((m: any) => m.id), [ids.get('hi')]); + }); + + it('note 4: wholeWord treats combining marks as word characters', async () => { + const { mod } = build([ + { text: 'café au lait', channel: 'z', ms: Date.now() - 2 * MIN }, + { text: 'नमस्ते', channel: 'z', ms: Date.now() - MIN }, + ]); + assert.equal((await call(mod, 'search', { query: 'cafe', wholeWord: true })).matches.length, 0); + assert.equal((await call(mod, 'search', { query: 'नम', wholeWord: true })).matches.length, 0); + }); +}); diff --git a/test/history-module-ux.test.ts b/test/history-module-ux.test.ts index f05ebc41..e3d10385 100644 --- a/test/history-module-ux.test.ts +++ b/test/history-module-ux.test.ts @@ -159,13 +159,15 @@ describe('HistoryModule UX: author filter', () => { assert.equal(all.scanned, 7); }); - it('extract author: maxScan bound reports truncated + scannedThrough, and a window of exactly maxScan is complete', async () => { + it('extract author: maxScan bound reports truncated + a positional resume, and a window of exactly maxScan is complete', async () => { const { mod } = stub(fixture); const d = data(await call(mod, 'extract', { author: 'antra', maxScan: 3 })); assert.equal(d.truncated, true); assert.equal(d.scanned, 3); - assert.equal(d.scannedThrough, new Date(min(3)).toISOString()); + assert.deepEqual(d.resume, { windowOffset: 3, offset: 0 }); assert.deepEqual(d.messages.map((m: any) => m.id), ['m1']); + const rest = data(await call(mod, 'extract', { author: 'antra', maxScan: 4, ...d.resume })); + assert.deepEqual(rest.messages.map((m: any) => m.id), ['m4', 'm5', 'm7']); const exact = data(await call(mod, 'extract', { author: 'antra', maxScan: 7 })); assert.equal(exact.truncated, false); assert.equal(exact.matchedCount, 4); From 63e82fd20aa2fda92ec0627320ec27afe81fa260 Mon Sep 17 00:00:00 2001 From: antra-tess Date: Tue, 29 Sep 2026 13:04:37 -0700 Subject: [PATCH 3/6] fix(history): address Greptile review on #174 - search: continuation carries `resume: {from|to, skip}`; `skip` steps past messages at the inclusive bound already scanned, so a limit stop or a maxScan stop inside a same-millisecond group no longer repeats forever. - extract author: a resumed scan reports matchedSinceWindowOffset, not a total; page-full counts stop at the resume point (additive across resumes); limit:0 no longer produces a non-advancing resume; maxScan:0 rejected. - extract aroundId: rejects `limit`; orders the anchor's whole millisecond by sequence before reaching outward, so large tie groups get true neighbours (cap 5000, loud error past it). - <@id> mentions accept alphanumeric ids; wholeWord reads full code points beside a match (astral letters); author descriptions match the declared array schema. - Real-store tests for all eight findings. Co-Authored-By: Claude Opus 5.5 --- changelog.d/history-author-around.added.md | 6 +- src/modules/history/index.ts | 167 ++++++++++++++------- test/history-module-realstore.test.ts | 83 ++++++++++ test/history-module-ux.test.ts | 3 +- 4 files changed, 203 insertions(+), 56 deletions(-) diff --git a/changelog.d/history-author-around.added.md b/changelog.d/history-author-around.added.md index 6f9352ef..c3c7b5f6 100644 --- a/changelog.d/history-author-around.added.md +++ b/changelog.d/history-author-around.added.md @@ -6,7 +6,9 @@ message id from `search` (its own channel by default, `allChannels` to interleave). `search` adds `wholeWord` (Unicode-aware) and `order: "newest"`, and when it stops early (at `limit` or `maxScan`) it - reports `scannedThrough` plus a hint naming the `from`/`to` to continue with. + reports `scannedThrough` and `resume: {from|to, skip}` to repeat the call + with (`skip` steps past messages at that exact instant already scanned). An author-filtered `extract` that stops early returns `resume: {windowOffset, offset}` to repeat the call with (position, not - timestamp, so late-appended backfill in a channel is not skipped). + timestamp, so late-appended backfill in a channel is not skipped); a + resumed call reports `matchedSinceWindowOffset` rather than a total. diff --git a/src/modules/history/index.ts b/src/modules/history/index.ts index d40f1c76..d2558a68 100644 --- a/src/modules/history/index.ts +++ b/src/modules/history/index.ts @@ -92,6 +92,8 @@ interface SearchInput { author?: AuthorSpec; excludeAuthor?: AuthorSpec; order?: 'oldest' | 'newest'; + /** Pool messages to skip at the start of the window (from a previous `resume`). */ + skip?: number; limit?: number; maxScan?: number; } @@ -139,10 +141,17 @@ const FILTER_SCAN_PAGE = 1000; /** `extract({aroundId})` window sizes, per side. */ const AROUND_DEFAULT = 10; const AROUND_MAX = 100; -/** Upper bound on same-millisecond neighbours fetched around an anchor — - * real stores see a handful at most; this only keeps a pathological burst - * from turning one call into an unbounded fetch. */ -const AROUND_TIE_CAP = 500; +/** Upper bound on messages sharing the anchor's millisecond that + * `extract({aroundId})` will fetch to order them by sequence — real stores + * see a handful; past this the call fails loudly rather than guessing. */ +const AROUND_TIE_CAP = 5000; + +function tieCapError(anchorId: string, n: number): Error { + return new Error( + `Message ${anchorId} shares its timestamp with ${n} messages (more than ${AROUND_TIE_CAP}), too many to order ` + + 'around it. Use extract with from/to set to that instant instead.', + ); +} /** Largest valid Date value (ms) — an open newest edge. */ const MAX_DATE_MS = 8.64e15; @@ -201,13 +210,13 @@ const AUTHOR_SPEC_JSON_SCHEMA = { type: 'array', items: { type: 'string' } }; const AUTHOR_SCHEMA = { ...AUTHOR_SPEC_JSON_SCHEMA, description: - 'Only messages written by one of these authors (a single string is fine too). Matches, case-insensitively and exactly, ' + + 'Only messages written by one of these authors (an array of names or ids). Matches, case-insensitively and exactly, ' + 'the author\'s display name or user id (a leading "@" or a <@id> mention is accepted), or the ' + 'stored participant for messages without author metadata — e.g. your own turns, under your name.', }; const EXCLUDE_AUTHOR_SCHEMA = { ...AUTHOR_SPEC_JSON_SCHEMA, - description: 'Drop messages written by any of these authors (a single string is fine too). Same matching as `author`.', + description: 'Drop messages written by any of these authors (an array of names or ids). Same matching as `author`.', }; export class HistoryModule implements Module { @@ -454,8 +463,8 @@ export class HistoryModule implements Module { 'format:"raw" returns the content blocks unmodified. ' + 'Two extra modes: (1) aroundId — pass a message id (e.g. from `search`) to get the conversation ' + 'around it: `before` messages before it, the message itself (anchor:true), and `after` messages ' + - 'after it, in the anchor\'s own channel unless channelId/allChannels says otherwise; from/to/offset ' + - 'do not apply in this mode. (2) author/excludeAuthor — keep only (or drop) messages by these ' + + 'after it, in the anchor\'s own channel unless channelId/allChannels says otherwise; from/to/offset/limit ' + + 'do not apply in this mode (at most before+after+1 messages). (2) author/excludeAuthor — keep only (or drop) messages by these ' + 'authors; this is filtered in-process over at most maxScan messages of the window, so a response ' + 'may report truncated:true with a `resume` object ({windowOffset, offset}) — repeat the call with those ' + 'fields added to continue exactly where it stopped.', @@ -509,6 +518,7 @@ export class HistoryModule implements Module { order: { type: 'string', enum: ['oldest', 'newest'], description: 'Scan and return oldest-first (default) or newest-first.' }, limit: { type: 'number', description: `Max matches to return (default ${SEARCH_DEFAULT_LIMIT}, hard cap ${SEARCH_MAX_LIMIT}). Must be a non-negative integer.` }, maxScan: { type: 'number', description: `Max candidate messages to scan (default ${SEARCH_DEFAULT_MAX_SCAN}, hard cap ${SEARCH_MAX_MAX_SCAN}). Must be a non-negative integer.` }, + skip: { type: 'number', description: 'Continuation only: messages at the start of the window already scanned by the previous call — copy it from that response\'s `resume`.' }, }, required: ['query'], }, @@ -673,6 +683,7 @@ export class HistoryModule implements Module { // appends old-timestamped messages late, so "continue from the last // timestamp seen" would skip them. const maxScan = clampCount(input.maxScan, EXTRACT_FILTER_DEFAULT_MAX_SCAN, EXTRACT_FILTER_MAX_MAX_SCAN, 'maxScan'); + if (maxScan === 0) throw new Error('"maxScan" must be at least 1 for an author-filtered extract.'); const windowOffset = clampCount(input.windowOffset, 0, NATIVE_OFFSET_MAX, 'windowOffset'); const page: StoredMessage[] = []; let kept = 0; @@ -699,7 +710,7 @@ export class HistoryModule implements Module { afterPage = windowOffset + scanned; } kept++; - if (page.length >= limit && kept > offset + limit) { + if (limit > 0 && page.length >= limit && kept > offset + limit) { // One filtered match beyond the page proves there is more; // no need to keep scanning. pageFull = true; @@ -725,7 +736,15 @@ export class HistoryModule implements Module { return { success: true, data: { - ...(exhausted ? { matchedCount: kept } : { matchedCountAtLeast: kept }), + // A resumed scan (windowOffset > 0) only saw the window from that + // position on, so its count is named for that — never a total. A + // page-full stop's lookahead match lies past the resume point and + // is left out, so counts add up across a chain of resumes. + ...(() => { + const n = pageFull ? kept - 1 : kept; + const base = windowOffset > 0 ? 'matchedSinceWindowOffset' : 'matchedCount'; + return { [exhausted ? base : `${base}AtLeast`]: n }; + })(), returned: page.length, scanned, truncated: !exhausted, @@ -779,15 +798,22 @@ export class HistoryModule implements Module { /** * `extract({aroundId})` — the conversation around one message, like - * Discord's fetch_around. Built from two index-backed time queries meeting - * at the anchor's millisecond: the newest of [.., t] and the oldest of - * [t, ..] (both via edgeWindow). Both include every message sharing t, so each side is - * over-fetched by that tie count; the union is deduped by id, ordered by - * (timestamp, sequence), and cut to `before`/`after` around the anchor. + * Discord's fetch_around. Neighbours are by (timestamp, sequence): first + * the anchor's own millisecond in sequence order, then, for whatever a + * side still needs, the nearest messages strictly before/after it (both + * via edgeWindow, which is exact by timestamp). */ private handleExtractAround(input: ExtractInput): ToolResult { - if (input.from !== undefined || input.to !== undefined || input.offset !== undefined || input.windowOffset !== undefined) { - throw new Error('"aroundId" cannot be combined with from/to/offset/windowOffset — it picks its own window.'); + if ( + input.from !== undefined || + input.to !== undefined || + input.offset !== undefined || + input.windowOffset !== undefined || + input.limit !== undefined + ) { + throw new Error( + '"aroundId" cannot be combined with from/to/offset/windowOffset/limit — it picks its own window; size it with before/after.', + ); } if (input.author !== undefined || input.excludeAuthor !== undefined) { throw new Error('"aroundId" returns the whole conversation around a message; author filters do not apply. Use search/extract with author instead.'); @@ -809,38 +835,39 @@ export class HistoryModule implements Module { const channelId = input.allChannels ? undefined : (this.resolveChannel(input.channelId) ?? anchorChannel); const t = anchor.timestamp.getTime(); - // Every message sharing the anchor's millisecond lands on BOTH sides - // below, so over-fetch each side by that count; the dedupe + sort then - // orders them by sequence around the anchor. - const ties = Math.min( - cm.queryMessagesByTimeAndChannel({ fromMs: t, toMs: t, channelId, limit: AROUND_TIE_CAP }).messages.length, - AROUND_TIE_CAP, - ); - const beforeSide = this.edgeWindow({ toMs: t, channelId, n: before + ties, side: 'newest' }).messages; - const afterSide = this.edgeWindow({ fromMs: t, channelId, n: after + ties, side: 'oldest' }).messages; - - const byId = new Map(); - for (const m of [...beforeSide, ...afterSide]) byId.set(String(m.id), m); - const ordered = [...byId.values()].sort( - (a, b) => a.timestamp.getTime() - b.timestamp.getTime() || a.sequence - b.sequence, - ); - const idx = ordered.findIndex((m) => String(m.id) === anchorId); + // The anchor's own millisecond first, whole and in sequence order — + // that is the (timestamp, sequence) order the neighbours are defined + // by, and it may hold many messages (a backfill burst). Then only if a + // side still needs more, the nearest messages strictly before / after + // that millisecond (edgeWindow is exact by timestamp). + let ties: StoredMessage[]; + if (channelId !== undefined) { + const n = cm.queryMessagesByTimeAndChannel({ fromMs: t, toMs: t, channelId, limit: 0 }).matchedCount; + if (n > AROUND_TIE_CAP) throw tieCapError(anchorId, n); + ties = cm.queryMessagesByTimeAndChannel({ fromMs: t, toMs: t, channelId, limit: n }).messages; + } else { + ties = cm.queryMessagesByTime({ fromMs: t, toMs: t, limit: AROUND_TIE_CAP + 1 }).messages; + if (ties.length > AROUND_TIE_CAP) throw tieCapError(anchorId, ties.length); + } + ties = [...ties].sort((x, y) => x.sequence - y.sequence); + const idx = ties.findIndex((m) => String(m.id) === anchorId); if (idx === -1) { - if (input.channelId === undefined || channelId === anchorChannel) { - // edgeWindow is exact, so the anchor can only be missed when more - // than AROUND_TIE_CAP messages share its millisecond. - throw new Error( - `Could not place message ${anchorId} among its neighbours: more than ${AROUND_TIE_CAP} messages share ` + - 'its timestamp. Use extract with from/to set to that instant instead.', - ); - } + // The anchor is in the store and every message of its millisecond in + // `channelId` was fetched, so only an explicit channelId excludes it. throw new Error( `Message ${anchorId} is not in channel ${JSON.stringify(input.channelId)}` + (anchorChannel ? ` (it is in ${anchorChannel}).` : '.') + ' Omit channelId to use the anchor\'s own channel.', ); } - const window = ordered.slice(Math.max(0, idx - before), idx + after + 1); + const tiesBefore = ties.slice(Math.max(0, idx - before), idx); + const tiesAfter = ties.slice(idx + 1, idx + 1 + after); + const needBefore = before - tiesBefore.length; + const needAfter = after - tiesAfter.length; + const earlier = + needBefore > 0 ? this.edgeWindow({ toMs: t - 1, channelId, n: needBefore, side: 'newest' }).messages.reverse() : []; + const later = needAfter > 0 ? this.edgeWindow({ fromMs: t + 1, channelId, n: needAfter, side: 'oldest' }).messages : []; + const window = [...earlier, ...tiesBefore, ties[idx]!, ...tiesAfter, ...later]; return { success: true, data: { @@ -900,12 +927,15 @@ export class HistoryModule implements Module { // detected up front rather than // silently scanning maxScan candidates and returning as if that were the // whole window. See the tool description's `truncated` contract. + const skip = clampCount(input.skip, 0, SEARCH_MAX_MAX_SCAN, 'skip'); let windowMessages: StoredMessage[]; + let skippedPool: StoredMessage[]; let truncated: boolean; { - const w = this.edgeWindow({ fromMs, toMs, channelId, n: maxScan, side: order }); + const w = this.edgeWindow({ fromMs, toMs, channelId, n: maxScan + skip, side: order }); truncated = w.more; - windowMessages = w.messages; + skippedPool = w.messages.slice(0, skip); + windowMessages = w.messages.slice(skip); } const candidatePoolSize = windowMessages.length; const candidates = authorFilter ? windowMessages.filter(authorFilter) : windowMessages; @@ -929,12 +959,26 @@ export class HistoryModule implements Module { if (!last) return {}; const at = last.timestamp.toISOString(); const bound = order === 'oldest' ? 'from' : 'to'; + // The continuation bound is inclusive, so the next window starts with + // every message at exactly `at` — skip the ones already scanned (the + // pool is (timestamp, sequence) ordered from the edge, so those are a + // prefix of the next window). Without this a `limit` stop, or a + // maxScan stop inside a same-millisecond group, repeats forever. + const pool = [...skippedPool, ...windowMessages]; + const stop = pool.indexOf(last); + const atMs = last.timestamp.getTime(); + let skipNext = 0; + for (let i = 0; i <= stop; i++) if (pool[i]!.timestamp.getTime() === atMs) skipNext++; + const resume = { [bound]: at, skip: skipNext }; return { scannedThrough: at, - hint: stoppedEarly - ? `Stopped at limit; more candidates remain. Continue with ${bound}:"${at}" (or raise limit).` - : `Only the ${order} ${candidatePoolSize} messages of this range were scanned (maxScan); continue with ` + - `${bound}:"${at}", raise maxScan, or narrow with channelId/author/dates.`, + resume, + hint: + (stoppedEarly + ? 'Stopped at limit; more candidates remain. ' + : `Only the ${order} ${candidatePoolSize} messages of this range were scanned (maxScan). `) + + `Continue by repeating the call with ${bound}:"${at}" and skip:${skipNext} (the fields of \`resume\`), ` + + 'or raise limit/maxScan, or narrow with channelId/author/dates.', }; }; @@ -1676,6 +1720,23 @@ function isWordChar(ch: string | undefined): boolean { return ch !== undefined && WORD_CHAR_RE.test(ch); } +/** The whole code point starting at `i` (astral letters are two UTF-16 units). */ +function codePointAt(s: string, i: number): string | undefined { + const cp = s.codePointAt(i); + return cp === undefined ? undefined : String.fromCodePoint(cp); +} + +/** The whole code point ending just before `i`. */ +function codePointBefore(s: string, i: number): string | undefined { + if (i <= 0) return undefined; + const lo = s.charCodeAt(i - 1); + if (lo >= 0xdc00 && lo <= 0xdfff && i >= 2) { + const hi = s.charCodeAt(i - 2); + if (hi >= 0xd800 && hi <= 0xdbff) return s.slice(i - 2, i); + } + return s[i - 1]; +} + /** * First occurrence of `needle` in `text`. With `wholeWord`, an occurrence * only counts when it isn't glued to a letter/digit/underscore on either @@ -1690,13 +1751,13 @@ function matchSubstring(text: string, needle: string, caseSensitive: boolean, wh const index = haystack.indexOf(needle); return index === -1 ? null : { index, length: needle.length }; } - const checkLeft = isWordChar(needle[0]); - const checkRight = isWordChar(needle[needle.length - 1]); + const checkLeft = isWordChar(codePointAt(needle, 0)); + const checkRight = isWordChar(codePointBefore(needle, needle.length)); for (let from = 0; ; ) { const index = haystack.indexOf(needle, from); if (index === -1) return null; const end = index + needle.length; - if ((!checkLeft || !isWordChar(haystack[index - 1])) && (!checkRight || !isWordChar(haystack[end]))) { + if ((!checkLeft || !isWordChar(codePointBefore(haystack, index))) && (!checkRight || !isWordChar(codePointAt(haystack, end)))) { return { index, length: needle.length }; } from = index + 1; @@ -1728,8 +1789,8 @@ function authorName(msg: StoredMessage): string | null { * <@id>/<@!id> mention to its id. */ function normalizeAuthorSpec(spec: string): string { const s = spec.trim(); - const mention = /^<@!?(\d+)>$/.exec(s); - if (mention) return mention[1]!; + const mention = /^<@!?([^\s<>@]+)>$/.exec(s); + if (mention) return mention[1]!.toLowerCase(); return s.replace(/^@/, '').toLowerCase(); } diff --git a/test/history-module-realstore.test.ts b/test/history-module-realstore.test.ts index 40b7943c..5dbaeb1e 100644 --- a/test/history-module-realstore.test.ts +++ b/test/history-module-realstore.test.ts @@ -170,4 +170,87 @@ describe('HistoryModule on a real store (PR #174 review)', () => { assert.equal((await call(mod, 'search', { query: 'cafe', wholeWord: true })).matches.length, 0); assert.equal((await call(mod, 'search', { query: 'नम', wholeWord: true })).matches.length, 0); }); + }); + +describe('Greptile review on 62b57f0 (real store)', () => { + const ids = (d: any, key = 'messages') => d[key].map((m: any) => m.id); + + it('G1: search limit continuation makes progress (distinct times and a same-ms group)', async () => { + const now = Date.now(); + const rows: Row[] = [1, 2, 3].map((i) => ({ text: `x${i}`, channel: 'c', ms: now - (10 - i) * MIN })); + for (let i = 0; i < 5; i++) rows.push({ text: `x-tie${i}`, channel: 'c', ms: now - MIN }); + const { mod } = build(rows); + for (const [order, extra] of [['oldest', { limit: 1 }], ['newest', { limit: 1 }], ['oldest', { maxScan: 2 }], ['newest', { maxScan: 2 }]] as const) { + const seen: string[] = []; + let d = await call(mod, 'search', { query: 'x', order, ...extra }); + seen.push(...ids(d, 'matches')); + for (let guard = 0; (d.resume || d.scannedThrough) && guard < 12; guard++) { + const cont = d.resume ?? { [order === 'oldest' ? 'from' : 'to']: d.scannedThrough }; + d = await call(mod, 'search', { query: 'x', order, ...extra, ...cont }); + seen.push(...ids(d, 'matches')); + } + assert.equal(new Set(seen).size, 8, `${order} ${JSON.stringify(extra)}: reached ${new Set(seen).size}/8`); + assert.equal(seen.length, 8, `${order} ${JSON.stringify(extra)}: no repeats`); + } + }); + + it('G2: a resumed, exhausted author extract does not claim the whole-window matchedCount', async () => { + const now = Date.now(); + const { mod } = build([1, 2, 3, 4].map((i) => ({ text: `a${i}`, channel: 'd', author: 'antra', ms: now - (10 - i) * MIN }))); + const first = await call(mod, 'extract', { author: 'antra', channelId: 'd', maxScan: 3 }); + const second = await call(mod, 'extract', { author: 'antra', channelId: 'd', maxScan: 3, ...first.resume }); + assert.equal(second.truncated, false); + assert.equal(second.matchedCount, undefined, 'resumed scan saw only part of the window'); + assert.equal(first.matchedCountAtLeast + second.matchedSinceWindowOffset, 4); + }); + + it('G3: aroundId rejects limit instead of silently ignoring it', async () => { + const { mod, ids: m } = build([{ text: 'a', channel: 'd', ms: Date.now() - MIN }]); + const r = await mod.handleToolCall({ id: 't', name: 'extract', input: { aroundId: m.get('a'), limit: 1 } } as ToolCall); + assert.equal(r.success, false); + }); + + it('G4: aroundId inside a same-millisecond group larger than the tie cap returns true neighbours', async () => { + const now = Date.now(); + const rows: Row[] = []; + for (let i = 0; i < 1100; i++) rows.push({ text: `t${i}`, channel: 'd', ms: now - MIN }); + const { mod, ids: m } = build(rows); + const d = await call(mod, 'extract', { aroundId: m.get('t590'), before: 3, after: 2 }); + assert.deepEqual(ids(d), [587, 588, 589, 590, 591, 592].map((i) => m.get(`t${i}`))); + }); + + it('G5: <@id> mention form works for alphanumeric ids', async () => { + const { mod, ids: m } = build([{ text: 'hi', channel: 'z', ms: Date.now() - MIN, author: 'Ann', authorId: 'U0ABC' }]); + assert.deepEqual(ids(await call(mod, 'search', { query: 'hi', author: '<@U0ABC>' }), 'matches'), [m.get('hi')]); + assert.equal((await call(mod, 'search', { query: 'hi', excludeAuthor: '<@U0ABC>' })).matches.length, 0); + }); + + it('G6: wholeWord sees astral-plane letters beside the match', async () => { + const { mod } = build([{ text: '\u{1D504}a b\u{1D504}', channel: 'z', ms: Date.now() - MIN }]); + assert.equal((await call(mod, 'search', { query: 'a', wholeWord: true })).matches.length, 0); + assert.equal((await call(mod, 'search', { query: 'b', wholeWord: true })).matches.length, 0); + }); + + it('G7: author schema does not advertise a form its declared type rejects', () => { + const tools = new HistoryModule().getTools(); + for (const t of tools) { + for (const k of ['author', 'excludeAuthor']) { + const p = (t.inputSchema as any).properties?.[k]; + if (!p) continue; + const allowsString = p.type === 'string' || (Array.isArray(p.type) && p.type.includes('string')); + assert.ok(allowsString || !/single string/i.test(p.description), `${t.name}.${k}`); + } + } + }); + + it('G8: limit:0 author extract resume advances; maxScan:0 is rejected', async () => { + const now = Date.now(); + const { mod } = build([1, 2, 3].map((i) => ({ text: `a${i}`, channel: 'd', author: 'antra', ms: now - (10 - i) * MIN }))); + const d = await call(mod, 'extract', { author: 'antra', channelId: 'd', limit: 0 }); + assert.ok(!d.truncated || d.resume.windowOffset > 0, JSON.stringify(d)); + const r = await mod.handleToolCall({ id: 't', name: 'extract', input: { author: 'antra', maxScan: 0 } } as ToolCall); + assert.equal(r.success, false); + }); +}); + diff --git a/test/history-module-ux.test.ts b/test/history-module-ux.test.ts index e3d10385..79ff5eee 100644 --- a/test/history-module-ux.test.ts +++ b/test/history-module-ux.test.ts @@ -152,7 +152,8 @@ describe('HistoryModule UX: author filter', () => { assert.equal(d.messages[0].author, 'antra'); // page full + one more antra message (m7) exists → not a total assert.equal(d.truncated, true); - assert.equal(d.matchedCountAtLeast, 4); + // counted up to the resume point; the lookahead (m7) is re-scanned by the resume + assert.equal(d.matchedCountAtLeast, 3); const all = data(await call(mod, 'extract', { author: 'antra' })); assert.equal(all.truncated, false); assert.equal(all.matchedCount, 4); From 03a0a19a9601818dcba26116ab7c4d29f41f2d61 Mon Sep 17 00:00:00 2001 From: antra-tess Date: Wed, 30 Sep 2026 09:47:44 -0700 Subject: [PATCH 4/6] fix(history): resume survives removals/insertions between calls (PR #174 re-review) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Blocker from the 09-30 re-review: `windowOffset` is a position in the LIVE window, and MessageStore.remove() splices (every discord:delete, host hide, /undo), so a removal before the cursor made the resumed page start one late and silently skip a message; an older-stamped insertion before it (a time-ordered window with backfill) made it repeat one. - extract author resume carries `afterId`, the message at windowOffset-1. On resume it is checked; if the window moved, it is re-found within 1000 positions and the scan re-anchors, reporting `windowChanged: {shift}` (a positive shift names messages this scan chain did not see). If it is gone, the call fails loudly. windowOffset without afterId is rejected. - search continuation: the `skip` count at the resume instant becomes `skipSequences: [lo, hi]`, the sequence range already scanned there. A removal can't shift the skip onto an unscanned message, and a message appended at that instant (always the highest sequence) is never mistaken for a scanned one, in either order. - search rejects maxScan:0 (truncated with no way forward). - aroundId accepts #173's `msg:` hit ids; `sum:` gets a clear error. Real-store tests for each; all 9 fail against 63e82fd (checked one by one). .Suite 955/950/1 fail: mcpl-awareness-barrier, timing under load, 25/25 alone ×3. Co-Authored-By: Claude Opus 5.5 (1M context) --- changelog.d/history-author-around.added.md | 18 +- src/modules/history/index.ts | 221 ++++++++++++++++++--- test/history-module-realstore.test.ts | 121 ++++++++++- test/history-module-ux.test.ts | 2 +- 4 files changed, 320 insertions(+), 42 deletions(-) diff --git a/changelog.d/history-author-around.added.md b/changelog.d/history-author-around.added.md index c3c7b5f6..b1585621 100644 --- a/changelog.d/history-author-around.added.md +++ b/changelog.d/history-author-around.added.md @@ -6,9 +6,17 @@ message id from `search` (its own channel by default, `allChannels` to interleave). `search` adds `wholeWord` (Unicode-aware) and `order: "newest"`, and when it stops early (at `limit` or `maxScan`) it - reports `scannedThrough` and `resume: {from|to, skip}` to repeat the call - with (`skip` steps past messages at that exact instant already scanned). + reports `scannedThrough` and `resume: {from|to, skipSequences}` to repeat + the call with (`skipSequences` names the messages at that exact instant + already scanned by sequence range, so a message added or removed there + between calls is neither lost nor shifts the skip). `search` rejects + `maxScan: 0`. An author-filtered `extract` that stops early returns - `resume: {windowOffset, offset}` to repeat the call with (position, not - timestamp, so late-appended backfill in a channel is not skipped); a - resumed call reports `matchedSinceWindowOffset` rather than a total. + `resume: {windowOffset, offset, afterId}` to repeat the call with (position, + not timestamp, so late-appended backfill in a channel is not skipped). On + resume, `afterId` pins the position: if messages were removed or inserted + before it since the previous call, the scan re-anchors and reports + `windowChanged: {shift}`; if the anchor itself is gone, it fails loudly + instead of silently skipping. A resumed call reports + `matchedSinceWindowOffset` rather than a total. + `aroundId` also accepts a `semantic_search` `msg:` hit id. diff --git a/src/modules/history/index.ts b/src/modules/history/index.ts index d2558a68..d1647fca 100644 --- a/src/modules/history/index.ts +++ b/src/modules/history/index.ts @@ -75,6 +75,8 @@ interface ExtractInput { maxScan?: number; /** Resume point of an author-filtered scan (from a previous `resume`). */ windowOffset?: number; + /** Id of the message just before `windowOffset` (from a previous `resume`). */ + afterId?: string | null; aroundId?: string; before?: number; after?: number; @@ -92,8 +94,11 @@ interface SearchInput { author?: AuthorSpec; excludeAuthor?: AuthorSpec; order?: 'oldest' | 'newest'; - /** Pool messages to skip at the start of the window (from a previous `resume`). */ - skip?: number; + /** Continuation only (from a previous `resume`): messages at exactly the + * resume bound's millisecond whose sequence lies in [lo, hi] were already + * scanned. A range, not a count, so a message added or removed at that + * millisecond between calls can't shift what gets skipped. */ + skipSequences?: [number, number]; limit?: number; maxScan?: number; } @@ -135,6 +140,9 @@ const NATIVE_OFFSET_MAX = 0xffffffff; // 4294967295 */ const EXTRACT_FILTER_DEFAULT_MAX_SCAN = 10000; const EXTRACT_FILTER_MAX_MAX_SCAN = 50000; +/** How far (in window positions) a filtered-extract resume looks for its + * `afterId` when the window moved since the previous call. */ +const RELOCATE_RADIUS = 1000; /** Native page size for the in-process filtered scans above. */ const FILTER_SCAN_PAGE = 1000; @@ -320,6 +328,66 @@ export class HistoryModule implements Module { return cm.queryMessagesByTime({ fromMs: ms, toMs: ms }).messages; } + /** Messages at exactly `ms` (in `channelId` when given). Bounded by + * AROUND_TIE_CAP: past it a same-millisecond group is too large to step + * through by sequence, and the caller should narrow instead. */ + private countAtMs(ms: number, channelId: string | undefined): number { + const cm = this.cm as ContextManager; + const n = + channelId !== undefined + ? cm.queryMessagesByTimeAndChannel({ fromMs: ms, toMs: ms, channelId, limit: 0 }).matchedCount + : cm.queryMessagesByTime({ fromMs: ms, toMs: ms, limit: AROUND_TIE_CAP + 1 }).messages.length; + if (n > AROUND_TIE_CAP) { + throw new Error( + `More than ${AROUND_TIE_CAP} messages share the instant ${new Date(ms).toISOString()}; a continuation can't step ` + + 'through them. Narrow the search with channelId/author instead.', + ); + } + return n; + } + + /** + * Validate a filtered-extract resume position against the window as it is + * NOW. `afterId` names the message the previous call saw at position + * `windowOffset - 1`. If it is still there, resume as asked. If the window + * moved (a removal or an insertion before that point since the previous + * call), look for it within RELOCATE_RADIUS positions either side and + * resume just past it, reporting the shift. If it can't be found — it was + * itself removed, or the window moved further than that — fail loudly: + * guessing is how a message gets silently skipped. + */ + private relocateWindowOffset( + windowOffset: number, + afterId: string | null | undefined, + q: { fromMs?: number; toMs?: number; channelId?: string }, + ): { windowOffset: number; shift: number } { + if (windowOffset === 0) { + if (afterId != null) throw new Error('"afterId" only applies together with a windowOffset > 0.'); + return { windowOffset, shift: 0 }; + } + if (afterId == null) { + throw new Error( + '"windowOffset" needs the "afterId" from the same `resume` object — without it a message added or removed ' + + 'since the previous call would be silently skipped. Pass the whole `resume` object back.', + ); + } + const cm = this.cm as ContextManager; + const at = cm.queryMessagesByTimeAndChannel({ ...q, limit: 1, offset: windowOffset - 1 }).messages[0]; + if (at && String(at.id) === afterId) return { windowOffset, shift: 0 }; + const lo = Math.max(0, windowOffset - 1 - RELOCATE_RADIUS); + const around = cm.queryMessagesByTimeAndChannel({ ...q, limit: 2 * RELOCATE_RADIUS + 1, offset: lo }).messages; + const i = around.findIndex((m) => String(m.id) === afterId); + if (i === -1) { + throw new Error( + `The window changed since the previous call: message ${afterId}, where that scan stopped, is no longer ` + + `within ${RELOCATE_RADIUS} positions of windowOffset ${windowOffset} (it may have been deleted). ` + + 'Restart the scan without windowOffset/afterId, or narrow it with from/to.', + ); + } + const relocated = lo + i + 1; + return { windowOffset: relocated, shift: relocated - windowOffset }; + } + /** * The `n` messages of a time/channel range nearest one of its edges, in * TIMESTAMP order walking away from that edge (newest-first for @@ -466,8 +534,9 @@ export class HistoryModule implements Module { 'after it, in the anchor\'s own channel unless channelId/allChannels says otherwise; from/to/offset/limit ' + 'do not apply in this mode (at most before+after+1 messages). (2) author/excludeAuthor — keep only (or drop) messages by these ' + 'authors; this is filtered in-process over at most maxScan messages of the window, so a response ' + - 'may report truncated:true with a `resume` object ({windowOffset, offset}) — repeat the call with those ' + - 'fields added to continue exactly where it stopped.', + 'may report truncated:true with a `resume` object ({windowOffset, offset, afterId}) — repeat the call with ' + + 'those fields added to continue exactly where it stopped. If messages were added or removed in the window ' + + 'meanwhile, the resume re-anchors on afterId and says so (windowChanged), or fails loudly when it can\'t.', inputSchema: { type: 'object' as const, properties: { @@ -476,7 +545,7 @@ export class HistoryModule implements Module { channelId: { type: 'string', description: 'Restrict to one channel. Accepts a channel label (e.g. "#general") or the raw internal channel id.' }, author: AUTHOR_SCHEMA, excludeAuthor: EXCLUDE_AUTHOR_SCHEMA, - aroundId: { type: 'string', description: 'Message id to center on (as returned by search/extract). Returns the surrounding conversation instead of a range.' }, + aroundId: { type: 'string', description: 'Message id to center on (as returned by search/extract; a "msg:" hit id works too). Returns the surrounding conversation instead of a range.' }, before: { type: 'number', description: `With aroundId: messages to include before the anchor (default ${AROUND_DEFAULT}, cap ${AROUND_MAX}).` }, after: { type: 'number', description: `With aroundId: messages to include after the anchor (default ${AROUND_DEFAULT}, cap ${AROUND_MAX}).` }, allChannels: { type: 'boolean', description: 'With aroundId: interleave every channel around the anchor\'s time instead of staying in its channel (default false).' }, @@ -484,6 +553,7 @@ export class HistoryModule implements Module { offset: { type: 'number', description: `Number of matching messages to skip (default 0, capped to ${NATIVE_OFFSET_MAX}). Must be a non-negative integer.` }, maxScan: { type: 'number', description: `With author/excludeAuthor: max window messages to examine (default ${EXTRACT_FILTER_DEFAULT_MAX_SCAN}, hard cap ${EXTRACT_FILTER_MAX_MAX_SCAN}).` }, windowOffset: { type: 'number', description: 'With author/excludeAuthor: resume a truncated scan at this position of the window — copy it from the previous response\'s `resume`.' }, + afterId: { type: 'string', description: 'With windowOffset: id of the last message the previous call scanned — copy it from `resume`. Lets a resume notice messages added or removed since, instead of silently skipping one.' }, format: { type: 'string', enum: ['text', 'raw'], description: 'Content rendering (default "text").' }, }, }, @@ -496,7 +566,8 @@ export class HistoryModule implements Module { 'to maxScan candidates — if the narrowed window is larger than maxScan, the response reports ' + 'truncated:true up front rather than silently missing later matches; narrow the filter or raise ' + 'maxScan and retry. When the scan stops early (truncated, or `limit` matches reached) the response ' + - 'carries scannedThrough: pass it as `from` (order:"oldest") or `to` (order:"newest") to continue. ' + + 'carries a `resume` object ({from|to, skipSequences}): repeat the call with its fields to continue ' + + '(scannedThrough is the same instant, for reading). ' + 'order:"newest" scans the most recent maxScan messages of the window first — usually what you want ' + 'for "when did this last come up". author/excludeAuthor narrow by who wrote the message; ' + 'wholeWord:true stops "mission" from matching inside "uncommissioned". Each match carries an `id` ' + @@ -518,7 +589,11 @@ export class HistoryModule implements Module { order: { type: 'string', enum: ['oldest', 'newest'], description: 'Scan and return oldest-first (default) or newest-first.' }, limit: { type: 'number', description: `Max matches to return (default ${SEARCH_DEFAULT_LIMIT}, hard cap ${SEARCH_MAX_LIMIT}). Must be a non-negative integer.` }, maxScan: { type: 'number', description: `Max candidate messages to scan (default ${SEARCH_DEFAULT_MAX_SCAN}, hard cap ${SEARCH_MAX_MAX_SCAN}). Must be a non-negative integer.` }, - skip: { type: 'number', description: 'Continuation only: messages at the start of the window already scanned by the previous call — copy it from that response\'s `resume`.' }, + skipSequences: { + type: 'array', + items: { type: 'number' }, + description: 'Continuation only: copy it from the previous response\'s `resume` together with its from/to. Marks the messages at exactly that instant that were already scanned.', + }, }, required: ['query'], }, @@ -667,8 +742,8 @@ export class HistoryModule implements Module { const authorFilter = buildAuthorFilter(input.author, input.excludeAuthor); const cm = this.cm as ContextManager; - if (!authorFilter && input.windowOffset !== undefined) { - throw new Error('"windowOffset" only applies together with author/excludeAuthor (it resumes a filtered scan). Use offset.'); + if (!authorFilter && (input.windowOffset !== undefined || input.afterId !== undefined)) { + throw new Error('"windowOffset"/"afterId" only apply together with author/excludeAuthor (they resume a filtered scan). Use offset.'); } if (authorFilter) { @@ -684,12 +759,25 @@ export class HistoryModule implements Module { // timestamp seen" would skip them. const maxScan = clampCount(input.maxScan, EXTRACT_FILTER_DEFAULT_MAX_SCAN, EXTRACT_FILTER_MAX_MAX_SCAN, 'maxScan'); if (maxScan === 0) throw new Error('"maxScan" must be at least 1 for an author-filtered extract.'); - const windowOffset = clampCount(input.windowOffset, 0, NATIVE_OFFSET_MAX, 'windowOffset'); + const requestedOffset = clampCount(input.windowOffset, 0, NATIVE_OFFSET_MAX, 'windowOffset'); + // A position is only meaningful against the window it was taken in: + // a removal before it (every discord:delete, host hide, /undo) shifts + // later positions left, and a message inserted before it (a + // time-ordered window with backfill) shifts them right. Either would + // silently skip or repeat a message. afterId pins the position to the + // message that was there; relocate it if the window moved. + const { windowOffset, shift } = this.relocateWindowOffset( + requestedOffset, + input.afterId, + { fromMs, toMs, channelId }, + ); const page: StoredMessage[] = []; let kept = 0; let scanned = 0; - // Window position just past the last message put on the page. + let lastScanned: StoredMessage | undefined; + // Window position just past the last message put on the page, and that message. let afterPage = windowOffset; + let lastOnPage: StoredMessage | undefined; let exhausted = false; let pageFull = false; outer: for (;;) { @@ -704,10 +792,12 @@ export class HistoryModule implements Module { }).messages; for (const m of chunk) { scanned++; + lastScanned = m; if (!authorFilter(m)) continue; if (kept >= offset && page.length < limit) { page.push(m); afterPage = windowOffset + scanned; + lastOnPage = m; } kept++; if (limit > 0 && page.length >= limit && kept > offset + limit) { @@ -748,14 +838,32 @@ export class HistoryModule implements Module { returned: page.length, scanned, truncated: !exhausted, + ...(shift !== 0 + ? { + windowChanged: { + shift, + note: + shift > 0 + ? `${shift} message(s) were added to the window before the resume point since the previous call; this scan chain did not see them.` + : `${-shift} message(s) before the resume point were removed since the previous call; the resume was re-anchored, nothing skipped.`, + }, + } + : {}), ...(!exhausted ? (() => { // pageFull: resume just past the last returned message, no // skip. maxScan: resume past everything scanned, still owing // whatever part of `offset` was not yet consumed. + // afterId = the message at position windowOffset-1 (see + // relocateWindowOffset). When nothing was scanned (no + // chunk came back), carry the incoming anchor forward. const resume = pageFull - ? { windowOffset: afterPage, offset: 0 } - : { windowOffset: windowOffset + scanned, offset: Math.max(0, offset - kept) }; + ? { windowOffset: afterPage, offset: 0, afterId: idOrNull(lastOnPage) ?? input.afterId ?? null } + : { + windowOffset: windowOffset + scanned, + offset: Math.max(0, offset - kept), + afterId: idOrNull(lastScanned) ?? input.afterId ?? null, + }; return { resume, hint: @@ -763,7 +871,7 @@ export class HistoryModule implements Module { ? 'More matching messages exist. ' : `Stopped after scanning ${scanned} messages of the window without reaching its end. `) + `Continue by repeating this call with the same from/to/channelId/author plus ` + - `windowOffset:${resume.windowOffset} and offset:${resume.offset} (the fields of \`resume\`).`, + `the fields of \`resume\` (windowOffset:${resume.windowOffset}, offset:${resume.offset}, afterId).`, }; })() : {}), @@ -826,7 +934,15 @@ export class HistoryModule implements Module { const format = input.format ?? 'text'; const cm = this.cm as ContextManager; - const anchor = cm.getMessage(String(input.aroundId)); + // semantic_search (#173) hands out `msg:` / `sum:`; accept the + // message form as-is, and say plainly what a summary id is. + const rawId = String(input.aroundId); + if (/^sum:/.test(rawId)) { + throw new Error( + `${rawId} is a compression summary, not a message — use overview for summaries, or aroundId with a msg: hit.`, + ); + } + const anchor = cm.getMessage(rawId.replace(/^msg:/, '')); if (!anchor) { throw new Error(`No message with id ${JSON.stringify(input.aroundId)} in this history.`); } @@ -927,15 +1043,30 @@ export class HistoryModule implements Module { // detected up front rather than // silently scanning maxScan candidates and returning as if that were the // whole window. See the tool description's `truncated` contract. - const skip = clampCount(input.skip, 0, SEARCH_MAX_MAX_SCAN, 'skip'); + if (maxScan === 0) throw new Error('"maxScan" must be at least 1 — a zero-message scan can never make progress.'); + // Continuation: `skipSequences` names what the previous call already + // scanned at exactly the resume bound's millisecond (the bound is + // inclusive). Identified by sequence range rather than a count, so a + // message removed at that millisecond, or appended there (a new + // message always gets the highest sequence), between calls neither + // shifts the skip onto an unscanned message nor hides the new one. + const skipRange = parseSkipSequences(input.skipSequences); + const edgeMs = order === 'oldest' ? fromMs : toMs; + if (skipRange && edgeMs === undefined) { + throw new Error(`"skipSequences" only applies together with "${order === 'oldest' ? 'from' : 'to'}" — pass the whole \`resume\` object back.`); + } + const isSkipped = (m: StoredMessage) => + !!skipRange && m.timestamp.getTime() === edgeMs && m.sequence >= skipRange[0] && m.sequence <= skipRange[1]; let windowMessages: StoredMessage[]; - let skippedPool: StoredMessage[]; let truncated: boolean; { - const w = this.edgeWindow({ fromMs, toMs, channelId, n: maxScan + skip, side: order }); - truncated = w.more; - skippedPool = w.messages.slice(0, skip); - windowMessages = w.messages.slice(skip); + // Over-fetch by the number of messages at the edge millisecond (the + // only ones that can be skipped), then drop the skipped ones. + const edgeTies = skipRange ? this.countAtMs(edgeMs!, channelId) : 0; + const w = this.edgeWindow({ fromMs, toMs, channelId, n: maxScan + edgeTies, side: order }); + const pool = skipRange ? w.messages.filter((m) => !isSkipped(m)) : w.messages; + truncated = w.more || pool.length > maxScan; + windowMessages = pool.slice(0, maxScan); } const candidatePoolSize = windowMessages.length; const candidates = authorFilter ? windowMessages.filter(authorFilter) : windowMessages; @@ -960,16 +1091,27 @@ export class HistoryModule implements Module { const at = last.timestamp.toISOString(); const bound = order === 'oldest' ? 'from' : 'to'; // The continuation bound is inclusive, so the next window starts with - // every message at exactly `at` — skip the ones already scanned (the - // pool is (timestamp, sequence) ordered from the edge, so those are a - // prefix of the next window). Without this a `limit` stop, or a - // maxScan stop inside a same-millisecond group, repeats forever. - const pool = [...skippedPool, ...windowMessages]; - const stop = pool.indexOf(last); + // every message at exactly `at`. Hand back the sequence range already + // scanned there: this call's messages at `at` up to the stop, plus — + // when `at` is still the incoming bound — the range the previous + // call had already covered. Within one millisecond the pool is in + // sequence order, so every message at `at` whose sequence lies in the + // union was scanned; anything appended later sorts outside it. const atMs = last.timestamp.getTime(); - let skipNext = 0; - for (let i = 0; i <= stop; i++) if (pool[i]!.timestamp.getTime() === atMs) skipNext++; - const resume = { [bound]: at, skip: skipNext }; + const stop = windowMessages.indexOf(last); + let lo = Infinity; + let hi = -Infinity; + for (let i = 0; i <= stop; i++) { + const m = windowMessages[i]!; + if (m.timestamp.getTime() !== atMs) continue; + lo = Math.min(lo, m.sequence); + hi = Math.max(hi, m.sequence); + } + if (skipRange && atMs === edgeMs) { + lo = Math.min(lo, skipRange[0]); + hi = Math.max(hi, skipRange[1]); + } + const resume = { [bound]: at, skipSequences: [lo, hi] as [number, number] }; return { scannedThrough: at, resume, @@ -977,7 +1119,7 @@ export class HistoryModule implements Module { (stoppedEarly ? 'Stopped at limit; more candidates remain. ' : `Only the ${order} ${candidatePoolSize} messages of this range were scanned (maxScan). `) + - `Continue by repeating the call with ${bound}:"${at}" and skip:${skipNext} (the fields of \`resume\`), ` + + `Continue by repeating the call with the fields of \`resume\` (${bound}:"${at}" and skipSequences), ` + 'or raise limit/maxScan, or narrow with channelId/author/dates.', }; }; @@ -1781,6 +1923,23 @@ function authorOf(msg: StoredMessage): { id?: string; name?: string } | undefine }; } +function parseSkipSequences(v: unknown): [number, number] | null { + if (v === undefined || v === null) return null; + if ( + !Array.isArray(v) || + v.length !== 2 || + !v.every((x) => typeof x === 'number' && Number.isInteger(x) && x >= 0) || + v[0] > v[1] + ) { + throw new Error('"skipSequences" must be [lo, hi], two non-negative integers with lo <= hi — copy it from `resume`.'); + } + return [v[0], v[1]]; +} + +function idOrNull(msg: StoredMessage | undefined): string | null { + return msg ? String(msg.id) : null; +} + function authorName(msg: StoredMessage): string | null { return authorOf(msg)?.name ?? null; } diff --git a/test/history-module-realstore.test.ts b/test/history-module-realstore.test.ts index 5dbaeb1e..7c562af0 100644 --- a/test/history-module-realstore.test.ts +++ b/test/history-module-realstore.test.ts @@ -39,14 +39,22 @@ function build(rows: Row[]) { } catch {} const ms = new MessageStore(store); const ids = new Map(); - rows.forEach((r, i) => { + // Live position of the next append: removals splice state items out. + let live = 0; + const append = (r: Row) => { const meta: Record = { channelId: r.channel }; if (r.author) meta.author = { id: r.authorId ?? `id-${r.author}`, name: r.author }; const m = ms.append('user', [{ type: 'text', text: r.text }], meta as never); ids.set(r.text, String(m.id)); - const item = store.getStateItemJson('messages', i) as Record; - store.editStateItem('messages', i, Buffer.from(JSON.stringify({ ...item, timestamp: r.ms }))); - }); + const item = store.getStateItemJson('messages', live) as Record; + store.editStateItem('messages', live, Buffer.from(JSON.stringify({ ...item, timestamp: r.ms }))); + live++; + }; + const remove = (text: string) => { + ms.remove(ids.get(text)! as never); + live--; + }; + rows.forEach(append); const cm = { queryMessagesByTime: (o: never) => ms.queryByTime(o), queryMessagesByTimeAndChannel: (o: never) => ms.queryByTimeAndChannel(o), @@ -54,7 +62,7 @@ function build(rows: Row[]) { } as unknown as ContextManager; const mod = new HistoryModule(); mod.bind(cm); - return { mod, ids }; + return { mod, ids, append, remove }; } async function call(mod: HistoryModule, name: string, input: Record): Promise { @@ -254,3 +262,106 @@ describe('Greptile review on 62b57f0 (real store)', () => { }); }); +describe('HistoryModule on a real store: window changes between resume calls (PR #174 re-review)', () => { + const rowsN = (n: number, channel = 'd') => { + const now = Date.now(); + return Array.from({ length: n }, (_, i) => ({ text: `m${i}`, channel, author: 'antra', ms: now - (100 - i) * MIN })); + }; + const texts2 = (d: any) => d.messages.map((m: any) => String(m.content).replace(/^.*?:\s*/, '')); + + for (const channelId of ['d', undefined]) { + it(`extract author resume after a removal before the cursor skips nothing (${channelId ? 'channel' : 'no channel'})`, async () => { + const { mod, remove } = build(rowsN(9)); + const q = { author: 'antra', limit: 4, ...(channelId ? { channelId } : {}) }; + const first = await call(mod, 'extract', q); + assert.deepEqual(texts2(first), ['m0', 'm1', 'm2', 'm3']); + remove('m1'); + const second = await call(mod, 'extract', { ...q, ...first.resume }); + assert.deepEqual(texts2(second), ['m4', 'm5', 'm6', 'm7']); + assert.equal(second.windowChanged?.shift, -1); + }); + } + + it('extract author resume after an older-stamped insertion before the cursor repeats nothing, and says so', async () => { + const { mod, append } = build(rowsN(9)); + const first = await call(mod, 'extract', { author: 'antra', limit: 4 }); + assert.deepEqual(texts2(first), ['m0', 'm1', 'm2', 'm3']); + // catch-up backfill: appended now, stamped between m0 and m1 → lands before the cursor in the time-ordered window + append({ text: 'late', channel: 'd', author: 'antra', ms: Date.now() - 99.5 * MIN }); + const second = await call(mod, 'extract', { author: 'antra', limit: 4, ...first.resume }); + assert.deepEqual(texts2(second), ['m4', 'm5', 'm6', 'm7']); + assert.equal(second.windowChanged?.shift, 1); + assert.match(second.windowChanged.note, /did not see/); + }); + + it('extract resume fails loudly when the anchor message itself was removed, or afterId is missing', async () => { + const { mod, remove } = build(rowsN(9)); + const first = await call(mod, 'extract', { author: 'antra', channelId: 'd', limit: 4 }); + remove('m3'); + const r = await mod.handleToolCall({ id: 't', name: 'extract', input: { author: 'antra', channelId: 'd', limit: 4, ...first.resume } } as ToolCall); + assert.equal(r.success, false); + assert.match(r.error!, /window changed/); + const noId = await mod.handleToolCall({ id: 't', name: 'extract', input: { author: 'antra', channelId: 'd', windowOffset: 4 } } as ToolCall); + assert.equal(noId.success, false); + assert.match(noId.error!, /afterId/); + }); + + const chain = async (mod: HistoryModule, base: Record, between?: (step: number) => void) => { + const seen: string[] = []; + let d = await call(mod, 'search', base); + seen.push(...d.matches.map((m: any) => m.id)); + for (let step = 0; d.resume && step < 20; step++) { + between?.(step); + d = await call(mod, 'search', { ...base, ...d.resume }); + seen.push(...d.matches.map((m: any) => m.id)); + } + return seen; + }; + + for (const channelId of ['c', undefined]) { + it(`search newest: a message appended at the resume instant between calls is not lost (${channelId ? 'channel' : 'no channel'})`, async () => { + const now = Date.now(); + const tieMs = now - MIN; + const rows: Row[] = [1, 2].map((i) => ({ text: `x-old${i}`, channel: 'c', ms: now - (10 - i) * MIN })); + for (let i = 0; i < 6; i++) rows.push({ text: `x-tie${i}`, channel: 'c', ms: tieMs }); + const { mod, append, ids: m } = build(rows); + const base = { query: 'x', order: 'newest', limit: 3, ...(channelId ? { channelId } : {}) }; + const seen = await chain(mod, base, (step) => { + if (step === 0) append({ text: 'x-tie-new', channel: 'c', ms: tieMs }); + }); + assert.ok(seen.includes(m.get('x-tie-new')!), 'appended tie reached'); + assert.equal(new Set(seen).size, 9); + assert.equal(seen.length, 9, 'no repeats'); + }); + } + + it('search oldest: removing an already-scanned message at the resume instant does not skip an unscanned one', async () => { + const now = Date.now(); + const tieMs = now - MIN; + const rows: Row[] = []; + for (let i = 0; i < 6; i++) rows.push({ text: `x-tie${i}`, channel: 'c', ms: tieMs }); + rows.push({ text: 'x-after', channel: 'c', ms: now - 0.5 * MIN }); + const { mod, remove, ids: m } = build(rows); + const seen = await chain(mod, { query: 'x', limit: 2 }, (step) => { + if (step === 0) remove('x-tie0'); + }); + const expected = ['x-tie0', 'x-tie1', 'x-tie2', 'x-tie3', 'x-tie4', 'x-tie5', 'x-after'].map((t) => m.get(t)); + assert.deepEqual(seen, expected); + }); + + it('search rejects maxScan:0 (it could never make progress)', async () => { + const { mod } = build(rowsN(2)); + const r = await mod.handleToolCall({ id: 't', name: 'search', input: { query: 'm', maxScan: 0 } } as ToolCall); + assert.equal(r.success, false); + assert.match(r.error!, /maxScan/); + }); + + it('aroundId accepts a semantic_search "msg:" hit and explains a "sum:" one', async () => { + const { mod, ids: m } = build(rowsN(5)); + const d = await call(mod, 'extract', { aroundId: `msg:${m.get('m2')}`, before: 1, after: 1 }); + assert.deepEqual(texts2(d), ['m1', 'm2', 'm3']); + const r = await mod.handleToolCall({ id: 't', name: 'extract', input: { aroundId: 'sum:42' } } as ToolCall); + assert.equal(r.success, false); + assert.match(r.error!, /summary/); + }); +}); diff --git a/test/history-module-ux.test.ts b/test/history-module-ux.test.ts index 79ff5eee..1ba4b53e 100644 --- a/test/history-module-ux.test.ts +++ b/test/history-module-ux.test.ts @@ -165,7 +165,7 @@ describe('HistoryModule UX: author filter', () => { const d = data(await call(mod, 'extract', { author: 'antra', maxScan: 3 })); assert.equal(d.truncated, true); assert.equal(d.scanned, 3); - assert.deepEqual(d.resume, { windowOffset: 3, offset: 0 }); + assert.deepEqual(d.resume, { windowOffset: 3, offset: 0, afterId: 'm3' }); assert.deepEqual(d.messages.map((m: any) => m.id), ['m1']); const rest = data(await call(mod, 'extract', { author: 'antra', maxScan: 4, ...d.resume })); assert.deepEqual(rest.messages.map((m: any) => m.id), ['m4', 'm5', 'm7']); From 26266002b58e271197eaba6bc367c8a684e0963a Mon Sep 17 00:00:00 2001 From: antra-tess Date: Fri, 2 Oct 2026 20:16:45 -0700 Subject: [PATCH 5/6] fix(history): catch balanced window changes behind a time-ordered resume Greptile (03a0a19): a removal plus an older-stamped insertion before the cursor leaves afterId at its original position, so the resume reported shift 0 and silently skipped the inserted message. Only a time-ordered (no channelId) window can take an insertion before the cursor; a channel window pages in append order, so appends land after it. The no-channel resume now carries `seqMark` (the store's highest sequence at that call). On resume, the append tail above it is checked directly; messages that landed behind the cursor are reported as windowChanged.addedBefore / missedIds even when shift is 0. More than RELOCATE_RADIUS appends since the mark fails loudly. Co-Authored-By: Claude Opus 5.5 --- changelog.d/history-author-around.added.md | 5 +- src/modules/history/index.ts | 85 ++++++++++++++++++++-- test/history-module-realstore.test.ts | 24 ++++++ test/history-module-ux.test.ts | 4 +- 4 files changed, 111 insertions(+), 7 deletions(-) diff --git a/changelog.d/history-author-around.added.md b/changelog.d/history-author-around.added.md index b1585621..1bc19f73 100644 --- a/changelog.d/history-author-around.added.md +++ b/changelog.d/history-author-around.added.md @@ -17,6 +17,9 @@ resume, `afterId` pins the position: if messages were removed or inserted before it since the previous call, the scan re-anchors and reports `windowChanged: {shift}`; if the anchor itself is gone, it fails loudly - instead of silently skipping. A resumed call reports + instead of silently skipping. Without a channel (a time-ordered window) the + resume also carries `seqMark`, and messages appended since then that landed + behind the cursor are reported in `windowChanged.missedIds` even when a + removal balanced them out (shift 0). A resumed call reports `matchedSinceWindowOffset` rather than a total. `aroundId` also accepts a `semantic_search` `msg:` hit id. diff --git a/src/modules/history/index.ts b/src/modules/history/index.ts index d1647fca..b5878805 100644 --- a/src/modules/history/index.ts +++ b/src/modules/history/index.ts @@ -77,6 +77,8 @@ interface ExtractInput { windowOffset?: number; /** Id of the message just before `windowOffset` (from a previous `resume`). */ afterId?: string | null; + /** Time-ordered (no channelId) resume only: the store's highest sequence when the previous call ran. */ + seqMark?: number; aroundId?: string; before?: number; after?: number; @@ -346,6 +348,36 @@ export class HistoryModule implements Module { return n; } + /** Highest sequence in the store right now (the last append), or -1 when empty. */ + private storeSeqMark(): number { + const cm = this.cm as ContextManager; + const count = cm.getMessageCount(); + if (count === 0) return -1; + return cm.getMessageWindow(count - 1, 1).messages[0]?.sequence ?? -1; + } + + /** + * Messages appended since `seqMark`, walking the store's append tail back + * until a sequence ≤ seqMark. Null when more than RELOCATE_RADIUS were + * appended — too many to check cheaply. + */ + private appendedSince(seqMark: number): StoredMessage[] | null { + const cm = this.cm as ContextManager; + const out: StoredMessage[] = []; + let end = cm.getMessageCount(); + while (end > 0) { + const start = Math.max(0, end - 100); + const w = cm.getMessageWindow(start, end - start).messages; + for (let i = w.length - 1; i >= 0; i--) { + if (w[i]!.sequence <= seqMark) return out; + out.push(w[i]!); + if (out.length > RELOCATE_RADIUS) return null; + } + end = start; + } + return out; + } + /** * Validate a filtered-extract resume position against the window as it is * NOW. `afterId` names the message the previous call saw at position @@ -553,6 +585,7 @@ export class HistoryModule implements Module { offset: { type: 'number', description: `Number of matching messages to skip (default 0, capped to ${NATIVE_OFFSET_MAX}). Must be a non-negative integer.` }, maxScan: { type: 'number', description: `With author/excludeAuthor: max window messages to examine (default ${EXTRACT_FILTER_DEFAULT_MAX_SCAN}, hard cap ${EXTRACT_FILTER_MAX_MAX_SCAN}).` }, windowOffset: { type: 'number', description: 'With author/excludeAuthor: resume a truncated scan at this position of the window — copy it from the previous response\'s `resume`.' }, + seqMark: { type: 'number', description: 'With windowOffset and no channelId: copy it from `resume`. Lets a resume find messages added behind it since the previous call.' }, afterId: { type: 'string', description: 'With windowOffset: id of the last message the previous call scanned — copy it from `resume`. Lets a resume notice messages added or removed since, instead of silently skipping one.' }, format: { type: 'string', enum: ['text', 'raw'], description: 'Content rendering (default "text").' }, }, @@ -742,7 +775,7 @@ export class HistoryModule implements Module { const authorFilter = buildAuthorFilter(input.author, input.excludeAuthor); const cm = this.cm as ContextManager; - if (!authorFilter && (input.windowOffset !== undefined || input.afterId !== undefined)) { + if (!authorFilter && (input.windowOffset !== undefined || input.afterId !== undefined || input.seqMark !== undefined)) { throw new Error('"windowOffset"/"afterId" only apply together with author/excludeAuthor (they resume a filtered scan). Use offset.'); } @@ -771,6 +804,38 @@ export class HistoryModule implements Module { input.afterId, { fromMs, toMs, channelId }, ); + // A time-ordered window (no channelId) can also take an older-stamped + // insertion BEFORE the cursor; paired with a removal there, positions + // balance out (shift 0) and afterId alone can't see it. Appends carry + // the highest sequences, so everything appended since the previous + // call is the store's append tail above its seqMark: check those + // directly. (A channel window pages in append order, so an append + // always lands after the cursor and is simply scanned.) + const seqMarkNow = channelId === undefined ? this.storeSeqMark() : undefined; + let missed: StoredMessage[] = []; + if (channelId === undefined && windowOffset > 0) { + if (typeof input.seqMark !== 'number') { + throw new Error( + '"windowOffset" without channelId needs the "seqMark" from the same `resume` object. Pass the whole `resume` object back.', + ); + } + const added = this.appendedSince(input.seqMark); + if (added === null) { + throw new Error( + `More than ${RELOCATE_RADIUS} messages were added since the previous call, too many to check against the ` + + 'resume point. Restart the scan without windowOffset, or narrow it with from/to/channelId.', + ); + } + const cursor = cm.queryMessagesByTimeAndChannel({ fromMs, toMs, limit: 1, offset: windowOffset - 1 }).messages[0]; + if (cursor) { + const ct = cursor.timestamp.getTime(); + missed = added.filter((m) => { + const ts = m.timestamp.getTime(); + if ((fromMs !== undefined && ts < fromMs) || (toMs !== undefined && ts > toMs)) return false; + return ts < ct || (ts === ct && m.sequence < cursor.sequence); + }); + } + } const page: StoredMessage[] = []; let kept = 0; let scanned = 0; @@ -838,13 +903,21 @@ export class HistoryModule implements Module { returned: page.length, scanned, truncated: !exhausted, - ...(shift !== 0 + ...(shift !== 0 || missed.length > 0 ? { windowChanged: { shift, + ...(channelId === undefined + ? { + addedBefore: missed.length, + // the ones this filter would have kept — fetch them with extract({aroundId}) + missedIds: missed.filter(authorFilter).map((m) => String(m.id)), + } + : {}), note: - shift > 0 - ? `${shift} message(s) were added to the window before the resume point since the previous call; this scan chain did not see them.` + missed.length > 0 || shift > 0 + ? `${missed.length || shift} message(s) were added to the window before the resume point since the previous call; this scan chain did not see them` + + (channelId === undefined ? ' (matching ones listed in missedIds).' : '.') : `${-shift} message(s) before the resume point were removed since the previous call; the resume was re-anchored, nothing skipped.`, }, } @@ -857,12 +930,14 @@ export class HistoryModule implements Module { // afterId = the message at position windowOffset-1 (see // relocateWindowOffset). When nothing was scanned (no // chunk came back), carry the incoming anchor forward. + const mark = seqMarkNow !== undefined ? { seqMark: seqMarkNow } : {}; const resume = pageFull - ? { windowOffset: afterPage, offset: 0, afterId: idOrNull(lastOnPage) ?? input.afterId ?? null } + ? { windowOffset: afterPage, offset: 0, afterId: idOrNull(lastOnPage) ?? input.afterId ?? null, ...mark } : { windowOffset: windowOffset + scanned, offset: Math.max(0, offset - kept), afterId: idOrNull(lastScanned) ?? input.afterId ?? null, + ...mark, }; return { resume, diff --git a/test/history-module-realstore.test.ts b/test/history-module-realstore.test.ts index 7c562af0..cb5cc78b 100644 --- a/test/history-module-realstore.test.ts +++ b/test/history-module-realstore.test.ts @@ -59,6 +59,8 @@ function build(rows: Row[]) { queryMessagesByTime: (o: never) => ms.queryByTime(o), queryMessagesByTimeAndChannel: (o: never) => ms.queryByTimeAndChannel(o), getMessage: (id: string) => ms.get(id as never), + getMessageCount: () => ms.length(), + getMessageWindow: (offset: number, limit: number) => ms.getWindow(offset, limit), } as unknown as ContextManager; const mod = new HistoryModule(); mod.bind(cm); @@ -294,6 +296,28 @@ describe('HistoryModule on a real store: window changes between resume calls (PR assert.match(second.windowChanged.note, /did not see/); }); + it('extract author resume notices a BALANCED change (one removal + one older-stamped insertion before the cursor)', async () => { + const { mod, append, remove, ids: m } = build(rowsN(9)); + const first = await call(mod, 'extract', { author: 'antra', limit: 4 }); + assert.deepEqual(texts2(first), ['m0', 'm1', 'm2', 'm3']); + remove('m1'); + append({ text: 'late', channel: 'd', author: 'antra', ms: Date.now() - 99.5 * MIN }); + const second = await call(mod, 'extract', { author: 'antra', limit: 4, ...first.resume }); + assert.deepEqual(texts2(second), ['m4', 'm5', 'm6', 'm7']); + assert.ok(second.windowChanged, 'a message landed behind the cursor; the response must say so'); + assert.deepEqual(second.windowChanged.missedIds, [m.get('late')]); + }); + + it('extract author resume in a channel window: a late append lands after the cursor and is simply returned', async () => { + const { mod, append, remove } = build(rowsN(9)); + const q = { author: 'antra', limit: 4, channelId: 'd' }; + const first = await call(mod, 'extract', q); + remove('m1'); + append({ text: 'late', channel: 'd', author: 'antra', ms: Date.now() - 99.5 * MIN }); + const second = await call(mod, 'extract', { ...q, limit: 10, ...first.resume }); + assert.deepEqual(texts2(second), ['m4', 'm5', 'm6', 'm7', 'm8', 'late']); + }); + it('extract resume fails loudly when the anchor message itself was removed, or afterId is missing', async () => { const { mod, remove } = build(rowsN(9)); const first = await call(mod, 'extract', { author: 'antra', channelId: 'd', limit: 4 }); diff --git a/test/history-module-ux.test.ts b/test/history-module-ux.test.ts index 1ba4b53e..cd0ae949 100644 --- a/test/history-module-ux.test.ts +++ b/test/history-module-ux.test.ts @@ -75,6 +75,8 @@ function stub(messages: StoredMessage[]) { getMessage(id: string): StoredMessage | null { return messages.find((m) => m.id === id) ?? null; }, + getMessageCount: () => bySeq.length, + getMessageWindow: (offset: number, limit: number) => ({ messages: bySeq.slice(offset, offset + limit), startIndex: offset, totalCount: bySeq.length }), } as unknown as ContextManager; const mod = new HistoryModule(); mod.bind(cm); @@ -165,7 +167,7 @@ describe('HistoryModule UX: author filter', () => { const d = data(await call(mod, 'extract', { author: 'antra', maxScan: 3 })); assert.equal(d.truncated, true); assert.equal(d.scanned, 3); - assert.deepEqual(d.resume, { windowOffset: 3, offset: 0, afterId: 'm3' }); + assert.deepEqual(d.resume, { windowOffset: 3, offset: 0, afterId: 'm3', seqMark: 7 }); assert.deepEqual(d.messages.map((m: any) => m.id), ['m1']); const rest = data(await call(mod, 'extract', { author: 'antra', maxScan: 4, ...d.resume })); assert.deepEqual(rest.messages.map((m: any) => m.id), ['m4', 'm5', 'm7']); From 9cc278b5bfb3b5ba68c460172019ba9806a90de5 Mon Sep 17 00:00:00 2001 From: antra-tess Date: Fri, 2 Oct 2026 20:37:07 -0700 Subject: [PATCH 6/6] fix(history): unrelated appends no longer break a no-channel resume Greptile (2626600): the seqMark check failed once more than 1000 messages of any kind were appended since the previous call, so a busy store broke no-channel resumes even when nothing behind the cursor changed. The append tail is now walked in pages and only messages in range and behind the cursor are kept; the walk is capped at 200,000 appends, and past it the response reports windowChanged.unverified instead of erroring. Co-Authored-By: Claude Opus 5.5 --- changelog.d/history-author-around.added.md | 4 +- src/modules/history/index.ts | 55 ++++++++++++++-------- test/history-module-realstore.test.ts | 9 ++++ 3 files changed, 48 insertions(+), 20 deletions(-) diff --git a/changelog.d/history-author-around.added.md b/changelog.d/history-author-around.added.md index 1bc19f73..1f23cb4e 100644 --- a/changelog.d/history-author-around.added.md +++ b/changelog.d/history-author-around.added.md @@ -20,6 +20,8 @@ instead of silently skipping. Without a channel (a time-ordered window) the resume also carries `seqMark`, and messages appended since then that landed behind the cursor are reported in `windowChanged.missedIds` even when a - removal balanced them out (shift 0). A resumed call reports + removal balanced them out (shift 0). Unrelated appends don't count; past + 200,000 appends since the mark the check is reported as + `windowChanged.unverified` instead of failing. A resumed call reports `matchedSinceWindowOffset` rather than a total. `aroundId` also accepts a `semantic_search` `msg:` hit id. diff --git a/src/modules/history/index.ts b/src/modules/history/index.ts index b5878805..176d4462 100644 --- a/src/modules/history/index.ts +++ b/src/modules/history/index.ts @@ -145,6 +145,10 @@ const EXTRACT_FILTER_MAX_MAX_SCAN = 50000; /** How far (in window positions) a filtered-extract resume looks for its * `afterId` when the window moved since the previous call. */ const RELOCATE_RADIUS = 1000; +/** Most appends a no-channel resume walks back through to check for + * messages added behind its cursor; past it the check is reported as + * unverified rather than failing. */ +const APPENDED_SCAN_CAP = 200_000; /** Native page size for the in-process filtered scans above. */ const FILTER_SCAN_PAGE = 1000; @@ -357,25 +361,31 @@ export class HistoryModule implements Module { } /** - * Messages appended since `seqMark`, walking the store's append tail back - * until a sequence ≤ seqMark. Null when more than RELOCATE_RADIUS were - * appended — too many to check cheaply. + * Of the messages appended since `seqMark` (the store's append tail above + * it, walked back in pages), the ones `relevant` keeps. Unrelated appends + * cost a look but never count. `complete` is false when the walk stopped + * at APPENDED_SCAN_CAP before reaching seqMark — the result is then a + * lower bound, not a verification. */ - private appendedSince(seqMark: number): StoredMessage[] | null { + private appendedSince( + seqMark: number, + relevant: (m: StoredMessage) => boolean, + ): { messages: StoredMessage[]; complete: boolean } { const cm = this.cm as ContextManager; const out: StoredMessage[] = []; + let looked = 0; let end = cm.getMessageCount(); while (end > 0) { - const start = Math.max(0, end - 100); + const start = Math.max(0, end - 1000); const w = cm.getMessageWindow(start, end - start).messages; for (let i = w.length - 1; i >= 0; i--) { - if (w[i]!.sequence <= seqMark) return out; - out.push(w[i]!); - if (out.length > RELOCATE_RADIUS) return null; + if (w[i]!.sequence <= seqMark) return { messages: out, complete: true }; + if (relevant(w[i]!)) out.push(w[i]!); + if (++looked >= APPENDED_SCAN_CAP) return { messages: out, complete: false }; } end = start; } - return out; + return { messages: out, complete: true }; } /** @@ -813,27 +823,23 @@ export class HistoryModule implements Module { // always lands after the cursor and is simply scanned.) const seqMarkNow = channelId === undefined ? this.storeSeqMark() : undefined; let missed: StoredMessage[] = []; + let unverified = false; if (channelId === undefined && windowOffset > 0) { if (typeof input.seqMark !== 'number') { throw new Error( '"windowOffset" without channelId needs the "seqMark" from the same `resume` object. Pass the whole `resume` object back.', ); } - const added = this.appendedSince(input.seqMark); - if (added === null) { - throw new Error( - `More than ${RELOCATE_RADIUS} messages were added since the previous call, too many to check against the ` + - 'resume point. Restart the scan without windowOffset, or narrow it with from/to/channelId.', - ); - } const cursor = cm.queryMessagesByTimeAndChannel({ fromMs, toMs, limit: 1, offset: windowOffset - 1 }).messages[0]; if (cursor) { const ct = cursor.timestamp.getTime(); - missed = added.filter((m) => { + const r = this.appendedSince(input.seqMark, (m) => { const ts = m.timestamp.getTime(); if ((fromMs !== undefined && ts < fromMs) || (toMs !== undefined && ts > toMs)) return false; return ts < ct || (ts === ct && m.sequence < cursor.sequence); }); + missed = r.messages; + unverified = !r.complete; } } const page: StoredMessage[] = []; @@ -903,10 +909,19 @@ export class HistoryModule implements Module { returned: page.length, scanned, truncated: !exhausted, - ...(shift !== 0 || missed.length > 0 + ...(shift !== 0 || missed.length > 0 || unverified ? { windowChanged: { shift, + ...(unverified + ? { + unverified: true, + unverifiedNote: + `More than ${APPENDED_SCAN_CAP} messages were appended since the previous call; only the newest ` + + 'were checked, so messages may have been added behind the resume point unseen. Restart without ' + + 'windowOffset, or narrow with from/to/channelId, if that matters.', + } + : {}), ...(channelId === undefined ? { addedBefore: missed.length, @@ -918,7 +933,9 @@ export class HistoryModule implements Module { missed.length > 0 || shift > 0 ? `${missed.length || shift} message(s) were added to the window before the resume point since the previous call; this scan chain did not see them` + (channelId === undefined ? ' (matching ones listed in missedIds).' : '.') - : `${-shift} message(s) before the resume point were removed since the previous call; the resume was re-anchored, nothing skipped.`, + : shift < 0 + ? `${-shift} message(s) before the resume point were removed since the previous call; the resume was re-anchored, nothing skipped.` + : 'No change found behind the resume point, but the check was incomplete (see unverifiedNote).', }, } : {}), diff --git a/test/history-module-realstore.test.ts b/test/history-module-realstore.test.ts index cb5cc78b..cdd17861 100644 --- a/test/history-module-realstore.test.ts +++ b/test/history-module-realstore.test.ts @@ -308,6 +308,15 @@ describe('HistoryModule on a real store: window changes between resume calls (PR assert.deepEqual(second.windowChanged.missedIds, [m.get('late')]); }); + it('extract author resume (no channel) is not broken by many unrelated appends since seqMark', async () => { + const { mod, append } = build(rowsN(9)); + const first = await call(mod, 'extract', { author: 'antra', limit: 4 }); + for (let i = 0; i < 1500; i++) append({ text: `busy${i}`, channel: 'other', author: 'bob', ms: Date.now() }); + const second = await call(mod, 'extract', { author: 'antra', limit: 4, ...first.resume }); + assert.deepEqual(texts2(second), ['m4', 'm5', 'm6', 'm7']); + assert.equal(second.windowChanged, undefined); + }); + it('extract author resume in a channel window: a late append lands after the cursor and is simply returned', async () => { const { mod, append, remove } = build(rowsN(9)); const q = { author: 'antra', limit: 4, channelId: 'd' };