mirror of
https://github.com/garrytan/gbrain.git
synced 2026-08-16 09:52:22 +00:00
Compare commits
1
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
08f2397615 |
@@ -1651,7 +1651,7 @@ async function extractTimelineFromDB(
|
||||
* make re-extraction idempotent). EVERY processed page is stamped, including
|
||||
* zero-link pages — they WERE processed.
|
||||
*/
|
||||
async function extractStaleFromDB(
|
||||
export async function extractStaleFromDB(
|
||||
engine: BrainEngine,
|
||||
opts: {
|
||||
dryRun: boolean;
|
||||
|
||||
+39
-1
@@ -1479,7 +1479,31 @@ export async function registerBuiltinHandlers(
|
||||
embedSkipReason = 'auto_embed_disabled';
|
||||
}
|
||||
|
||||
return { ...result, embed_job_id: embedJobId, embed_skip_reason: embedSkipReason };
|
||||
// #2849: large-sync extract deferral follow-up. performSync skips inline
|
||||
// link/timeline extraction when totalChanges > 100, leaving
|
||||
// links_extracted_at unstamped. A standalone sync job (webhook push,
|
||||
// sync trigger) has no autopilot extract phase behind it, so the pages
|
||||
// would stay extraction-stale until a manual `gbrain extract --stale`.
|
||||
// Queue a source-scoped stale sweep instead. Best-effort + idempotent:
|
||||
// a duplicate sweep finds 0 stale pages and no-ops.
|
||||
let extractJobId: number | null = null;
|
||||
if (result.extractDeferred) {
|
||||
try {
|
||||
const { MinionQueue } = await import('../core/minions/queue.ts');
|
||||
const queue = new MinionQueue(engine);
|
||||
const followUp = await queue.add(
|
||||
'extract',
|
||||
{ stale: true, ...(sourceId ? { sourceId } : {}) },
|
||||
{
|
||||
idempotency_key: `sync-extract-stale:${sourceId ?? 'default'}:${Math.floor(Date.now() / 30_000)}`,
|
||||
maxWaiting: 1,
|
||||
},
|
||||
);
|
||||
extractJobId = followUp.id;
|
||||
} catch { /* best-effort: extract --stale sweeps it later */ }
|
||||
}
|
||||
|
||||
return { ...result, embed_job_id: embedJobId, embed_skip_reason: embedSkipReason, extract_stale_job_id: extractJobId };
|
||||
});
|
||||
|
||||
registerBuiltinJob(worker, engine, 'embed', async (job) => {
|
||||
@@ -1652,6 +1676,20 @@ export async function registerBuiltinHandlers(
|
||||
});
|
||||
|
||||
worker.register('extract', async (job) => {
|
||||
// #2849: stale-sweep mode — the sync handler's large-sync deferral
|
||||
// follow-up. DB-source (reads page content from the DB, so it runs on
|
||||
// checkout-less brains), source-scopable, idempotent. Same core as
|
||||
// `gbrain extract --stale`.
|
||||
if (job.data.stale === true) {
|
||||
const { extractStaleFromDB } = await import('./extract.ts');
|
||||
return await extractStaleFromDB(engine, {
|
||||
dryRun: !!job.data.dryRun,
|
||||
jsonMode: false,
|
||||
includeFrontmatter: false,
|
||||
sourceIdFilter: typeof job.data.sourceId === 'string' ? job.data.sourceId : undefined,
|
||||
catchUp: false,
|
||||
});
|
||||
}
|
||||
const { runExtractCore } = await import('./extract.ts');
|
||||
const mode = (typeof job.data.mode === 'string' && ['links', 'timeline', 'all'].includes(job.data.mode))
|
||||
? (job.data.mode as 'links' | 'timeline' | 'all')
|
||||
|
||||
@@ -2146,8 +2146,13 @@ export async function runServeHttp(engine: BrainEngine, options: ServeHttpOption
|
||||
// Other event types (ping, pull_request, etc.) return 202 'ignored'
|
||||
// so GitHub doesn't retry.
|
||||
// D15.5: HMAC compare uses the shared safeHexEqual helper.
|
||||
// D18: submits 'sync' job with auto_embed_backfill=true and priority -10
|
||||
// (above autopilot's 0).
|
||||
// D18: submits 'sync' job with extraction + auto_embed_backfill enabled and
|
||||
// priority -10 (above autopilot's 0). noExtract:false opts normal
|
||||
// incremental pushes into sync's inline link/timeline extraction (#2849
|
||||
// — the standalone sync handler defaults noExtract to TRUE, which left
|
||||
// webhook-imported pages permanently stale). Large (>100 file) pushes
|
||||
// defer inline extract; the sync handler queues an extract --stale
|
||||
// follow-up job for that branch.
|
||||
// ---------------------------------------------------------------------------
|
||||
const githubWebhookLimiter = rateLimit({
|
||||
windowMs: 60_000,
|
||||
@@ -2267,6 +2272,7 @@ export async function runServeHttp(engine: BrainEngine, options: ServeHttpOption
|
||||
'sync',
|
||||
{
|
||||
sourceId: source.id,
|
||||
noExtract: false,
|
||||
auto_embed_backfill: true,
|
||||
embed_reason: 'webhook',
|
||||
},
|
||||
|
||||
+19
-1
@@ -222,6 +222,14 @@ export interface SyncResult {
|
||||
* everything," the exact misdiagnosis in the #1794 recurrence report.
|
||||
*/
|
||||
bankedFiles?: number;
|
||||
/**
|
||||
* #2849: true when extraction was REQUESTED (noExtract false) but this sync
|
||||
* skipped inline link/timeline extraction because totalChanges > 100 (the
|
||||
* #1794 large-sync deferral). links_extracted_at stays unstamped for the
|
||||
* imported pages. The standalone `sync` job handler queues a source-scoped
|
||||
* `extract --stale` follow-up when set; CLI runs print the manual hint.
|
||||
*/
|
||||
extractDeferred?: boolean;
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -1379,6 +1387,10 @@ See also:
|
||||
{
|
||||
sourceId: sourceIdArg,
|
||||
repoPath: source.local_path,
|
||||
// #2849: opt in to inline extraction — the standalone sync handler
|
||||
// defaults noExtract to TRUE (dedupe for doctor's [sync, extract]
|
||||
// remediation plan), which would leave triggered syncs extraction-stale.
|
||||
noExtract: false,
|
||||
auto_embed_backfill: true,
|
||||
embed_reason: 'sync_trigger',
|
||||
},
|
||||
@@ -3287,11 +3299,16 @@ async function performSyncInner(engine: BrainEngine, opts: SyncOpts): Promise<Sy
|
||||
// the stale sweep scans the whole source, so banked-across-runs pages are
|
||||
// covered regardless.
|
||||
const extractOpts = opts.sourceId ? { sourceId: opts.sourceId } : undefined;
|
||||
let extractDeferred = false;
|
||||
if (!opts.noExtract && totalChanges > 100 && pagesAffected.length > 0) {
|
||||
// #2849: surface the deferral to callers. A standalone sync job (webhook
|
||||
// push, sync trigger) has no autopilot extract phase behind it, so the
|
||||
// job handler queues an `extract --stale` follow-up off this flag.
|
||||
extractDeferred = true;
|
||||
slog(
|
||||
` Large sync: deferring link/timeline extraction. ` +
|
||||
`Run 'gbrain extract --stale${opts.sourceId ? ` --source-id ${opts.sourceId}` : ''}' ` +
|
||||
`(or let the autopilot cycle's extract phase sweep it).`,
|
||||
`(sync jobs queue this follow-up automatically).`,
|
||||
);
|
||||
}
|
||||
if (!opts.noExtract && totalChanges <= 100 && pagesAffected.length > 0) {
|
||||
@@ -3400,6 +3417,7 @@ async function performSyncInner(engine: BrainEngine, opts: SyncOpts): Promise<Sy
|
||||
chunksCreated,
|
||||
embedded,
|
||||
pagesAffected,
|
||||
extractDeferred,
|
||||
};
|
||||
}
|
||||
|
||||
|
||||
+20
-101
@@ -428,13 +428,8 @@ export async function runPhaseSynthesize(
|
||||
|
||||
const queue = new MinionQueue(engine);
|
||||
const childIds: number[] = [];
|
||||
/**
|
||||
* 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>();
|
||||
/** Map child job_id → chunk metadata for D6 orchestrator-side slug rewrite. */
|
||||
const chunkInfo = new Map<number, { idx: number; hash6: string }>();
|
||||
/** Skip reasons for the cycle report (D5 cap hits, D8 legacy-key skips). */
|
||||
const skipReports: Array<{ filePath: string; reason: string }> = [];
|
||||
|
||||
@@ -518,14 +513,9 @@ export async function runPhaseSynthesize(
|
||||
{ allowProtectedSubmit: true },
|
||||
);
|
||||
childIds.push(child.id);
|
||||
childMeta.set(child.id, {
|
||||
idx: i,
|
||||
hash6,
|
||||
chunkTotal: chunks.length,
|
||||
transcriptSource: t.transcriptSource,
|
||||
transcriptId: stripContentVersionSuffix(t.basename),
|
||||
inferredDate: t.inferredDate,
|
||||
});
|
||||
if (isChunked) {
|
||||
chunkInfo.set(child.id, { idx: i, hash6 });
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -554,14 +544,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: childMeta drives post-hoc rewrite of
|
||||
// D6 orchestrator slug rewrite: chunkInfo 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, childMeta, cycleSourceId);
|
||||
const writtenRefs = await collectChildPutPageSlugs(engine, childIds, chunkInfo, cycleSourceId);
|
||||
|
||||
const summaryDate = opts.date ?? today();
|
||||
|
||||
@@ -569,12 +559,7 @@ 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.
|
||||
// #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);
|
||||
await stampDreamProvenance(engine, writtenRefs, summaryDate);
|
||||
|
||||
// Dual-write: reverse-render each DB row → markdown file.
|
||||
const reverseWriteCount = await reverseWriteRefs(engine, opts.brainDir, writtenRefs, cycleSourceId);
|
||||
@@ -1110,17 +1095,15 @@ function sanitizeForSlug(s: string): string {
|
||||
* fake"): we no longer need detection because the rewrite enforces
|
||||
* uniqueness at slug-write time.
|
||||
*
|
||||
* `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.
|
||||
* `chunkInfo` maps child job_id → { chunk_index, hash6 }. Single-chunk
|
||||
* children are absent from the map and pass through unchanged.
|
||||
*/
|
||||
async function collectChildPutPageSlugs(
|
||||
engine: BrainEngine,
|
||||
childIds: number[],
|
||||
childMeta: Map<number, ChildMeta>,
|
||||
chunkInfo: Map<number, { idx: number; hash6: string }>,
|
||||
sourceId = 'default',
|
||||
): Promise<Array<{ slug: string; source_id: string; jobId: number }>> {
|
||||
): Promise<Array<{ slug: string; source_id: string }>> {
|
||||
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
|
||||
@@ -1139,73 +1122,16 @@ async function collectChildPutPageSlugs(
|
||||
FROM subagent_tool_executions
|
||||
WHERE job_id = ANY($1::int[])
|
||||
AND tool_name = 'brain_put_page'
|
||||
AND status = 'complete'
|
||||
ORDER BY id`,
|
||||
AND status = 'complete'`,
|
||||
[childIds],
|
||||
);
|
||||
const rewritten = new Map<string, number>();
|
||||
const rewritten = new Set<string>();
|
||||
for (const r of rows) {
|
||||
if (typeof r.slug !== 'string' || r.slug.length === 0) continue;
|
||||
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);
|
||||
const ci = chunkInfo.get(r.job_id);
|
||||
rewritten.add(ci ? rewriteChunkedSlug(r.slug, ci.hash6, ci.idx) : r.slug);
|
||||
}
|
||||
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;
|
||||
return Array.from(rewritten).sort().map(slug => ({ slug, source_id: sourceId }));
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -1251,19 +1177,12 @@ async function hasLegacySingleChunkCompletion(
|
||||
*/
|
||||
async function stampDreamProvenance(
|
||||
engine: BrainEngine,
|
||||
refs: Array<{ slug: string; source_id: string; jobId?: number }>,
|
||||
refs: Array<{ slug: string; source_id: string }>,
|
||||
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, 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 };
|
||||
for (const { slug, source_id } of refs) {
|
||||
try {
|
||||
await executeRawJsonb(
|
||||
engine,
|
||||
@@ -1271,7 +1190,7 @@ async function stampDreamProvenance(
|
||||
SET frontmatter = COALESCE(frontmatter, '{}'::jsonb) || $3::jsonb
|
||||
WHERE slug = $1 AND source_id = $2`,
|
||||
[slug, source_id],
|
||||
[stamp],
|
||||
[{ dream_generated: true, dream_cycle_date: cycleDate }],
|
||||
);
|
||||
} 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, dirname } from 'node:path';
|
||||
import { join, basename } from 'node:path';
|
||||
import { createHash } from 'node:crypto';
|
||||
import { pruneDir } from '../sync.ts';
|
||||
|
||||
@@ -23,22 +23,8 @@ export interface DiscoveredTranscript {
|
||||
content: string;
|
||||
/** Filename basename without extension; used as a topic-slug seed. */
|
||||
basename: string;
|
||||
/**
|
||||
* 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.
|
||||
*/
|
||||
/** Inferred date if the basename matches `YYYY-MM-DD...` (or null). */
|
||||
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 {
|
||||
@@ -175,36 +161,6 @@ 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
|
||||
@@ -269,11 +225,8 @@ 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 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;
|
||||
const inferredDate = dateMatch ? dateMatch[1] : null;
|
||||
if (!isInDateRange(inferredDate, opts)) continue;
|
||||
|
||||
let content: string;
|
||||
try {
|
||||
@@ -288,17 +241,12 @@ 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),
|
||||
});
|
||||
}
|
||||
}
|
||||
@@ -342,7 +290,6 @@ export function readSingleTranscript(
|
||||
contentHash: hashContent(content),
|
||||
content,
|
||||
basename: baseName,
|
||||
inferredDate: inferContentDate(content) ?? (dateMatch ? dateMatch[1] : null),
|
||||
transcriptSource: deriveTranscriptSource(filePath),
|
||||
inferredDate: dateMatch ? dateMatch[1] : null,
|
||||
};
|
||||
}
|
||||
|
||||
@@ -26,7 +26,6 @@ const transcript: DiscoveredTranscript = {
|
||||
content: 'User: hello world',
|
||||
contentHash: 'abcdef0123456789',
|
||||
inferredDate: '2026-07-17',
|
||||
transcriptSource: null,
|
||||
} as DiscoveredTranscript;
|
||||
|
||||
describe('#2415: buildSynthesisPrompt output root', () => {
|
||||
|
||||
@@ -152,94 +152,3 @@ 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,7 +302,6 @@ describe('judgeSignificance', () => {
|
||||
content: 'A short conversation about something interesting.',
|
||||
basename: 'x',
|
||||
inferredDate: null,
|
||||
transcriptSource: null,
|
||||
};
|
||||
}
|
||||
|
||||
@@ -416,7 +415,6 @@ describe('judgeSignificance — UTF-16 safety (v0.41.13)', () => {
|
||||
content,
|
||||
basename: 'long',
|
||||
inferredDate: null,
|
||||
transcriptSource: null,
|
||||
};
|
||||
}
|
||||
|
||||
|
||||
@@ -45,7 +45,6 @@ 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', () => {
|
||||
|
||||
@@ -1,109 +0,0 @@
|
||||
/**
|
||||
* #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');
|
||||
});
|
||||
});
|
||||
@@ -216,21 +216,17 @@ 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, 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).
|
||||
// Baseline: one handler already installed by prior tests in this file.
|
||||
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 NET leaked live entries.
|
||||
// After 50 reporter lifecycles, still exactly one handler and zero leaked live entries.
|
||||
expect(__signalHandlerInstalledForTest()).toBe(installedBefore || true);
|
||||
expect(__liveReporterCountForTest()).toBe(liveBefore);
|
||||
expect(__liveReporterCountForTest()).toBe(0);
|
||||
});
|
||||
|
||||
test('startHeartbeat() fires heartbeats and stop() clears', async () => {
|
||||
|
||||
@@ -16,6 +16,7 @@
|
||||
*/
|
||||
import { describe, test, expect } from 'bun:test';
|
||||
import { createHmac } from 'node:crypto';
|
||||
import { readFileSync } from 'node:fs';
|
||||
import { safeHexEqual } from '../src/core/timing-safe.ts';
|
||||
|
||||
const GITHUB_SECRET = 'super-secret-webhook-key';
|
||||
@@ -123,3 +124,25 @@ describe('Branch ref construction (D5)', () => {
|
||||
expect(pushedRef === `refs/heads/${trackedBranch}`).toBe(false);
|
||||
});
|
||||
});
|
||||
|
||||
describe('Webhook sync job extraction contract (#2849)', () => {
|
||||
test('opts into extraction before the pushed commit is consumed', () => {
|
||||
const serveSource = readFileSync(
|
||||
new URL('../src/commands/serve-http.ts', import.meta.url),
|
||||
'utf8',
|
||||
);
|
||||
const routeStart = serveSource.indexOf("'/webhooks/github'");
|
||||
const queueStart = serveSource.indexOf('const job = await queue.add(', routeStart);
|
||||
const responseStart = serveSource.indexOf('res.status(202)', queueStart);
|
||||
expect(routeStart).toBeGreaterThanOrEqual(0);
|
||||
expect(queueStart).toBeGreaterThan(routeStart);
|
||||
expect(responseStart).toBeGreaterThan(queueStart);
|
||||
|
||||
const routeSource = serveSource.slice(queueStart, responseStart);
|
||||
const payload = routeSource.match(
|
||||
/queue\.add\(\s*'sync',\s*\{([\s\S]*?)\}\s*,\s*\{/,
|
||||
);
|
||||
expect(payload).not.toBeNull();
|
||||
expect(payload?.[1]).toMatch(/\bnoExtract:\s*false\b/);
|
||||
});
|
||||
});
|
||||
|
||||
@@ -0,0 +1,130 @@
|
||||
/**
|
||||
* #2849 — large-sync extract deferral queues an `extract --stale` follow-up.
|
||||
*
|
||||
* performSync's incremental path skips inline link/timeline extraction when
|
||||
* totalChanges > 100 (the #1794 large-sync deferral), leaving
|
||||
* links_extracted_at unstamped. Pre-fix, a standalone sync job (webhook push,
|
||||
* `gbrain sync trigger`) had NOTHING behind it to sweep those pages — the
|
||||
* autopilot cycle's extract phase only walks that cycle's changedSlugs — so a
|
||||
* large webhook push left extraction permanently stale until a manual
|
||||
* `gbrain extract --stale`.
|
||||
*
|
||||
* Pins:
|
||||
* (a) performSync surfaces `extractDeferred: true` on the >100 branch and
|
||||
* leaves the pages unstamped/unlinked.
|
||||
* (b) the `sync` job handler queues an `extract` job with
|
||||
* { stale: true, sourceId? } when extractDeferred is set.
|
||||
* (c) the `extract` handler's stale mode actually sweeps: links created +
|
||||
* watermark stamped (end-to-end recovery, no manual step).
|
||||
*
|
||||
* Marked .serial.test.ts — spawns git subprocesses + shares one PGLite engine.
|
||||
*/
|
||||
|
||||
import { describe, test, expect, beforeAll, afterAll } from 'bun:test';
|
||||
import { mkdtempSync, writeFileSync, rmSync, mkdirSync } from 'fs';
|
||||
import { execSync } from 'child_process';
|
||||
import { tmpdir } from 'os';
|
||||
import { join } from 'path';
|
||||
import { PGLiteEngine } from '../src/core/pglite-engine.ts';
|
||||
import { MinionWorker } from '../src/core/minions/worker.ts';
|
||||
import { MinionQueue } from '../src/core/minions/queue.ts';
|
||||
import { registerBuiltinHandlers } from '../src/commands/jobs.ts';
|
||||
|
||||
let engine: PGLiteEngine;
|
||||
let worker: MinionWorker;
|
||||
let repoPath: string;
|
||||
|
||||
function git(cmd: string): void { execSync(cmd, { cwd: repoPath, stdio: 'pipe' }); }
|
||||
|
||||
describe('#2849 — large sync defers extract and queues a stale sweep', () => {
|
||||
beforeAll(async () => {
|
||||
engine = new PGLiteEngine();
|
||||
await engine.connect({});
|
||||
await engine.initSchema();
|
||||
worker = new MinionWorker(engine, { queue: 'test' });
|
||||
await registerBuiltinHandlers(worker, engine, { quiet: true });
|
||||
|
||||
repoPath = mkdtempSync(join(tmpdir(), 'gbrain-large-defer-'));
|
||||
git('git init');
|
||||
git('git config user.email "t@t.com"');
|
||||
git('git config user.name "T"');
|
||||
mkdirSync(join(repoPath, 'people'), { recursive: true });
|
||||
mkdirSync(join(repoPath, 'notes'), { recursive: true });
|
||||
writeFileSync(join(repoPath, 'people/alice.md'), [
|
||||
'---', 'type: person', 'title: Alice', '---', '', 'Alice is a founder.',
|
||||
].join('\n'));
|
||||
git('git add -A && git commit -m "initial"');
|
||||
|
||||
// Seed: full first sync imports the anchor page + sets last_commit.
|
||||
const { performSync } = await import('../src/commands/sync.ts');
|
||||
await performSync(engine, { repoPath, full: true, noPull: true, noEmbed: true });
|
||||
|
||||
// Second commit: 101 new pages → incremental totalChanges > 100.
|
||||
for (let i = 0; i < 101; i++) {
|
||||
writeFileSync(join(repoPath, `notes/n${i}.md`), [
|
||||
'---', 'type: note', `title: Note ${i}`, '---', '',
|
||||
`[Alice](people/alice) appears in note ${i}.`,
|
||||
].join('\n'));
|
||||
}
|
||||
git('git add -A && git commit -m "add 101 pages"');
|
||||
}, 120_000);
|
||||
|
||||
afterAll(async () => {
|
||||
if (repoPath) rmSync(repoPath, { recursive: true, force: true });
|
||||
if (engine) await engine.disconnect();
|
||||
}, 60_000);
|
||||
|
||||
test('sync handler defers inline extract and queues extract{stale} follow-up; stale sweep recovers', async () => {
|
||||
const syncHandler = (worker as unknown as { handlers: Map<string, (job: unknown) => Promise<unknown>> })
|
||||
.handlers.get('sync');
|
||||
expect(syncHandler).toBeDefined();
|
||||
|
||||
// Same payload shape the webhook submits (minus embed backfill noise).
|
||||
const result = await syncHandler!({
|
||||
data: { repoPath, noExtract: false, noPull: true, auto_embed_backfill: false },
|
||||
signal: { aborted: false },
|
||||
updateProgress: async () => {},
|
||||
}) as { status: string; extractDeferred?: boolean; extract_stale_job_id?: number | null };
|
||||
|
||||
expect(result.status).toBe('synced');
|
||||
// (a) inline extract was deferred, pages left stale.
|
||||
expect(result.extractDeferred).toBe(true);
|
||||
const staleBefore = await engine.countStalePagesForExtraction();
|
||||
expect(staleBefore).toBeGreaterThan(100);
|
||||
expect(await engine.getLinks('notes/n0')).toHaveLength(0);
|
||||
|
||||
// (b) a follow-up extract job with stale:true was queued.
|
||||
expect(result.extract_stale_job_id).toBeGreaterThan(0);
|
||||
const queue = new MinionQueue(engine);
|
||||
const extractJobs = await queue.getJobs({ name: 'extract', limit: 5 });
|
||||
expect(extractJobs.length).toBe(1);
|
||||
expect((extractJobs[0].data as { stale: boolean }).stale).toBe(true);
|
||||
|
||||
// (c) running the extract handler's stale mode recovers: links + stamps.
|
||||
const extractHandler = (worker as unknown as { handlers: Map<string, (job: unknown) => Promise<unknown>> })
|
||||
.handlers.get('extract');
|
||||
await extractHandler!({
|
||||
data: extractJobs[0].data,
|
||||
signal: { aborted: false },
|
||||
updateProgress: async () => {},
|
||||
});
|
||||
const links = await engine.getLinks('notes/n0');
|
||||
expect(links.some(l => l.to_slug === 'people/alice')).toBe(true);
|
||||
const rows = await engine.executeRaw<{ links_extracted_at: string | null }>(
|
||||
`SELECT links_extracted_at FROM pages WHERE slug = 'notes/n0'`,
|
||||
);
|
||||
expect(rows[0]?.links_extracted_at).not.toBeNull();
|
||||
}, 180_000);
|
||||
|
||||
test('sub-threshold sync does NOT set extractDeferred (no spurious follow-up)', async () => {
|
||||
// One more small commit → inline extract path, no deferral.
|
||||
writeFileSync(join(repoPath, 'notes/small.md'), [
|
||||
'---', 'type: note', 'title: Small', '---', '', 'No big deal.',
|
||||
].join('\n'));
|
||||
git('git add -A && git commit -m "one small page"');
|
||||
const { performSync } = await import('../src/commands/sync.ts');
|
||||
const result = await performSync(engine, { repoPath, noPull: true, noEmbed: true });
|
||||
expect(result.status).toBe('synced');
|
||||
expect(result.extractDeferred).toBeFalsy();
|
||||
}, 60_000);
|
||||
});
|
||||
@@ -100,6 +100,10 @@ describe('runSyncTrigger', () => {
|
||||
const job = jobs[0];
|
||||
expect(job.priority).toBe(-10);
|
||||
expect((job.data as { sourceId: string }).sourceId).toBe('default');
|
||||
// #2849: opt in to inline extraction — the standalone sync handler
|
||||
// defaults noExtract to TRUE, which would leave triggered syncs
|
||||
// extraction-stale.
|
||||
expect((job.data as { noExtract: boolean }).noExtract).toBe(false);
|
||||
expect((job.data as { auto_embed_backfill: boolean }).auto_embed_backfill).toBe(true);
|
||||
});
|
||||
|
||||
|
||||
Reference in New Issue
Block a user