diff --git a/changelog.d/history-author-around.added.md b/changelog.d/history-author-around.added.md new file mode 100644 index 00000000..1f23cb4e --- /dev/null +++ b/changelog.d/history-author-around.added.md @@ -0,0 +1,27 @@ +- 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; 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` 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, 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. 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). 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/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/src/modules/history/index.ts b/src/modules/history/index.ts index 8ea36121..176d4462 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,37 @@ interface ExtractInput { limit?: number; offset?: number; format?: 'text' | 'raw'; + author?: AuthorSpec; + excludeAuthor?: AuthorSpec; + 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; + /** Time-ordered (no channelId) resume only: the store's highest sequence when the previous call ran. */ + seqMark?: 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'; + /** 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; } @@ -109,6 +134,42 @@ 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; +/** 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; + +/** `extract({aroundId})` window sizes, per side. */ +const AROUND_DEFAULT = 10; +const AROUND_MAX = 100; +/** 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; + const SEARCH_DEFAULT_LIMIT = 20; const SEARCH_MAX_LIMIT = 100; const SEARCH_DEFAULT_MAX_SCAN = 5000; @@ -157,6 +218,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 (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 (an array of names or ids). Same matching as `author`.', +}; + export class HistoryModule implements Module { readonly name = 'history'; @@ -258,6 +334,209 @@ 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; + } + + /** 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; + } + + /** + * 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, + 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 - 1000); + const w = cm.getMessageWindow(start, end - start).messages; + for (let i = w.length - 1; i >= 0; i--) { + 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 { messages: out, complete: true }; + } + + /** + * 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 + * 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. 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; + 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 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(newest ? (a, b) => cmp(b, a) : cmp); + + if (channelId === undefined) { + 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 }; + 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 fetchCap = Math.max(4 * n, n + 1000); + 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 { this.ctx = ctx; } @@ -291,15 +570,33 @@ 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/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, 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: { 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; 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).' }, 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`.' }, + 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").' }, }, }, @@ -311,7 +608,14 @@ 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 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` ' + + '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,11 +623,20 @@ 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.` }, + 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'], }, @@ -459,28 +772,320 @@ 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 && (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.'); + } + + 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. + // + // 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'); + if (maxScan === 0) throw new Error('"maxScan" must be at least 1 for an author-filtered extract.'); + 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 }, + ); + // 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[] = []; + 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 cursor = cm.queryMessagesByTimeAndChannel({ fromMs, toMs, limit: 1, offset: windowOffset - 1 }).messages[0]; + if (cursor) { + const ct = cursor.timestamp.getTime(); + 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[] = []; + let kept = 0; + let scanned = 0; + 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 (;;) { + const want = Math.min(FILTER_SCAN_PAGE, maxScan - scanned); + if (want <= 0) break; + 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); + afterPage = windowOffset + scanned; + lastOnPage = m; + } + kept++; + 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; + 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: Math.min(windowOffset + scanned, NATIVE_OFFSET_MAX), + }).messages; + if (next.length === 0) exhausted = true; + } + return { + success: true, + data: { + // 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, + ...(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, + // the ones this filter would have kept — fetch them with extract({aroundId}) + missedIds: missed.filter(authorFilter).map((m) => String(m.id)), + } + : {}), + note: + 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 < 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).', + }, + } + : {}), + ...(!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 mark = seqMarkNow !== undefined ? { seqMark: seqMarkNow } : {}; + const resume = pageFull + ? { 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, + 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 ` + + `the fields of \`resume\` (windowOffset:${resume.windowOffset}, offset:${resume.offset}, afterId).`, + }; + })() + : {}), + 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. 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 || + 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.'); + } + 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; + + // 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.`); + } + const anchorId = String(anchor.id); + const anchorChannel = getChannelId(anchor); + const channelId = input.allChannels ? undefined : (this.resolveChannel(input.channelId) ?? anchorChannel); + const t = anchor.timestamp.getTime(); + + // 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) { + // 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 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: { - 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 +1105,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 +1129,92 @@ 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; + 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 truncated: boolean; + { + // 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; + 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'; + // The continuation bound is inclusive, so the next window starts with + // 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(); + 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, + 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 the fields of \`resume\` (${bound}:"${at}" and skipSequences), ` + + 'or raise limit/maxScan, or narrow with channelId/author/dates.', + }; + }; if (input.regex) { // Regex matching against caller-supplied patterns is ReDoS-shaped: a @@ -541,7 +1227,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 +1245,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 +1259,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 +1319,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 +1888,7 @@ function projectMessage(msg: StoredMessage, format: 'text' | 'raw'): Record= 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 + * side (Unicode-aware, so Cyrillic words work too) — "mission" no longer + * hits inside "uncommissioned". Word-ness is checked only on the sides + * where the needle itself starts/ends with a word character, so a needle + * like "#general" or "v2." still matches the way a person would expect. + */ +function matchSubstring(text: string, needle: string, caseSensitive: boolean, wholeWord = false): MatchHit | null { const haystack = caseSensitive ? text : text.toLowerCase(); - const index = haystack.indexOf(needle); - if (index === -1) return null; - return { index, length: needle.length }; + if (!wholeWord) { + const index = haystack.indexOf(needle); + return index === -1 ? null : { index, length: needle.length }; + } + 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(codePointBefore(haystack, index))) && (!checkRight || !isWordChar(codePointAt(haystack, end)))) { + return { index, length: needle.length }; + } + from = index + 1; + } +} + +// ============================================================================ +// author filtering +// ============================================================================ + +/** metadata.author as written by MCPL channel ingestion ({id, name}). Every + * incoming channel message is stored with participant "user" — the real + * author lives only here. */ +function authorOf(msg: StoredMessage): { id?: string; name?: string } | undefined { + const a = (msg.metadata as { author?: unknown } | undefined)?.author; + if (!a || typeof a !== 'object') return undefined; + const { id, name } = a as { id?: unknown; name?: unknown }; + return { + ...(typeof id === 'string' || typeof id === 'number' ? { id: String(id) } : {}), + ...(typeof name === 'string' ? { name } : {}), + }; +} + +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; +} + +/** Normalize one author spec: trim, case-fold, strip a leading "@", unwrap a + * <@id>/<@!id> mention to its id. */ +function normalizeAuthorSpec(spec: string): string { + const s = spec.trim(); + const mention = /^<@!?([^\s<>@]+)>$/.exec(s); + if (mention) return mention[1]!.toLowerCase(); + 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); + // 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; + }; } /** ~SNIPPET_CONTEXT_CHARS of surrounding context on each side of a match, or diff --git a/test/history-module-realstore.test.ts b/test/history-module-realstore.test.ts new file mode 100644 index 00000000..cdd17861 --- /dev/null +++ b/test/history-module-realstore.test.ts @@ -0,0 +1,400 @@ +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(); + // 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', 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), + 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); + return { mod, ids, append, remove }; +} + +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); + }); + +}); + +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); + }); +}); + +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 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 (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' }; + 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 }); + 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 new file mode 100644 index 00000000..cd0ae949 --- /dev/null +++ b/test/history-module-ux.test.ts @@ -0,0 +1,356 @@ +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; + }, + 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); + 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); + // 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); + assert.equal(all.scanned, 7); + }); + + 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.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']); + 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 () => {