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
Original file line number Diff line number Diff line change
Expand Up @@ -257,19 +257,31 @@ describe('retiring dormant Search embeddings', () => {
await sql`INSERT INTO embedding VALUES ('26003', 'second-search', 'second-search-doc')`
await sql`CREATE TABLE deletion_blocker (id text REFERENCES embedding(id))`
await sql`INSERT INTO deletion_blocker VALUES ('26002')`
/**
* Pages split by mutated rows under an adaptive limit, so how far each failed run gets depends
* on page geometry. The geometry-free invariant: every committed delete sits at or behind the
* cursor, and a failed page leaves every row past it (IDs 00001-26003 are contiguous).
*/
async function expectRolledBackPastCursor() {
const [progress] = await sql`SELECT phase, after_id FROM search_embedding_cleanup_progress`
expect(progress.phase).toBe('embeddings')
const [beyond] =
await sql`SELECT count(*)::int AS n FROM embedding WHERE id > ${progress.after_id}`
expect(beyond.n).toBe(26003 - Number(progress.after_id))
expect(progress.after_id < '26002').toBe(true)
}
await expect(pass()).rejects.toThrow()
const before = await sql`SELECT * FROM search_embedding_cleanup_progress`
const [remaining] = await sql`SELECT count(*)::int AS n FROM embedding`
expect(remaining.n).toBeGreaterThan(501)
expect(remaining.n).toBeLessThan(26002)
expect(await sql`SELECT name FROM script_migrations`).toHaveLength(0)
await expectRolledBackPastCursor()
await sql`UPDATE knowledge_base SET is_search_index = false WHERE id = 'second-search'`
await expect(pass()).rejects.toThrow('no longer a Search knowledge base')
expect((await sql`SELECT count(*)::int AS n FROM embedding`)[0].n).toBe(remaining.n)
await expectRolledBackPastCursor()
await sql`UPDATE knowledge_base SET is_search_index = true WHERE id = 'second-search'`
await expect(pass()).rejects.toThrow()
expect(await sql`SELECT * FROM search_embedding_cleanup_progress`).toEqual(before)
expect((await sql`SELECT count(*)::int AS n FROM embedding`)[0].n).toBe(remaining.n)
await expectRolledBackPastCursor()
await sql`DROP TABLE deletion_blocker`
await sql`INSERT INTO document (id, knowledge_base_id) VALUES ('aaa-late-document', 'search')`
await sql`INSERT INTO embedding VALUES ('00000', 'search', 'aaa-late-document')`
Expand Down Expand Up @@ -309,6 +321,225 @@ describe('retiring dormant Search embeddings', () => {
).toBe(0)
}, 60_000)

it('splits pages whose writes would outrun the statement timeout and resumes at the split', async () => {
const bound = 300
await sql`INSERT INTO document (id, knowledge_base_id, user_excluded, enabled)
SELECT 'doc-' || lpad(i::text, 5, '0'), CASE WHEN i % 3 = 0 THEN 'ordinary' ELSE 'search' END,
i % 7 = 0, i % 7 <> 0
FROM generate_series(1, 4000) i`
await sql`INSERT INTO embedding
SELECT lpad(i::text, 5, '0'), CASE WHEN i % 3 = 0 THEN 'ordinary' ELSE 'search' END,
CASE WHEN i % 3 = 0 THEN 'ordinary-doc' ELSE 'search-doc' END
FROM generate_series(1003, 5002) i`
await sql`CREATE TABLE committed_statement (rows integer NOT NULL)`
/** Sequences are not transactional, so this counts the timed-out statements that rolled back. */
await sql`CREATE SEQUENCE timed_out_statement`
/** Stands in for write cost: a statement touching more than `bound` rows times out and rolls back. */
await sql.unsafe(`CREATE FUNCTION bound_statement_rows() RETURNS trigger LANGUAGE plpgsql AS $$
DECLARE touched integer;
BEGIN
SELECT count(*) INTO touched FROM changed_rows;
IF touched > ${bound} THEN
PERFORM nextval('timed_out_statement');
RAISE EXCEPTION 'canceling statement due to statement timeout' USING ERRCODE = 'query_canceled';
END IF;
INSERT INTO committed_statement VALUES (touched);
RETURN NULL;
END $$`)
await sql`CREATE TRIGGER bound_document_update AFTER UPDATE ON document
REFERENCING NEW TABLE AS changed_rows FOR EACH STATEMENT EXECUTE FUNCTION bound_statement_rows()`
await sql`CREATE TRIGGER bound_embedding_delete AFTER DELETE ON embedding
REFERENCING OLD TABLE AS changed_rows FOR EACH STATEMENT EXECUTE FUNCTION bound_statement_rows()`
try {
expect(await pass()).toBe(true)
const [{ largest, total }] = await sql`SELECT max(rows)::int AS largest,
sum(rows)::int AS total FROM committed_statement`
expect(largest).toBeLessThanOrEqual(bound)
expect(largest).toBeGreaterThan(0)
/**
* The limit halves 2,000 → 1,000 → 500 → 250 on the first document page and never grows back
* to a size that timed out, in either phase: exactly three rolled-back statements. A page
* that read past its row limit in either phase would time out again.
*/
const [{ timeouts }] = await sql`SELECT last_value::int AS timeouts FROM timed_out_statement`
expect(timeouts).toBe(3)
/**
* Documents: i % 3 <> 0 gives 2,667 Search rows, of which i % 7 = 0 leaves 381 retired, so
* 2,286 are updated, plus `search-doc`. Chunks: 501 Search rows from the fixture plus the
* 2,667 with i % 3 <> 0 among 1003-5002. Equality proves no row was mutated twice.
*/
expect(total).toBe(2287 + 3168)
expect(
(
await sql`SELECT count(*)::int AS n FROM document
WHERE knowledge_base_id = 'search' AND (NOT user_excluded OR enabled)`
)[0].n
).toBe(0)
expect(
(
await sql`SELECT count(*)::int AS n FROM document
WHERE knowledge_base_id = 'ordinary' AND NOT user_excluded AND enabled`
)[0].n
/** `ordinary-doc` plus the 1,333 i % 3 = 0 rows, less the 190 of them seeded retired. */
).toBe(1 + 1333 - 190)
expect(
(await sql`SELECT count(*)::int AS n FROM embedding WHERE knowledge_base_id = 'search'`)[0]
.n
).toBe(0)
expect(
(
await sql`SELECT count(*)::int AS n FROM embedding WHERE knowledge_base_id = 'ordinary'`
)[0].n
/** 501 fixture chunks plus the 1,333 i % 3 = 0 rows among 1003-5002. */
).toBe(501 + 1333)
} finally {
await sql`DROP TRIGGER IF EXISTS bound_document_update ON document`
await sql`DROP TRIGGER IF EXISTS bound_embedding_delete ON embedding`
await sql`DROP FUNCTION bound_statement_rows()`
await sql`DROP TABLE committed_statement`
await sql`DROP SEQUENCE timed_out_statement`
}
}, 60_000)

it('bounds the IDs each page reads by the row limit once the limit shrinks', async () => {
const docs = 3000
const bound = 25
await sql`INSERT INTO document (id, knowledge_base_id)
SELECT 'doc-' || lpad(i::text, 5, '0'), 'search' FROM generate_series(1, ${docs}) i`
/** Statements over `bound` rows time out, which pins the row limit at its 25-row floor. */
await sql.unsafe(`CREATE FUNCTION bound_document_update() RETURNS trigger LANGUAGE plpgsql AS $$
BEGIN
IF (SELECT count(*) FROM changed_rows) > ${bound} THEN
RAISE EXCEPTION 'canceling statement due to statement timeout' USING ERRCODE = 'query_canceled';
END IF;
RETURN NULL;
END $$`)
await sql`CREATE TRIGGER bound_document_update AFTER UPDATE ON document
REFERENCING NEW TABLE AS changed_rows FOR EACH STATEMENT EXECUTE FUNCTION bound_document_update()`
/**
* On a table this small the planner may answer any page with a sequential scan, which reads
* every row whatever the window; production pages use the primary key, so the test does too.
*/
await sql`SET enable_seqscan = off`
/** Document rows read by any scan, counted across committed and rolled-back pages alike. */
async function documentReads() {
await sql`SELECT pg_stat_force_next_flush()`
await admin`SELECT pg_stat_clear_snapshot()`
const [row] = await admin`SELECT (seq_tup_read + coalesce(idx_tup_fetch, 0))::int AS n
FROM pg_stat_user_tables WHERE schemaname = ${schema} AND relname = 'document'`
return row.n
}
try {
const before = await documentReads()
expect(await pass()).toBe(true)
const reads = (await documentReads()) - before
/**
* About 120 pages retire the 3,001 documents 25 at a time once the limit has halved down to
* its floor. A window of four IDs per row reads about 100 IDs and 25 update lookups per page,
* roughly 5 reads per document, plus the chunk phase and completion rechecks. A fixed
* 25,000-ID window re-reads the rest of the table on every attempt, over 100 per document.
*/
expect(reads).toBeLessThan(30 * docs)
expect(
(
await sql`SELECT count(*)::int AS n FROM document
WHERE knowledge_base_id = 'search' AND (NOT user_excluded OR enabled)`
)[0].n
).toBe(0)
} finally {
await sql`RESET enable_seqscan`
await sql`DROP TRIGGER IF EXISTS bound_document_update ON document`
await sql`DROP FUNCTION bound_document_update()`
}
}, 120_000)

it('gives the completed-retirement recheck the completion timeout when it resumes', async () => {
expect(await pass()).toBe(true)
await sql`DELETE FROM script_migrations`
/** Stands in for a recheck that outlasts the two-minute page timeout on a large target set. */
await sql`CREATE FUNCTION require_recheck_timeout() RETURNS boolean LANGUAGE plpgsql AS $$
BEGIN
IF current_setting('statement_timeout') <> '30min' THEN
RAISE EXCEPTION 'canceling statement due to statement timeout' USING ERRCODE = 'query_canceled';
END IF;
RETURN true;
END $$`
await sql`ALTER TABLE search_embedding_cleanup_targets RENAME TO captured_targets`
await sql`CREATE VIEW search_embedding_cleanup_targets AS
SELECT knowledge_base_id FROM captured_targets WHERE require_recheck_timeout()`
try {
expect(await sql`SELECT phase FROM search_embedding_cleanup_progress`).toEqual([
{ phase: 'done' },
])
expect(await pass()).toBe(true)
} finally {
await sql`DROP VIEW IF EXISTS search_embedding_cleanup_targets`
await sql`ALTER TABLE IF EXISTS captured_targets RENAME TO search_embedding_cleanup_targets`
await sql`DROP FUNCTION require_recheck_timeout()`
}
})

it('scans past a run of already-retired documents longer than the row limit in one page', async () => {
/** 6,000 retired Search documents exceed the initial 2,000-row limit but fit one 25,000-ID scan. */
await sql`INSERT INTO document (id, knowledge_base_id, user_excluded, enabled)
SELECT 'doc-' || lpad(i::text, 5, '0'), 'search', true, false FROM generate_series(1, 6000) i`
await sql`INSERT INTO document (id, knowledge_base_id)
SELECT 'doc-' || lpad(i::text, 5, '0'), 'search' FROM generate_series(6001, 6010) i`
await sql`CREATE SEQUENCE document_page_statements`
await sql`CREATE FUNCTION count_document_page() RETURNS trigger LANGUAGE plpgsql AS $$
BEGIN
PERFORM nextval('document_page_statements');
RETURN NULL;
END $$`
/** Statement triggers fire even for zero rows, so this counts every documents-phase page. */
await sql`CREATE TRIGGER count_document_page AFTER UPDATE ON document
FOR EACH STATEMENT EXECUTE FUNCTION count_document_page()`
try {
expect(await pass()).toBe(true)
/** One page retires the 11 unretired rows (10 bulk plus `search-doc`); one more finds the end. */
expect((await sql`SELECT last_value::int AS n FROM document_page_statements`)[0].n).toBe(2)
Comment thread
waleedlatif1 marked this conversation as resolved.
expect(
(
await sql`SELECT count(*)::int AS n FROM document
WHERE knowledge_base_id = 'search' AND (NOT user_excluded OR enabled)`
)[0].n
).toBe(0)
} finally {
await sql`DROP TRIGGER IF EXISTS count_document_page ON document`
await sql`DROP FUNCTION count_document_page()`
await sql`DROP SEQUENCE document_page_statements`
}
})

it('fails at once on a timeout outside the page mutation instead of shrinking the page', async () => {
await sql`CREATE SEQUENCE completion_attempts`
/** Times out the completion checkpoint, a statement no smaller row limit can speed up. */
await sql`CREATE FUNCTION time_out_completion() RETURNS trigger LANGUAGE plpgsql AS $$
BEGIN
PERFORM nextval('completion_attempts');
RAISE EXCEPTION 'canceling statement due to statement timeout' USING ERRCODE = 'query_canceled';
END $$`
await sql`CREATE TABLE search_embedding_cleanup_progress (
id integer PRIMARY KEY CHECK (id = 1), knowledge_base_id text NOT NULL,
phase text NOT NULL CHECK (phase IN ('documents', 'embeddings', 'done')),
after_id text NOT NULL)`
await sql`CREATE TRIGGER time_out_completion BEFORE UPDATE ON search_embedding_cleanup_progress
FOR EACH ROW WHEN (NEW.phase = 'done') EXECUTE FUNCTION time_out_completion()`
try {
await expect(pass()).rejects.toMatchObject({ code: '57014' })
/** The sequence is not transactional, so it counts rolled-back attempts too. */
expect((await sql`SELECT last_value::int AS n FROM completion_attempts`)[0].n).toBe(1)
expect(await sql`SELECT name FROM script_migrations`).toHaveLength(0)
expect((await sql`SELECT count(*)::int AS n FROM embedding`)[0].n).toBe(501)
await sql`DROP TRIGGER time_out_completion ON search_embedding_cleanup_progress`
expect(await pass()).toBe(true)
} finally {
await sql`DROP TRIGGER IF EXISTS time_out_completion ON search_embedding_cleanup_progress`
await sql`DROP FUNCTION time_out_completion()`
await sql`DROP SEQUENCE completion_attempts`
}
})

it('rejects inconsistent document ownership before deleting any chunk in the page', async () => {
await sql`UPDATE embedding SET document_id = 'ordinary-doc' WHERE id = '00002'`
await expect(pass()).rejects.toThrow('Search content changed after retirement')
Expand Down
Loading
Loading