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
128 changes: 83 additions & 45 deletions ci/journey.sh
Original file line number Diff line number Diff line change
Expand Up @@ -137,6 +137,19 @@ ice_files() {
q "$HOST" "SELECT coldfront.ensure_attached(); SELECT r['n'] FROM duckdb.query('SELECT count(*) AS n FROM iceberg_metadata(''$1'') WHERE status <> ''DELETED'' AND content NOT LIKE ''%DELETES'' AND regexp_matches(file_path, ''$2'')') AS t(r);" | tail -1
}

# ice_deletes <ice_ref>: live delete files, counted the same way.
ice_deletes() {
q "$HOST" "SELECT coldfront.ensure_attached(); SELECT r['n'] FROM duckdb.query('SELECT count(*) AS n FROM iceberg_metadata(''$1'') WHERE status <> ''DELETED'' AND content LIKE ''%DELETES''') AS t(r);" | tail -1
}

# vec_misplaced <table> <column>: cold rows whose cluster is not the nearest live
# centroid, scored in PostgreSQL with the pgvector operator, independently of the
# generator the write paths use. A row with no embedding is skipped; a row with
# no assignment counts. Zero is the invariant every probe relies on.
vec_misplaced() {
q "$HOST" "SELECT coldfront.ensure_attached(); SELECT count(*) FROM iceberg_scan('ice.public.$1') r WHERE r['$2'] IS NOT NULL AND r['_cf_vec_list_$2']::int IS DISTINCT FROM (SELECT c.centroid_id FROM coldfront.vector_centroids c WHERE c.table_name = '$1' AND c.column_name = '$2' AND c.generation = (SELECT generation FROM coldfront.vector_config WHERE table_name = '$1' AND column_name = '$2') ORDER BY c.centroid <=> r['$2']::real[] LIMIT 1);" | tail -1
}

# ice_spec <ice_ref>: the table's partition spec as transform:column terms in
# field order (e.g. "identity:region,month:ts"), empty for an unpartitioned
# table or one that has no data yet.
Expand Down Expand Up @@ -1138,6 +1151,12 @@ EOSQL
)
assert_eq "training clamps nprobe to the trained count" "2/2" "$(extract CLAMP "$CL")"

# A row another engine appends straight to Iceberg carries no assignment. One
# such row, written under the claim as a foreign committer would and holding
# coldins's vector, is what the null arm is asserted with below.
local fts; fts=$(q "$HOST" "SELECT (date_trunc('month',now()) - interval '4 months' + interval '25 days')::text;")
q "$HOST" "SET duckdb.unsafe_allow_mixed_transactions = on; SELECT coldfront.ensure_attached(); SELECT coldfront._exec_iceberg_with_claim('\"ice\".\"public\".\"chunks\"', 'INSERT INTO \"ice\".\"public\".\"chunks\" (id, ts, body, embedding) VALUES (990001, TIMESTAMPTZ ''$fts'', ''foreign'', [4.0, 5.5, -6.0]::FLOAT[])');" >/dev/null

# With a generation live, every path that writes a vector into the cold tier
# stamps the cluster in the same statement. Two paths, one vector, read back
# from Iceberg: the view cannot project the column, by design.
Expand Down Expand Up @@ -1171,15 +1190,17 @@ EOSQL
)
assert_eq "the stamped cluster is the nearest centroid" "$(extract EXPECT "$E")" "$(extract ACTUAL "$E")"

# Rows archived before any training stay unassigned, which a probe reads
# through the null arm of its predicate rather than missing.
# Rows archived before any training are assigned by the training itself, to
# the centroid nearest them, the same as a row written after it.
local U; U=$(qf "$HOST" <<'EOSQL'
SELECT coldfront.ensure_attached();
SELECT 'PRE:' || count(*) FROM iceberg_scan('ice.public.chunks') r
WHERE r['body'] = 'cold' AND r['_cf_vec_list_embedding'] IS NULL;
SELECT 'EXPECT:' || (SELECT centroid_id FROM coldfront.vector_centroids
WHERE table_name = 'chunks' AND generation = (SELECT generation FROM coldfront.vector_config WHERE table_name = 'chunks')
ORDER BY centroid <=> ARRAY[1.5,-2,3]::real[] LIMIT 1);
SELECT 'PRE:' || min(r['_cf_vec_list_embedding']::int) FROM iceberg_scan('ice.public.chunks') r WHERE r['body'] = 'cold';
EOSQL
)
assert_eq "a row written before training stays unassigned" "1" "$(extract PRE "$U")"
assert_eq "a row archived before training is assigned to its nearest centroid" "$(extract EXPECT "$U")" "$(extract PRE "$U")"

