mirror of
https://github.com/garrytan/gbrain.git
synced 2026-08-15 17:32:37 +00:00
Compare commits
3
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
bcdb435d73 | ||
|
|
093d693502 | ||
|
|
9268552d70 |
-43
@@ -466,11 +466,6 @@ async function main() {
|
||||
const result = JSON.parse(JSON.stringify(rawResult, bigintToStringReplacer));
|
||||
const output = formatResult(op.name, result);
|
||||
if (output) process.stdout.write(output);
|
||||
// #1484 — invisible-miss hint: a bare query/search that hit zero results
|
||||
// on a multi-source brain tells the user (stderr) which source it
|
||||
// actually searched and how to widen the scope.
|
||||
const hint = await sourceScopeHint(op.name, params, ctx.sourceId, engine, result);
|
||||
if (hint) console.error(hint);
|
||||
} catch (e: unknown) {
|
||||
// v0.42.20.0 (codex D4): on error, set exitCode + return so the `finally`
|
||||
// STILL runs (drains every background-work sink + disconnects). A bare
|
||||
@@ -842,44 +837,6 @@ async function makeContext(engine: BrainEngine, params: Record<string, unknown>)
|
||||
};
|
||||
}
|
||||
|
||||
/**
|
||||
* #1484 — a bare `gbrain query`/`search` silently scopes to the resolved
|
||||
* source (usually 'default'); on a multi-source brain a zero-hit run looks
|
||||
* identical to "the brain doesn't know this" even when the answer lives in
|
||||
* another source. Returns a stderr hint when (a) the op is query/search,
|
||||
* (b) it returned zero results, (c) the caller did NOT scope explicitly
|
||||
* (--source / --source-id / --all-sources), and (d) the brain has >1
|
||||
* registered source. Best-effort: any lookup failure returns null.
|
||||
*
|
||||
* Exported for tests (same import-safety contract as formatResult).
|
||||
*/
|
||||
export async function sourceScopeHint(
|
||||
opName: string,
|
||||
params: Record<string, unknown>,
|
||||
sourceId: string,
|
||||
engine: BrainEngine,
|
||||
result: unknown,
|
||||
): Promise<string | null> {
|
||||
if (opName !== 'query' && opName !== 'search') return null;
|
||||
if (!Array.isArray(result) || result.length > 0) return null;
|
||||
// Explicit scoping (flag tier) = user intent; don't second-guess it.
|
||||
if (params.source || params.source_id || params.all_sources) return null;
|
||||
if (sourceId === '__all__') return null;
|
||||
try {
|
||||
const rows = await engine.executeRaw<{ n: number }>(
|
||||
`SELECT count(*)::int AS n FROM sources`,
|
||||
);
|
||||
const n = Number(rows[0]?.n ?? 0);
|
||||
if (n <= 1) return null;
|
||||
return (
|
||||
`Hint: this brain has ${n} sources; you searched only "${sourceId}". ` +
|
||||
`Retry with --source-id __all__ (all sources) or --source-id <id>.`
|
||||
);
|
||||
} catch {
|
||||
return null; // hint is best-effort; never fail the query over it
|
||||
}
|
||||
}
|
||||
|
||||
// Exported for tests (same import-safety contract as cliAliases/printOpHelp).
|
||||
export function formatResult(opName: string, result: unknown): string {
|
||||
switch (opName) {
|
||||
|
||||
+20
-49
@@ -1000,24 +1000,9 @@ const voyageCompatFetch = (async (input: RequestInfo | URL, init?: RequestInit)
|
||||
// Voyage diverges from OpenAI in two places that break the parser:
|
||||
// - `embedding` is a base64 string (SDK schema expects `number[]`)
|
||||
// - `usage` lacks `prompt_tokens` (SDK schema requires it when usage present)
|
||||
//
|
||||
// #1610: read the body ONCE via text() and JSON.parse it. The pre-fix
|
||||
// `await resp.clone().json()` truncated large bodies on bun < 1.1.27
|
||||
// (oven-sh/bun#6348) — the parse threw, the catch fell back to the raw
|
||||
// response, and multi-chunk pages died with "Invalid JSON response".
|
||||
// Every JSON return path below rebuilds the Response so a stale
|
||||
// Content-Length/Content-Encoding header from the original can't lie
|
||||
// about the rewritten body.
|
||||
const bodyText = await resp.text();
|
||||
const rebuild = (body: string) => {
|
||||
const headers = new Headers(resp.headers);
|
||||
headers.delete('content-length');
|
||||
headers.delete('content-encoding');
|
||||
return new Response(body, { status: resp.status, statusText: resp.statusText, headers });
|
||||
};
|
||||
try {
|
||||
const json: any = JSON.parse(bodyText);
|
||||
if (!json || typeof json !== 'object') return rebuild(bodyText);
|
||||
const json: any = await resp.clone().json();
|
||||
if (!json || typeof json !== 'object') return resp;
|
||||
let modified = false;
|
||||
if (Array.isArray(json.data)) {
|
||||
for (const item of json.data) {
|
||||
@@ -1052,19 +1037,22 @@ const voyageCompatFetch = (async (input: RequestInfo | URL, init?: RequestInit)
|
||||
: 0;
|
||||
modified = true;
|
||||
}
|
||||
if (!modified) return rebuild(bodyText);
|
||||
return rebuild(JSON.stringify(json));
|
||||
if (!modified) return resp;
|
||||
return new Response(JSON.stringify(json), {
|
||||
status: resp.status,
|
||||
statusText: resp.statusText,
|
||||
headers: resp.headers,
|
||||
});
|
||||
} catch (err) {
|
||||
// OOM-cap throws MUST propagate. The catch is here for "Voyage returned
|
||||
// JSON I can't reshape" (parse error, unexpected schema) — falling back
|
||||
// to the original body is correct in that case. Letting the
|
||||
// to the original response is correct in that case. Letting the
|
||||
// too-large response through here would defeat the entire purpose of
|
||||
// Layer 2 (the per-embedding cap that fires when Content-Length wasn't
|
||||
// available to Layer 1).
|
||||
if (err instanceof VoyageResponseTooLargeError) throw err;
|
||||
// If parsing/transformation fails, pass the original body through
|
||||
// (rebuilt — resp's body stream is already consumed by text()).
|
||||
return rebuild(bodyText);
|
||||
// If parsing/transformation fails, fall back to the original response.
|
||||
return resp;
|
||||
}
|
||||
}) as unknown as typeof fetch;
|
||||
|
||||
@@ -1204,21 +1192,9 @@ const zeroEntropyCompatFetch = (async (input: RequestInfo | URL, init?: RequestI
|
||||
// validates. Also map usage.total_tokens → prompt_tokens (SDK requires
|
||||
// prompt_tokens when `usage` is present — same divergence Voyage hit at
|
||||
// gateway.ts:655).
|
||||
//
|
||||
// #1610: read the body ONCE via text() + JSON.parse — `resp.clone().json()`
|
||||
// truncated large bodies on bun < 1.1.27 (oven-sh/bun#6348), so the parse
|
||||
// threw and the catch fell back to the RAW ZE `{results: ...}` shape, which
|
||||
// the AI SDK schema rejects → "Invalid JSON response" on multi-chunk pages.
|
||||
const bodyText = await resp.text();
|
||||
const rebuild = (body: string) => {
|
||||
const headers = new Headers(resp.headers);
|
||||
headers.delete('content-length');
|
||||
headers.delete('content-encoding');
|
||||
return new Response(body, { status: resp.status, statusText: resp.statusText, headers });
|
||||
};
|
||||
try {
|
||||
const json: any = JSON.parse(bodyText);
|
||||
if (!json || typeof json !== 'object') return rebuild(bodyText);
|
||||
const json: any = await resp.clone().json();
|
||||
if (!json || typeof json !== 'object') return resp;
|
||||
let modified = false;
|
||||
if (Array.isArray(json.results) && !Array.isArray(json.data)) {
|
||||
// Layer 2 OOM cap — per-embedding size. ZE returns float[] arrays,
|
||||
@@ -1252,25 +1228,20 @@ const zeroEntropyCompatFetch = (async (input: RequestInfo | URL, init?: RequestI
|
||||
// SDK also expects total_tokens; ZE provides it directly.
|
||||
modified = true;
|
||||
}
|
||||
if (!modified) return rebuild(bodyText);
|
||||
return rebuild(JSON.stringify(json));
|
||||
if (!modified) return resp;
|
||||
return new Response(JSON.stringify(json), {
|
||||
status: resp.status,
|
||||
statusText: resp.statusText,
|
||||
headers: resp.headers,
|
||||
});
|
||||
} catch (err) {
|
||||
// OOM-cap throws MUST propagate. Voyage's pattern: instanceof check on
|
||||
// its own tagged class. Same here — only rethrow our own cap class.
|
||||
if (err instanceof ZeroEntropyResponseTooLargeError) throw err;
|
||||
return rebuild(bodyText);
|
||||
return resp;
|
||||
}
|
||||
}) as unknown as typeof fetch;
|
||||
|
||||
/**
|
||||
* Test-only seams (#1610): the compat shims are module-private closures;
|
||||
* exporting them lets tests drive the response-rewrite paths behaviorally
|
||||
* (truncating clone(), stale Content-Length) without a live provider.
|
||||
* Same pattern as __getShrinkStateForTests.
|
||||
*/
|
||||
export const __voyageCompatFetchForTests = voyageCompatFetch;
|
||||
export const __zeroEntropyCompatFetchForTests = zeroEntropyCompatFetch;
|
||||
|
||||
/**
|
||||
* Generic asymmetric-embedding shim for openai-compatible recipes that
|
||||
* ship no compat fetch of their own (llama-server, litellm, ollama, ...).
|
||||
|
||||
+101
-20
@@ -428,8 +428,13 @@ export async function runPhaseSynthesize(
|
||||
|
||||
const queue = new MinionQueue(engine);
|
||||
const childIds: number[] = [];
|
||||
/** Map child job_id → chunk metadata for D6 orchestrator-side slug rewrite. */
|
||||
const chunkInfo = new Map<number, { idx: number; hash6: string }>();
|
||||
/**
|
||||
* Map child job_id → transcript metadata. Drives D6 orchestrator-side
|
||||
* slug rewrite for chunked transcripts AND the deterministic frontmatter
|
||||
* stampDreamProvenance merges into each written page. Populated for
|
||||
* every child (single-chunk children carry chunkTotal=1).
|
||||
*/
|
||||
const childMeta = new Map<number, ChildMeta>();
|
||||
/** Skip reasons for the cycle report (D5 cap hits, D8 legacy-key skips). */
|
||||
const skipReports: Array<{ filePath: string; reason: string }> = [];
|
||||
|
||||
@@ -513,9 +518,14 @@ export async function runPhaseSynthesize(
|
||||
{ allowProtectedSubmit: true },
|
||||
);
|
||||
childIds.push(child.id);
|
||||
if (isChunked) {
|
||||
chunkInfo.set(child.id, { idx: i, hash6 });
|
||||
}
|
||||
childMeta.set(child.id, {
|
||||
idx: i,
|
||||
hash6,
|
||||
chunkTotal: chunks.length,
|
||||
transcriptSource: t.transcriptSource,
|
||||
transcriptId: stripContentVersionSuffix(t.basename),
|
||||
inferredDate: t.inferredDate,
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
@@ -544,14 +554,14 @@ export async function runPhaseSynthesize(
|
||||
|
||||
// Collect slugs from put_page tool executions across the children
|
||||
// (codex finding #2: deterministic provenance, NOT pages.updated_at).
|
||||
// D6 orchestrator slug rewrite: chunkInfo drives post-hoc rewrite of
|
||||
// D6 orchestrator slug rewrite: childMeta drives post-hoc rewrite of
|
||||
// bare-hash slugs to `<hash6>-c<idx>` so chunked siblings can't collide
|
||||
// even if Sonnet drops the chunk suffix.
|
||||
// v0.32.8: refs carry source_id so reverseWriteRefs picks the correct
|
||||
// (source, slug) row. #1586: refs are stamped with the cycle's resolved
|
||||
// source (children write there via SubagentHandlerData.source_id).
|
||||
const cycleSourceId = opts.sourceId ?? 'default';
|
||||
const writtenRefs = await collectChildPutPageSlugs(engine, childIds, chunkInfo, cycleSourceId);
|
||||
const writtenRefs = await collectChildPutPageSlugs(engine, childIds, childMeta, cycleSourceId);
|
||||
|
||||
const summaryDate = opts.date ?? today();
|
||||
|
||||
@@ -559,7 +569,12 @@ export async function runPhaseSynthesize(
|
||||
// of every child-written page BEFORE reverse-rendering, so generated pages
|
||||
// are queryable (`frontmatter->>'dream_generated'`) and a later put_page
|
||||
// write-through (which re-renders from the DB row) can't erase the stamp.
|
||||
await stampDreamProvenance(engine, writtenRefs, summaryDate);
|
||||
// #2285: the stamp also carries the orchestrator-owned deterministic
|
||||
// frontmatter (transcript_id, transcript_source, transcript_hash, date,
|
||||
// chunk) derived from childMeta — subagent drift on those fields can't
|
||||
// leak, and reverseWriteRefs below re-reads the row so the same fields
|
||||
// land in the on-disk markdown.
|
||||
await stampDreamProvenance(engine, writtenRefs, summaryDate, childMeta);
|
||||
|
||||
// Dual-write: reverse-render each DB row → markdown file.
|
||||
const reverseWriteCount = await reverseWriteRefs(engine, opts.brainDir, writtenRefs, cycleSourceId);
|
||||
@@ -1095,15 +1110,17 @@ function sanitizeForSlug(s: string): string {
|
||||
* fake"): we no longer need detection because the rewrite enforces
|
||||
* uniqueness at slug-write time.
|
||||
*
|
||||
* `chunkInfo` maps child job_id → { chunk_index, hash6 }. Single-chunk
|
||||
* children are absent from the map and pass through unchanged.
|
||||
* `childMeta` maps child job_id → per-child transcript metadata. Chunked
|
||||
* children (chunkTotal > 1) get the slug rewrite; single-chunk children
|
||||
* pass through unchanged. Each returned ref carries the job_id that wrote
|
||||
* it so stampDreamProvenance can pair the slug back to its childMeta entry.
|
||||
*/
|
||||
async function collectChildPutPageSlugs(
|
||||
engine: BrainEngine,
|
||||
childIds: number[],
|
||||
chunkInfo: Map<number, { idx: number; hash6: string }>,
|
||||
childMeta: Map<number, ChildMeta>,
|
||||
sourceId = 'default',
|
||||
): Promise<Array<{ slug: string; source_id: string }>> {
|
||||
): Promise<Array<{ slug: string; source_id: string; jobId: number }>> {
|
||||
if (childIds.length === 0) return [];
|
||||
// Raw fetch — NO SELECT DISTINCT. Preserves per-child slug duplicates so
|
||||
// the orchestrator sees what each child wrote. COALESCE handles both
|
||||
@@ -1122,16 +1139,73 @@ async function collectChildPutPageSlugs(
|
||||
FROM subagent_tool_executions
|
||||
WHERE job_id = ANY($1::int[])
|
||||
AND tool_name = 'brain_put_page'
|
||||
AND status = 'complete'`,
|
||||
AND status = 'complete'
|
||||
ORDER BY id`,
|
||||
[childIds],
|
||||
);
|
||||
const rewritten = new Set<string>();
|
||||
const rewritten = new Map<string, number>();
|
||||
for (const r of rows) {
|
||||
if (typeof r.slug !== 'string' || r.slug.length === 0) continue;
|
||||
const ci = chunkInfo.get(r.job_id);
|
||||
rewritten.add(ci ? rewriteChunkedSlug(r.slug, ci.hash6, ci.idx) : r.slug);
|
||||
const meta = childMeta.get(r.job_id);
|
||||
const finalSlug = meta && meta.chunkTotal > 1
|
||||
? rewriteChunkedSlug(r.slug, meta.hash6, meta.idx)
|
||||
: r.slug;
|
||||
// Last writer wins, in execution-row order (ORDER BY id): if two children
|
||||
// collide on a final slug, the pages row holds the LAST put_page write, so
|
||||
// the stamp must attribute that child's transcript — not an arbitrary one.
|
||||
rewritten.set(finalSlug, r.job_id);
|
||||
}
|
||||
return Array.from(rewritten).sort().map(slug => ({ slug, source_id: sourceId }));
|
||||
return [...rewritten.entries()]
|
||||
.sort(([a], [b]) => a.localeCompare(b))
|
||||
.map(([slug, jobId]) => ({ slug, source_id: sourceId, jobId }));
|
||||
}
|
||||
|
||||
/**
|
||||
* Per-child orchestrator state. Drives D6 chunked-slug rewrite (idx + hash6)
|
||||
* AND the deterministic frontmatter stampDreamProvenance merges into each
|
||||
* written page. Populated for every child, not just chunked ones.
|
||||
*/
|
||||
interface ChildMeta {
|
||||
idx: number;
|
||||
hash6: string;
|
||||
chunkTotal: number;
|
||||
transcriptSource: string | null;
|
||||
transcriptId: string;
|
||||
inferredDate: string | null;
|
||||
}
|
||||
|
||||
/**
|
||||
* Strip the content-version suffix that claude-code-archive appends when a
|
||||
* conversation is edited (`<uuid>--<contentHash>.md`). The session UUID is
|
||||
* the stable transcript identifier; the suffix changes with content. Used to
|
||||
* populate `transcript_id` so edits of the same session collapse to one id.
|
||||
*/
|
||||
function stripContentVersionSuffix(basename: string): string {
|
||||
return basename.replace(/--[a-f0-9]+$/i, '');
|
||||
}
|
||||
|
||||
/**
|
||||
* Deterministic frontmatter for one synthesized page (#2285). Every field
|
||||
* here is owned by the orchestrator — the subagent's value for any of these
|
||||
* is overwritten. The subagent retains authority over type / title / tags /
|
||||
* body. `date` feeds the effective-date precedence chain
|
||||
* (src/core/effective-date.ts) so re-imports keep the conversation date even
|
||||
* when sync tools re-stamp file mtimes.
|
||||
*/
|
||||
function buildDeterministicFrontmatter(
|
||||
meta: ChildMeta,
|
||||
cycleDate: string,
|
||||
): Record<string, unknown> {
|
||||
const overrides: Record<string, unknown> = {
|
||||
dream_generated: true,
|
||||
dream_cycle_date: cycleDate,
|
||||
transcript_id: meta.transcriptId,
|
||||
transcript_hash: meta.hash6,
|
||||
};
|
||||
if (meta.transcriptSource) overrides.transcript_source = meta.transcriptSource;
|
||||
if (meta.chunkTotal > 1) overrides.chunk = `${meta.idx + 1}/${meta.chunkTotal}`;
|
||||
if (meta.inferredDate) overrides.date = meta.inferredDate;
|
||||
return overrides;
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -1177,12 +1251,19 @@ async function hasLegacySingleChunkCompletion(
|
||||
*/
|
||||
async function stampDreamProvenance(
|
||||
engine: BrainEngine,
|
||||
refs: Array<{ slug: string; source_id: string }>,
|
||||
refs: Array<{ slug: string; source_id: string; jobId?: number }>,
|
||||
cycleDate: string,
|
||||
childMeta?: Map<number, ChildMeta>,
|
||||
): Promise<void> {
|
||||
if (refs.length === 0) return;
|
||||
const { executeRawJsonb } = await import('../sql-query.ts');
|
||||
for (const { slug, source_id } of refs) {
|
||||
for (const { slug, source_id, jobId } of refs) {
|
||||
// #2285: when the ref pairs back to a child, the stamp also carries the
|
||||
// orchestrator-owned deterministic frontmatter for that transcript.
|
||||
const meta = jobId !== undefined ? childMeta?.get(jobId) : undefined;
|
||||
const stamp = meta
|
||||
? buildDeterministicFrontmatter(meta, cycleDate)
|
||||
: { dream_generated: true, dream_cycle_date: cycleDate };
|
||||
try {
|
||||
await executeRawJsonb(
|
||||
engine,
|
||||
@@ -1190,7 +1271,7 @@ async function stampDreamProvenance(
|
||||
SET frontmatter = COALESCE(frontmatter, '{}'::jsonb) || $3::jsonb
|
||||
WHERE slug = $1 AND source_id = $2`,
|
||||
[slug, source_id],
|
||||
[{ dream_generated: true, dream_cycle_date: cycleDate }],
|
||||
[stamp],
|
||||
);
|
||||
} catch (e) {
|
||||
const msg = e instanceof Error ? e.message : String(e);
|
||||
|
||||
@@ -10,7 +10,7 @@
|
||||
*/
|
||||
|
||||
import { readFileSync, readdirSync, statSync } from 'node:fs';
|
||||
import { join, basename } from 'node:path';
|
||||
import { join, basename, dirname } from 'node:path';
|
||||
import { createHash } from 'node:crypto';
|
||||
import { pruneDir } from '../sync.ts';
|
||||
|
||||
@@ -23,8 +23,22 @@ export interface DiscoveredTranscript {
|
||||
content: string;
|
||||
/** Filename basename without extension; used as a topic-slug seed. */
|
||||
basename: string;
|
||||
/** Inferred date if the basename matches `YYYY-MM-DD...` (or null). */
|
||||
/**
|
||||
* Inferred conversation date (YYYY-MM-DD) or null. Precedence: the
|
||||
* `| First message | <ISO> |` row in the transcript's `## Metadata`
|
||||
* table (stable across mtime-restamping re-syncs) wins; a leading
|
||||
* `YYYY-MM-DD` in the basename is the fallback.
|
||||
*/
|
||||
inferredDate: string | null;
|
||||
/**
|
||||
* Transcript source archive name, derived from the path's grandparent
|
||||
* directory (the immediate parent of the date directory). For the
|
||||
* canonical layout `<corpus>/<source>/<date>/<id>.md` this yields the
|
||||
* source-name segment — e.g. `claude-code` for the claude-code-archive
|
||||
* output, `meetings` for meeting recordings. Null when the file does
|
||||
* not live under a `<source>/<date>/` pair (ad-hoc inputs).
|
||||
*/
|
||||
transcriptSource: string | null;
|
||||
}
|
||||
|
||||
export interface DiscoverOpts {
|
||||
@@ -161,6 +175,36 @@ function matchesAnyExclude(text: string, patterns: RegExp[]): boolean {
|
||||
return false;
|
||||
}
|
||||
|
||||
/**
|
||||
* Content-based conversation date: the `| First message | <ISO timestamp> |`
|
||||
* row claude-code-archive writes into the transcript's `## Metadata` table.
|
||||
* Stable across rsync/Dropbox/Syncthing/B2 re-syncs that re-stamp mtime,
|
||||
* unlike anything derived from file metadata. Returns YYYY-MM-DD or null.
|
||||
*/
|
||||
const FIRST_MESSAGE_RE = /^\|\s*First message\s*\|\s*(\d{4}-\d{2}-\d{2})/im;
|
||||
|
||||
export function inferContentDate(content: string): string | null {
|
||||
const m = FIRST_MESSAGE_RE.exec(content);
|
||||
return m ? m[1] : null;
|
||||
}
|
||||
|
||||
/**
|
||||
* Derive the archive source name from a transcript path. Returns the basename
|
||||
* of the directory two levels above the file when the immediate parent is a
|
||||
* date directory and the grandparent looks like a source-name slug (lowercase
|
||||
* alphanumeric segments separated by hyphens); otherwise null. This pins the
|
||||
* canonical claude-code-archive layout `<corpus>/<source>/<date>/<id>.md`
|
||||
* without claiming a source for ad-hoc inputs that don't match.
|
||||
*/
|
||||
export function deriveTranscriptSource(filePath: string): string | null {
|
||||
const parentName = basename(dirname(filePath));
|
||||
if (!/^\d{4}-\d{2}-\d{2}/.test(parentName)) return null;
|
||||
const grandparentName = basename(dirname(dirname(filePath)));
|
||||
if (!grandparentName) return null;
|
||||
if (!/^[a-z0-9]+(-[a-z0-9]+)*$/.test(grandparentName)) return null;
|
||||
return grandparentName;
|
||||
}
|
||||
|
||||
function listTextFiles(dir: string): string[] {
|
||||
// Recursive walk with descent-time pruning (closes codex C12/C13 spec gap).
|
||||
// Accepts BOTH .txt and .md per transcript-discovery's domain rules — does
|
||||
@@ -225,8 +269,11 @@ export function discoverTranscripts(opts: DiscoverOpts): DiscoveredTranscript[]
|
||||
const ext = filePath.endsWith('.md') ? '.md' : '.txt';
|
||||
const baseName = basename(filePath, ext);
|
||||
const dateMatch = DATE_RE.exec(baseName);
|
||||
const inferredDate = dateMatch ? dateMatch[1] : null;
|
||||
if (!isInDateRange(inferredDate, opts)) continue;
|
||||
const filenameDate = dateMatch ? dateMatch[1] : null;
|
||||
// Fast path: date-named files outside the window skip before the read.
|
||||
// ponytail: a date-named file whose content date differs is filtered on
|
||||
// its filename date — acceptable; archive layouts use UUID basenames.
|
||||
if (filenameDate && !isInDateRange(filenameDate, opts)) continue;
|
||||
|
||||
let content: string;
|
||||
try {
|
||||
@@ -241,12 +288,17 @@ export function discoverTranscripts(opts: DiscoverOpts): DiscoveredTranscript[]
|
||||
}
|
||||
if (matchesAnyExclude(content, excludeRes)) continue;
|
||||
|
||||
// Content-metadata date wins (survives mtime restamps); filename next.
|
||||
const inferredDate = inferContentDate(content) ?? filenameDate;
|
||||
if (!isInDateRange(inferredDate, opts)) continue;
|
||||
|
||||
results.push({
|
||||
filePath,
|
||||
contentHash: hashContent(content),
|
||||
content,
|
||||
basename: baseName,
|
||||
inferredDate,
|
||||
transcriptSource: deriveTranscriptSource(filePath),
|
||||
});
|
||||
}
|
||||
}
|
||||
@@ -290,6 +342,7 @@ export function readSingleTranscript(
|
||||
contentHash: hashContent(content),
|
||||
content,
|
||||
basename: baseName,
|
||||
inferredDate: dateMatch ? dateMatch[1] : null,
|
||||
inferredDate: inferContentDate(content) ?? (dateMatch ? dateMatch[1] : null),
|
||||
transcriptSource: deriveTranscriptSource(filePath),
|
||||
};
|
||||
}
|
||||
|
||||
+2
-13
@@ -499,7 +499,7 @@ export function linkReadScopeOpts(ctx: OperationContext): { sourceId?: string; s
|
||||
* FAIL-CLOSED: anything not strictly `ctx.remote === false` is untrusted.
|
||||
*
|
||||
* This is the SINGLE resolver for every read op that accepts a per-call
|
||||
* `source_id` / `all_sources` parameter (query, search, code_callers, code_callees,
|
||||
* `source_id` / `all_sources` parameter (query, code_callers, code_callees,
|
||||
* get_page, search_by_image, code_blast, code_flow). Inlining the `__all__`
|
||||
* branch per handler is the bug class that leaked cross-source reads (#1924,
|
||||
* #1371): a remote client could pass `source_id: '__all__'` to opt out of its
|
||||
@@ -1442,24 +1442,13 @@ const search: Operation = {
|
||||
limit: { type: 'number', description: 'Max results (default 20)' },
|
||||
offset: { type: 'number', description: 'Skip first N results (for pagination)' },
|
||||
mode: { type: 'string', description: 'Search mode (conservative|balanced|tokenmax). Local callers only.' },
|
||||
source_id: {
|
||||
type: 'string',
|
||||
description:
|
||||
"Scope search to a single source. Defaults to OperationContext.sourceId. Pass '__all__' to span every source for trusted local callers; for remote callers '__all__' spans only your granted sources.",
|
||||
},
|
||||
all_sources: { type: 'boolean', description: "Span sources (equivalent to source_id=__all__): every source locally, your grant remotely." },
|
||||
},
|
||||
handler: async (ctx, p) => {
|
||||
const startedAt = Date.now();
|
||||
const queryText = p.query as string;
|
||||
const limit = (p.limit as number) || 20;
|
||||
const offset = (p.offset as number) || 0;
|
||||
// #1484 follow-up: route through the canonical fail-closed resolver so
|
||||
// `--source-id __all__` / `all_sources` behave the same as on `query`
|
||||
// (the zero-hit CLI hint advises exactly that retry). Without a per-call
|
||||
// param, `search` silently ignored --source-id — the retry looked like
|
||||
// a genuine miss.
|
||||
const scope = resolveRequestedScope(ctx, p.source_id as string | undefined, p.all_sources === true);
|
||||
const scope = sourceScopeOpts(ctx);
|
||||
|
||||
// T4/D5 — per-call mode honored ONLY for trusted/local callers so a remote
|
||||
// OAuth client can't escalate to the costly tokenmax bundle. Local + unknown
|
||||
|
||||
@@ -1323,18 +1323,8 @@ export async function hybridSearch(
|
||||
if (effectiveModality === 'both' && imageVectorList !== null) {
|
||||
vectorLists = [...vectorLists, imageVectorList];
|
||||
}
|
||||
} catch (err) {
|
||||
// Embedding/vector failure is non-fatal — fall back to keyword-only —
|
||||
// but say WHY (#1626): this arm only runs when the embedding provider
|
||||
// probed available, so a throw here is a real failure (embed timeout,
|
||||
// transient pooler error on the searchVector fan-out). Pre-fix the bare
|
||||
// catch made a cross-source `--source __all__` run silently collapse to
|
||||
// keyword-only/"No results" with zero diagnostics.
|
||||
warnOncePerProcess(
|
||||
'hybrid-vector-arm-failed',
|
||||
`[gbrain] vector arm failed (fail-open, keyword-only fallback): ` +
|
||||
`${err instanceof Error ? err.message : String(err)}`,
|
||||
);
|
||||
} catch {
|
||||
// Embedding failure is non-fatal, fall back to keyword-only
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -1,115 +0,0 @@
|
||||
/**
|
||||
* #1610 — Voyage/ZeroEntropy compat shims must read the response body ONCE
|
||||
* via text() instead of `resp.clone().json()`.
|
||||
*
|
||||
* On bun < 1.1.27, Response.clone() truncates large bodies (oven-sh/bun#6348):
|
||||
* the clone().json() parse threw, the shim's catch fell back to the ORIGINAL
|
||||
* response — whose wire shape (ZE `{results: ...}`, Voyage base64 embeddings)
|
||||
* the AI SDK's openai-compatible Zod schema rejects — and multi-chunk pages
|
||||
* failed with "Invalid JSON response".
|
||||
*
|
||||
* These tests simulate the truncating clone() and assert the shims still
|
||||
* return the fully rewritten body. They also pin that the rewritten Response
|
||||
* does NOT carry the original (now stale) Content-Length header, which lied
|
||||
* about the rewritten body's size (gateway.ts previously copied
|
||||
* `headers: resp.headers` verbatim).
|
||||
*/
|
||||
|
||||
import { afterEach, describe, expect, test } from 'bun:test';
|
||||
import {
|
||||
__voyageCompatFetchForTests,
|
||||
__zeroEntropyCompatFetchForTests,
|
||||
} from '../../src/core/ai/gateway.ts';
|
||||
|
||||
const origFetch = globalThis.fetch;
|
||||
afterEach(() => {
|
||||
globalThis.fetch = origFetch;
|
||||
});
|
||||
|
||||
/** Build a Response whose clone() truncates the body (bun < 1.1.27 behavior). */
|
||||
function truncatingCloneResponse(body: string): Response {
|
||||
const headers = {
|
||||
'content-type': 'application/json',
|
||||
// Deliberately stale after any rewrite: the original wire body's length.
|
||||
'content-length': String(Buffer.byteLength(body)),
|
||||
};
|
||||
const resp = new Response(body, { status: 200, headers });
|
||||
(resp as any).clone = () =>
|
||||
new Response(body.slice(0, 32), { status: 200, headers });
|
||||
return resp;
|
||||
}
|
||||
|
||||
describe('voyageCompatFetch — single body read (#1610)', () => {
|
||||
test('rewrites base64 embeddings even when clone() truncates the body', async () => {
|
||||
const floats = new Float32Array([0.5, 0.25, -1]);
|
||||
const b64 = Buffer.from(floats.buffer).toString('base64');
|
||||
const wireBody = JSON.stringify({
|
||||
object: 'list',
|
||||
data: [{ object: 'embedding', embedding: b64, index: 0 }],
|
||||
model: 'voyage-3',
|
||||
usage: { total_tokens: 7 },
|
||||
});
|
||||
globalThis.fetch = (async () => truncatingCloneResponse(wireBody)) as unknown as typeof fetch;
|
||||
|
||||
const out = await __voyageCompatFetchForTests('https://api.voyageai.com/v1/embeddings', {
|
||||
method: 'POST',
|
||||
body: JSON.stringify({ input: ['hello'], model: 'voyage-3' }),
|
||||
headers: { 'content-type': 'application/json' },
|
||||
});
|
||||
|
||||
const json: any = await out.json();
|
||||
expect(Array.from(json.data[0].embedding)).toEqual([0.5, 0.25, -1]);
|
||||
expect(json.usage.prompt_tokens).toBe(7);
|
||||
// Stale Content-Length from the wire body must not survive the rewrite.
|
||||
expect(out.headers.get('content-length')).toBeNull();
|
||||
expect(out.headers.get('content-encoding')).toBeNull();
|
||||
});
|
||||
});
|
||||
|
||||
describe('zeroEntropyCompatFetch — single body read (#1610)', () => {
|
||||
test('rewrites {results} → {data} even when clone() truncates the body', async () => {
|
||||
const wireBody = JSON.stringify({
|
||||
results: [{ embedding: [0.1, 0.2] }, { embedding: [0.3, 0.4] }],
|
||||
usage: { total_bytes: 42, total_tokens: 9 },
|
||||
});
|
||||
let fetchedUrl = '';
|
||||
globalThis.fetch = (async (url: string | URL | Request) => {
|
||||
fetchedUrl = String(url);
|
||||
return truncatingCloneResponse(wireBody);
|
||||
}) as unknown as typeof fetch;
|
||||
|
||||
const out = await __zeroEntropyCompatFetchForTests('https://api.zeroentropy.dev/v1/embeddings', {
|
||||
method: 'POST',
|
||||
body: JSON.stringify({ input: ['hello'], model: 'zembed-1' }),
|
||||
headers: { 'content-type': 'application/json' },
|
||||
});
|
||||
|
||||
expect(fetchedUrl.endsWith('/v1/models/embed')).toBe(true);
|
||||
const json: any = await out.json();
|
||||
// The AI SDK schema requires {data: [{embedding, index}]} — the raw ZE
|
||||
// {results} fallback is exactly the pre-fix "Invalid JSON response".
|
||||
expect(json.results).toBeUndefined();
|
||||
expect(json.data).toHaveLength(2);
|
||||
expect(json.data[0]).toEqual({ object: 'embedding', embedding: [0.1, 0.2], index: 0 });
|
||||
expect(json.data[1].index).toBe(1);
|
||||
expect(json.usage.prompt_tokens).toBe(9);
|
||||
expect(out.headers.get('content-length')).toBeNull();
|
||||
});
|
||||
|
||||
test('non-JSON body falls back to the original bytes (rebuilt, still readable)', async () => {
|
||||
const wireBody = 'plain text, not json';
|
||||
globalThis.fetch = (async () =>
|
||||
new Response(wireBody, {
|
||||
status: 200,
|
||||
headers: { 'content-type': 'application/json' },
|
||||
})) as unknown as typeof fetch;
|
||||
|
||||
const out = await __zeroEntropyCompatFetchForTests('https://api.zeroentropy.dev/v1/embeddings', {
|
||||
method: 'POST',
|
||||
body: JSON.stringify({ input: ['hello'] }),
|
||||
});
|
||||
// Body was consumed by the shim's single read; the fallback must
|
||||
// rebuild a readable Response rather than return the drained original.
|
||||
expect(await out.text()).toBe(wireBody);
|
||||
});
|
||||
});
|
||||
@@ -98,18 +98,16 @@ describe('zeroEntropyCompatFetch — OOM caps', () => {
|
||||
expect(src).toMatch(/MAX_ZEROENTROPY_RESPONSE_BYTES\s*=\s*256\s*\*\s*1024\s*\*\s*1024/);
|
||||
});
|
||||
|
||||
test('Layer 1: Content-Length pre-check before the body is read', async () => {
|
||||
test('Layer 1: Content-Length pre-check before resp.clone().json()', async () => {
|
||||
const src = await Bun.file(GATEWAY_PATH).text();
|
||||
// Find the zeroEntropyCompatFetch block bounds, then assert ordering
|
||||
// within it (mirroring the voyage cap test pattern). #1610 moved the
|
||||
// body read from `resp.clone().json()` to a single `resp.text()` (bun
|
||||
// < 1.1.27 truncates clone()d bodies, oven-sh/bun#6348).
|
||||
// within it (mirroring the voyage cap test pattern).
|
||||
const zeFetchStart = src.indexOf('const zeroEntropyCompatFetch');
|
||||
expect(zeFetchStart).toBeGreaterThan(0);
|
||||
const block = src.slice(zeFetchStart, zeFetchStart + 9000);
|
||||
const block = src.slice(zeFetchStart, zeFetchStart + 8000);
|
||||
|
||||
const preCheckIdx = block.indexOf("resp.headers.get('content-length')");
|
||||
const jsonParseIdx = block.indexOf('const bodyText = await resp.text()');
|
||||
const jsonParseIdx = block.indexOf('await resp.clone().json()');
|
||||
expect(preCheckIdx).toBeGreaterThan(0);
|
||||
expect(jsonParseIdx).toBeGreaterThan(0);
|
||||
// The pre-check MUST appear before the JSON parse — Voyage's lesson
|
||||
|
||||
@@ -1,61 +0,0 @@
|
||||
/**
|
||||
* #1484 — invisible-miss hint. A bare `gbrain query` resolves to a single
|
||||
* source (usually 'default'); on a multi-source brain a zero-hit run gave no
|
||||
* signal that the answer might live in another source. sourceScopeHint
|
||||
* returns the stderr hint exactly when: query/search op + zero results +
|
||||
* no explicit scoping param + >1 registered source.
|
||||
*/
|
||||
|
||||
import { describe, expect, test } from 'bun:test';
|
||||
import { sourceScopeHint } from '../src/cli.ts';
|
||||
import type { BrainEngine } from '../src/core/engine.ts';
|
||||
|
||||
function fakeEngine(sourceCount: number, fail = false): BrainEngine {
|
||||
return {
|
||||
executeRaw: async () => {
|
||||
if (fail) throw new Error('sources table missing');
|
||||
return [{ n: sourceCount }];
|
||||
},
|
||||
} as unknown as BrainEngine;
|
||||
}
|
||||
|
||||
describe('sourceScopeHint (#1484)', () => {
|
||||
test('fires on a bare zero-hit query against a multi-source brain', async () => {
|
||||
const hint = await sourceScopeHint('query', {}, 'default', fakeEngine(3), []);
|
||||
expect(hint).toContain('3 sources');
|
||||
expect(hint).toContain('"default"');
|
||||
expect(hint).toContain('--source-id __all__');
|
||||
});
|
||||
|
||||
test('fires for search too', async () => {
|
||||
const hint = await sourceScopeHint('search', {}, 'wiki', fakeEngine(2), []);
|
||||
expect(hint).toContain('"wiki"');
|
||||
});
|
||||
|
||||
test('silent when results were found', async () => {
|
||||
expect(await sourceScopeHint('query', {}, 'default', fakeEngine(3), [{ slug: 'a' }])).toBeNull();
|
||||
});
|
||||
|
||||
test('silent when the caller scoped explicitly', async () => {
|
||||
expect(await sourceScopeHint('query', { source_id: 'wiki' }, 'wiki', fakeEngine(3), [])).toBeNull();
|
||||
expect(await sourceScopeHint('query', { source: 'wiki' }, 'wiki', fakeEngine(3), [])).toBeNull();
|
||||
expect(await sourceScopeHint('query', { all_sources: true }, '__all__', fakeEngine(3), [])).toBeNull();
|
||||
});
|
||||
|
||||
test('silent when the resolved scope is already __all__', async () => {
|
||||
expect(await sourceScopeHint('query', {}, '__all__', fakeEngine(3), [])).toBeNull();
|
||||
});
|
||||
|
||||
test('silent on a single-source brain', async () => {
|
||||
expect(await sourceScopeHint('query', {}, 'default', fakeEngine(1), [])).toBeNull();
|
||||
});
|
||||
|
||||
test('silent for non-search ops and non-array results', async () => {
|
||||
expect(await sourceScopeHint('get_stats', {}, 'default', fakeEngine(3), [])).toBeNull();
|
||||
expect(await sourceScopeHint('query', {}, 'default', fakeEngine(3), { rows: [] })).toBeNull();
|
||||
});
|
||||
|
||||
test('best-effort: sources lookup failure returns null, never throws', async () => {
|
||||
expect(await sourceScopeHint('query', {}, 'default', fakeEngine(3, true), [])).toBeNull();
|
||||
});
|
||||
});
|
||||
@@ -26,6 +26,7 @@ const transcript: DiscoveredTranscript = {
|
||||
content: 'User: hello world',
|
||||
contentHash: 'abcdef0123456789',
|
||||
inferredDate: '2026-07-17',
|
||||
transcriptSource: null,
|
||||
} as DiscoveredTranscript;
|
||||
|
||||
describe('#2415: buildSynthesisPrompt output root', () => {
|
||||
|
||||
@@ -152,3 +152,94 @@ describe('#2569: stampDreamProvenance persists the marker into DB frontmatter',
|
||||
await stampDreamProvenance(engine as any, refs, '2026-07-17'); // idempotent
|
||||
});
|
||||
});
|
||||
|
||||
describe('#2285: orchestrator-owned deterministic transcript frontmatter', () => {
|
||||
const meta = {
|
||||
idx: 1,
|
||||
hash6: 'abc123',
|
||||
chunkTotal: 3,
|
||||
transcriptSource: 'claude-code',
|
||||
transcriptId: 'session-uuid',
|
||||
inferredDate: '2026-05-15',
|
||||
};
|
||||
|
||||
test('collectChildPutPageSlugs pairs each ref back to the writing job', async () => {
|
||||
const refs = await collectChildPutPageSlugs(
|
||||
engine as any, [1001], new Map([[1001, { ...meta, chunkTotal: 1 }]]), 'mybrain',
|
||||
);
|
||||
expect(refs.length).toBeGreaterThan(0);
|
||||
for (const r of refs) {
|
||||
expect(r.jobId).toBe(1001);
|
||||
expect(r.source_id).toBe('mybrain'); // #1586: cycle source, never hardcoded 'default'
|
||||
}
|
||||
});
|
||||
|
||||
test('slug collision across children attributes the LAST writer (matches surviving putPage)', async () => {
|
||||
const db = (engine as any).db;
|
||||
// Jobs 1001 then 1002 write the same slug; the pages row would hold
|
||||
// 1002's content (last put_page wins), so the ref must carry jobId 1002.
|
||||
await db.query(
|
||||
`INSERT INTO subagent_tool_executions (job_id, message_idx, tool_use_id, tool_name, status, input)
|
||||
VALUES (1001, 9, 'tool_dup_a', 'brain_put_page', 'complete', $1::jsonb)`,
|
||||
[JSON.stringify({ slug: 'wiki/agents/test/collision', body: 'first' })],
|
||||
);
|
||||
await db.query(
|
||||
`INSERT INTO subagent_tool_executions (job_id, message_idx, tool_use_id, tool_name, status, input)
|
||||
VALUES (1002, 9, 'tool_dup_b', 'brain_put_page', 'complete', $1::jsonb)`,
|
||||
[JSON.stringify({ slug: 'wiki/agents/test/collision', body: 'second' })],
|
||||
);
|
||||
const refs = await collectChildPutPageSlugs(engine as any, [1001, 1002], new Map());
|
||||
const hit = refs.find((r: { slug: string }) => r.slug === 'wiki/agents/test/collision');
|
||||
expect(hit?.jobId).toBe(1002);
|
||||
});
|
||||
|
||||
test('stampDreamProvenance merges the transcript metadata into DB frontmatter', async () => {
|
||||
const slug = 'wiki/originals/ideas/2026-07-17-transcript-meta-abc123';
|
||||
await engine.putPage(slug, {
|
||||
type: 'note',
|
||||
title: 'Meta stamp',
|
||||
compiled_truth: 'body',
|
||||
timeline: '',
|
||||
frontmatter: { keep_me: 'yes', transcript_id: 'subagent-drift' },
|
||||
});
|
||||
await stampDreamProvenance(
|
||||
engine as any,
|
||||
[{ slug, source_id: 'default', jobId: 42 }],
|
||||
'2026-07-17',
|
||||
new Map([[42, meta]]),
|
||||
);
|
||||
const rows = await engine.executeRaw<{ fm: Record<string, unknown> }>(
|
||||
`SELECT frontmatter AS fm FROM pages WHERE slug = $1`, [slug],
|
||||
);
|
||||
const fm = rows[0].fm as Record<string, unknown>;
|
||||
expect(fm.dream_generated).toBe(true);
|
||||
expect(fm.dream_cycle_date).toBe('2026-07-17');
|
||||
expect(fm.transcript_id).toBe('session-uuid'); // orchestrator wins over subagent drift
|
||||
expect(fm.transcript_hash).toBe('abc123');
|
||||
expect(fm.transcript_source).toBe('claude-code');
|
||||
expect(fm.chunk).toBe('2/3');
|
||||
expect(fm.date).toBe('2026-05-15');
|
||||
expect(fm.keep_me).toBe('yes'); // subagent-owned keys survive
|
||||
});
|
||||
|
||||
test('single-chunk children with no inferredDate stamp only the applicable fields', async () => {
|
||||
const slug = 'wiki/originals/ideas/2026-07-17-minimal-meta-abc123';
|
||||
await engine.putPage(slug, {
|
||||
type: 'note', title: 'Minimal', compiled_truth: 'b', timeline: '', frontmatter: {},
|
||||
});
|
||||
await stampDreamProvenance(
|
||||
engine as any,
|
||||
[{ slug, source_id: 'default', jobId: 43 }],
|
||||
'2026-07-17',
|
||||
new Map([[43, { ...meta, chunkTotal: 1, transcriptSource: null, inferredDate: null }]]),
|
||||
);
|
||||
const rows = await engine.executeRaw<{ fm: Record<string, unknown> }>(
|
||||
`SELECT frontmatter AS fm FROM pages WHERE slug = $1`, [slug],
|
||||
);
|
||||
const fm = rows[0].fm as Record<string, unknown>;
|
||||
expect(fm.transcript_id).toBe('session-uuid');
|
||||
expect(fm.chunk).toBeUndefined();
|
||||
expect(fm.transcript_source).toBeUndefined();
|
||||
expect(fm.date).toBeUndefined();
|
||||
});
|
||||
});
|
||||
|
||||
@@ -302,6 +302,7 @@ describe('judgeSignificance', () => {
|
||||
content: 'A short conversation about something interesting.',
|
||||
basename: 'x',
|
||||
inferredDate: null,
|
||||
transcriptSource: null,
|
||||
};
|
||||
}
|
||||
|
||||
@@ -415,6 +416,7 @@ describe('judgeSignificance — UTF-16 safety (v0.41.13)', () => {
|
||||
content,
|
||||
basename: 'long',
|
||||
inferredDate: null,
|
||||
transcriptSource: null,
|
||||
};
|
||||
}
|
||||
|
||||
|
||||
@@ -45,6 +45,7 @@ const FIXTURE_TRANSCRIPT: DiscoveredTranscript = {
|
||||
content: 'Synthetic transcript content for gateway-adapter parity tests.',
|
||||
contentHash: 'sha-fixture-1',
|
||||
inferredDate: '2026-05-24',
|
||||
transcriptSource: null,
|
||||
};
|
||||
|
||||
describe('makeJudgeClient — construction-time provider probe', () => {
|
||||
|
||||
@@ -0,0 +1,109 @@
|
||||
/**
|
||||
* #2285 — transcript metadata discovery.
|
||||
*
|
||||
* Pins the two discovery-side additions:
|
||||
* 1. `transcriptSource` — derived from the `<source>/<date>/<file>` path
|
||||
* layout; null for ad-hoc inputs that don't match.
|
||||
* 2. Content-based date inference — the `| First message | <ISO> |` row in
|
||||
* the transcript's `## Metadata` table wins over the filename-regex
|
||||
* date (stable across mtime-restamping re-syncs); filename is the
|
||||
* fallback.
|
||||
*
|
||||
* Pure filesystem; no engine, no LLM.
|
||||
*/
|
||||
|
||||
import { describe, test, expect, beforeEach, afterEach } from 'bun:test';
|
||||
import { mkdtempSync, rmSync, writeFileSync, mkdirSync } from 'node:fs';
|
||||
import { tmpdir } from 'node:os';
|
||||
import { join, dirname } from 'node:path';
|
||||
import {
|
||||
discoverTranscripts,
|
||||
readSingleTranscript,
|
||||
deriveTranscriptSource,
|
||||
inferContentDate,
|
||||
} from '../../src/core/cycle/transcript-discovery.ts';
|
||||
|
||||
let tmpDir: string;
|
||||
|
||||
beforeEach(() => {
|
||||
tmpDir = mkdtempSync(join(tmpdir(), 'gbrain-transcript-meta-'));
|
||||
});
|
||||
|
||||
afterEach(() => {
|
||||
rmSync(tmpDir, { recursive: true, force: true });
|
||||
});
|
||||
|
||||
function write(relPath: string, body: string): string {
|
||||
const full = join(tmpDir, relPath);
|
||||
mkdirSync(dirname(full), { recursive: true });
|
||||
writeFileSync(full, body);
|
||||
return full;
|
||||
}
|
||||
|
||||
const FILLER = 'User: hello world. '.repeat(200);
|
||||
const METADATA_BLOCK =
|
||||
'## Metadata\n\n| Key | Value |\n| --- | --- |\n| First message | 2026-05-15T03:51:11.584Z |\n\n';
|
||||
|
||||
describe('deriveTranscriptSource', () => {
|
||||
test('extracts the source slug from <source>/<date>/<file> layout', () => {
|
||||
expect(deriveTranscriptSource('/corpus/claude-code/2026-06-12/abc.md')).toBe('claude-code');
|
||||
expect(deriveTranscriptSource('/corpus/voice-notes/2026-06-12/xyz.md')).toBe('voice-notes');
|
||||
});
|
||||
|
||||
test('null when the parent dir is not a date dir or grandparent is not a slug', () => {
|
||||
expect(deriveTranscriptSource('/corpus/flat-file.md')).toBeNull();
|
||||
expect(deriveTranscriptSource('/corpus/claude-code/not-a-date/abc.md')).toBeNull();
|
||||
expect(deriveTranscriptSource('/corpus/Not A Slug/2026-06-12/abc.md')).toBeNull();
|
||||
});
|
||||
});
|
||||
|
||||
describe('inferContentDate', () => {
|
||||
test('parses the | First message | row', () => {
|
||||
expect(inferContentDate(METADATA_BLOCK)).toBe('2026-05-15');
|
||||
});
|
||||
test('null when absent', () => {
|
||||
expect(inferContentDate(FILLER)).toBeNull();
|
||||
});
|
||||
});
|
||||
|
||||
describe('discoverTranscripts — transcriptSource + date cascade', () => {
|
||||
test('populates transcriptSource per file; null for flat files', () => {
|
||||
write('claude-code/2026-06-12/aaaa.md', FILLER);
|
||||
write('2026-06-12-flat.md', FILLER);
|
||||
const out = discoverTranscripts({ corpusDir: tmpDir, minChars: 100 });
|
||||
const byBase = new Map(out.map(t => [t.basename, t.transcriptSource]));
|
||||
expect(byBase.get('aaaa')).toBe('claude-code');
|
||||
expect(byBase.get('2026-06-12-flat')).toBeNull();
|
||||
});
|
||||
|
||||
test('content First-message date wins over the filename date', () => {
|
||||
write('2026-01-01-named.md', METADATA_BLOCK + FILLER);
|
||||
const out = discoverTranscripts({ corpusDir: tmpDir, minChars: 100 });
|
||||
expect(out).toHaveLength(1);
|
||||
expect(out[0].inferredDate).toBe('2026-05-15');
|
||||
});
|
||||
|
||||
test('filename date remains the fallback when content has no metadata row', () => {
|
||||
write('2026-01-01-named.md', FILLER);
|
||||
const out = discoverTranscripts({ corpusDir: tmpDir, minChars: 100 });
|
||||
expect(out[0].inferredDate).toBe('2026-01-01');
|
||||
});
|
||||
|
||||
test('date filter matches on the content date for UUID-named transcripts', () => {
|
||||
write('claude-code/2026-05-15/uuid-basename.md', METADATA_BLOCK + FILLER);
|
||||
const hit = discoverTranscripts({ corpusDir: tmpDir, minChars: 100, date: '2026-05-15' });
|
||||
expect(hit).toHaveLength(1);
|
||||
const miss = discoverTranscripts({ corpusDir: tmpDir, minChars: 100, date: '2026-05-16' });
|
||||
expect(miss).toHaveLength(0);
|
||||
});
|
||||
});
|
||||
|
||||
describe('readSingleTranscript — same metadata surface', () => {
|
||||
test('carries transcriptSource and prefers the content date', () => {
|
||||
const p = write('claude-code/2026-05-15/2026-01-01-single.md', METADATA_BLOCK + FILLER);
|
||||
const t = readSingleTranscript(p, { minChars: 100 });
|
||||
expect(t).not.toBeNull();
|
||||
expect(t!.transcriptSource).toBe('claude-code');
|
||||
expect(t!.inferredDate).toBe('2026-05-15');
|
||||
});
|
||||
});
|
||||
@@ -1,77 +0,0 @@
|
||||
/**
|
||||
* #1626 — hybridSearch's text-vector arm must not fail DARK.
|
||||
*
|
||||
* The arm only runs when the embedding provider probed available, so a throw
|
||||
* inside it (embed timeout, transient pooler error on searchVector) is a real
|
||||
* failure. Pre-fix, a bare `catch {}` swallowed it and the run silently
|
||||
* collapsed to keyword-only — under `--source __all__` on a strained pooler
|
||||
* that read as a non-deterministic "No results". The fix logs the swallowed
|
||||
* reason via warnOncePerProcess while keeping the keyword fallback.
|
||||
*/
|
||||
|
||||
import { afterAll, beforeAll, describe, expect, test } from 'bun:test';
|
||||
import { PGLiteEngine } from '../src/core/pglite-engine.ts';
|
||||
import { hybridSearch } from '../src/core/search/hybrid.ts';
|
||||
import {
|
||||
__setEmbedTransportForTests,
|
||||
configureGateway,
|
||||
resetGateway,
|
||||
} from '../src/core/ai/gateway.ts';
|
||||
import { _resetWarnOnceForTests } from '../src/core/utils.ts';
|
||||
|
||||
let engine: PGLiteEngine;
|
||||
const origWarn = console.warn;
|
||||
|
||||
beforeAll(async () => {
|
||||
// Pin the gateway to OpenAI with a stub key (put-page-provenance pattern):
|
||||
// embed() runs instantiateEmbedding — which requires OPENAI_API_KEY — BEFORE
|
||||
// the stubbed transport is reached. Without this, a keyless CI environment
|
||||
// throws the config error instead of the transport's, and the assertion on
|
||||
// the swallowed reason fails. The key never leaves the process.
|
||||
configureGateway({
|
||||
embedding_model: 'openai:text-embedding-3-large',
|
||||
embedding_dimensions: 1536,
|
||||
env: { ...process.env, OPENAI_API_KEY: process.env.OPENAI_API_KEY || 'sk-test-stub' },
|
||||
});
|
||||
engine = new PGLiteEngine();
|
||||
await engine.connect({});
|
||||
await engine.initSchema();
|
||||
await engine.putPage('people/alice-example', {
|
||||
type: 'person',
|
||||
title: 'Alice Example',
|
||||
compiled_truth: 'Alice Example is a test person for the vector-arm warn test.',
|
||||
});
|
||||
});
|
||||
|
||||
afterAll(async () => {
|
||||
console.warn = origWarn;
|
||||
__setEmbedTransportForTests(null);
|
||||
resetGateway();
|
||||
await engine.disconnect();
|
||||
});
|
||||
|
||||
describe('hybridSearch vector-arm failure telemetry (#1626)', () => {
|
||||
test('embed failure logs the swallowed reason and falls back to keyword', async () => {
|
||||
_resetWarnOnceForTests();
|
||||
// Installing a transport makes isAvailable('embedding') true (test-seam
|
||||
// fast path), so the vector arm RUNS — and then throws.
|
||||
__setEmbedTransportForTests(() => {
|
||||
throw new Error('pooler exploded mid-fanout');
|
||||
});
|
||||
const warnings: string[] = [];
|
||||
console.warn = (...args: unknown[]) => {
|
||||
warnings.push(args.map(String).join(' '));
|
||||
};
|
||||
try {
|
||||
const results = await hybridSearch(engine, 'alice');
|
||||
// Keyword fallback still returns results — fail-open preserved.
|
||||
expect(results.some((r) => r.slug === 'people/alice-example')).toBe(true);
|
||||
} finally {
|
||||
console.warn = origWarn;
|
||||
__setEmbedTransportForTests(null);
|
||||
}
|
||||
const armWarnings = warnings.filter((w) => w.includes('vector arm failed'));
|
||||
expect(armWarnings).toHaveLength(1);
|
||||
expect(armWarnings[0]).toContain('pooler exploded mid-fanout');
|
||||
});
|
||||
});
|
||||
@@ -216,17 +216,21 @@ describe('progress reporter', () => {
|
||||
});
|
||||
|
||||
test('only one process-level signal handler installed across many reporters', () => {
|
||||
// Baseline: one handler already installed by prior tests in this file.
|
||||
// Baseline: one handler already installed by prior tests in this file, and
|
||||
// possibly live reporters from OTHER test files sharing this bun process
|
||||
// (shard composition is not this test's invariant — assert the delta, not
|
||||
// an absolute zero, or shard reshuffles make this fail spuriously).
|
||||
const installedBefore = __signalHandlerInstalledForTest();
|
||||
const liveBefore = __liveReporterCountForTest();
|
||||
const { stream } = sink(false);
|
||||
for (let i = 0; i < 50; i++) {
|
||||
const p = createProgress({ mode: 'json', stream, minIntervalMs: 0, minItems: 1 });
|
||||
p.start(`phase_${i}`, 1);
|
||||
p.finish();
|
||||
}
|
||||
// After 50 reporter lifecycles, still exactly one handler and zero leaked live entries.
|
||||
// After 50 reporter lifecycles, still exactly one handler and zero NET leaked live entries.
|
||||
expect(__signalHandlerInstalledForTest()).toBe(installedBefore || true);
|
||||
expect(__liveReporterCountForTest()).toBe(0);
|
||||
expect(__liveReporterCountForTest()).toBe(liveBefore);
|
||||
});
|
||||
|
||||
test('startHeartbeat() fires heartbeats and stop() clears', async () => {
|
||||
|
||||
@@ -1,87 +0,0 @@
|
||||
/**
|
||||
* #1484 follow-up — the `search` op must honor per-call `source_id` /
|
||||
* `all_sources` through the canonical fail-closed resolver
|
||||
* (resolveRequestedScope), exactly like `query` does.
|
||||
*
|
||||
* Pre-fix, `search` had no source_id param at all: the zero-hit CLI hint
|
||||
* advised "retry with --source-id __all__", the flag parsed into params,
|
||||
* NOTHING consumed it, and the retry silently re-ran the same single-source
|
||||
* search — an invisible false negative (and the retry's params.source_id
|
||||
* suppressed the hint, so the user got no second warning).
|
||||
*/
|
||||
|
||||
import { describe, expect, test } from 'bun:test';
|
||||
import { operationsByName } from '../src/core/operations.ts';
|
||||
import type { OperationContext } from '../src/core/operations.ts';
|
||||
import type { BrainEngine } from '../src/core/engine.ts';
|
||||
|
||||
const searchOp = operationsByName['search'];
|
||||
|
||||
/** Fake engine: keyword-only config so the handler's scope goes straight to
|
||||
* searchKeyword, where we capture the opts it was called with. */
|
||||
function makeCtx(remote: boolean, allowedSources?: string[]) {
|
||||
const captured: { opts?: Record<string, unknown> } = {};
|
||||
const engine = {
|
||||
getConfig: async (key: string) => (key === 'search.mcp_keyword_only' ? 'true' : null),
|
||||
searchKeyword: async (_q: string, opts: Record<string, unknown>) => {
|
||||
captured.opts = opts;
|
||||
return [];
|
||||
},
|
||||
} as unknown as BrainEngine;
|
||||
const ctx = {
|
||||
engine,
|
||||
config: { engine: 'pglite' },
|
||||
logger: { info: () => {}, warn: () => {}, error: () => {} },
|
||||
dryRun: false,
|
||||
remote,
|
||||
sourceId: 'default',
|
||||
...(allowedSources ? { auth: { allowedSources } } : {}),
|
||||
} as unknown as OperationContext;
|
||||
return { ctx, captured };
|
||||
}
|
||||
|
||||
describe('search op per-call source scope (#1484 follow-up)', () => {
|
||||
test('op declares source_id + all_sources params (the CLI hint advises them)', () => {
|
||||
expect(searchOp.params.source_id).toBeDefined();
|
||||
expect(searchOp.params.all_sources).toBeDefined();
|
||||
});
|
||||
|
||||
test('default: scopes to ctx.sourceId', async () => {
|
||||
const { ctx, captured } = makeCtx(false);
|
||||
await searchOp.handler(ctx, { query: 'x' });
|
||||
expect(captured.opts?.sourceId).toBe('default');
|
||||
});
|
||||
|
||||
test("local + source_id '__all__' spans the whole brain (no source filter)", async () => {
|
||||
const { ctx, captured } = makeCtx(false);
|
||||
await searchOp.handler(ctx, { query: 'x', source_id: '__all__' });
|
||||
expect(captured.opts?.sourceId).toBeUndefined();
|
||||
expect(captured.opts?.sourceIds).toBeUndefined();
|
||||
});
|
||||
|
||||
test('local + all_sources=true spans the whole brain', async () => {
|
||||
const { ctx, captured } = makeCtx(false);
|
||||
await searchOp.handler(ctx, { query: 'x', all_sources: true });
|
||||
expect(captured.opts?.sourceId).toBeUndefined();
|
||||
expect(captured.opts?.sourceIds).toBeUndefined();
|
||||
});
|
||||
|
||||
test('explicit source_id wins over ctx.sourceId', async () => {
|
||||
const { ctx, captured } = makeCtx(false);
|
||||
await searchOp.handler(ctx, { query: 'x', source_id: 'wiki' });
|
||||
expect(captured.opts?.sourceId).toBe('wiki');
|
||||
});
|
||||
|
||||
test("remote + '__all__' collapses to the caller's grant (fail-closed)", async () => {
|
||||
const { ctx, captured } = makeCtx(true, ['wiki', 'essays']);
|
||||
await searchOp.handler(ctx, { query: 'x', source_id: '__all__' });
|
||||
expect(captured.opts?.sourceIds).toEqual(['wiki', 'essays']);
|
||||
});
|
||||
|
||||
test('remote + out-of-grant source_id is denied', async () => {
|
||||
const { ctx } = makeCtx(true, ['wiki']);
|
||||
await expect(searchOp.handler(ctx, { query: 'x', source_id: 'secrets' })).rejects.toThrow(
|
||||
/outside your granted sources/,
|
||||
);
|
||||
});
|
||||
});
|
||||
@@ -34,7 +34,7 @@ describe('v0.31.8 — voyage Content-Length pre-check + per-item cap', () => {
|
||||
expect(source).toMatch(/MAX_VOYAGE_RESPONSE_BYTES\s*=\s*256\s*\*\s*1024\s*\*\s*1024/);
|
||||
});
|
||||
|
||||
test('Layer 1: Content-Length pre-check fires BEFORE the body is read (D10 OOM defense)', async () => {
|
||||
test('Layer 1: Content-Length pre-check fires BEFORE resp.clone().json() (D10 OOM defense)', async () => {
|
||||
const source = await Bun.file(new URL('../src/core/ai/gateway.ts', import.meta.url)).text();
|
||||
// Anchor relative to the post-fetch handler block. The function declaration
|
||||
// contains an OUTBOUND request body section earlier; we want to verify
|
||||
@@ -47,10 +47,8 @@ describe('v0.31.8 — voyage Content-Length pre-check + per-item cap', () => {
|
||||
// doesn't pin to comment text.
|
||||
const preCheckIdx = inboundBlock.indexOf("resp.headers.get('content-length')");
|
||||
// Use the full lvalue assignment so the match doesn't accidentally hit
|
||||
// comment text that mentions the body read for context. (#1610 moved the
|
||||
// read from `resp.clone().json()` to a single `resp.text()` — bun <
|
||||
// 1.1.27 truncates clone()d bodies, oven-sh/bun#6348.)
|
||||
const jsonParseIdx = inboundBlock.indexOf('const bodyText = await resp.text()');
|
||||
// comment text that mentions `await resp.clone().json()` for context.
|
||||
const jsonParseIdx = inboundBlock.indexOf('const json: any = await resp.clone().json()');
|
||||
expect(preCheckIdx).toBeGreaterThan(0);
|
||||
expect(jsonParseIdx).toBeGreaterThan(0);
|
||||
// The pre-check MUST appear before the JSON parse — otherwise the OOM
|
||||
|
||||
Reference in New Issue
Block a user