Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions .env.example
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
36 changes: 36 additions & 0 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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
2 changes: 2 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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:<id>]` 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
Expand Down
41 changes: 41 additions & 0 deletions api/scripts/ptah/embedding-v2.json
Original file line number Diff line number Diff line change
@@ -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
}
}
8 changes: 8 additions & 0 deletions api/scripts/ptah/source.sql
Original file line number Diff line number Diff line change
@@ -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;
152 changes: 152 additions & 0 deletions api/scripts/ptah/verify.js
Original file line number Diff line number Diff line change
@@ -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 });
}
}
4 changes: 4 additions & 0 deletions api/scripts/reembed.js
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down
Loading