# A cold UPDATE that sets a new embedding must re-derive the cluster in the
# same statement. The row is found by its new cluster and no longer by its
Expand Down Expand Up @@ -1244,7 +1265,7 @@ EOSQL
ONE_ROWS=$(q "$HOST" "SET coldfront.vector_nprobe = 1; $topk LIMIT 100;" | grep -vx SET)
assert_eq "a probe reads the cluster it aimed at" "1" "$(printf '%s\n' "$ONE_ROWS" | grep -c '^asg_trig$' || true)"
assert_eq "a probe skips the cluster it did not look in" "0" "$(printf '%s\n' "$ONE_ROWS" | grep -c '^asg_bulk1$' || true)"
assert_eq "a probe keeps every unassigned row" "1" "$(printf '%s\n' "$ONE_ROWS" | grep -c '^coldins$' || true)"
assert_eq "a probe keeps every unassigned row" "1" "$(printf '%s\n' "$ONE_ROWS" | grep -c '^foreign$' || true)"
assert_eq "a probe keeps every hot row" "1" "$(printf '%s\n' "$ONE_ROWS" | grep -c '^hot$' || true)"

# The predicate reaches the Iceberg scan itself rather than filtering above it,
Expand All @@ -1256,12 +1277,11 @@ EOSQL
assert_eq "a declined shape carries no predicate" "0" "$(printf '%s\n' "$PLAN_PLAIN" | grep -c '_cf_vec_list_embedding' || true)"

# The null arm carrying its weight, stated as an answer rather than a count:
# coldins holds this exact vector and was written before there was a generation,
# so a probe returns it only through that disjunct. Asserted as membership in the
# top few rather than as the single winner: asg_trig carries a centroid, and a
# centroid of a cluster with one member IS that member's vector, so the two can
# tie at distance zero and either is then a correct answer to LIMIT 1.
local NEAR; NEAR=$(q "$HOST" "SET coldfront.vector_nprobe = 1; SELECT body FROM chunks ORDER BY embedding <=> ARRAY[4,5.5,-6]::real[] LIMIT 3;" | grep -vx SET | grep -c '^coldins$' || true)
# the foreign row holds this exact vector and no assignment, so a probe returns
# it only through that disjunct. Asserted as membership in the top few rather
# than as the single winner: coldins holds the same vector, so the two tie at
# distance zero and either is a correct answer to LIMIT 1.
local NEAR; NEAR=$(q "$HOST" "SET coldfront.vector_nprobe = 1; SELECT body FROM chunks ORDER BY embedding <=> ARRAY[4,5.5,-6]::real[] LIMIT 3;" | grep -vx SET | grep -c '^foreign$' || true)
assert_eq "a probed search still finds an unassigned row" "1" "$NEAR"

# Grouping, aggregates, windows and DISTINCT ride the same narrowed scan: the
Expand Down Expand Up @@ -1333,39 +1353,7 @@ EOSQL
assert_eq "clusters under one row group are counted" "true" "$(extract FLOOR "$S")"
assert_eq "a fixture this small draws advice" "true" "$(extract ADVICE "$S")"

