Skip to content
Merged
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
56 changes: 56 additions & 0 deletions db/migrations/018-store-embeddings-as-blobs.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1,56 @@
-- Stores document embeddings as float32 blobs instead of JSON text.
--
-- JSON text took 3-5x the space of the binary vector (about 20-32 KB per chunk
-- at 1536 dimensions versus 6 KB), on top of the binary copy in documents_vec.
-- sqlite-vec accepts both forms, so the vector triggers pass the blob through.
--
-- The full-text update trigger is narrowed first: it fired on every column,
-- so converting embeddings would have rewritten the whole full-text index.

-- @migration-step scope full-text update trigger
-- Only content, metadata (path) and page_id (title, url) feed documents_fts.
DROP TRIGGER IF EXISTS documents_fts_after_update;

CREATE TRIGGER documents_fts_after_update
AFTER UPDATE OF content, metadata, page_id ON documents
BEGIN
DELETE FROM documents_fts WHERE rowid = old.id;
INSERT INTO documents_fts(rowid, content, title, url, path)
SELECT new.id, new.content, p.title, p.url, json_extract(new.metadata, '$.path')
FROM pages p WHERE p.id = new.page_id;
END;

-- @migration-step convert embeddings to float32 blobs
-- The vector triggers are dropped so the conversion does not rewrite
-- documents_vec, which already holds the same vectors in binary form.
-- Every stored value already went through sqlite-vec's parser when it was
-- written (the triggers below, and the backfill in 011), so vec_f32 accepts it.
DROP TRIGGER IF EXISTS documents_vec_after_insert;
DROP TRIGGER IF EXISTS documents_vec_after_update;

UPDATE documents
SET embedding = vec_f32(embedding)
WHERE typeof(embedding) = 'text';

-- @migration-step recreate vector triggers
CREATE TRIGGER documents_vec_after_insert
AFTER INSERT ON documents
WHEN NEW.embedding IS NOT NULL
BEGIN
INSERT OR REPLACE INTO documents_vec (rowid, library_id, version_id, embedding)
SELECT NEW.id, v.library_id, v.id, NEW.embedding
FROM pages p
JOIN versions v ON p.version_id = v.id
WHERE p.id = NEW.page_id;
END;

CREATE TRIGGER documents_vec_after_update
AFTER UPDATE OF embedding, page_id ON documents
BEGIN
DELETE FROM documents_vec WHERE rowid = OLD.id;
INSERT OR REPLACE INTO documents_vec (rowid, library_id, version_id, embedding)
SELECT NEW.id, v.library_id, v.id, NEW.embedding
FROM pages p
JOIN versions v ON p.version_id = v.id
WHERE p.id = NEW.page_id AND NEW.embedding IS NOT NULL;
END;
12 changes: 8 additions & 4 deletions docs/concepts/data-storage.md
Original file line number Diff line number Diff line change
Expand Up @@ -174,6 +174,10 @@ Sequential SQL migrations in `db/migrations/`:
12. `011-add-vector-triggers.sql` - FTS and vector table trigger maintenance
13. `012-add-source-content-type.sql` - Source content type tracking on pages
14. `013-create-metadata-table.sql` - Key-value metadata table for embedding model tracking
15. `014-rebuild-vector-partition-keys.sql` - Partition the vector table by library and version
16. `015-add-progress-pages-indexed.sql` - Indexed page count in version progress
17. `016-add-content-url-to-pages.sql` - Retrieval location when it differs from the page URL
18. `018-store-embeddings-as-blobs.sql` - Float32 blob embeddings and a full-text update trigger limited to indexed columns

**Code Reference:** All migration files in `db/migrations/` directory

Expand Down Expand Up @@ -239,12 +243,12 @@ Handles document lifecycle operations with normalized schema access.

### Vector Storage

Embeddings stored as BLOB in documents table:
Each embedding is stored as a float32 blob in `documents.embedding`, 4 bytes per dimension:

- 1536-dimensional vectors by default (configurable via `embeddings.vectorDimension`)
- Provider-agnostic binary serialization
- NULL handling for documents without embeddings
- Direct storage eliminates need for separate vector table
- Triggers on `documents` copy the blob into the `documents_vec` sqlite-vec table, which serves KNN search
- When `documents_vec` is rebuilt, vectors missing from it are backfilled from `documents.embedding`
- NULL for documents without embeddings

**Code Reference:** `src/store/types.ts` line 4 (EMBEDDINGS_VECTOR_DIMENSION constant)

Expand Down
8 changes: 7 additions & 1 deletion openspec/specs/database-migrations/spec.md
Original file line number Diff line number Diff line change
Expand Up @@ -26,13 +26,19 @@ The system SHALL execute database migrations with rollback-capable SQLite journa

### Requirement: Production SQLite settings after migrations

The system SHALL configure production SQLite settings after migration execution completes, including WAL mode, bounded WAL checkpointing, busy timeout, foreign key enforcement, and `synchronous = NORMAL`.
The system SHALL configure production SQLite settings after migration execution completes, including WAL mode, bounded WAL checkpointing, a 64 MB `journal_size_limit`, busy timeout, foreign key enforcement, and `synchronous = NORMAL`.

#### Scenario: Post-migration settings are applied

- **WHEN** migrations complete successfully or the schema is already up to date
- **THEN** the database connection MUST be configured for WAL mode, bounded autocheckpointing, busy timeout, foreign keys, and normal synchronous durability

#### Scenario: WAL file shrinks after a large write

- **WHEN** a large write grows the WAL file beyond 64 MB
- **AND** a later checkpoint resets the WAL
- **THEN** the WAL file MUST be truncated to at most 64 MB instead of keeping its peak size

### Requirement: Visible migration progress

The system SHALL emit diagnostic progress for each pending migration. Progress MUST include the migration index and total pending migration count, migration identifier, a visible marker for each completed execution block, total elapsed time, and a completion or failure outcome.
Expand Down
45 changes: 45 additions & 0 deletions src/store/DocumentStore.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -142,6 +142,51 @@ describe("DocumentStore - With Embeddings", () => {
});

