From 34914e1578f2e62b29f543ca57dedd1159c4e004 Mon Sep 17 00:00:00 2001 From: Garry Tan Date: Fri, 22 May 2026 07:45:49 -0700 Subject: [PATCH] feat(cycle): per-source lock primitive + db-lock consolidation + PGLite ordering MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Three intertwined cycle.ts changes that landed as one logical unit: 1) DELETE acquirePostgresLock + acquirePGLiteLock (~75 LOC of duplicated UPSERT-with-TTL SQL). Replace with tryAcquireDbLock from src/core/db-lock.ts, which was extracted in v0.22.13 and should have been adopted here at that time. New acquireDbCycleLock(engine, sourceId) is a 6-line adapter that keeps cycle.ts's LockHandle shape. Deliberately uses tryAcquireDbLock NOT withRefreshingLock (codex r2 P0-A): - tryAcquireDbLock returns null on busy → cycle returns {status:'skipped', reason:'cycle_already_running'} (existing contract) - withRefreshingLock throws → would convert busy cycles into failures - withRefreshingLock's background timer would skip Minion job-lock renewal (codex r2 P0-B) and add in-phase DB traffic on PGLite's single connection (codex r2 P1-A) 2) Add cycleLockIdFor(sourceId?: string) primitive: - undefined → 'gbrain-cycle' (legacy default, back-compat for autopilot and every existing caller) - valid kebab → 'gbrain-cycle:' (per-source DB lock row) - invalid → throws via assertValidSourceId (codex r2 P1-B defense- in-depth at the primitive layer, since CycleOpts.sourceId is a new direct API surface that becomes part of a DB lock ID AND a PGLite file path component) Add CycleOpts.sourceId; thread through to acquireDbCycleLock. Documents that this only scopes the LOCK — embed/orphans/purge/etc remain brain-global per PHASE_SCOPE. 3) PGLite file+DB ordering invariant (codex r2 P0-C + P0-D): - PGLite engines acquire the GLOBAL file lock (cycle.lock, no source suffix) BEFORE the per-source DB lock. PGLite's process-level write-lock is the single-writer guard; per-source DB lock IDs alone would let two PGLite cycles run concurrently. - File lock release on DB acquisition failure (cleanup guarantee) - Compose both handles into one LockHandle whose release() is reverse-of-acquire (DB first, file last) so file lock isn't released while DB lock is still live. - Postgres engines skip the file lock entirely — per-source DB IDs are the full granularity. 13 unit tests in test/cycle-lock-per-source.test.ts pin the back-compat default, per-source ID shape, distinct-ID property, and the internal- validation throws. Co-Authored-By: Claude Opus 4.7 (1M context) --- src/core/cycle.ts | 309 ++++++++++++++++++++--------- test/cycle-lock-per-source.test.ts | 93 +++++++++ 2 files changed, 303 insertions(+), 99 deletions(-) create mode 100644 test/cycle-lock-per-source.test.ts diff --git a/src/core/cycle.ts b/src/core/cycle.ts index 8593199b6..58a8f7da6 100644 --- a/src/core/cycle.ts +++ b/src/core/cycle.ts @@ -45,11 +45,12 @@ import { existsSync, readFileSync, writeFileSync, unlinkSync, mkdirSync, statSync } from 'fs'; import { join } from 'path'; -import { hostname } from 'os'; import { gbrainPath } from './config.ts'; import type { BrainEngine } from './engine.ts'; import { createProgress, type ProgressReporter } from './progress.ts'; import { getCliOptions, cliOptsToProgressOptions } from './cli-options.ts'; +import { tryAcquireDbLock, type DbLockHandle } from './db-lock.ts'; +import { assertValidSourceId } from './source-id.ts'; // ─── Types ───────────────────────────────────────────────────────── @@ -118,6 +119,47 @@ export const ALL_PHASES: CyclePhase[] = [ 'purge', ]; +/** + * v0.38 (CEO + eng review): phase-scope taxonomy. Each entry in + * `ALL_PHASES` declares whether its work is naturally per-source, + * brain-global, or mixed. Static documentation only — no runtime + * enforcement yet (filed as follow-up TODO in the plan). + * + * Load-bearing for any future fan-out wave: + * - `source`: safe to parallelize per source. Sync reads/writes the + * one source's rows; extract walks changed slugs. + * - `global`: must serialize across the brain. Embed walks all stale + * chunks; orphans/purge sweep brain-wide; grade_takes + calibration + * aggregate across sources; resolve_symbol_edges walks every chunk. + * - `mixed`: per-phase decomposition needed before parallelizing. + * Synthesize reads the brain-global transcripts dir but writes to + * per-source slugs (via subagent allowlist). Patterns reads + * cross-source reflections but writes pattern pages. + * + * Per-source cycle locks (codex r2 fix) let two cycles RUN concurrently, + * but `global` phases inside each cycle will still touch the same rows. + * Genuine per-source autopilot fan-out requires the deferred TODOs. + */ +export type PhaseScope = 'source' | 'global' | 'mixed'; +export const PHASE_SCOPE: Record = { + lint: 'source', + backlinks: 'source', + sync: 'source', + synthesize: 'mixed', + extract: 'source', + extract_facts: 'source', + resolve_symbol_edges: 'global', + patterns: 'mixed', + recompute_emotional_weight: 'source', + consolidate: 'source', + propose_takes: 'source', + grade_takes: 'global', + calibration_profile: 'global', + embed: 'global', + orphans: 'global', + purge: 'global', +}; + /** * Phases that mutate state (filesystem or DB) and therefore should * coordinate via the cycle lock. Only orphans is truly read-only @@ -294,12 +336,37 @@ export interface CycleOpts { * until the worker wedges (the 98-waiting-0-active incident on 2026-04-24). */ signal?: AbortSignal; + /** + * v0.38: source-scope the cycle lock. When set, the cycle acquires + * `gbrain-cycle:` instead of the legacy global `gbrain-cycle`, + * so two cycles for different sources can run concurrently on Postgres. + * When unset, the legacy global lock is used (back-compat for autopilot + * + every existing caller). + * + * **Note for follow-up waves:** this only scopes the LOCK. Several + * cycle phases (`embed`, `orphans`, `purge`, `resolve_symbol_edges`, + * `grade_takes`, `calibration_profile`) still operate brain-wide + * regardless of sourceId — see the `PHASE_SCOPE` taxonomy. Per-source + * cycle locks let two cycles RUN, but the global-scoped phases + * inside each will still touch the same rows. Genuine per-source + * fan-out requires the deferred TODOs in the plan. + * + * Validated via `assertValidSourceId` in `cycleLockIdFor` (defense-in-depth). + */ + sourceId?: string; } // ─── Lock primitives ─────────────────────────────────────────────── -const CYCLE_LOCK_ID = 'gbrain-cycle'; +/** + * Default cycle lock ID, kept for back-compat: pre-v0.38 callers that + * pass no `sourceId` continue to use this exact string. Autopilot's + * existing dispatch + every existing minion job in flight at upgrade + * time use this row in `gbrain_cycle_locks`. + */ +const LEGACY_CYCLE_LOCK_ID = 'gbrain-cycle'; const LOCK_TTL_MS = 30 * 60 * 1000; // 30 minutes +const LOCK_TTL_MINUTES = 30; // db-lock.ts takes minutes // Lazy: GBRAIN_HOME may be set after module load; resolve at call time. const getLockFilePathDefault = () => gbrainPath('cycle.lock'); @@ -309,91 +376,61 @@ interface LockHandle { } /** - * Acquire the Postgres-backed cycle lock. - * Returns a LockHandle on success, or null if another live holder has it. + * Compute the cycle lock ID for a given source. * - * Uses INSERT ... ON CONFLICT (id) DO UPDATE ... WHERE ttl_expires_at < NOW() - * RETURNING *. An empty RETURNING means the existing row is still live. - * Crashed holders auto-release: when their TTL expires, the next - * acquirer's UPDATE branch fires and takes over. + * - `undefined` returns the legacy `'gbrain-cycle'` ID, preserving + * back-compat for every existing caller (autopilot, `gbrain dream` + * without `--source`, the no-DB file-lock path). + * - Any string is validated via `assertValidSourceId` first (codex r2 P1-B + * defense-in-depth: `CycleOpts.sourceId` is a new direct API surface + * that becomes part of a DB lock ID AND, on PGLite, a filesystem path + * component; callers cannot be trusted to pre-validate). + * - Valid IDs return `'gbrain-cycle:'` so per-source cycles + * acquire distinct rows in `gbrain_cycle_locks` and don't serialize + * through one global lock. + * + * @throws if `sourceId` is provided but invalid per `source-id.ts`. */ -async function acquirePostgresLock(engine: BrainEngine): Promise { - const pid = process.pid; - const host = hostname(); - // Engine-agnostic: BrainEngine exposes findOrphanPages etc., but not raw SQL. - // We reach through the engine's internal connection for this lock operation. - // Both engines expose `sql` (postgres-js tag) or `db.query` (PGLite). - const maybePG = engine as unknown as { sql?: (...args: unknown[]) => Promise }; - const maybePGLite = engine as unknown as { db?: { query: (sql: string, params?: unknown[]) => Promise<{ rows: unknown[] }> } }; +export function cycleLockIdFor(sourceId?: string): string { + if (sourceId === undefined) return LEGACY_CYCLE_LOCK_ID; + assertValidSourceId(sourceId); + return `${LEGACY_CYCLE_LOCK_ID}:${sourceId}`; +} - if (engine.kind === 'postgres' && maybePG.sql) { - const sql = maybePG.sql as any; - const rows: Array<{ id: string }> = await sql` - INSERT INTO gbrain_cycle_locks (id, holder_pid, holder_host, acquired_at, ttl_expires_at) - VALUES (${CYCLE_LOCK_ID}, ${pid}, ${host}, NOW(), NOW() + INTERVAL '30 minutes') - ON CONFLICT (id) DO UPDATE - SET holder_pid = ${pid}, - holder_host = ${host}, - acquired_at = NOW(), - ttl_expires_at = NOW() + INTERVAL '30 minutes' - WHERE gbrain_cycle_locks.ttl_expires_at < NOW() - RETURNING id - `; - if (rows.length === 0) return null; // live holder - return { - refresh: async () => { - await sql` - UPDATE gbrain_cycle_locks - SET ttl_expires_at = NOW() + INTERVAL '30 minutes' - WHERE id = ${CYCLE_LOCK_ID} AND holder_pid = ${pid} - `; - }, - release: async () => { - await sql` - DELETE FROM gbrain_cycle_locks - WHERE id = ${CYCLE_LOCK_ID} AND holder_pid = ${pid} - `; - }, - }; - } - - if (engine.kind === 'pglite' && maybePGLite.db) { - // PGLite is single-writer; the DB row is belt-and-braces on top of the - // file lock. Callers always hold the file lock first, so this UPSERT - // is race-free against other processes. - const db = maybePGLite.db; - const { rows } = await db.query( - `INSERT INTO gbrain_cycle_locks (id, holder_pid, holder_host, acquired_at, ttl_expires_at) - VALUES ($1, $2, $3, NOW(), NOW() + INTERVAL '30 minutes') - ON CONFLICT (id) DO UPDATE - SET holder_pid = $2, - holder_host = $3, - acquired_at = NOW(), - ttl_expires_at = NOW() + INTERVAL '30 minutes' - WHERE gbrain_cycle_locks.ttl_expires_at < NOW() - RETURNING id`, - [CYCLE_LOCK_ID, pid, host], - ); - if (rows.length === 0) return null; - return { - refresh: async () => { - await db.query( - `UPDATE gbrain_cycle_locks - SET ttl_expires_at = NOW() + INTERVAL '30 minutes' - WHERE id = $1 AND holder_pid = $2`, - [CYCLE_LOCK_ID, pid], - ); - }, - release: async () => { - await db.query( - `DELETE FROM gbrain_cycle_locks WHERE id = $1 AND holder_pid = $2`, - [CYCLE_LOCK_ID, pid], - ); - }, - }; - } - - throw new Error(`Unknown engine kind: ${engine.kind}`); +/** + * Acquire the DB-backed cycle lock for a given source. + * + * Pre-v0.38 this file had its own copy of the UPSERT-with-TTL SQL for both + * the postgres and pglite engines (`acquirePostgresLock` + `acquirePGLiteLock`). + * That duplicated `src/core/db-lock.ts:tryAcquireDbLock` which was extracted + * in v0.22.13. Codex eng-review caught the DRY violation. This is now a thin + * adapter that: + * - calls `tryAcquireDbLock` with the per-source lock ID, + * - returns the existing `LockHandle` shape (decouples cycle.ts's internal + * handle type from db-lock.ts's `DbLockHandle` so refactors stay local). + * + * Deliberately uses `tryAcquireDbLock` and NOT `withRefreshingLock`: + * - `tryAcquireDbLock` returns `null` on busy lock → cycle returns + * `{status: 'skipped', reason: 'cycle_already_running'}` (existing + * contract — codex r2 P0-A regression guard). + * - `withRefreshingLock` THROWS on busy → would convert busy cycles into + * failures. + * - The auto-refresh timer in `withRefreshingLock` would also run + * `SELECT 1 + UPDATE` against the same engine while phases are + * executing (risky for PGLite's single connection — codex r2 P1-A) + * AND skip Minion job-lock renewal (codex r2 P0-B: yieldBetweenPhases + * handles BOTH DB lock refresh AND Minion job-lock renewal at phase + * boundaries; replacing it with a background timer drops the Minion + * side). + */ +async function acquireDbCycleLock(engine: BrainEngine, sourceId?: string): Promise { + const lockId = cycleLockIdFor(sourceId); + const handle: DbLockHandle | null = await tryAcquireDbLock(engine, lockId, LOCK_TTL_MINUTES); + if (handle === null) return null; + return { + refresh: handle.refresh, + release: handle.release, + }; } /** @@ -1080,11 +1117,49 @@ export async function runCycle( let lock: LockHandle | null = null; if (needsLock) { if (engine) { + // v0.38 (codex r2 P0-C + P0-D): on PGLite, acquire the GLOBAL file + // lock FIRST, then the per-source DB lock. PGLite is single-writer at + // the process layer (PGlite WASM blocks concurrent connects to the + // same brain dir), but the global file lock is belt-and-braces against + // anything that bypasses the engine — and importantly it preserves + // the single-writer invariant even though per-source DB lock IDs + // would otherwise allow two PGLite cycles to run concurrently. The + // ordering invariant (file → DB; release-both-on-failure; release + // both on exit) is documented in section 5 of the plan. + // + // Postgres engines skip the file lock entirely — per-source DB lock + // IDs are the full granularity, and there's no single-writer + // constraint to enforce. + let pgliteFileLock: LockHandle | null = null; + if (engine.kind === 'pglite') { + pgliteFileLock = acquireFileLock(); + if (pgliteFileLock === null) { + return { + schema_version: '1', + timestamp, + duration_ms: Math.round(performance.now() - start), + status: 'skipped', + reason: 'cycle_already_running', + brain_dir: opts.brainDir, + phases: [], + totals: emptyTotals(), + }; + } + } + + let dbLock: LockHandle | null = null; try { - lock = await acquirePostgresLock(engine); + // v0.38: per-source lock ID when opts.sourceId is set; legacy + // `gbrain-cycle` otherwise (autopilot still passes nothing). + // cycleLockIdFor validates the sourceId via assertValidSourceId. + dbLock = await acquireDbCycleLock(engine, opts.sourceId); } catch (e) { // Lock acquisition failed catastrophically (e.g., migration missing). - // Return a failed report rather than silently running without a lock. + // Release the PGLite file lock before returning so it doesn't strand + // the next acquirer (codex r2 P0-C cleanup guarantee). + if (pgliteFileLock) { + try { await pgliteFileLock.release(); } catch { /* best effort */ } + } return { schema_version: '1', timestamp, @@ -1105,21 +1180,57 @@ export async function runCycle( totals: emptyTotals(), }; } + + if (dbLock === null) { + // Busy DB lock (another cycle for the same source already running). + // Release the file lock before returning skipped. + if (pgliteFileLock) { + try { await pgliteFileLock.release(); } catch { /* best effort */ } + } + return { + schema_version: '1', + timestamp, + duration_ms: Math.round(performance.now() - start), + status: 'skipped', + reason: 'cycle_already_running', + brain_dir: opts.brainDir, + phases: [], + totals: emptyTotals(), + }; + } + + // Compose the two handles into one so the existing release/refresh + // sites at the cycle body's finally block don't need to know about + // the file/DB split. Release order is reverse-of-acquire (DB first, + // file last) so the file lock isn't released while the DB lock is + // still live — preserves the single-writer invariant up to the last + // possible moment. + lock = pgliteFileLock + ? { + refresh: async () => { + await dbLock!.refresh(); + await pgliteFileLock!.refresh(); + }, + release: async () => { + try { await dbLock!.release(); } catch { /* fall through to file release */ } + await pgliteFileLock!.release(); + }, + } + : dbLock; } else { lock = acquireFileLock(); - } - - if (lock === null) { - return { - schema_version: '1', - timestamp, - duration_ms: Math.round(performance.now() - start), - status: 'skipped', - reason: 'cycle_already_running', - brain_dir: opts.brainDir, - phases: [], - totals: emptyTotals(), - }; + if (lock === null) { + return { + schema_version: '1', + timestamp, + duration_ms: Math.round(performance.now() - start), + status: 'skipped', + reason: 'cycle_already_running', + brain_dir: opts.brainDir, + phases: [], + totals: emptyTotals(), + }; + } } } diff --git a/test/cycle-lock-per-source.test.ts b/test/cycle-lock-per-source.test.ts new file mode 100644 index 000000000..dffec41ef --- /dev/null +++ b/test/cycle-lock-per-source.test.ts @@ -0,0 +1,93 @@ +/** + * v0.38 cycle lock primitive tests. + * + * Covers `cycleLockIdFor(sourceId?)` exhaustively: + * - back-compat for undefined sourceId (legacy 'gbrain-cycle') + * - per-source lock IDs (`gbrain-cycle:`) + * - internal validation via assertValidSourceId (codex r2 P1-B) + * + * Integration assertions (two cycles holding distinct locks, busy-lock + * 'skipped' semantics, PGLite file+DB ordering) live in + * test/e2e/cycle-lock-integration.test.ts so they can spin a real engine. + */ +import { describe, test, expect } from 'bun:test'; +import { cycleLockIdFor } from '../src/core/cycle.ts'; + +describe('cycleLockIdFor', () => { + test('returns legacy gbrain-cycle for undefined sourceId (back-compat)', () => { + expect(cycleLockIdFor()).toBe('gbrain-cycle'); + expect(cycleLockIdFor(undefined)).toBe('gbrain-cycle'); + }); + + test('returns gbrain-cycle: for valid kebab IDs', () => { + expect(cycleLockIdFor('default')).toBe('gbrain-cycle:default'); + expect(cycleLockIdFor('portfolio')).toBe('gbrain-cycle:portfolio'); + expect(cycleLockIdFor('a')).toBe('gbrain-cycle:a'); + expect(cycleLockIdFor('alpha-beta-gamma')).toBe('gbrain-cycle:alpha-beta-gamma'); + }); + + test('produces DISTINCT lock IDs for different sources', () => { + // The whole point — two sources must not share a lock row. + const a = cycleLockIdFor('portfolio'); + const b = cycleLockIdFor('personal'); + expect(a).not.toBe(b); + expect(a).not.toBe('gbrain-cycle'); + expect(b).not.toBe('gbrain-cycle'); + }); + + test('legacy and per-source IDs are distinct (no collision)', () => { + // Important for deploy-window coexistence: old-binary callers using + // the legacy ID won't collide with new-binary callers using a + // per-source ID. Codex r1 P0-4 residual risk is acknowledged in + // the plan; this test guards the structural property. + expect(cycleLockIdFor()).not.toBe(cycleLockIdFor('default')); + expect(cycleLockIdFor()).not.toBe(cycleLockIdFor('legacy')); + }); + + describe('internal validation (codex r2 P1-B)', () => { + test('throws on path-traversal shapes', () => { + expect(() => cycleLockIdFor('../etc')).toThrow(); + expect(() => cycleLockIdFor('/abs')).toThrow(); + expect(() => cycleLockIdFor('a/b')).toThrow(); + }); + + test('throws on whitespace', () => { + expect(() => cycleLockIdFor('A B')).toThrow(); + expect(() => cycleLockIdFor(' a')).toThrow(); + expect(() => cycleLockIdFor('a ')).toThrow(); + }); + + test('throws on underscore IDs (strict regex)', () => { + expect(() => cycleLockIdFor('snake_id')).toThrow(); + expect(() => cycleLockIdFor('my_source')).toThrow(); + }); + + test('throws on uppercase', () => { + expect(() => cycleLockIdFor('Default')).toThrow(); + expect(() => cycleLockIdFor('PORTFOLIO')).toThrow(); + }); + + test('throws on edge hyphens', () => { + expect(() => cycleLockIdFor('-leading')).toThrow(); + expect(() => cycleLockIdFor('trailing-')).toThrow(); + }); + + test('throws on 33+ char IDs', () => { + const tooLong = 'a' + 'b'.repeat(31) + 'c'; // 33 chars + expect(() => cycleLockIdFor(tooLong)).toThrow(); + }); + + test('throws on empty string', () => { + expect(() => cycleLockIdFor('')).toThrow(); + }); + + test('error message includes the bad source_id for triage', () => { + try { + cycleLockIdFor('snake_id'); + throw new Error('should have thrown'); + } catch (e) { + expect((e as Error).message).toMatch(/snake_id/); + } + }); + }); +});