# Assigning the rows that predate training. These are the rows vector_status
# just counted as unassigned, and they are read by every probe whatever clusters
# it looks in, so on a corpus tiered before training they are the entire cost of
# a search. One claimed UPDATE, reusing the SET item a cold UPDATE emits.
local unasg_before; unasg_before=$(extract REPORTS "$S" | cut -d/ -f2)
local A; A=$(qf "$HOST" <<'EOSQL'
CALL coldfront.vector_assign('public', 'chunks', 'embedding');
SELECT coldfront.ensure_attached();
SELECT 'LEFT:' || count(*) FROM iceberg_scan('ice.public.chunks') r
WHERE r['_cf_vec_list_embedding'] IS NULL;
SELECT 'ROWS:' || count(*) FROM iceberg_scan('ice.public.chunks') r;
SELECT 'MATCHES:' || count(*) FROM iceberg_scan('ice.public.chunks') r
WHERE r['_cf_vec_list_embedding']::int = (
SELECT centroid_id FROM coldfront.vector_centroids c
WHERE c.table_name = 'chunks'
AND c.generation = (SELECT generation FROM coldfront.vector_config WHERE table_name = 'chunks')
ORDER BY c.centroid <=> r['embedding']::real[] LIMIT 1);
EOSQL
)
assert_gt "there were rows to assign" "0" "$unasg_before"
assert_eq "no cold row is left without a cluster" "0" "$(extract LEFT "$A")"
assert_eq "assigning rewrote rows, it did not add any" "7" "$(extract ROWS "$A")"
assert_eq "every row sits in its nearest cluster" "$(extract ROWS "$A")" "$(extract MATCHES "$A")"

# Which is the point: with nothing unassigned, a narrow probe stops reading the
# whole table.
local narrowed
narrowed=$(q "$HOST" "SET coldfront.vector_nprobe = 1; $topk LIMIT 100;" | grep -vx SET | grep -c . || true)
assert_gt "a narrow probe now reads less than everything" "$narrowed" "$ALL"

# Idempotent: nothing to do the second time, and no rewrite to pay for.
local A2; A2=$(q_may "$HOST" "CALL coldfront.vector_assign('public','chunks','embedding');")
assert_contains "a second assign has nothing to do" "already assigned" "$A2"
story_vector_retrain

# A NULL embedding through the identity-omitted slow path. The Iceberg schema
# declares one cluster column per vector column unconditionally, so the row's
Expand Down Expand Up @@ -1394,6 +1382,56 @@ EOSQL
story_vector_compaction
}

# ───────────────────────────────────────────────────────────────────────────
# Story 5b, retraining. A retrain replaces the centroids and, in the same
# transaction, re-assigns every cold row whose nearest centroid changed, so no
# probe ever reads a row against centroids it was not assigned under. At the
# same nlist the centroids keep their ids and move with their data, so the rows
# that stay put are not rewritten; at a new nlist every row moves. Every probe
# is asserted against the exact scan, which is the answer a probe may narrow
# but never lose.
# ───────────────────────────────────────────────────────────────────────────
story_vector_retrain() {
step "5b. Retraining re-assigns the rows that moved"
local topk="SELECT body FROM chunks ORDER BY embedding <=> ARRAY[4,5.5,-6]::real[] LIMIT 100"
local live="SELECT generation FROM coldfront.vector_config WHERE table_name = 'chunks'"
local ids="SELECT string_agg(centroid_id::text, ',' ORDER BY centroid_id) FROM coldfront.vector_centroids WHERE table_name = 'chunks' AND generation = ($live)"
local exact; exact=$(q "$HOST" "SET coldfront.vector_probe = off; $topk;" | grep -vx SET | sort | tr '\n' ' ')
local gen; gen=$(q "$HOST" "$live;")

# TC-195: a smaller nlist. Every row moves, the foreign row included, and an
# exhaustive probe (nprobe = nlist) must still return the exact answer.
local R; R=$(q_may "$HOST" "CALL coldfront.vector_train('public','chunks','embedding', 1);")
assert_contains "TC-195: the retrain wrote one centroid" "trained 1 centroids" "$R"
assert_eq "TC-195: an exhaustive probe after the retrain returns the exact answer" "$exact" \
"$(q "$HOST" "SET coldfront.vector_nprobe = 1; $topk;" | grep -vx SET | sort | tr '\n' ' ')"
assert_eq "TC-195: every cold row sits in its nearest cluster" "0" "$(vec_misplaced chunks embedding)"

# TC-196: a different nlist starts afresh, and a narrow probe narrows again.
R=$(q_may "$HOST" "CALL coldfront.vector_train('public','chunks','embedding', 2);")
assert_contains "TC-196: the retrain wrote two centroids" "trained 2 centroids" "$R"
assert_eq "TC-196: an exhaustive probe returns the exact answer" "$exact" \
"$(q "$HOST" "SET coldfront.vector_nprobe = 2; $topk;" | grep -vx SET | sort | tr '\n' ' ')"
assert_eq "TC-196: every cold row sits in its nearest cluster" "0" "$(vec_misplaced chunks embedding)"
local all one
all=$(q "$HOST" "SET coldfront.vector_nprobe = 2; $topk;" | grep -vx SET | grep -c . || true)
one=$(q "$HOST" "SET coldfront.vector_nprobe = 1; $topk;" | grep -vx SET | grep -c . || true)
assert_gt "TC-196: a narrow probe reads less than everything" "$one" "$all"

# TC-197: the same nlist keeps each centroid's identity. Nothing has changed
# since the last training, so no row moves and no delete file is written.
local ids_before deletes_before
ids_before=$(q "$HOST" "$ids;")
deletes_before=$(ice_deletes ice.public.chunks)
R=$(q_may "$HOST" "CALL coldfront.vector_train('public','chunks','embedding');")
assert_eq "TC-197: the retrain moved the pointer" "$((gen + 3))" "$(q "$HOST" "$live;")"
assert_eq "TC-197: the centroids kept their ids" "$ids_before" "$(q "$HOST" "$ids;")"
assert_eq "TC-197: no row moved, so no delete file was written" "$deletes_before" "$(ice_deletes ice.public.chunks)"
assert_eq "TC-197: every cold row sits in its nearest cluster" "0" "$(vec_misplaced chunks embedding)"
assert_eq "TC-197: an exhaustive probe returns the exact answer" "$exact" \
"$(q "$HOST" "SET coldfront.vector_nprobe = 2; $topk;" | grep -vx SET | sort | tr '\n' ' ')"
}

