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
66 changes: 10 additions & 56 deletions bun.lock

Large diffs are not rendered by default.

12 changes: 2 additions & 10 deletions package.json
Original file line number Diff line number Diff line change
Expand Up @@ -12,7 +12,7 @@
"homepage": "https://workglow.dev",
"description": "A worked example of the Workglow libraries: SEC EDGAR and Form ADV ingestion, filing-to-markdown conversion, and retrieval-grounded question answering over what filings say.",
"scripts": {
"release-checks": "bun run format && bun run lint && bun run typecheck && bun run build && bun run prepack-check",
"release-checks": "bun run format-check && bun run lint && bun run typecheck && bun run build && bun run test && bun run prepack-check",
"release": "bun run release-checks && bunset --auto --push --commit --tag",
"prepack-check": "bun ./scripts/checkPackedContents.ts",
"dev": "concurrently -c 'auto' -n 'sec:' 'bun:dev-*'",
Expand Down Expand Up @@ -52,27 +52,19 @@
},
"dependencies": {
"@huggingface/transformers": "^4.3.0",
"@huggingface/transformers-structured-output": "^4.3.0",
"@modelcontextprotocol/sdk": "^1.30.0",
"@workglow/cli": "0.6.7",
"cheerio": "^1.2.0",
"cheerio-json-mapper": "^1.0.4",
"commander": "^15.0.0",
"compromise": "^14.17.0",
"csv-parse": "^7.0.2",
"fast-xml-parser": "^5.11.1",
"html-entities": "^2.6.0",
"pdf2json": "^4.1.0",
"pg": "^8.23.0",
"typebox": "1.3.34",
"words-to-numbers": "^1.5.1",
"workglow": "0.6.7",
"xml2js": "^0.6.2"
"workglow": "0.6.7"
},
"devDependencies": {
"@types/bun": "1.4.2",
"@types/pg": "^8.23.1",
"@types/xml2js": "^0.4.14",
"bunset": "1.1.2",
"concurrently": "^10.0.5",
"oxfmt": "0.68.0",
Expand Down
142 changes: 142 additions & 0 deletions src/config/kbIndexedStamp.sqlite.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,142 @@
/**
* @license
* Copyright 2026 Steven Roussey <sroussey@gmail.com>
* SPDX-License-Identifier: Apache-2.0
*/

import { afterEach, beforeEach, describe, expect, it } from "vitest";
import type { DocumentNode } from "workglow";
import { Document, globalServiceRegistry, NodeKind } from "workglow";
import { getSecKnowledgeBase, resetSecKnowledgeBaseForTesting } from "../kb/secKnowledgeBase";
import {
FILING_DOCUMENT_REPOSITORY_TOKEN,
type FilingDocument,
} from "../storage/document/FilingDocumentSchema";
import { kbDocIdFor } from "../task/kb/selectDocumentsToIndex";
import { selectDocumentsToIndex } from "../task/kb/selectDocumentsToIndex";
import { getDb } from "../util/db";
import { syncKbIndexedStamp } from "./kbIndexedStamp";
import { withSqliteDb } from "./testing/withSqliteDb";

const accession = (index: number) => `0000320193-26-${String(index).padStart(6, "0")}`;

const doc = (index: number, over: Partial<FilingDocument> = {}): FilingDocument => ({
cik: 320193,
accession_number: accession(index),
doc_file: "primary.htm",
doc_type: "10-K",
description: null,
sequence: 1,
is_primary: true,
form: "10-K",
filing_date: "2026-03-01",
title: `Filing ${index}`,
section_count: 1,
char_count: 100,
converter_version: "1",
converted_at: "2026-01-01T00:00:00.000Z",
kb_indexed_at: null,
...over,
});

/**
* The stamp is a cache of "is there a `kb_document` row for this", and the three
* properties below are what make it safe to read: it never claims more than the
* anti-join, `db setup` never leaves it empty, and the steady state stops
* touching the table.
*/
describe("kb_indexed_at stamp (sqlite)", () => {
withSqliteDb("kb_indexed_stamp", [FILING_DOCUMENT_REPOSITORY_TOKEN]);

beforeEach(async () => {
await resetSecKnowledgeBaseForTesting();
delete process.env.SEC_EMBEDDING_MODEL;
});
afterEach(async () => {
await resetSecKnowledgeBaseForTesting();
});

async function seed(count: number) {
const repo = globalServiceRegistry.get(FILING_DOCUMENT_REPOSITORY_TOKEN);
for (let i = 1; i <= count; i += 1) await repo.put(doc(i) as never);
}

async function putInKb(index: number) {
const kb = await getSecKnowledgeBase();
const title = `Filing ${index}`;
const root = { kind: NodeKind.DOCUMENT, title, children: [] } as unknown as DocumentNode;
await kb.upsertDocument(
new Document(root, { title } as never, [], kbDocIdFor(accession(index), "primary.htm"))
);
}

const stampOf = (index: number): string | null =>
(
getDb()
.prepare("SELECT `kb_indexed_at` AS s FROM `filing_document` WHERE `accession_number` = ?")
.get(accession(index)) as { s: string | null } | undefined
)?.s ?? null;

it("backfills a database whose documents were indexed before the column existed", async () => {
// The state every existing deployment is in: rows in the knowledge base,
// no stamp. Left un-backfilled the partial index would cover the whole
// table and the selection would walk it, which is the plan this replaces.
await seed(3);
await putInKb(1);
await putInKb(2);

syncKbIndexedStamp(getDb());

expect(stampOf(1)).not.toBeNull();
expect(stampOf(2)).not.toBeNull();
expect(stampOf(3)).toBeNull();
});

it("is idempotent, and does not restamp a row on a second run", async () => {
await seed(1);
await putInKb(1);

syncKbIndexedStamp(getDb());
const first = stampOf(1);
syncKbIndexedStamp(getDb());

expect(stampOf(1)).toBe(first);
});

it("creates the partial index, and the steady state reads it instead of the table", async () => {
await seed(3);
for (const i of [1, 2, 3]) await putInKb(i);
syncKbIndexedStamp(getDb());

const plan = getDb()
.prepare(
`EXPLAIN QUERY PLAN
SELECT d.* FROM \`filing_document\` d
WHERE d.\`section_count\` > 0 AND d.\`kb_indexed_at\` IS NULL
ORDER BY d.\`filing_date\` DESC, d.\`accession_number\` DESC`
)
.all() as Array<{ detail?: string }>;
const detail = plan.map((row) => row.detail ?? "").join(" | ");

expect(detail).toContain("filing_document_kb_unindexed");
// The point of the partial index: no table scan, and no temp B-tree to
// sort what a scan would have produced.
expect(detail).not.toContain("SCAN d\n");
expect(detail).not.toContain("TEMP B-TREE");
});

it("never widens the selection: an unstamped document already in the kb is still skipped", async () => {
// The stamp is a narrowing of the anti-join, never a replacement for it.
// If a stamp goes missing — a crash between the upsert and the stamp, a
// database restored from before the backfill — the anti-join still has to
// keep that document out, because indexing it twice spends embedding calls.
await seed(2);
await putInKb(1);
// Deliberately NOT stamped.
expect(stampOf(1)).toBeNull();

const picked = await selectDocumentsToIndex({ limit: 10 });

expect(picked.map((d) => d.accession_number)).toEqual([accession(2)]);
});
});
88 changes: 88 additions & 0 deletions src/config/kbIndexedStamp.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,88 @@
/**
* @license
* Copyright 2026 Steven Roussey <sroussey@gmail.com>
* SPDX-License-Identifier: Apache-2.0
*/

import { Sqlite } from "workglow";
import { KB_DOCUMENT_TABLE } from "../kb/secKbTables";

/** The partial index the steady-state selection walks. */
const PARTIAL_INDEX = "filing_document_kb_unindexed";

/**
* Brings `filing_document.kb_indexed_at` up to date and indexes the nulls.
*
* Two statements that have to travel together, which is why they are one
* function rather than a registry entry. The column is a cache of "is there a
* `kb_document` row for this", and a cache that is added empty is worse than no
* cache: every row reads as unindexed, so the partial index below covers the
* whole table and the planner walks it in `filing_date` order, issuing a random
* probe per row — the access path that measured 11x slower than the plain scan
* it replaced. Backfilling in the same pass that creates the index is what
* keeps that state from existing.
*
* Both statements are idempotent and cheap on a database already in step: the
* `UPDATE` matches nothing once every indexed document is stamped, and
* `CREATE INDEX IF NOT EXISTS` is a catalog read.
*
* Callers must have applied the schema's columns first — the column is added
* generically from the registry, not here — so this no-ops when it is absent
* rather than assuming an order it cannot see.
*/
export function syncKbIndexedStamp(db: Sqlite.Database): void {
if (!hasColumn(db, "filing_document", "kb_indexed_at")) return;
if (!tableExists(db, KB_DOCUMENT_TABLE)) {
// No knowledge base yet, so nothing is indexed and every row's null is
// already the truth. The index still pays off on the first `ask`.
createPartialIndex(db);
return;
}

// `converted_at` rather than a timestamp taken here: the row's own history is
// a better answer than "whenever someone ran db setup", and it keeps the
// backfill deterministic, so two runs cannot disagree about a row.
db.prepare(
`UPDATE \`filing_document\`
SET \`kb_indexed_at\` = \`converted_at\`
WHERE \`kb_indexed_at\` IS NULL
AND EXISTS (
SELECT 1 FROM \`${KB_DOCUMENT_TABLE}\` k
WHERE k.\`doc_id\` = \`filing_document\`.\`accession_number\`
|| ':' || \`filing_document\`.\`doc_file\`
)`
).run();

createPartialIndex(db);
}

/**
* Indexed on the order the selection reads — newest first, ties by accession —
* so the same walk serves both the filter and the `ORDER BY`. Partial, because
* the rows it must not contain are the ones that accumulate: a corpus that is
* fully indexed leaves it empty, which is the point.
*/
function createPartialIndex(db: Sqlite.Database): void {
db.prepare(
`CREATE INDEX IF NOT EXISTS \`${PARTIAL_INDEX}\`
ON \`filing_document\` (\`filing_date\` DESC, \`accession_number\` DESC)
WHERE \`kb_indexed_at\` IS NULL`
).run();
}

function hasColumn(db: Sqlite.Database, table: string, column: string): boolean {
const rows = db.prepare(`PRAGMA table_info(\`${table}\`)`).all() as Array<{ name?: unknown }>;
return rows.some((row) => row.name === column);
}

function tableExists(db: Sqlite.Database, table: string): boolean {
return (
db.prepare("SELECT name FROM sqlite_master WHERE type = 'table' AND name = ?").all(table)
.length > 0
);
}

/** Whether the fast path is available — the column exists, so setup backfilled it. */
export function kbIndexedStampAvailable(db: Sqlite.Database): boolean {
return hasColumn(db, "filing_document", "kb_indexed_at");
}
15 changes: 9 additions & 6 deletions src/config/setupAllDatabases.ts
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@ import {
} from "./addMissingColumns";
import { alignPostgresColumnTypes } from "./alignPostgresColumnTypes";
import { dropStaleCheckConstraints } from "./dropStaleCheckConstraints";
import { syncKbIndexedStamp } from "./kbIndexedStamp";
import { SEC_STORAGE_REGISTRY } from "./storageRegistry";
import { SEC_DB_FOLDER, SEC_DB_TYPE } from "./tokens";

Expand Down Expand Up @@ -49,12 +50,14 @@ export async function setupAllDatabases(): Promise<void> {

// Skipped under --dry-run: these passes issue DDL through raw SQL, which the
// repositories' ReadOnlyTabularStorage wrapper cannot intercept.
if (
dbType === "sqlite" &&
globalServiceRegistry.has(SEC_DB_FOLDER) &&
shouldAddMissingColumns("sqlite")
) {
addMissingColumnsSqlite(getDb());
if (dbType === "sqlite" && globalServiceRegistry.has(SEC_DB_FOLDER)) {
if (shouldAddMissingColumns("sqlite")) {
addMissingColumnsSqlite(getDb());
}
// After the column pass, because the stamp it maintains is one of the
// columns that pass adds — and it must be backfilled in the same run that
// creates its index, never left empty.
syncKbIndexedStamp(getDb());
}
if (dbType === "postgres" && !isDryRun()) {
if (shouldAddMissingColumns("postgres")) await addMissingColumnsPostgres();
Expand Down
18 changes: 18 additions & 0 deletions src/storage/document/FilingDocumentSchema.ts
Original file line number Diff line number Diff line change
Expand Up @@ -90,6 +90,24 @@ export const FilingDocumentSchema = Type.Object({
*/
converter_version: Type.String({ maxLength: 32, description: "Converter version stamp" }),
converted_at: Type.String({ description: "ISO 8601 timestamp" }),
/**
* When this document entered the knowledge base, or null while it has not.
*
* A denormalisation of "is there a `kb_document` row for this", and it earns
* its keep in the one regime the index selection spends its life in. `ask`
* pre-indexes before every question, so after the first question the
* selection matches nothing — and a `LIMIT` that never fills has to visit
* every candidate to learn that. Against a partial index over the nulls, the
* steady state is an empty index scan instead of the whole table.
*
* It is a fast path, never the authority: the anti-join against
* `kb_document` still runs, so a stamp that is missing or stale costs a
* traversal rather than a document indexed twice. That direction matters,
* because indexing twice spends embedding calls.
*/
kb_indexed_at: TypeNullable(
Type.String({ description: "ISO 8601 timestamp the document entered the knowledge base" })
),
});

export type FilingDocument = Static<typeof FilingDocumentSchema>;
Expand Down
5 changes: 5 additions & 0 deletions src/task/document/ConvertFilingDocumentTask.ts
Original file line number Diff line number Diff line change
Expand Up @@ -387,6 +387,11 @@ export class ConvertFilingDocumentTask extends Task<
char_count: doc.charCount,
converter_version: FILING_CONVERTER_VERSION,
converted_at: convertedAt,
// Freshly converted markdown is not in the knowledge base yet, and a
// RE-conversion at a newer converter version has replaced the text the
// knowledge base holds — so both start unindexed. Narrower than the
// anti-join beside it, never wider, which is the safe direction.
kb_indexed_at: null,
});
}

Expand Down
1 change: 1 addition & 0 deletions src/task/document/selectFilingsToConvert.sqlite.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,7 @@ const documentRow = (accession: string, docFile: string, isPrimary: boolean) =>
char_count: 900,
converter_version: VERSION,
converted_at: "2026-03-02T00:00:00.000Z",
kb_indexed_at: null,
});

