mirror of
https://github.com/garrytan/gbrain.git
synced 2026-07-30 03:12:32 +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>
441 lines
16 KiB
TypeScript
441 lines
16 KiB
TypeScript
/**
|
|
* Connection Manager — route Postgres queries by query type (v0.30.1, Fix 1).
|
|
*
|
|
* Three pools, one decision: read() goes to the pooler (port 6543, fast,
|
|
* many connections); ddl() and bulk() go to a direct connection (port 5432,
|
|
* 30min statement_timeout, capped at 3 conns) so DDL doesn't time out on
|
|
* the Supabase pooler's 2-min statement_timeout.
|
|
*
|
|
* The connection-manager is the URL-routing layer. It layers on top of
|
|
* postgres.js's existing pool primitives + PostgresEngine.withReservedConnection.
|
|
*
|
|
* ┌─────────────────────────────┐
|
|
* │ GBRAIN_DATABASE_URL │ GBRAIN_DIRECT_DATABASE_URL (override)
|
|
* │ (pooler, port 6543) │
|
|
* └────────┬────────────────────┘
|
|
* │
|
|
* ▼ auto-detect Supabase
|
|
* ┌──────────────┐ ┌──────────────┐
|
|
* │ read pool │ │ direct pool │
|
|
* │ size 10 │ │ size 3 │
|
|
* │ prepare:no │ │ stmt 30min │
|
|
* │ stmt 5min │ │ idle 5min │
|
|
* └──────────────┘ │ mwm 256MB │
|
|
* └──────────────┘
|
|
*
|
|
* Architectural notes:
|
|
* - INSTANCE-owned (T5 / X1 amendment): each PostgresEngine constructs its
|
|
* own ConnectionManager. Worker engines (cycle, sync) inherit the parent's
|
|
* via constructor option `parent`. transaction() clones share the parent's.
|
|
* - Lazy direct pool init via cached Promise<Sql> (A1): concurrent first
|
|
* callers await the same Promise, so no double-init.
|
|
* - Kill-switch (F1): GBRAIN_DISABLE_DIRECT_POOL=1 falls back to single-pool
|
|
* legacy path. With parent set, inherit parent's kill-switch state (A2).
|
|
* - Audit (F8): every acquire/release/error logs to connection-events.jsonl.
|
|
* - Non-Supabase passthrough: if URL isn't a Supabase pooler and no
|
|
* GBRAIN_DIRECT_DATABASE_URL override, ddl()/bulk() share the read pool.
|
|
*/
|
|
|
|
import postgres from 'postgres';
|
|
import { resolvePrepare, resolveSessionTimeouts, resolvePoolSize } from './db.ts';
|
|
import { redactPgUrl } from './url-redact.ts';
|
|
import { logConnectionEvent } from './connection-audit.ts';
|
|
|
|
export type Sql = ReturnType<typeof postgres>;
|
|
|
|
export interface ConnectionManagerOpts {
|
|
/** Primary URL — usually the pooler (port 6543) on Supabase. */
|
|
url: string;
|
|
/**
|
|
* Override for the direct URL. When set, takes precedence over auto-derivation.
|
|
* Sourced from GBRAIN_DIRECT_DATABASE_URL or explicit caller config.
|
|
*/
|
|
directUrl?: string | null;
|
|
/**
|
|
* Inherit pools + kill-switch state from a parent manager (worker engines,
|
|
* transaction clones). When set, this manager is a thin reference holder
|
|
* and does NOT open its own pools.
|
|
*/
|
|
parent?: ConnectionManager;
|
|
/**
|
|
* Read pool size override (defaults to resolvePoolSize() — 10 normally).
|
|
*/
|
|
readPoolSize?: number;
|
|
/**
|
|
* Direct pool size override (defaults to GBRAIN_DIRECT_POOL_SIZE env or 3).
|
|
*/
|
|
directPoolSize?: number;
|
|
/**
|
|
* When true, the read pool is owned by some other code (e.g. db.ts:connect's
|
|
* module singleton). The connection manager will USE it via getReadPool but
|
|
* not call .end() on disconnect(). Default false (we own both pools).
|
|
*/
|
|
readPoolOwnedExternally?: boolean;
|
|
}
|
|
|
|
/** Default direct-pool size (P1 raised from 2 to 3). Override via env. */
|
|
export const DEFAULT_DIRECT_POOL_SIZE = 3;
|
|
|
|
/** Search statement timeout (F5 consolidation) — was 8s scattered. */
|
|
export const SEARCH_STMT_TIMEOUT_MS = 8000;
|
|
|
|
/** DDL pool default statement_timeout (Fix 1). */
|
|
const DDL_STMT_TIMEOUT_MS = 30 * 60 * 1000; // 30min
|
|
|
|
/** DDL pool default idle-in-transaction timeout (Fix 1). */
|
|
const DDL_IDLE_TX_TIMEOUT_MS = 5 * 60 * 1000; // 5min
|
|
|
|
/** Bulk pool default maintenance_work_mem (P1). 256MB safe on Supabase. */
|
|
const BULK_MAINTENANCE_WORK_MEM = '256MB';
|
|
|
|
/**
|
|
* Hostname patterns that indicate a Supabase pooler. Used for auto-detection
|
|
* of the dual-pool topology. Adding more patterns is safe; mis-detection
|
|
* just means we open a "direct" pool against the same URL — wasteful but
|
|
* not broken (the kill-switch is the operator's escape hatch).
|
|
*/
|
|
const SUPABASE_POOLER_HOSTNAME_PATTERNS = [
|
|
/\.pooler\.supabase\.com$/i,
|
|
/^pooler\.supabase\.com$/i,
|
|
];
|
|
|
|
const SUPABASE_POOLER_PORTS = new Set(['6543']);
|
|
|
|
/**
|
|
* True if the URL looks like a Supabase pooler endpoint. Used for kill-switch
|
|
* activation and dual-pool routing.
|
|
*/
|
|
export function isSupabasePoolerUrl(url: string): boolean {
|
|
try {
|
|
const parsed = new URL(url.replace(/^postgres(ql)?:\/\//, 'http://'));
|
|
if (SUPABASE_POOLER_PORTS.has(parsed.port)) return true;
|
|
if (SUPABASE_POOLER_HOSTNAME_PATTERNS.some(re => re.test(parsed.hostname))) return true;
|
|
return false;
|
|
} catch {
|
|
return false;
|
|
}
|
|
}
|
|
|
|
/**
|
|
* Derive a direct (non-pooler) URL from a Supabase pooler URL. Two known shapes:
|
|
*
|
|
* Pooler hostname: aws-N-region.pooler.supabase.com on port 6543
|
|
* → swap to db.<project-ref>.supabase.co on port 5432
|
|
* (project-ref encoded in the user component as postgres.<ref>)
|
|
* Direct hostname: db.<ref>.supabase.co already on port 5432 → returned as-is
|
|
*
|
|
* For the modern shape, we try to extract project-ref from the user component.
|
|
* If we cannot, we fall back to swapping port-only and the caller may warn.
|
|
*
|
|
* Returns null when the URL isn't a recognized Supabase pooler.
|
|
*/
|
|
export function deriveDirectUrl(url: string): string | null {
|
|
try {
|
|
const parsed = new URL(url.replace(/^postgres(ql)?:\/\//, 'http://'));
|
|
const port = parsed.port;
|
|
const hostname = parsed.hostname;
|
|
const isPoolerHost = SUPABASE_POOLER_HOSTNAME_PATTERNS.some(re => re.test(hostname));
|
|
if (port !== '6543' && !isPoolerHost) return null;
|
|
// User part on Supabase pooler is typically `postgres.<project-ref>`.
|
|
// Extract <project-ref> for the direct hostname.
|
|
const user = parsed.username || '';
|
|
const decodedUser = decodeURIComponent(user);
|
|
const refMatch = decodedUser.match(/^postgres\.([a-z0-9]+)$/i);
|
|
let directHost = hostname;
|
|
if (refMatch && refMatch[1] && isPoolerHost) {
|
|
directHost = `db.${refMatch[1]}.supabase.co`;
|
|
}
|
|
// Compose direct URL by swapping host + port. Preserve auth, db, query.
|
|
parsed.hostname = directHost;
|
|
parsed.port = '5432';
|
|
// Reconstruct with the original scheme.
|
|
const scheme = url.match(/^postgres(?:ql)?:\/\//i)?.[0] ?? 'postgres://';
|
|
const auth = parsed.username
|
|
? `${parsed.username}${parsed.password ? `:${parsed.password}` : ''}@`
|
|
: '';
|
|
const search = parsed.search ?? '';
|
|
const path = parsed.pathname ?? '';
|
|
return `${scheme}${auth}${directHost}:5432${path}${search}`;
|
|
} catch {
|
|
return null;
|
|
}
|
|
}
|
|
|
|
/**
|
|
* Read kill-switch state from env. Subordinate to parent manager's state
|
|
* when present (A2 inheritance).
|
|
*/
|
|
export function readKillSwitchEnv(): boolean {
|
|
return process.env.GBRAIN_DISABLE_DIRECT_POOL === '1' ||
|
|
process.env.GBRAIN_DISABLE_DIRECT_POOL === 'true';
|
|
}
|
|
|
|
/**
|
|
* Resolve direct pool size: explicit > env > default.
|
|
*/
|
|
export function resolveDirectPoolSize(explicit?: number): number {
|
|
if (typeof explicit === 'number' && explicit > 0) return explicit;
|
|
const raw = process.env.GBRAIN_DIRECT_POOL_SIZE;
|
|
if (raw) {
|
|
const parsed = parseInt(raw, 10);
|
|
if (Number.isFinite(parsed) && parsed > 0 && parsed <= 20) return parsed;
|
|
}
|
|
return DEFAULT_DIRECT_POOL_SIZE;
|
|
}
|
|
|
|
export class ConnectionManager {
|
|
private readonly opts: ConnectionManagerOpts;
|
|
private _readPool: Sql | null = null;
|
|
private _readPoolOwnedExternally: boolean;
|
|
private _directInit: Promise<Sql | null> | null = null;
|
|
private _directPool: Sql | null = null;
|
|
private _killSwitch: boolean;
|
|
private _directUrl: string | null;
|
|
private _isSupabase: boolean;
|
|
|
|
constructor(opts: ConnectionManagerOpts) {
|
|
this.opts = opts;
|
|
this._readPoolOwnedExternally = opts.readPoolOwnedExternally === true;
|
|
|
|
// A2: kill-switch resolution. Parent overrides env when present.
|
|
if (opts.parent) {
|
|
this._killSwitch = opts.parent.isKillSwitchActive();
|
|
this._isSupabase = opts.parent.isSupabase();
|
|
this._directUrl = opts.parent.resolveDirectUrl();
|
|
this._readPool = opts.parent.peekReadPool();
|
|
this._readPoolOwnedExternally = true; // never end the parent's pool
|
|
} else {
|
|
this._killSwitch = readKillSwitchEnv();
|
|
this._isSupabase = isSupabasePoolerUrl(opts.url);
|
|
// Direct URL: explicit override > env > derive > null
|
|
const envOverride = process.env.GBRAIN_DIRECT_DATABASE_URL;
|
|
this._directUrl = opts.directUrl ?? envOverride ?? deriveDirectUrl(opts.url);
|
|
}
|
|
}
|
|
|
|
/** Whether dual-pool routing is active (false on non-Supabase or kill-switch). */
|
|
isDualPoolActive(): boolean {
|
|
return this._isSupabase && !this._killSwitch && !!this._directUrl;
|
|
}
|
|
|
|
isSupabase(): boolean { return this._isSupabase; }
|
|
isKillSwitchActive(): boolean { return this._killSwitch; }
|
|
resolveDirectUrl(): string | null { return this._directUrl; }
|
|
|
|
/**
|
|
* Internal: peek at the read pool without forcing init. Used by parent
|
|
* inheritance to share the same instance.
|
|
*/
|
|
peekReadPool(): Sql | null { return this._readPool; }
|
|
|
|
/**
|
|
* Set the read pool. Used by db.ts:connect or PostgresEngine.connect when
|
|
* they own the pool externally (and connection-manager is just routing).
|
|
*/
|
|
setReadPool(sql: Sql): void {
|
|
this._readPool = sql;
|
|
this._readPoolOwnedExternally = true;
|
|
}
|
|
|
|
/**
|
|
* Get or lazily create the read pool. Honors `readPoolOwnedExternally` so
|
|
* we don't double-create when db.ts:connect already owns the singleton.
|
|
*/
|
|
async getReadPool(): Promise<Sql> {
|
|
if (this._readPool) return this._readPool;
|
|
if (this._readPoolOwnedExternally) {
|
|
throw new Error('connection-manager: read pool marked as externally-owned but not provided');
|
|
}
|
|
const opts: Record<string, unknown> = {
|
|
max: resolvePoolSize(this.opts.readPoolSize),
|
|
idle_timeout: 20,
|
|
connect_timeout: 10,
|
|
types: { bigint: postgres.BigInt },
|
|
};
|
|
const timeouts = resolveSessionTimeouts();
|
|
if (Object.keys(timeouts).length > 0) opts.connection = timeouts;
|
|
const prepare = resolvePrepare(this.opts.url);
|
|
if (typeof prepare === 'boolean') opts.prepare = prepare;
|
|
this._readPool = postgres(this.opts.url, opts);
|
|
logConnectionEvent({ pool: 'read', op: 'init' });
|
|
return this._readPool;
|
|
}
|
|
|
|
/**
|
|
* Acquire the read connection. Synchronous accessor — assumes read pool
|
|
* is already initialized (matches existing engine.sql semantics).
|
|
* Throws if pool not ready.
|
|
*/
|
|
read(): Sql {
|
|
if (!this._readPool) {
|
|
throw new Error('connection-manager: read pool not initialized; call getReadPool() first or set externally');
|
|
}
|
|
return this._readPool;
|
|
}
|
|
|
|
/**
|
|
* Acquire (and lazy-init) the direct DDL pool. When kill-switch is active
|
|
* or non-Supabase, returns the read pool (single-pool fallback).
|
|
*
|
|
* A1: lazy init wraps in a cached Promise<Sql> so concurrent first-callers
|
|
* await the same init instead of racing two pool constructions.
|
|
*/
|
|
async ddl(): Promise<Sql> {
|
|
if (!this.isDualPoolActive()) {
|
|
return this.getReadPool();
|
|
}
|
|
return this.getDirectPool();
|
|
}
|
|
|
|
/**
|
|
* Acquire the direct pool for a long-running BULK operation. Caller can
|
|
* override the per-op timeout via SET LOCAL inside a sql.begin block.
|
|
* Same pool as ddl(); the distinction is callsite intent (used by audit
|
|
* + caller-side timeout SET LOCAL).
|
|
*/
|
|
async bulk(_timeoutSeconds?: number): Promise<Sql> {
|
|
if (!this.isDualPoolActive()) {
|
|
return this.getReadPool();
|
|
}
|
|
return this.getDirectPool();
|
|
}
|
|
|
|
private async getDirectPool(): Promise<Sql> {
|
|
if (this._directPool) return this._directPool;
|
|
// A1: cache the Promise so concurrent first callers await the same init.
|
|
if (!this._directInit) {
|
|
this._directInit = this.initDirectPool().then(pool => {
|
|
this._directPool = pool;
|
|
return pool;
|
|
}).catch(err => {
|
|
// Reset cache on failure so next caller can retry.
|
|
this._directInit = null;
|
|
throw err;
|
|
});
|
|
}
|
|
const pool = await this._directInit;
|
|
if (!pool) {
|
|
// Defensive — initDirectPool should have thrown.
|
|
throw new Error('connection-manager: direct pool init returned null');
|
|
}
|
|
return pool;
|
|
}
|
|
|
|
private async initDirectPool(): Promise<Sql> {
|
|
if (!this._directUrl) {
|
|
throw new Error('connection-manager: cannot init direct pool — no direct URL');
|
|
}
|
|
const size = resolveDirectPoolSize(this.opts.directPoolSize);
|
|
const opts: Record<string, unknown> = {
|
|
max: size,
|
|
idle_timeout: 20,
|
|
connect_timeout: 10,
|
|
types: { bigint: postgres.BigInt },
|
|
// Always use prepared statements on the direct pool — no PgBouncer
|
|
// here, so the prepare-cache invalidation issue doesn't apply.
|
|
prepare: true,
|
|
// Apply DDL session GUCs as connection startup parameters (durable
|
|
// through any intermediary pooling layer, same trick as
|
|
// resolveSessionTimeouts).
|
|
connection: {
|
|
statement_timeout: String(DDL_STMT_TIMEOUT_MS),
|
|
idle_in_transaction_session_timeout: String(DDL_IDLE_TX_TIMEOUT_MS),
|
|
maintenance_work_mem: BULK_MAINTENANCE_WORK_MEM,
|
|
},
|
|
};
|
|
const t0 = Date.now();
|
|
try {
|
|
const pool = postgres(this._directUrl, opts);
|
|
// Probe to validate connectivity early.
|
|
await pool`SELECT 1`;
|
|
logConnectionEvent({
|
|
pool: 'ddl',
|
|
op: 'init',
|
|
duration_ms: Date.now() - t0,
|
|
host: this._directUrl ? this.hostOnly(this._directUrl) : undefined,
|
|
});
|
|
return pool;
|
|
} catch (err) {
|
|
logConnectionEvent({
|
|
pool: 'ddl',
|
|
op: 'error',
|
|
duration_ms: Date.now() - t0,
|
|
error: { message: err instanceof Error ? err.message : String(err) },
|
|
});
|
|
throw err;
|
|
}
|
|
}
|
|
|
|
/**
|
|
* SELECT 1 latency probe on each pool. Used by doctor's connection_routing
|
|
* check + healthCheck() in the gateway.
|
|
*/
|
|
async healthCheck(): Promise<{ read: number | null; direct: number | null }> {
|
|
const result: { read: number | null; direct: number | null } = { read: null, direct: null };
|
|
try {
|
|
const t0 = Date.now();
|
|
const pool = await this.getReadPool();
|
|
await pool`SELECT 1`;
|
|
result.read = Date.now() - t0;
|
|
} catch { /* leave null */ }
|
|
if (this.isDualPoolActive()) {
|
|
try {
|
|
const t0 = Date.now();
|
|
const pool = await this.getDirectPool();
|
|
await pool`SELECT 1`;
|
|
result.direct = Date.now() - t0;
|
|
} catch { /* leave null */ }
|
|
}
|
|
return result;
|
|
}
|
|
|
|
/**
|
|
* Disconnect pools we own. Read pool stays alive if marked externally owned
|
|
* (db.ts singleton path). Direct pool is always ours.
|
|
*/
|
|
async disconnect(): Promise<void> {
|
|
if (this._directPool) {
|
|
try { await this._directPool.end(); } catch { /* idempotent */ }
|
|
this._directPool = null;
|
|
this._directInit = null;
|
|
}
|
|
if (this._readPool && !this._readPoolOwnedExternally) {
|
|
try { await this._readPool.end(); } catch { /* idempotent */ }
|
|
this._readPool = null;
|
|
}
|
|
}
|
|
|
|
/**
|
|
* Diagnostic snapshot for doctor / get_health surfaces.
|
|
*/
|
|
describeMode(): {
|
|
mode: 'split' | 'single (kill-switch)' | 'single (non-supabase)' | 'single (no-direct-url)';
|
|
direct_host?: string;
|
|
kill_switch_active: boolean;
|
|
direct_pool_size: number;
|
|
} {
|
|
let mode: 'split' | 'single (kill-switch)' | 'single (non-supabase)' | 'single (no-direct-url)';
|
|
if (!this._isSupabase) mode = 'single (non-supabase)';
|
|
else if (this._killSwitch) mode = 'single (kill-switch)';
|
|
else if (!this._directUrl) mode = 'single (no-direct-url)';
|
|
else mode = 'split';
|
|
return {
|
|
mode,
|
|
direct_host: this._directUrl ? this.hostOnly(this._directUrl) : undefined,
|
|
kill_switch_active: this._killSwitch,
|
|
direct_pool_size: resolveDirectPoolSize(this.opts.directPoolSize),
|
|
};
|
|
}
|
|
|
|
private hostOnly(url: string): string {
|
|
// Redact creds first, then strip everything except host:port for doctor display.
|
|
const redacted = redactPgUrl(url);
|
|
try {
|
|
const parsed = new URL(redacted.replace(/^postgres(ql)?:\/\//, 'http://'));
|
|
return `${parsed.hostname}${parsed.port ? `:${parsed.port}` : ''}`;
|
|
} catch {
|
|
return redacted;
|
|
}
|
|
}
|
|
}
|