# ───────────────────────────────────────────────────────────────────────────
# Story 5c — Compaction on a clustered table. Each cold write leaves a file
# sorted within itself, so what accumulates is a set of files with overlapping
Expand Down
6 changes: 3 additions & 3 deletions docs/architecture_decoupled.md
Original file line number Diff line number Diff line change
Expand Up @@ -337,9 +337,9 @@ ERROR: coldfront: "public.orders" is adopted read-only
HINT: Release it with coldfront.release_iceberg_table() and adopt again with p_writable => true to arm INSERT/UPDATE/DELETE.
```

Reads never consult the flag. `vector_train()`, `vector_assign()` and
`drop_iceberg_table()` refuse a read-only relation too, since each
rewrites or destroys the Iceberg table. The archiver and
Reads never consult the flag. `vector_train()` and `drop_iceberg_table()`
refuse a read-only relation too, since each rewrites or destroys the Iceberg
table. The archiver and
`create_iceberg_table()` set the flag; adoption defaults it to false.

### One relation per Iceberg table
Expand Down
62 changes: 30 additions & 32 deletions docs/architecture_vectors.md
Original file line number Diff line number Diff line change
Expand Up @@ -128,8 +128,9 @@ spellings of the name.
## Assignment

`coldfront._vec_list_expr(schema, table, column, vec_expr)` is the only place a
cluster assignment is defined. Given the text of an expression that yields the
vector as DuckDB sees it, it returns:
cluster assignment is defined for a write; its formula, `_vec_nearest_expr`, is
shared with the assignment at the end of training. Given the text of an
expression that yields the vector as DuckDB sees it, it returns:

```sql
(SELECT arg_min(c.centroid_id, list_cosine_distance(c.centroid, <vec_expr>))
Expand Down Expand Up @@ -223,8 +224,10 @@ Lloyd iterations run as DuckDB statements over a reservoir sample, with fixed
`REPEATABLE` seeds so a retrain on unchanged data reproduces the same centroids.
The mean recompute unnests the vector against a matching `range` so the two lists
advance together, because DuckDB has no `WITH ORDINALITY`. The sample and the
working tables live in `memory.main`, which is session-scoped, so the whole loop
must run in one call.
working tables are DuckDB temporary tables, session-scoped, so the whole loop
must run in one call; temporary rather than `memory.main` because a DuckDB
transaction may write one attached database, the temporary database is exempt,
and the assignment at the end writes `ice`.

The centroids return through a temporary heap table. A single
`INSERT … SELECT FROM duckdb.query(…)` fails with `DuckDB does not support
Expand All @@ -235,6 +238,26 @@ Empty clusters do not come back from the mean, so the stored count can be below
`nlist`; the actual count is recorded rather than padded. A `vector_config` row
must exist first, since it holds the generation pointer the procedure writes.

A retrain at the same `nlist` (a call without `p_nlist`, or with the trained
count) starts the iterations from the live centroids instead of fresh seeds, so
each centroid keeps its id and moves with its data. A first training, or a
changed `nlist`, has no set to start from and seeds with k-means++.

Training ends with the loop's final step applied to the table rather than the
sample: every cold row is assigned to its nearest centroid of the generation just
written, in one claimed `UPDATE` whose `WHERE` is `cluster IS DISTINCT FROM
nearest`, so only a row whose cluster changed is rewritten and a row with no
cluster counts as changed. It scores against the session's copy of the new
centroids (`temp.main.cf_gen`), because the rows inserted into
`vector_centroids` in this transaction are invisible over `pglocal` until commit;
the formula is `_vec_nearest_expr`, the same one every write path uses. The
pointer and the assignments commit together, so no search ever reads a row
against centroids it was not assigned under. The claim is held from the sample
to the commit, as the compactor holds it across its read and rewrite, so a cold
write cannot land between the two with an assignment against the replaced set.
What the pass leaves behind is a merge-on-read delete per rewritten row, and
those rows in update order; compaction resolves both.

Seeds come from k-means++: one sample row at random, then each next drawn with
probability proportional to its squared distance from the nearest seed already
chosen, by an exponential race (the minimum of `-ln(u)/w` is a weighted draw, in one
Expand Down Expand Up @@ -398,33 +421,6 @@ Two session knobs, both `PGC_USERSET`:
| `coldfront.vector_probe` | `on` | `off` gives the exact scan a recall measurement compares against |
| `coldfront.vector_nprobe` | `0` | `0` uses the column's configured `nprobe`; a value at or above `nlist` is exhaustive |

## Assigning what predates training

`CALL coldfront.vector_assign(schema, table, column)` gives a cluster to the cold
rows that have none, which are exactly those written before the generation existed.
It is one claimed UPDATE whose SET item is the same generator a cold UPDATE uses
when a caller changes an embedding, applied to the embedding already stored, so it
serialises through the bakery like every other cold write and adds no new way to
compute an assignment.

It refuses without a live generation rather than reporting success. The lookup would
resolve to NULL for every row, so the table would be rewritten in full and left
exactly as it was.

Two constraints shape the body, and both bite anything else written here:

- **A DuckDB transaction may write one attached database.** Staging the row count in
a `memory.main` table would spend this transaction's one database on `memory` and
leave the UPDATE unable to write `ice` at all.
- **`duckdb.query` needs a constant at plan time, not a literal in the source.**
Built through `EXECUTE format(...)` the argument is a constant again, which is how
a dynamic table name is read without staging anything. `vector_status` reads its
distribution the same way.

What it leaves behind is a merge-on-read delete per rewritten row, and rows in
update order rather than cluster order. Compaction resolves both. The sequence is
train, assign, compact.

## Reporting

`CALL coldfront.vector_status([schema, table])` fills a session-lifetime temporary
Expand All @@ -437,7 +433,9 @@ because a bare `CALL` is its own transaction.

Nothing is staged on the way. DuckDB groups the table by cluster and aggregates
that grouping in one query, so a row of scalars crosses back per table, read through
the same `EXECUTE format(...)` form `vector_assign` uses. The work list is a pair of
`EXECUTE format(...)`: `duckdb.query` needs a constant at plan time, not a literal
in the source, and a dynamic table name built into the executed text is one
again. The work list is a pair of
key arrays walked by index, because a `FOR` over a query would hold a portal open
for its body and pg_duckdb refuses a DuckDB read while one is. `vector_train`'s
`memory.main` tables are its algorithm's own state, not a way to move a result
Expand Down
Loading
Loading