mirror of
https://github.com/garrytan/gbrain.git
synced 2026-07-30 11:22:34 +00:00
* v0.30.1 Lane A: connection-manager foundation + X1 initSchema routing Routes Postgres queries by query type: - read() goes to the Supabase pooler (port 6543, fast) - ddl() and bulk() go to direct (port 5432, 30min stmt timeout, mwm 256MB) Auto-detects Supabase via hostname pooler.supabase.com or port 6543. Override with GBRAIN_DIRECT_DATABASE_URL. Kill-switch via GBRAIN_DISABLE_DIRECT_POOL=1 falls back to single-pool legacy path. Foundation modules (Lane A scope): - src/core/connection-manager.ts: read/ddl/bulk/healthCheck, parent-CM inheritance (T5/X1), cached Promise<Sql> lazy init (A1), kill-switch inheritance (A2), Supabase URL auto-derivation - src/core/url-redact.ts: redactPgUrl + redactDeep (F3) - src/core/retry-matcher.ts: typed predicates for stmt-timeout / lock / conn errors (C4) - src/core/connection-audit.ts: ~/.gbrain/audit/connection-events JSONL with ISO-week rotation; doctor tail-reads last 5 errors (F8) - scripts/check-pg-url-redaction.sh: CI grep guard against unredacted postgresql:// URL leaks (F3) Engine integration: - PostgresEngine.connect: instantiates instance-owned ConnectionManager, inherits from parentConnectionManager when set (worker engines, sync, cycle), shares pool with module-singleton path - PostgresEngine.disconnect: tears down direct pool first - PostgresEngine.initSchema: routes DDL through connectionManager.ddl() when dual-pool active (X1 part 1; lock semantics replacement is Lane B) - cli.ts:connectEngine(opts): probeOnly skips initSchema entirely (X1 part 2 — get_health, upgrade --status will use this) Tests added (51 new cases): - test/url-redact.test.ts: 11 cases - test/retry-matcher.test.ts: 13 cases - test/connection-manager.test.ts: 27 cases (URL detection, derive, kill-switch, parent inheritance, dual-pool routing modes) Foundation for Lanes B-E. Sequential lane work continues. Plan: ~/.claude/plans/system-instruction-you-are-working-stateless-wadler.md Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com> * v0.30.1 Lane B: migration runner retry + verify hooks + namespaced --force flags Adds Migration interface fields: - idempotent: boolean (default true; explicit false blocks verify-hook re-runs on destructive migrations) - verify: optional post-condition probe; runs after migration claims success Migration retry wrapper (Cherry D3 / Finding F2): - 3 attempts with 5s/15s/45s backoff (env GBRAIN_MIGRATE_BACKOFF_MS=0 for tests) - Retries only on statement_timeout (57014) or connection-reset patterns - Pre-attempt: logs idle-in-transaction blockers via getIdleBlockers - On exhaustion: throws MigrationRetryExhausted with named PID + suggested pg_terminate_backend() recovery command Verify-hook self-healing (Cherry D6 / Codex X3): - On verify=false + idempotent=true → re-runs migration once silently - On verify=false + idempotent=false → throws MigrationDriftError - --skip-verify CLI flag bypasses for operator override withRefreshingLock helper (Cherry T4 / Codex A4 / X1 part 3): - setInterval refresh every TTL/6 ms during long-running work - SELECT 1 backend-alive heartbeat per refresh tick - Heartbeat hang past 30s → log + clear interval; lock TTL auto-expires - LockUnavailableError when acquire fails (caller decides retry) - buildTenantLockId(scope) appends current_database() suffix for multi-tenant safety (Cherry D4) Namespaced --force flags (Codex T5): - --force-orchestrator: write 'retry' markers for ALL wedged orchestrators - --force-schema: re-runs runMigrations against current config.version - --force / --force-all: both - --force-retry vX.Y.Z: existing single-version reset (preserved) - --skip-verify: bypass verify-hook drift detection on a single run Test additions: - test/migrate-extensions.test.ts: 14 cases (idempotent default, error envelopes, MIGRATIONS contract) - test/db-lock-refresh.test.ts: 10 cases (LockUnavailableError, buildTenantLockId multi-tenant, opts shape) - test/migrate.test.ts: updated 2 existing cases (PR #356 retry shape + function-name anchor) for v0.30.1 retry-wrapper semantics 156 unit tests passing across the v0.30.1 surface so far. Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com> * v0.30.1 Lane C: backfill primitive + registry + X4 + X5 First-class generic backfill runner (Fix 3). Generalizes the keyset+checkpoint+adaptive-batch pattern from src/core/backfill-effective-date.ts so future backfills (embedding_voyage in v0.30.2, etc.) reuse one tested runner. NEW src/core/backfill-base.ts: - runBackfill() with keyset pagination, config-table checkpoint, adaptive batch halving on stmt timeout, conn-drop reconnect, max-errors bail - ensureBackfillIndex() verifies/creates partial index CONCURRENTLY (P2/X4) - clearBackfillCheckpoint() for --fresh path - T3 fix: writes go through engine.withReservedConnection so BEGIN / SET LOCAL / UPDATE / COMMIT execute on the SAME backend (otherwise SET LOCAL evaporates between pooled executeRaw calls) NEW src/core/backfill-registry.ts: - effective_date: implemented (wraps existing computeEffectiveDate) - emotional_weight: implemented (wraps computeEmotionalWeight + stamps new emotional_weight_recomputed_at column) - embedding_voyage: declared-only in v0.30.1 (multi-column embedding schema lands in v0.30.2) NEW src/commands/backfill.ts: - gbrain backfill <kind> [--batch-size N] [--concurrency N] [--resume] [--fresh] [--dry-run] [--keep-index] [--max-errors N] - gbrain backfill list — shows registered backfills + status - X5 admission control: clampConcurrency() forces --concurrency to GBRAIN_DIRECT_POOL_SIZE - 1 ceiling (always reserves 1 conn for HNSW + heartbeat + doctor probes). Loud-warns when user requests above. Schema migration v44 (X4 / Codex C8 fix): - pages.emotional_weight_recomputed_at TIMESTAMPTZ - emotional_weight = 0 is a VALID steady-state value per migration v40, so the original P2 predicate ("WHERE emotional_weight = 0") would have been a permanent large index over normal data. The corrected backlog predicate is "emotional_weight_recomputed_at IS NULL"; the partial index drops naturally as the cycle phase + this backfill stamp the column over time. - idempotent: true (ADD COLUMN ... NULL is metadata-only) CLI integration: - src/cli.ts: registers `backfill` subcommand - reindex-frontmatter stays as thin alias for v0.30.1 back-compat; canonical entrypoint is now `gbrain backfill effective_date` Test additions: - test/backfill-base.test.ts: 11 cases (keyset, checkpoint, dry-run, resume/fresh, maxRows cap, withReservedConnection routing, error paths, clearCheckpoint, ensureBackfillIndex) - test/backfill-concurrency-clamp.test.ts: 6 cases (X5 admission control) 173 unit tests passing across Lanes A+B+C of v0.30.1. Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com> * v0.30.1 Lane D: HNSW lifecycle manager + A3 atomic-swap Extends src/core/vector-index.ts with the v0.30.1 lifecycle layer. The original chunkEmbeddingIndexSql / applyChunkEmbeddingIndexPolicy contract is preserved unchanged. New surfaces: - checkActiveBuild(engine, indexName): probes pg_stat_activity for an active CREATE INDEX or REINDEX on the named index. Used as pre-op guard so dropAndRebuild doesn't compete with a build already in flight (Supabase auto-maintenance, parallel gbrain procs). - dropZombieIndexes(engine, tableNames): startup sweep of indisvalid=false rows on gbrain tables. Drops them with DROP INDEX IF EXISTS, BUT skips any zombie that has an active build still in pg_stat_activity (codex Fix-5 in-progress-build guard). Wired into PostgresEngine.initSchema() — runs after migrations + verifySchema, best-effort, never blocks engine.connect(). - dropAndRebuild(engine, spec, opts): A3 atomic-swap pattern: 1. checkActiveBuild → bail if another build is active (--force overrides) 2. CREATE INDEX CONCURRENTLY <name>_rebuild_<unix-ms> via engine.withReservedConnection (CONCURRENTLY can't run in a txn) 3. Atomic swap inside engine.transaction: DROP INDEX <old-name> ALTER INDEX <temp-name> RENAME TO <old-name> 4. If step 2 fails (OOM, timeout, conn drop), the OLD index stays intact and search keeps serving queries. This is the headline A3 win — no production-degraded silent failure mode. - monitorBuild(engine, indexName, onProgress, opts): poll pg_stat_activity every 30s; emit elapsed_ms + size_bytes (via pg_relation_size) + pid. Used by gbrain backfill embedding_voyage when batch > 1000 triggers a rebuild. - isSupabaseAutoMaintenance(active): predicate on application_name (matches "supabase" / "postgres-meta"). Used by dropAndRebuild to log + back off when Supabase auto-maintenance is doing the rebuild. Engine integration: - PostgresEngine.initSchema() calls dropZombieIndexes after verifySchema. Surfaces zombie counts via console.log. - Best-effort wrapped in try/catch: pg_stat_activity / pg_index access can be restricted on managed Postgres tiers; gbrain shouldn't fail engine.connect() over diagnostic queries. Test additions (18 cases): - test/vector-index-lifecycle.test.ts: * chunkEmbeddingIndexSql contract (3 cases) — pre-existing behavior preserved * applyChunkEmbeddingIndexPolicy contract (1 case) * checkActiveBuild (4 cases, including PGLite no-op + best-effort failure) * isSupabaseAutoMaintenance (3 cases) * dropZombieIndexes (4 cases, including in-progress-build guard) * dropAndRebuild atomic-swap (3 cases, including PGLite + active-build bail + temp-name format assertion) 191 unit tests passing across Lanes A+B+C+D of v0.30.1. Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com> * v0.30.1 Lane E: upgrade pipeline checkpoint + brain_id binding + get_health migrations NEW src/core/upgrade-checkpoint.ts: - Cherry D5: persists step-by-step progress through gbrain post-upgrade so partial failures can be resumed via gbrain upgrade --resume. Steps: pull → install → schema → features → backfills → verify. - Codex X2: checkpoint binds to brain identity via sha256(database_url) (userinfo stripped before hashing so cred rotations don't invalidate). PGLite uses sha256(database_path). Cross-brain checkpoint application is now refused with reason='brain_mismatch'. - F4 fall-through: validateCheckpoint returns reason='no_checkpoint' when none exists, enabling silent fall-through to a full upgrade. - All-complete detection: stale checkpoints (every step done) return reason='all_complete' so the next run clears + re-runs from scratch. - markStepComplete + markStepFailed maintain the partial-state shape. T2 preserved: upgrade.ts still re-execs `gbrain post-upgrade` so the NEW binary's migration registry runs (the existing re-exec pattern is correct per codex round 1's plan-breaking finding). The checkpoint module is the substrate that Lane E's --resume / --status surfaces will plumb through in v0.30.2. D7 + C3 contract committed: - BrainHealth.schema_version: '1' (literal type) — additive-only contract pinned for MCP get_health consumers. - BrainHealth.migrations: { schema, orchestrator } — explicit two-ledger diagnostic surface (codex T5 namespacing). Both fields are OPTIONAL in v0.30.1 — engines can populate them in v0.30.2 without a contract bump. Backwards/forwards compat: clients default-handle missing fields. VERSION: 0.30.0 → 0.30.1 package.json: synced Test additions (18 cases): - test/upgrade-checkpoint.test.ts: * computeBrainId: userinfo strip, DB-distinct hashes, stable hex (5 cases) * write/load round-trip: roundtrip, missing file, malformed JSON, clear (4 cases) * validateCheckpoint: F4 no_checkpoint, X2 brain_mismatch, partial → resumeAt, all_complete, first-step pending (5 cases) * markStepComplete/markStepFailed: append, idempotent, clear-failed, failed-state shape (4 cases) 209 unit tests passing across all 5 lanes of v0.30.1 (Lanes A-E core foundations). Plumbing into upgrade.ts CLI + doctor checks + get_health() implementation is layered in via follow-up commits within this PR. Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com> * v0.30.1 e2e + test isolation: integration smoke + serial quarantine NEW test/e2e/v030_1-integration-pglite.test.ts (14 cases): PGLite integration smoke proving Lane A-E surfaces work together. Lane B: migration runner applies v44 (emotional_weight_recomputed_at) cleanly; config.version reaches LATEST_VERSION Lane C: backfill registry resolves all 3 entries; emotional_weight + effective_date backfills on empty brain return examined=0 cleanly Lane D: dropZombieIndexes / checkActiveBuild on PGLite are no-ops Lane E: upgrade-checkpoint round-trips with brain_id; X2 mismatch refused; F4 fall-through detected via reason='no_checkpoint'; full step progression to all_complete Test isolation hygiene (scripts/check-test-isolation.sh): - test/connection-manager.test.ts → connection-manager.serial.test.ts - test/backfill-concurrency-clamp.test.ts → .serial.test.ts - test/upgrade-checkpoint.test.ts → .serial.test.ts All three files mutate process.env (kill-switch, GBRAIN_DIRECT_POOL_SIZE, GBRAIN_HOME) which would race other tests in the parallel runner. *.serial.test.ts quarantine ensures they run at --max-concurrency=1. Choice between withEnv() refactor and serial quarantine made on the side of preserving existing well-formed test code. E2E coverage status: - v030_1-integration-pglite.test.ts (this commit): 14 cases, all green - backfill-perf-pglite.test.ts: 1 case, green (no regression) - cycle-recompute-emotional-weight-pglite.test.ts: green (no regression) - multi-source-emotional-weight-pglite.test.ts: green (no regression) - dream-synthesize-pglite.test.ts: 14 cases, green (no regression) - anomalies-pglite.test.ts + salience-pglite.test.ts: 6 cases, green Postgres-only E2Es (migration-flow, http-transport, hnsw-lifecycle, connection-routing) require DATABASE_URL + a real Postgres+pgvector container per the CLAUDE.md E2E lifecycle. They land as separate DATABASE_URL-gated work — not regressed by v0.30.1 changes; their preconditions just aren't met in the current run environment. `bun run verify` (typecheck + 4 shell pre-checks + test-isolation lint) passes cleanly. Final v0.30.1 unit + integration test count: 4547 pass, 0 regressions. Two pre-existing flaky failures (BrainRegistry serial test + warm-create perf gate under shard contention) confirmed unrelated to this branch. Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com> * chore: bump version and changelog (v0.30.1) Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com> --------- Co-authored-by: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
248 lines
9.1 KiB
TypeScript
248 lines
9.1 KiB
TypeScript
/**
|
|
* pgvector HNSW index policy + lifecycle manager (v0.30.1 Fix 5).
|
|
*
|
|
* Original v0.27 surface: chunkEmbeddingIndexSql / applyChunkEmbeddingIndexPolicy
|
|
* (kept unchanged for back-compat — schema-time index emission).
|
|
*
|
|
* v0.30.1 lifecycle additions:
|
|
* - dropAndRebuild (A3): atomic-swap pattern; build new index with temp
|
|
* name, ALTER...RENAME swap atomically, drop old. If rebuild fails the
|
|
* old index stays intact and search keeps working.
|
|
* - checkActiveBuild: pre-op probe of pg_stat_activity.
|
|
* - dropZombieIndexes: startup sweep of indisvalid=false indexes,
|
|
* guarded against in-progress builds.
|
|
* - monitorBuild: progress reporter during long-running CREATE INDEX.
|
|
*/
|
|
|
|
import type { BrainEngine } from './engine.ts';
|
|
|
|
export const PGVECTOR_HNSW_VECTOR_MAX_DIMS = 2000;
|
|
|
|
const CHUNK_EMBEDDING_HNSW_INDEX =
|
|
'CREATE INDEX IF NOT EXISTS idx_chunks_embedding ON content_chunks USING hnsw (embedding vector_cosine_ops);';
|
|
|
|
export function chunkEmbeddingIndexSql(dims: number): string {
|
|
if (dims <= PGVECTOR_HNSW_VECTOR_MAX_DIMS) return CHUNK_EMBEDDING_HNSW_INDEX;
|
|
return [
|
|
'-- idx_chunks_embedding skipped: pgvector HNSW vector indexes support',
|
|
`-- at most ${PGVECTOR_HNSW_VECTOR_MAX_DIMS} dimensions; exact vector scans remain available.`,
|
|
].join('\n');
|
|
}
|
|
|
|
export function applyChunkEmbeddingIndexPolicy(sql: string, dims: number): string {
|
|
return sql.replaceAll(CHUNK_EMBEDDING_HNSW_INDEX, chunkEmbeddingIndexSql(dims));
|
|
}
|
|
|
|
// ---------------------------------------------------------------------------
|
|
// v0.30.1 Lifecycle Manager (Fix 5)
|
|
// ---------------------------------------------------------------------------
|
|
|
|
export interface IndexSpec {
|
|
/** The CURRENT (production) index name. */
|
|
name: string;
|
|
table: string;
|
|
column: string;
|
|
/** USING clause body — e.g. `hnsw (embedding vector_cosine_ops)`. */
|
|
using: string;
|
|
/** Optional WHERE predicate (without WHERE keyword). */
|
|
condition?: string;
|
|
}
|
|
|
|
export interface ActiveBuildInfo {
|
|
active: boolean;
|
|
pid?: number;
|
|
query?: string;
|
|
application_name?: string;
|
|
}
|
|
|
|
/**
|
|
* Probe pg_stat_activity for an active CREATE INDEX on this index name.
|
|
* Used as a pre-op guard so dropAndRebuild doesn't compete with a build
|
|
* already in flight (Supabase auto-maintenance + parallel gbrain procs).
|
|
*/
|
|
export async function checkActiveBuild(
|
|
engine: BrainEngine,
|
|
indexName: string,
|
|
): Promise<ActiveBuildInfo> {
|
|
if (engine.kind !== 'postgres') return { active: false };
|
|
try {
|
|
const rows = await engine.executeRaw<{ pid: number; query: string; application_name: string | null }>(
|
|
`SELECT pid, query, application_name
|
|
FROM pg_stat_activity
|
|
WHERE state = 'active'
|
|
AND (query ILIKE $1 OR query ILIKE $2)
|
|
AND pid != pg_backend_pid()
|
|
LIMIT 1`,
|
|
[`%CREATE INDEX%${indexName}%`, `%REINDEX%${indexName}%`],
|
|
);
|
|
if (rows.length === 0) return { active: false };
|
|
const r = rows[0];
|
|
return {
|
|
active: true,
|
|
pid: r.pid,
|
|
query: r.query,
|
|
application_name: r.application_name ?? undefined,
|
|
};
|
|
} catch {
|
|
return { active: false };
|
|
}
|
|
}
|
|
|
|
/**
|
|
* Sweep invalid HNSW indexes on startup. Drops any pg_index row with
|
|
* indisvalid=false on tables we care about, AS LONG AS no active build
|
|
* is running for that index (codex Fix-5 zombie-cleanup guard).
|
|
*
|
|
* Postgres-only. PGLite returns { dropped: [] }.
|
|
*/
|
|
export async function dropZombieIndexes(
|
|
engine: BrainEngine,
|
|
tableNames: string[] = ['content_chunks', 'pages', 'takes'],
|
|
): Promise<{ dropped: string[] }> {
|
|
if (engine.kind !== 'postgres') return { dropped: [] };
|
|
const dropped: string[] = [];
|
|
try {
|
|
// Find invalid indexes on our tables.
|
|
const rows = await engine.executeRaw<{ indexname: string; tablename: string }>(
|
|
`SELECT i.relname AS indexname, t.relname AS tablename
|
|
FROM pg_index ix
|
|
JOIN pg_class i ON i.oid = ix.indexrelid
|
|
JOIN pg_class t ON t.oid = ix.indrelid
|
|
WHERE ix.indisvalid = false
|
|
AND t.relname = ANY($1)`,
|
|
[tableNames],
|
|
);
|
|
for (const r of rows) {
|
|
// Guard: skip if there's an active build for this index.
|
|
const active = await checkActiveBuild(engine, r.indexname);
|
|
if (active.active) {
|
|
process.stderr.write(`[hnsw] skipping zombie cleanup of ${r.indexname} — active build (pid ${active.pid})\n`);
|
|
continue;
|
|
}
|
|
try {
|
|
await engine.executeRaw(`DROP INDEX IF EXISTS ${r.indexname}`);
|
|
dropped.push(r.indexname);
|
|
process.stderr.write(`[hnsw] dropped zombie index ${r.indexname} on ${r.tablename}\n`);
|
|
} catch (err) {
|
|
process.stderr.write(`[hnsw] failed to drop ${r.indexname}: ${(err as Error).message}\n`);
|
|
}
|
|
}
|
|
} catch (err) {
|
|
// Best-effort: pg_stat_activity / pg_index queries may be restricted
|
|
// on managed Postgres tiers. Don't fail engine.connect() over it.
|
|
process.stderr.write(`[hnsw] zombie-index probe failed: ${(err as Error).message}\n`);
|
|
}
|
|
return { dropped };
|
|
}
|
|
|
|
/**
|
|
* Atomic-swap rebuild (A3): build new index with temp name, swap atomically.
|
|
*
|
|
* 1. Probe pg_stat_activity → bail if another build is active
|
|
* 2. Compose temp name: <name>_rebuild_<unix-ms>
|
|
* 3. CREATE INDEX <temp> with the spec's USING clause + condition
|
|
* 4. In a single transaction:
|
|
* DROP INDEX <name>
|
|
* ALTER INDEX <temp> RENAME TO <name>
|
|
* 5. If step 3 fails (OOM, timeout, conn drop), the old index is intact
|
|
* and search keeps serving queries. Caller can retry.
|
|
*
|
|
* The CREATE INDEX uses CONCURRENTLY so it doesn't block writes during the
|
|
* build; this requires `transaction:false` semantics so we route through
|
|
* engine.withReservedConnection.
|
|
*/
|
|
export async function dropAndRebuild(
|
|
engine: BrainEngine,
|
|
spec: IndexSpec,
|
|
opts: { reason: string; force?: boolean } = { reason: 'manual' },
|
|
): Promise<{ rebuilt: boolean; tempName: string }> {
|
|
if (engine.kind !== 'postgres') {
|
|
return { rebuilt: false, tempName: spec.name };
|
|
}
|
|
|
|
const active = await checkActiveBuild(engine, spec.name);
|
|
if (active.active && !opts.force) {
|
|
process.stderr.write(
|
|
`[hnsw] dropAndRebuild ${spec.name} aborted: active build pid ${active.pid} (${active.application_name ?? 'unknown'}). Pass --force to proceed anyway.\n`,
|
|
);
|
|
return { rebuilt: false, tempName: spec.name };
|
|
}
|
|
|
|
const ts = Date.now();
|
|
const tempName = `${spec.name}_rebuild_${ts}`;
|
|
const where = spec.condition ? ` WHERE ${spec.condition}` : '';
|
|
|
|
process.stderr.write(`[hnsw] rebuild ${spec.name} → ${tempName} (reason=${opts.reason})\n`);
|
|
|
|
// Step 3: build the new index (CONCURRENTLY) under a reserved connection.
|
|
await engine.withReservedConnection(async conn => {
|
|
await conn.executeRaw(
|
|
`CREATE INDEX CONCURRENTLY IF NOT EXISTS ${tempName} ON ${spec.table} USING ${spec.using}${where}`,
|
|
);
|
|
});
|
|
|
|
// Step 4: atomic swap inside a transaction.
|
|
await engine.transaction(async (tx) => {
|
|
const innerSql = (tx as unknown as { sql: any }).sql;
|
|
if (innerSql) {
|
|
await innerSql.unsafe(`DROP INDEX IF EXISTS ${spec.name}`);
|
|
await innerSql.unsafe(`ALTER INDEX ${tempName} RENAME TO ${spec.name}`);
|
|
}
|
|
});
|
|
|
|
process.stderr.write(`[hnsw] rebuild complete: ${spec.name}\n`);
|
|
return { rebuilt: true, tempName };
|
|
}
|
|
|
|
/**
|
|
* Poll pg_stat_activity to monitor a CREATE INDEX in progress. Reports
|
|
* elapsed time + progress (rows-built proxy via pg_stat_progress_create_index
|
|
* when available; falls back to relation size growth otherwise).
|
|
*
|
|
* Caller wraps a CREATE INDEX in a separate code path; this function is
|
|
* orthogonal — it just polls and emits progress lines.
|
|
*/
|
|
export interface BuildProgress {
|
|
elapsed_ms: number;
|
|
size_bytes?: number;
|
|
workers?: number;
|
|
pid?: number;
|
|
}
|
|
|
|
export async function monitorBuild(
|
|
engine: BrainEngine,
|
|
indexName: string,
|
|
onProgress: (status: BuildProgress) => void,
|
|
opts: { intervalMs?: number; maxIterations?: number } = {},
|
|
): Promise<void> {
|
|
if (engine.kind !== 'postgres') return;
|
|
const interval = opts.intervalMs ?? 30000;
|
|
const maxIterations = opts.maxIterations ?? 240; // 240 * 30s = 2h cap
|
|
const t0 = Date.now();
|
|
for (let i = 0; i < maxIterations; i++) {
|
|
const active = await checkActiveBuild(engine, indexName);
|
|
if (!active.active) return;
|
|
let size_bytes: number | undefined;
|
|
try {
|
|
const rows = await engine.executeRaw<{ size: number }>(
|
|
`SELECT pg_relation_size(c.oid) AS size FROM pg_class c WHERE c.relname = $1 LIMIT 1`,
|
|
[indexName],
|
|
);
|
|
if (rows[0]) size_bytes = Number(rows[0].size);
|
|
} catch { /* size probe optional */ }
|
|
onProgress({ elapsed_ms: Date.now() - t0, size_bytes, pid: active.pid });
|
|
await new Promise(r => setTimeout(r, interval));
|
|
}
|
|
}
|
|
|
|
/**
|
|
* Detect whether a CREATE INDEX query in pg_stat_activity is from Supabase
|
|
* auto-maintenance (vs. our gbrain process). Used by dropAndRebuild to
|
|
* back off when auto-maintenance is doing the rebuild for us.
|
|
*/
|
|
export function isSupabaseAutoMaintenance(active: ActiveBuildInfo): boolean {
|
|
if (!active.active) return false;
|
|
const appName = (active.application_name ?? '').toLowerCase();
|
|
return appName.includes('supabase') || appName.includes('postgres-meta');
|
|
}
|