describe("Document Storage and Retrieval", () => {
it("stores each embedding as a float32 blob identical to its indexed vector", async () => {
const originalApiKey = process.env.OPENAI_API_KEY;
try {
process.env.OPENAI_API_KEY = "test-key-for-blob-storage";
await store.shutdown();
const cfg = loadConfig();
cfg.app.embeddingModel = "openai:text-embedding-3-small";
store = new DocumentStore(":memory:", cfg);
await store.initialize();

await store.addDocuments(
"bloblib",
"1.0.0",
1,
createScrapeResult(
"Blob Storage",
"https://example.com/blob-storage",
"Embeddings are stored in their binary form",
),
);

// @ts-expect-error Accessing private property for testing
const db = store.db;
const rows = db
.prepare(`
SELECT typeof(d.embedding) AS type, length(d.embedding) AS bytes,
(SELECT dv.embedding FROM documents_vec dv WHERE dv.rowid = d.id)
= d.embedding AS matchesIndex
FROM documents d
`)
.all() as Array<{ type: string; bytes: number; matchesIndex: number }>;

expect(rows.length).toBeGreaterThan(0);
for (const row of rows) {
expect(row).toEqual({ type: "blob", bytes: 1536 * 4, matchesIndex: 1 });
}
} finally {
if (originalApiKey === undefined) {
delete process.env.OPENAI_API_KEY;
} else {
process.env.OPENAI_API_KEY = originalApiKey;
}
}
});

it("should store and retrieve documents with proper metadata", async () => {
// Add two pages separately
await store.addDocuments(
Expand Down
24 changes: 12 additions & 12 deletions src/store/DocumentStore.ts
Original file line number Diff line number Diff line change
Expand Up @@ -149,7 +149,7 @@ export class DocumentStore {
// Updated for new schema - documents table now uses page_id
insertDocument: Database.Statement<[number, string, string, number]>;
// Updated for new schema - embeddings stored directly in documents table
insertEmbedding: Database.Statement<[string, bigint]>;
insertEmbedding: Database.Statement<[Buffer, bigint]>;
// New statement for pages table
insertPage: Database.Statement<
[
Expand Down Expand Up @@ -301,7 +301,7 @@ export class DocumentStore {
private prepareStatements(): void {
const statements = {
getById: this.db.prepare<[bigint]>(
`SELECT d.id, d.page_id, d.content, json(d.metadata) as metadata, d.sort_order, d.embedding, d.created_at, p.url, p.title, p.source_content_type, p.content_type, p.content_url
`SELECT d.id, d.page_id, d.content, json(d.metadata) as metadata, d.sort_order, d.created_at, p.url, p.title, p.source_content_type, p.content_type, p.content_url
FROM documents d
JOIN pages p ON d.page_id = p.id
WHERE d.id = ?`,
Expand All @@ -310,7 +310,7 @@ export class DocumentStore {
insertDocument: this.db.prepare<[number, string, string, number]>(
"INSERT INTO documents (page_id, content, metadata, sort_order) VALUES (?, ?, ?, ?)",
),
insertEmbedding: this.db.prepare<[string, bigint]>(
insertEmbedding: this.db.prepare<[Buffer, bigint]>(
"UPDATE documents SET embedding = ? WHERE id = ?",
),
insertPage: this.db.prepare<
Expand Down Expand Up @@ -427,7 +427,7 @@ export class DocumentStore {
getChildChunks: this.db.prepare<
[string, string, string, number, string, bigint, number]
>(`
SELECT d.id, d.page_id, d.content, json(d.metadata) as metadata, d.sort_order, d.embedding, d.created_at, p.url, p.title, p.source_content_type, p.content_type, p.content_url FROM documents d
SELECT d.id, d.page_id, d.content, json(d.metadata) as metadata, d.sort_order, d.created_at, p.url, p.title, p.source_content_type, p.content_type, p.content_url FROM documents d
JOIN pages p ON d.page_id = p.id
JOIN versions v ON p.version_id = v.id
JOIN libraries l ON v.library_id = l.id
Expand All @@ -443,7 +443,7 @@ export class DocumentStore {
getPrecedingSiblings: this.db.prepare<
[string, string, string, bigint, string, number]
>(`
SELECT d.id, d.page_id, d.content, json(d.metadata) as metadata, d.sort_order, d.embedding, d.created_at, p.url, p.title, p.source_content_type, p.content_type, p.content_url FROM documents d
SELECT d.id, d.page_id, d.content, json(d.metadata) as metadata, d.sort_order, d.created_at, p.url, p.title, p.source_content_type, p.content_type, p.content_url FROM documents d
JOIN pages p ON d.page_id = p.id
JOIN versions v ON p.version_id = v.id
JOIN libraries l ON v.library_id = l.id
Expand All @@ -458,7 +458,7 @@ export class DocumentStore {
getSubsequentSiblings: this.db.prepare<
[string, string, string, bigint, string, number]
>(`
SELECT d.id, d.page_id, d.content, json(d.metadata) as metadata, d.sort_order, d.embedding, d.created_at, p.url, p.title, p.source_content_type, p.content_type, p.content_url FROM documents d
SELECT d.id, d.page_id, d.content, json(d.metadata) as metadata, d.sort_order, d.created_at, p.url, p.title, p.source_content_type, p.content_type, p.content_url FROM documents d
JOIN pages p ON d.page_id = p.id
JOIN versions v ON p.version_id = v.id
JOIN libraries l ON v.library_id = l.id
Expand All @@ -471,7 +471,7 @@ export class DocumentStore {
LIMIT ?
`),
getParentChunk: this.db.prepare<[string, string, string, string, bigint]>(`
SELECT d.id, d.page_id, d.content, json(d.metadata) as metadata, d.sort_order, d.embedding, d.created_at, p.url, p.title, p.source_content_type, p.content_type, p.content_url FROM documents d
SELECT d.id, d.page_id, d.content, json(d.metadata) as metadata, d.sort_order, d.created_at, p.url, p.title, p.source_content_type, p.content_type, p.content_url FROM documents d
JOIN pages p ON d.page_id = p.id
JOIN versions v ON p.version_id = v.id
JOIN libraries l ON v.library_id = l.id
Expand Down Expand Up @@ -1260,12 +1260,12 @@ export class DocumentStore {
d.id,
v.library_id,
v.id,
json_extract(d.embedding, '$')
d.embedding
FROM documents d
JOIN pages p ON d.page_id = p.id
JOIN versions v ON p.version_id = v.id
WHERE d.embedding IS NOT NULL
AND vec_length(json_extract(d.embedding, '$')) = ?
AND vec_length(d.embedding) = ?
AND NOT EXISTS (
SELECT 1 FROM documents_vec existing WHERE existing.rowid = d.id
)
Expand Down Expand Up @@ -2109,7 +2109,7 @@ export class DocumentStore {
// Insert into vector table only if vector search is enabled
if (this.isVectorSearchEnabled && paddedEmbeddings.length > 0) {
this.statements.insertEmbedding.run(
JSON.stringify(paddedEmbeddings[docIndex]),
Buffer.from(new Float32Array(paddedEmbeddings[docIndex]).buffer),
BigInt(rowId),
);
}
Expand Down Expand Up @@ -2655,7 +2655,7 @@ export class DocumentStore {
// Use parameterized query for variable number of IDs
const placeholders = ids.map(() => "?").join(",");
const stmt = this.db.prepare(
`SELECT d.id, d.page_id, d.content, json(d.metadata) as metadata, d.sort_order, d.embedding, d.created_at, p.url, p.title, p.source_content_type, p.content_type, p.content_url FROM documents d
`SELECT d.id, d.page_id, d.content, json(d.metadata) as metadata, d.sort_order, d.created_at, p.url, p.title, p.source_content_type, p.content_type, p.content_url FROM documents d
JOIN pages p ON d.page_id = p.id
JOIN versions v ON p.version_id = v.id
JOIN libraries l ON v.library_id = l.id
Expand Down Expand Up @@ -2687,7 +2687,7 @@ export class DocumentStore {
try {
const normalizedVersion = normalizeVersionLabel(version);
const stmt = this.db.prepare(
`SELECT d.id, d.page_id, d.content, json(d.metadata) as metadata, d.sort_order, d.embedding, d.created_at, p.url, p.title, p.source_content_type, p.content_type, p.content_url FROM documents d
`SELECT d.id, d.page_id, d.content, json(d.metadata) as metadata, d.sort_order, d.created_at, p.url, p.title, p.source_content_type, p.content_type, p.content_url FROM documents d
JOIN pages p ON d.page_id = p.id
JOIN versions v ON p.version_id = v.id
JOIN libraries l ON v.library_id = l.id
Expand Down
91 changes: 91 additions & 0 deletions src/store/applyMigrations.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1086,4 +1086,95 @@ describe("Database Migrations", () => {
expect(bestExactScore).toBeLessThanOrEqual(bestPartialScore);
}
});

describe("embedding storage", () => {
const vector = new Array(1536).fill(0).map((_, index) => (index === 2 ? 1 : 0));

/** Inserts one document on a fresh library, version and page. */
function insertDocument(content: string, embedding: string | null = null) {
db.prepare("INSERT INTO libraries (name) VALUES (?)").run("blob-lib");
const { id: libraryId } = db
.prepare("SELECT id FROM libraries WHERE name = ?")
.get("blob-lib") as { id: number };
db.prepare("INSERT INTO versions (library_id, name) VALUES (?, ?)").run(
libraryId,
"1.0.0",
);
const { id: versionId } = db
.prepare("SELECT id FROM versions WHERE library_id = ?")
.get(libraryId) as { id: number };
const pageId = db
.prepare("INSERT INTO pages (version_id, url, title) VALUES (?, ?, ?)")
.run(versionId, "https://example.com/blob", "Blob").lastInsertRowid;
const docId = db
.prepare(
"INSERT INTO documents (page_id, content, metadata, sort_order, embedding) VALUES (?, ?, ?, ?, ?)",
)
.run(
pageId,
content,
JSON.stringify({ path: "/blob" }),
0,
embedding,
).lastInsertRowid;
return { libraryId, versionId, docId: Number(docId) };
}

it("should convert JSON embeddings to float32 blobs and keep them searchable", async () => {
await expect(applyMigrations(db)).resolves.toBeUndefined();
// Rows written before migration 018 hold JSON text
const { libraryId, versionId, docId } = insertDocument(
"Stored as JSON",
JSON.stringify(vector),
);
db.prepare("DELETE FROM _schema_migrations WHERE id = ?").run(
"018-store-embeddings-as-blobs.sql",
);

await expect(applyMigrations(db)).resolves.toBeUndefined();

const { embedding } = db
.prepare("SELECT embedding FROM documents WHERE id = ?")
.get(docId) as { embedding: Buffer };
expect(Buffer.isBuffer(embedding)).toBe(true);
expect(Array.from(new Float32Array(Uint8Array.from(embedding).buffer))).toEqual(
vector,
);

const nearest = db
.prepare(`
SELECT rowid, distance
FROM documents_vec
WHERE library_id = ? AND version_id = ? AND embedding MATCH ? AND k = 1
`)
.get(libraryId, versionId, JSON.stringify(vector)) as
| { rowid: number; distance: number }
| undefined;
expect(nearest?.rowid).toBe(docId);
expect(nearest?.distance).toBeCloseTo(0, 6);
});

it("should reindex full text when document content changes", async () => {
await expect(applyMigrations(db)).resolves.toBeUndefined();
const { docId } = insertDocument("original wording");

db.prepare("UPDATE documents SET content = ? WHERE id = ?").run(
"revised wording",
docId,
);

const matches = (term: string) =>
db
.prepare("SELECT rowid FROM documents_fts WHERE documents_fts MATCH ?")
.all(term);
expect(matches("revised")).toEqual([{ rowid: docId }]);
expect(matches("original")).toEqual([]);
});
});

it("should cap the WAL file size kept after checkpoints", async () => {
await expect(applyMigrations(db)).resolves.toBeUndefined();

expect(db.pragma("journal_size_limit", { simple: true })).toBe(64 * 1024 * 1024);
});
});
8 changes: 7 additions & 1 deletion src/store/applyMigrations.ts
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@ const MIGRATION_STEP_MARKER = /^\s*--\s*@migration-step\s+(.+?)\s*$/;
const VECTOR_PARTITION_MIGRATION = "014-rebuild-vector-partition-keys.sql";
const DEFAULT_VECTOR_DIMENSION = 1536;
const VECTOR_DIMENSION_TOKEN = "__DOCUMENTS_VEC_DIMENSION__";
const WAL_SIZE_LIMIT_BYTES = 64 * 1024 * 1024;

interface MigrationStep {
label: string;
Expand Down Expand Up @@ -303,6 +304,11 @@ export async function applyMigrations(
// Configure WAL autocheckpoint to prevent unbounded growth
db.pragma("wal_autocheckpoint = 1000"); // Checkpoint every 1000 pages (~4MB)

// SQLite reuses the WAL file after a checkpoint but never shrinks it, so one
// large write (a bulk scrape, a migration, VACUUM) would leave it at its peak
// size until the last connection closes. Trim it back once the WAL resets.
db.pragma(`journal_size_limit = ${WAL_SIZE_LIMIT_BYTES}`);

// Set busy timeout for better handling of concurrent access
db.pragma("busy_timeout = 30000"); // 30 seconds

Expand All @@ -313,7 +319,7 @@ export async function applyMigrations(
db.pragma("synchronous = NORMAL");

logger.debug(
"Applied production database configuration (WAL mode, autocheckpoint, foreign keys, busy timeout)",
"Applied production database configuration (WAL mode, autocheckpoint, WAL size limit, foreign keys, busy timeout)",
);
} catch (_error) {
logger.warn("⚠️ Could not apply all production database settings");
Expand Down
Loading
Loading