From 93afcf8da691100f98c9d9a027e14362c0393df2 Mon Sep 17 00:00:00 2001 From: Denis V <5462781+denisvmedia@users.noreply.github.com> Date: Fri, 2 Oct 2026 16:50:43 +0200 Subject: [PATCH] feat: prepare embedding replacements beside active vectors with Ptah --- .env.example | 3 + .github/workflows/ci.yml | 36 +++++ README.md | 2 + api/scripts/ptah/embedding-v2.json | 41 ++++++ api/scripts/ptah/source.sql | 8 ++ api/scripts/ptah/verify.js | 152 ++++++++++++++++++++ api/scripts/reembed.js | 4 + api/src/services/pgvector.js | 45 ++++-- api/tests/storage-pgvector-helpers.test.js | 18 +++ docs/embedding-generations.md | 156 +++++++++++++++++++++ 10 files changed, 450 insertions(+), 15 deletions(-) create mode 100644 api/scripts/ptah/embedding-v2.json create mode 100644 api/scripts/ptah/source.sql create mode 100644 api/scripts/ptah/verify.js create mode 100644 docs/embedding-generations.md diff --git a/.env.example b/.env.example index 095b032..4dcc6ee 100644 --- a/.env.example +++ b/.env.example @@ -50,6 +50,9 @@ GEMINI_EMBEDDING_DIMS=1536 # 1536 is the v4 default — sweet spot for qualit # EMBED_DOC_PREFIX= # Rarely needed; documents usually embed raw. # Swapping encoders on an existing corpus? Re-embed in place (dry-run first): # node api/scripts/reembed.js [--commit] +# To build beside the active vectors instead, see docs/embedding-generations.md. +# Set this together with the matching encoder only after verification: +# PGVECTOR_COLUMN=embedding_v2 # --- Structured Storage Backend --- # v4: Postgres is required. The postgres container uses the pgvector/pgvector:pg16 image diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 607b807..cccd09e 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -58,3 +58,39 @@ jobs: timeout 5 node src/index.js 2>&1 | head -1 || test $? -eq 124 env: BRAIN_API_KEY: test-key + + embedding-generation: + runs-on: ubuntu-latest + timeout-minutes: 10 + services: + postgres: + image: pgvector/pgvector:pg16 + env: + POSTGRES_DB: zengram_test + POSTGRES_PASSWORD: test + ports: + - 5432:5432 + options: >- + --health-cmd "pg_isready -U postgres -d zengram_test" + --health-interval 5s --health-timeout 5s --health-retries 10 + steps: + - uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7 + with: + persist-credentials: false + - uses: actions/setup-node@820762786026740c76f36085b0efc47a31fe5020 # v7 + with: + node-version: 22 + - run: npm ci + working-directory: api + - name: Install Ptah 0.11.4 + run: | + curl --fail --location --retry 3 https://github.com/stokaro/ptah/releases/download/v0.11.4/ptah_0.11.4_linux_amd64.tar.gz -o "$RUNNER_TEMP/ptah.tar.gz" + echo "de0f2d49570ed57545b50b6d4de83e5891847a79921d0ecef262467bdb38a53b $RUNNER_TEMP/ptah.tar.gz" | sha256sum -c - + mkdir "$RUNNER_TEMP/ptah-bin" + tar -xzf "$RUNNER_TEMP/ptah.tar.gz" -C "$RUNNER_TEMP/ptah-bin" ptah LICENSE + echo "$RUNNER_TEMP/ptah-bin" >> "$GITHUB_PATH" + - name: Verify embedding generation handover + working-directory: api + run: node scripts/ptah/verify.js + env: + POSTGRES_URL: postgres://postgres:test@localhost:5432/zengram_test?sslmode=disable diff --git a/README.md b/README.md index ddceca0..f98c01f 100644 --- a/README.md +++ b/README.md @@ -183,6 +183,8 @@ Copy [`adapters/claude-code/sessionend/`](adapters/claude-code/sessionend/) to y ## Roadmap +Changing embedding models on a live corpus? [Build a replacement beside the active vectors with Ptah](docs/embedding-generations.md), then deploy the matching encoder and column after verification. + **Recently shipped**: cross-encoder reranking stage, entity-graph retrieval path, weighted RRF, self-hosted encoder support (local endpoints + instruction prefixes + in-place re-embed), agentic iterate-until-sufficient retrieval (`brain_research`) with grounded `[mem:]` citations, pgvector migration (single-Postgres storage), multi-collection support, on-demand LLM reflection, temporal validity — [full changelog](CHANGELOG.md) **Coming next**: Automatic memory capture, hosted docs, LangChain/LlamaIndex integration diff --git a/api/scripts/ptah/embedding-v2.json b/api/scripts/ptah/embedding-v2.json new file mode 100644 index 0000000..ded0b14 --- /dev/null +++ b/api/scripts/ptah/embedding-v2.json @@ -0,0 +1,41 @@ +{ + "version": 1, + "name": "zengram-memories-v2", + "source": { + "schema": "public", + "table": "memories", + "key_fields": ["id"], + "input_fields": ["embedding_text"], + "version_strategy": "input_hash", + "mutable": true + }, + "preprocessing": { + "prefix": "", + "null_policy": "empty", + "empty_policy": "skip", + "unicode_normalization": "none", + "truncate": "refuse" + }, + "model": { + "provider": "openai-compatible", + "endpoint_class": "local", + "endpoint": "http://127.0.0.1:8000/v1", + "identifier": "BAAI/bge-small-en-v1.5", + "reported_dimension": 384, + "normalization": "none" + }, + "target": { + "schema": "public", + "table": "memories", + "column": "embedding_v2", + "representation": "vector", + "metric": "cosine", + "index_method": "hnsw" + }, + "consistency": { "mode": "outbox" }, + "policy": { + "require_exact_approval": true, + "require_consistency_mode": true, + "min_source_rows": 1 + } +} diff --git a/api/scripts/ptah/source.sql b/api/scripts/ptah/source.sql new file mode 100644 index 0000000..69e6501 --- /dev/null +++ b/api/scripts/ptah/source.sql @@ -0,0 +1,8 @@ +-- Expose the same text fallback as reembed.js to Ptah's column-based source. +-- A generated column follows payload edits without another application write. +-- Adding a STORED column rewrites the table: schedule this one-time setup. +SET lock_timeout = '5s'; +ALTER TABLE public.memories ADD COLUMN embedding_text text GENERATED ALWAYS AS ( + COALESCE(NULLIF(payload->>'text', ''), NULLIF(payload->>'content', ''), + NULLIF(payload->>'note', ''), payload->>'title', '') +) STORED; diff --git a/api/scripts/ptah/verify.js b/api/scripts/ptah/verify.js new file mode 100644 index 0000000..36f8fa5 --- /dev/null +++ b/api/scripts/ptah/verify.js @@ -0,0 +1,152 @@ +// Run only against a disposable PostgreSQL database with pgvector installed. +// The deterministic provider tests migration mechanics, not retrieval quality. +import assert from 'node:assert/strict'; +import { fork, spawn } from 'node:child_process'; +import { createServer } from 'node:http'; +import { mkdtemp, readFile, writeFile, rm } from 'node:fs/promises'; +import { tmpdir } from 'node:os'; +import { join } from 'node:path'; +import pg from 'pg'; + +if (process.env.PTAH_TEST_ACTOR) { + const { initEmbeddings, embed } = await import('../../src/services/embedders/interface.js'); + const { initPgvector, upsertPoint, supersedeAndInsert, searchPoints } = await import('../../src/services/pgvector.js'); + await initEmbeddings(); + await initPgvector(); + process.send({ ready: true }); + process.on('message', async ({ command, id, text, tenant = 'a', collection = 'shared_memories' }) => { + try { + const vector = await embed(text, command === 'search' ? 'search' : 'store'); + let result; + if (command === 'search') result = await searchPoints(vector, { client_id: tenant, active: true }, 10, [], [], collection); + else { + const payload = { text, type: 'fact', client_id: tenant, key: 'favorite', active: true }; + if (command === 'supersede') result = await supersedeAndInsert('key', 'favorite', id, vector, payload, {}, collection); + else await upsertPoint(id, vector, payload, collection); + } + process.send({ result: result ?? null }); + } catch (error) { process.send({ error: error.message }); } + }); +} else { + assert.ok(process.env.POSTGRES_URL, 'Set POSTGRES_URL to a disposable database'); + const db = new pg.Pool({ connectionString: process.env.POSTGRES_URL }); + assert.equal((await db.query("SELECT to_regclass('public.memories') AS table")).rows[0].table, null, + 'Refusing to test against a database that already contains memories'); + const work = await mkdtemp(join(tmpdir(), 'zengram-ptah-')); + const actors = []; + let failProvider = false; + let duringBackfill; + const provider = createServer(async (req, res) => { + let body = ''; + for await (const chunk of req) body += chunk; + const { input, model, encoding_format } = JSON.parse(body); + if (failProvider && model === 'candidate') { res.writeHead(503).end('test outage'); return; } + if (model === 'candidate' && duringBackfill) { + const check = duringBackfill; + duringBackfill = null; + await check(); + } + const dims = model === 'candidate' ? 384 : 1536; + const inputs = Array.isArray(input) ? input : [input]; + res.setHeader('Content-Type', 'application/json'); + res.end(JSON.stringify({ data: inputs.map((text, index) => { + const values = Array(dims).fill(0); + values.splice(0, 4, text.includes('apple') ? 1 : 0.1, text.includes('pear') ? 1 : 0.1, 0.2, 0.3); + return { index, embedding: encoding_format === 'base64' + ? Buffer.from(new Float32Array(values).buffer).toString('base64') : values }; + }) })); + }); + await new Promise(resolve => provider.listen(0, '127.0.0.1', resolve)); + const endpoint = `http://127.0.0.1:${provider.address().port}/v1`; + const actor = async (column, model) => { + const child = fork(new URL(import.meta.url), [], { env: { ...process.env, + PTAH_TEST_ACTOR: '1', PGVECTOR_COLUMN: column, EMBEDDING_PROVIDER: 'openai', + OPENAI_BASE_URL: endpoint, OPENAI_EMBEDDING_MODEL: model, OPENAI_EMBEDDING_DIMS: '', + EMBED_DOC_PREFIX: '', EMBED_QUERY_PREFIX: '', + }, stdio: ['ignore', 'inherit', 'inherit', 'ipc'] }); + actors.push(child); + await new Promise((resolve, reject) => { + child.once('message', resolve); + child.once('exit', code => reject(new Error(`API actor exited: ${code}`))); + }); + return message => new Promise((resolve, reject) => { + child.once('message', reply => reply.error ? reject(new Error(reply.error)) : resolve(reply.result)); + child.send(message); + }); + }; + const specPath = join(work, 'candidate.json'); + const ptah = async (verb, extra = [], ok = true) => { + const child = spawn(process.env.PTAH_BIN || 'ptah', ['inference', verb, '--spec', specPath, + '--db-url', process.env.POSTGRES_URL, '--run-id', 'zengram-proof', ...extra]); + let output = ''; + child.stdout.on('data', chunk => { output += chunk; }); + child.stderr.on('data', chunk => { output += chunk; }); + const code = await new Promise((resolve, reject) => { child.on('close', resolve); child.on('error', reject); }); + console.log(`ptah ${verb} (exit ${code}):\n${output}`); + if (ok) assert.equal(code, 0, output); else assert.notEqual(code, 0, output); + return output; + }; + try { + const old = await actor('vector', 'legacy'); + for (const [id, text, tenant] of [['keep', 'apple', 'a'], ['update', 'pear', 'a'], ['delete', 'apple', 'a'], ['other', 'apple', 'b']]) { + await old({ command: 'store', id, text, tenant }); + } + await old({ command: 'store', id: 'private', text: 'apple', collection: 'private' }); + const before = (await db.query('SELECT id, vector::text, payload FROM memories ORDER BY id')).rows; + const oldIndex = (await db.query("SELECT oid FROM pg_class WHERE relname = 'idx_memories_vector_hnsw'")).rows[0].oid; + await db.query(await readFile(new URL('./source.sql', import.meta.url), 'utf8')); + const spec = JSON.parse(await readFile(new URL('./embedding-v2.json', import.meta.url), 'utf8')); + Object.assign(spec.model, { endpoint, identifier: 'candidate', reported_dimension: 384 }); + await writeFile(specPath, JSON.stringify(spec)); + await ptah('prepare'); + failProvider = true; + await ptah('backfill', [], false); + assert.deepEqual((await db.query('SELECT id, vector::text, payload FROM memories ORDER BY id')).rows, before); + assert.ok((await old({ command: 'search', text: 'apple' })).some(r => r.id === 'keep')); + failProvider = false; + duringBackfill = async () => { + assert.ok((await old({ command: 'search', text: 'apple' })).some(r => r.id === 'keep')); + }; + await ptah('backfill'); + assert.equal(duringBackfill, null, 'old-model search must run while candidate embedding is in flight'); + assert.deepEqual((await db.query('SELECT id, vector::text, payload FROM memories ORDER BY id')).rows, before); + await old({ command: 'store', id: 'update', text: 'apple' }); + await old({ command: 'store', id: 'late', text: 'pear' }); + await db.query("DELETE FROM memories WHERE id = 'delete'"); + await ptah('catchup'); + for (const [id, expected] of [['update', [1, 0.1, 0.2, 0.3]], ['late', [0.1, 1, 0.2, 0.3]]]) { + const { rows } = await db.query('SELECT embedding_v2::text AS vector FROM memories WHERE id = $1', [id]); + assert.deepEqual(JSON.parse(rows[0].vector).slice(0, 4), expected); + } + await ptah('index'); + await ptah('verify'); + const refusal = await ptah('cutover', [], false); + const digest = refusal.match(/plan ([a-f0-9]{12,64})/)?.[1]; + assert.ok(digest, 'cutover must show the exact approval digest'); + await ptah('cutover', ['--approve', digest, '--approver', 'integration test']); + assert.equal((await db.query("SELECT oid FROM pg_class WHERE relname = 'idx_memories_vector_hnsw'")).rows[0].oid, oldIndex); + assert.equal((await db.query("SELECT vector::text FROM memories WHERE id = 'keep'")).rows[0].vector, + before.find(r => r.id === 'keep').vector); + const next = await actor('embedding_v2', 'candidate'); + const results = await next({ command: 'search', text: 'apple' }); + assert.ok(results.some(r => r.id === 'keep')); + assert.ok(results.some(r => r.id === 'update')); + assert.ok(results.some(r => r.id === 'late')); + for (const id of ['delete', 'other', 'private']) assert.ok(!results.some(r => r.id === id)); + await next({ command: 'store', id: 'after', text: 'apple' }); + await next({ command: 'store', id: 'after', text: 'pear' }); + assert.equal((await db.query("SELECT vector_dims(embedding_v2) AS dims, vector FROM memories WHERE id = 'after'")).rows[0].dims, 384); + assert.equal((await db.query("SELECT vector FROM memories WHERE id = 'after'")).rows[0].vector, null); + await next({ command: 'supersede', id: 'successor', text: 'pear' }); + await next({ command: 'supersede', id: 'successor', text: 'apple' }); + assert.equal((await db.query("SELECT vector_dims(embedding_v2) AS dims FROM memories WHERE id = 'successor'")).rows[0].dims, 384); + await ptah('catchup'); + await ptah('verify'); + console.log('PASS: failed/retried backfill, old search, live edits, catch-up, approval, dimension switch, both write paths, and tenant/collection filters'); + } finally { + for (const child of actors) child.kill(); + await new Promise(resolve => provider.close(resolve)); + await db.end(); + await rm(work, { recursive: true, force: true }); + } +} diff --git a/api/scripts/reembed.js b/api/scripts/reembed.js index a8a3d89..c643149 100644 --- a/api/scripts/reembed.js +++ b/api/scripts/reembed.js @@ -25,6 +25,10 @@ import pg from 'pg'; import { initEmbeddings, embed, getEmbeddingDimensions } from '../src/services/embedders/interface.js'; import { errorSummary } from '../src/lib/log.js'; +if (process.env.PGVECTOR_COLUMN && process.env.PGVECTOR_COLUMN !== 'vector') { + throw new Error('reembed.js only manages the legacy vector column. Use the Ptah generation workflow for PGVECTOR_COLUMN.'); +} + const COMMIT = process.argv.includes('--commit'); const argVal = (name, dflt) => { const i = process.argv.indexOf(name); diff --git a/api/src/services/pgvector.js b/api/src/services/pgvector.js index 31f54e3..94ccc16 100644 --- a/api/src/services/pgvector.js +++ b/api/src/services/pgvector.js @@ -8,6 +8,15 @@ import pg from 'pg'; import { getEmbeddingDimensions } from './embedders/interface.js'; const POSTGRES_URL = process.env.POSTGRES_URL; +// Select the prepared column together with its matching encoder at deployment. +// The default keeps the existing storage path; Ptah owns candidate indexes. +export function embeddingColumn(value = 'vector') { + if (!/^(vector|embedding_[a-z0-9_]+)$/.test(value) || value.length > 63) { + throw new Error('PGVECTOR_COLUMN must be vector or embedding_ (at most 63 characters)'); + } + return value; +} +const VECTOR_COLUMN = embeddingColumn(process.env.PGVECTOR_COLUMN); // Effective-confidence decay applied on read to fact- and status-type memories. const DECAY_FACTOR = parseFloat(process.env.DECAY_FACTOR) || 0.98; @@ -66,11 +75,12 @@ export function halfvecMode(dims) { // SQL distance expression for the active vector mode. In halfvec mode both the // stored column and the query param are cast to halfvec(dims) so the operator // matches the halfvec HNSW index. -export function vectorDistanceExpr(mode, dims, param = '$1') { +export function vectorDistanceExpr(mode, dims, param = '$1', column = 'vector') { + column = embeddingColumn(column); if (mode === 'halfvec') { - return `(vector::halfvec(${dims})) <=> ${param}::halfvec(${dims})`; + return `(${column}::halfvec(${dims})) <=> ${param}::halfvec(${dims})`; } - return `vector <=> ${param}::vector`; + return `${column} <=> ${param}::vector`; } // Startup dims-guard decision: does the existing vector column's declared @@ -108,7 +118,7 @@ export async function initPgvector() { const pgvectorVersion = parseVectorVersion(verRes.rows[0]?.extversion); iterativeScanSupported = supportsIterativeScan(pgvectorVersion); - await pool.query(` + if (VECTOR_COLUMN === 'vector') await pool.query(` CREATE TABLE IF NOT EXISTS memories ( id TEXT PRIMARY KEY, vector vector(${dims}), @@ -134,13 +144,17 @@ export async function initPgvector() { // declared dimension than the provider now reports, every write would fail // with an opaque dimension-mismatch. Fail fast at startup with a fix instead. const colRes = await pool.query( - `SELECT atttypmod FROM pg_attribute - WHERE attrelid = 'memories'::regclass AND attname = 'vector'` + `SELECT atttypmod, atttypid = 'vector'::regtype AS is_vector FROM pg_attribute + WHERE attrelid = 'memories'::regclass AND attname = $1 AND NOT attisdropped`, + [VECTOR_COLUMN] ); + if (!colRes.rows[0]?.is_vector) { + throw new Error(`Prepare memories.${VECTOR_COLUMN} as a vector column before starting the API`); + } const atttypmod = colRes.rows[0]?.atttypmod; if (dimsGuardShouldExit(atttypmod, dims)) { throw new Error( - `[pgvector] FATAL: existing 'memories.vector' column is vector(${atttypmod}) but the ` + + `[pgvector] FATAL: existing 'memories.${VECTOR_COLUMN}' column is vector(${atttypmod}) but the ` + `embedding provider reports ${dims} dims. Fix the mismatch: set the provider's dims env ` + `(e.g. GEMINI_EMBEDDING_DIMS/OPENAI_EMBEDDING_DIMS) back to ${atttypmod}, or re-embed the ` + `corpus into a fresh column at ${dims} dims. Refusing to start with a column that would ` + @@ -151,9 +165,10 @@ export async function initPgvector() { // Indexes — HNSW for vector ANN, btree for hot-path filters, GIN for JSONB entity filter. // HNSW creation is idempotent via IF NOT EXISTS but takes a moment on first create. // >2000 dims exceeds the `vector`-type HNSW cap, so index the halfvec cast instead. - if (vectorMode === 'halfvec') { + // Ptah builds and verifies candidate indexes before deployment. + if (VECTOR_COLUMN === 'vector' && vectorMode === 'halfvec') { await pool.query(`CREATE INDEX IF NOT EXISTS idx_memories_vector_hnsw ON memories USING hnsw ((vector::halfvec(${dims})) halfvec_cosine_ops)`); - } else { + } else if (VECTOR_COLUMN === 'vector') { await pool.query(`CREATE INDEX IF NOT EXISTS idx_memories_vector_hnsw ON memories USING hnsw (vector vector_cosine_ops)`); } await pool.query(`CREATE INDEX IF NOT EXISTS idx_memories_type ON memories(type) WHERE active = true`); @@ -170,7 +185,7 @@ export async function initPgvector() { await pool.query(`CREATE INDEX IF NOT EXISTS idx_memories_collection ON memories(collection)`); console.log( - `[pgvector] Table 'memories' ready (vector dims: ${dims}, mode: ${vectorMode}, ` + + `[pgvector] Table 'memories' ready (column: ${VECTOR_COLUMN}, vector dims: ${dims}, mode: ${vectorMode}, ` + `pgvector: ${verRes.rows[0]?.extversion || 'unknown'}, iterative_scan: ${iterativeScanSupported ? 'relaxed_order' : 'off'})` ); } @@ -209,14 +224,14 @@ export async function upsertPoint(id, vector, payload, collection) { const vecLit = toVectorLiteral(vector); await pool.query(` INSERT INTO memories ( - id, vector, type, source_agent, client_id, content_hash, + id, ${VECTOR_COLUMN}, type, source_agent, client_id, content_hash, key, subject, active, consolidated, importance, confidence, access_count, created_at, last_accessed_at, payload, collection ) VALUES ( $1, $2::vector, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, $16, $17 ) ON CONFLICT (id) DO UPDATE SET - vector = EXCLUDED.vector, + ${VECTOR_COLUMN} = EXCLUDED.${VECTOR_COLUMN}, type = EXCLUDED.type, source_agent = EXCLUDED.source_agent, client_id = EXCLUDED.client_id, @@ -288,14 +303,14 @@ export async function supersedeAndInsert(keyField, keyValue, newId, vector, payl const storedPayload = { ...payload, supersedes: supersededId }; await client.query( `INSERT INTO memories ( - id, vector, type, source_agent, client_id, content_hash, + id, ${VECTOR_COLUMN}, type, source_agent, client_id, content_hash, key, subject, active, consolidated, importance, confidence, access_count, created_at, last_accessed_at, payload, collection ) VALUES ( $1, $2::vector, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, $16, $17 ) ON CONFLICT (id) DO UPDATE SET - vector = EXCLUDED.vector, + ${VECTOR_COLUMN} = EXCLUDED.${VECTOR_COLUMN}, payload = EXCLUDED.payload, active = EXCLUDED.active`, [ @@ -380,7 +395,7 @@ export async function searchPoints(vector, filter = {}, limit = 10, nestedFilter // pgvector '<=>' is cosine distance (0 = identical, 2 = opposite). We map it to // a [0,1] similarity: score = 1 - distance/2, which is exactly 0.5 + cosine_sim/2. // SEARCH_SCORE_FLOOR (default 0.55 ≈ cosine 0.1) drops near-orthogonal matches. - const distExpr = vectorDistanceExpr(vectorMode, vectorDims); + const distExpr = vectorDistanceExpr(vectorMode, vectorDims, '$1', VECTOR_COLUMN); const scoreExpr = `1 - (${distExpr}) / 2`; const sql = ` SELECT id, payload, ${scoreExpr} AS score diff --git a/api/tests/storage-pgvector-helpers.test.js b/api/tests/storage-pgvector-helpers.test.js index f845496..9dfe31e 100644 --- a/api/tests/storage-pgvector-helpers.test.js +++ b/api/tests/storage-pgvector-helpers.test.js @@ -1,6 +1,7 @@ import { describe, it } from 'node:test'; import assert from 'node:assert/strict'; import { + embeddingColumn, parseVectorVersion, supportsIterativeScan, clampEfSearch, @@ -112,3 +113,20 @@ describe('dimsGuardShouldExit', () => { assert.equal(dimsGuardShouldExit(undefined, 1536), false); }); }); + +// A column cannot be a bind parameter: reject SQL and unrelated payload columns. +describe('embeddingColumn', () => { + it('keeps the default and accepts candidate names', () => { + assert.equal(embeddingColumn(), 'vector'); + assert.equal(embeddingColumn('embedding_v2'), 'embedding_v2'); + assert.equal(vectorDistanceExpr('vector', 384, '$1', 'embedding_v2'), + 'embedding_v2 <=> $1::vector'); + assert.equal(vectorDistanceExpr('halfvec', 3072, '$1', 'embedding_v2'), + '(embedding_v2::halfvec(3072)) <=> $1::halfvec(3072)'); + }); + it('rejects empty, unrelated, oversized, and SQL names', () => { + for (const column of ['', 'payload', 'embedding_V2', 'embedding_x; DROP TABLE memories', 'embedding_' + 'x'.repeat(54)]) { + assert.throws(() => embeddingColumn(column), /PGVECTOR_COLUMN/); + } + }); +}); diff --git a/docs/embedding-generations.md b/docs/embedding-generations.md new file mode 100644 index 0000000..193ecd2 --- /dev/null +++ b/docs/embedding-generations.md @@ -0,0 +1,156 @@ +# Change encoders without overwriting active vectors + +`reembed.js` replaces `memories.vector` in place. On a dimension change it drops +HNSW and clears that column before filling it again. For larger corpora, use +[Ptah inference migrations](https://docs.ptah.run/v0.11.4/inference/overview/) to build +and verify a separate column while the existing API keeps serving the old one. + +```mermaid +flowchart LR + API[API + current encoder] --> V[memories.vector + existing HNSW] + P[Ptah + replacement encoder] --> C[memories.embedding_v2 + new HNSW] + C --> Check[Catch up and verify] + Check --> Deploy[Deploy matching encoder + PGVECTOR_COLUMN] + Deploy --> C +``` + +This is an opt-in operator workflow, tested with Ptah **0.11.4** and PostgreSQL +16. It does not add Ptah to the API image or change the default storage path. +The example uses an OpenAI-compatible replacement encoder with at most 2,000 +dimensions and a `vector` HNSW index. The existing encoder may use any Zengram +provider. Native Gemini is not an OpenAI-compatible endpoint. + +## Prepare once + +Install `ptah` from its [release](https://github.com/stokaro/ptah/releases/tag/v0.11.4). +Keep the current API configuration and encoder running. Run these commands from +the repository root, with an explicit database URL for the intended database: + +```sh +export POSTGRES_URL='postgres://user:password@localhost:5433/shared_brain?sslmode=disable' +psql "$POSTGRES_URL" -v ON_ERROR_STOP=1 -f api/scripts/ptah/source.sql +``` + +`source.sql` exposes the text fallback used by `reembed.js` as a generated +`embedding_text` column. Ptah reads columns, while Zengram stores text inside +JSON payloads. The generated column also follows subsequent payload edits. +Adding a stored generated column rewrites the table and takes a table lock: +schedule this one-time setup separately, with enough disk space. Its five-second +lock timeout bounds waiting for a lock, not the rewrite duration. Do not rerun +it once the column exists. + +## Build the candidate + +Copy the example, then configure the new model's endpoint, identifier, reported +dimension, and document prefix: + +```sh +cp api/scripts/ptah/embedding-v2.json embedding-v2.json +``` + +Ptah accepts JSON specifications as well as YAML. The prefix must equal the decoded +`EMBED_DOC_PREFIX` you will deploy. Do not trim, normalize, or truncate text +differently between backfill and API writes. + +For a hosted endpoint, set `endpoint_class` to `hosted` and use a credential +reference such as `"credential": "env:OPENAI_API_KEY"`; do not put a key in the +file. If you request shortened dimensions in Zengram, set the same +`requested_dimension` in the specification. Leave it absent when the server +expects native dimensions. Keep credentials and private corpus text out of Git. + +```sh +export PTAH_DB_URL="$POSTGRES_URL" +export PTAH_SPEC="$PWD/embedding-v2.json" +export PTAH_RUN_ID=memories-v2 +ptah inference plan +ptah inference prepare +ptah inference backfill +ptah inference catchup +ptah inference index +ptah inference verify +``` + +Use a new column and run ID for every model change, including changes that keep +the same dimension. `prepare` adds the candidate column, metadata, and outbox +triggers. `catchup` brings inserts, text changes, and deletes into the candidate. +The original column and index remain in use. A failed backfill can be retried +with the same specification and run ID; it does not clear the active vectors. + +Verification checks coverage, freshness, dimensions, and index readiness. It +does not establish that the new encoder retrieves better answers. Evaluate +representative queries with that model's query prefix before switching; see +[Ptah evaluation](https://docs.ptah.run/v0.11.4/inference/reference/evaluation-corpus/). + +## Switch the encoder and column together + +Briefly pause **all writers**, including consolidation jobs and imports. Drain +in-flight writes, catch up again, and verify before approving the switch: + +```sh +ptah inference catchup +ptah inference verify +ptah inference cutover +# The unapproved command refuses and prints a plan digest. Review it, then: +ptah inference cutover --approve --approver 'operator name' +``` + +Ptah records the active generation. Zengram does **not** read that pointer: +its encoder and column are fixed for the lifetime of each API process. Deploy +them as one configuration change, for example: + +```dotenv +PGVECTOR_COLUMN=embedding_v2 +EMBEDDING_PROVIDER=openai +OPENAI_BASE_URL=http://encoder:8000/v1 +OPENAI_EMBEDDING_MODEL=BAAI/bge-small-en-v1.5 +# Unset OPENAI_EMBEDDING_DIMS when the server uses native dimensions. +# Set EMBED_DOC_PREFIX and EMBED_QUERY_PREFIX to this model's expected values. +``` + +Point the API at the same model used by Ptah, even when its network address +differs. Startup checks the selected column's dimension and refuses a missing +or incompatible column. Dimension equality alone does not prove model equality. +Check search with the new API, retire old API processes, then resume writers. +Old processes can serve reads during preparation and this deployment, but must +not resume writing with the old configuration. This is not an automatic, +zero-downtime deployment controller. + +Both normal writes and fact/status supersession now write the selected column. +The API leaves candidate index creation to Ptah. Keyword search, tenant and +collection filters, payloads, and entity extraction keep their existing paths. +Do not use `reembed.js` with a selected candidate column; it refuses that setup. + +## Keep or discard the previous column + +Before writers resume, the original vectors still provide a way to cancel the +application switch. Afterward they become stale: keeping a column is **not** a +complete rollback strategy. Do not switch back without rebuilding/catching up +that model's vectors and verifying them while writers are paused. This example +does not register the legacy column as a Ptah generation or provide automatic +rollback. Keep a database backup and both encoder configurations. + +Outbox triggers remain installed after cutover. Schedule `ptah inference +catchup` with this specification and run ID to keep metadata current and drain +captured changes. Keep the active generation's specification; do not retire it +while the API still reads its column. Use Ptah's +[retirement guide](https://docs.ptah.run/v0.11.4/inference/guides/rollback-and-retire/) +when a generation is no longer needed. + +## Reproduce the migration test + +With Node.js, Ptah 0.11.4, and an **empty, disposable** pgvector database: + +```sh +cd api +npm ci +POSTGRES_URL='postgres://postgres:test@localhost:5432/zengram_test?sslmode=disable' \ + node scripts/ptah/verify.js +``` + +The test runs real Zengram storage code and Ptah against PostgreSQL. A local, +deterministic embeddings endpoint simulates different dimensions and an outage. +It checks retry, old-column search during the build, catch-up of edits and +deletes, explicit cutover approval, writes after switching, and tenant/collection +filters. It does not call a paid provider or measure model quality. The test +refuses a database that already contains `memories`; discard the test database +afterward.