/**
Expand Down
1 change: 1 addition & 0 deletions src/task/document/selectFilingsToConvert.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -61,6 +61,7 @@ async function markConverted(
char_count: 100,
converter_version: version,
converted_at: "2026-08-01T00:00:00.000Z",
kb_indexed_at: null,
});
}

Expand Down
25 changes: 25 additions & 0 deletions src/task/kb/IndexFilingSectionsTask.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,7 @@ const header = (index: number, filingDate: string): FilingDocument => ({
char_count: 100,
converter_version: "1",
converted_at: "2026-01-01T00:00:00.000Z",
kb_indexed_at: null,
});

/**
Expand Down Expand Up @@ -213,4 +214,28 @@ describe("IndexFilingSectionsTask under --dry-run", () => {
expect(upsert).toHaveBeenCalledTimes(1);
expect(out).toMatchObject({ success: true, indexed: 1, sections: 1 });
});

it("stamps kb_indexed_at on what it indexed, and leaves it alone on a dry run", async () => {
// The stamp is what lets the next run answer "anything to index?" from a
// partial index rather than a traversal. A dry run indexes nothing, so it
// must claim nothing either.
const kb = await getSecKnowledgeBase();
vi.spyOn(kb, "upsert").mockResolvedValue({
doc_id: "0000320193-26-000000:primary.htm",
} as never);
const stamp = () =>
(
getDb().prepare("SELECT `kb_indexed_at` AS s FROM `filing_document` LIMIT 1").get() as
| { s: string | null }
| undefined
)?.s ?? null;

globalServiceRegistry.registerInstance(SEC_DRY_RUN, true);
await new IndexFilingSectionsTask().run({});
expect(stamp()).toBeNull();

globalServiceRegistry.registerInstance(SEC_DRY_RUN, false);
await new IndexFilingSectionsTask().run({});
expect(stamp()).not.toBeNull();
});
});
Loading
Loading