From f138ac0abe6bcf2040e511d62f955f7b0f6cb07f Mon Sep 17 00:00:00 2001 From: "shoney.arickathil" Date: Sun, 30 Aug 2026 10:22:42 +0200 Subject: [PATCH] feat(db2-delta): replay and compaction fold delta chains MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - apply_delta: DELTA replay arm — fold pre-delta state via back_off, overlay the field, remove-then-recreate so indexes stay correct - apply_record/replay loop: dispatch DELTA to apply_delta, drop its payload back to the log same as INSERT/UPDATE - stage_flattened_row: compaction's delta-chain path — fold + re-encode as one fresh INSERT instead of copying the chain - wo_wal_compact: peek the row's current record kind, flatten deltas, keep the byte-for-byte copy for chains already at length zero - test_wal: three new tests — chain-of-three replay incl. secondary index, compaction flattens to chain length zero (asserts the record is a full row, not a delta), and the commit-before-repoint crash window replays the update without ever re-pointing the map Co-Authored-By: Claude Opus 5 (1M context) (cherry picked from commit 7e4ae70d7beb2a2e396a002184bc06f41253decb) --- database/src/wal.c | 138 +++++++++++++++++++++- runtime/test/test_wal.c | 250 ++++++++++++++++++++++++++++++++++++++++ 2 files changed, 382 insertions(+), 6 deletions(-) diff --git a/database/src/wal.c b/database/src/wal.c index fc5e352..1c41fd6 100644 --- a/database/src/wal.c +++ b/database/src/wal.c @@ -705,6 +705,45 @@ static int copy_record(wo_wal *nw, wo_wal *ow, uint64_t off, uint32_t want_cid, return rc; } +/* Task 5 (replay and compaction fold the same way): copy_record's + * counterpart for a row whose chain is NOT empty — a delta cannot be moved + * byte-for-byte the way copy_record moves a base row, because the whole + * point of compaction is to reset every chain to length zero. Fold the row + * (THE fold, same one reads and replay use) and re-encode it as a fresh + * INSERT via enc_val, exactly the representation the fold's own contract + * promises (ENGINE values, not the VM's) — so this, unlike copy_record, + * never touches wo_row_ptr/a slab at all. */ +static int stage_flattened_row(wo_wal *nw, wo_wal *ow, wo_db *db, uint64_t off, + uint32_t class_id, uint64_t id, const char **why) { + const wo_classdesc *c = &db->classes[class_id]; + uint64_t *vals = c->field_cnt ? calloc(c->field_cnt, sizeof *vals) : NULL; + if (c->field_cnt && !vals) { + *why = "out of memory flattening a delta chain"; + return -1; + } + uint32_t got_cid = 0; + uint64_t got_id = 0; + const char *fmsg = ""; + if (wo_wal_fold_row_at(ow, db, off, &got_cid, &got_id, vals, &fmsg) != 0 || + got_cid != class_id || got_id != id) { + for (uint32_t i = 0; i < c->field_cnt; i++) wo_db_val_free(db, c->kinds[i], vals[i]); + free(vals); + *why = "a live row's delta chain does not fold cleanly"; + return -1; + } + wbuf p = {0}; + wput_u8(&p, WO_WAL_INSERT); + wput_u32(&p, class_id); + wput_u64(&p, id); + for (uint32_t i = 0; i < c->field_cnt; i++) enc_val(&p, db->classes, c->kinds[i], vals[i]); + int rc = stage(nw, &p); + free(p.b); + for (uint32_t i = 0; i < c->field_cnt; i++) wo_db_val_free(db, c->kinds[i], vals[i]); + free(vals); + if (rc != 0) *why = "staging a flattened row failed"; + return rc; +} + int wo_wal_compact(wo_wal *w, wo_db *db) { /* staged records would be written into a file about to be replaced */ if (!w->path || w->len != 0) return -1; @@ -757,7 +796,25 @@ int wo_wal_compact(wo_wal *w, wo_db *db) { goto fail; } uint64_t at = nw.off + (uint64_t)nw.len; - if (copy_record(&nw, w, o1 - 1, cid, id, &why) != 0) goto fail; + /* Task 5: a chain flattens to ONE full row on every + * checkpoint — the bound the whole design relies on + * (without it, chains grow without limit). Peek the kind + * byte at the row's current record (position fixed by the + * frame header, valid whatever the record turns out to be) + * to choose: an empty chain keeps the existing byte-for-byte + * copy, faster and unchanged; a delta chain folds instead of + * being copied — copying it would just move the chain, not + * flatten it. */ + uint8_t kind_byte = 0; + if (pread(w->fd, &kind_byte, 1, (off_t)(o1 - 1 + 8)) != 1) { + why = "cannot read a row's recorded offset"; + goto fail; + } + if (kind_byte == WO_WAL_DELTA) { + if (stage_flattened_row(&nw, w, db, o1 - 1, cid, id, &why) != 0) goto fail; + } else if (copy_record(&nw, w, o1 - 1, cid, id, &why) != 0) { + goto fail; + } /* value-only update: cannot rehash, so `cur` stays valid */ if (wo_row_set_offset(db, cid, id, at) != 0) { why = "row vanished from the id map mid-compaction"; @@ -822,6 +879,72 @@ fail: return -1; /* the live log is untouched and still usable */ } +/* Task 5 (replay and compaction fold the same way): the DELTA replay arm. + * + * A delta cannot be applied like an insert/update — its body is one field, + * not the row — so it borrows WO_WAL_UPDATE's remove-then-recreate SHAPE + * instead: fold the row's PRE-delta state (THE fold, the same one reads and + * compaction use) via [back_off], overlay this delta's field onto it, then + * remove-and-recreate through the ordinary choke points. Recreating, rather + * than patching an index in place, is what re-indexes a changed INDEXED + * column for free — wo_row_remove takes the OLD value's entries out, + * wo_row_raw_commit puts the NEW value's entries in, exactly as a live + * update would have. The caller (wo_wal_replay_ex) still owns dropping the + * recreated row's payload back to the log, the same way it does for + * INSERT/UPDATE. */ +static int apply_delta(wo_db *db, uint32_t cid, uint64_t id, rbuf *r) { + uint32_t field_idx = rd_u32(r); + uint64_t back_off = rd_u64(r); + const wo_classdesc *c = &db->classes[cid]; + if (r->bad || field_idx >= c->field_cnt) return -1; + uint64_t nv; + if (dec_val(r, db, c->kinds[field_idx], &nv) != 0 || (size_t)(r->end - r->p) != 0) + return -1; + + uint32_t field_cnt = c->field_cnt; + uint64_t *vals = field_cnt ? calloc(field_cnt, sizeof *vals) : NULL; + if (field_cnt && !vals) { + wo_db_val_free(db, c->kinds[field_idx], nv); + return -1; + } + uint32_t got_cid = 0; + uint64_t got_id = 0; + const char *fmsg = ""; + if (wo_wal_fold_row_at((wo_wal *)db->rt->wal, db, back_off, &got_cid, &got_id, vals, + &fmsg) != 0 || + got_cid != cid || got_id != id) { + for (uint32_t i = 0; i < field_cnt; i++) wo_db_val_free(db, c->kinds[i], vals[i]); + free(vals); + wo_db_val_free(db, c->kinds[field_idx], nv); + return -1; + } + wo_db_val_free(db, c->kinds[field_idx], vals[field_idx]); + vals[field_idx] = nv; + + /* the row must exist (its base record precedes its deltas in a correct + log); anything else is corruption */ + if (wo_row_remove(db, cid, id) != 0) { + for (uint32_t i = 0; i < field_cnt; i++) wo_db_val_free(db, c->kinds[i], vals[i]); + free(vals); + return -1; + } + db_row *row = wo_row_create_raw(db, cid, id); + if (!row) { + for (uint32_t i = 0; i < field_cnt; i++) wo_db_val_free(db, c->kinds[i], vals[i]); + free(vals); + return -1; + } + memcpy(row->slots, vals, (size_t)field_cnt * sizeof *vals); + free(vals); + /* a unique violation here is corruption, same reasoning as INSERT/UPDATE + below: the live update that produced this delta would have refused it */ + if (wo_row_raw_commit(db, cid, row) != 0) { + wo_row_remove(db, cid, id); + return -1; + } + return 0; +} + static int apply_record(wo_db *db, const uint8_t *payload, uint32_t len) { rbuf r = {payload, payload + len, 0}; uint8_t kind = rd_u8(&r); @@ -835,6 +958,7 @@ static int apply_record(wo_db *db, const uint8_t *payload, uint32_t len) { * any. -2 so the caller can say which of the two it is. */ if (db->classes[cid].flags & WO_CLASSF_VOLATILE) return -2; if (kind == WO_WAL_REMOVE) return wo_row_remove(db, cid, id); + if (kind == WO_WAL_DELTA) return apply_delta(db, cid, id, &r); if (kind != WO_WAL_INSERT && kind != WO_WAL_UPDATE) return -1; if (kind == WO_WAL_UPDATE) { /* replace: the row must exist (its insert precedes its update in a @@ -1120,9 +1244,10 @@ int64_t wo_wal_replay_ex(const char *path, wo_db *db, uint32_t *volatile_cid) { int64_t applied = 0; /* databasev2 2: replay must be able to READ rows back, not only write them. - * A REMOVE (or an UPDATE, which replays as remove-then-recreate) on a - * keys-resident table reaches wo_row_remove, whose keys arm borrows the row - * out of the log to find its index entries — and a borrow reads through + * A REMOVE (or an UPDATE or a DELTA, both of which replay as + * remove-then-recreate — see apply_delta, Task 5) on a keys-resident + * table reaches wo_row_remove, whose keys arm borrows the row out of the + * log to find its index entries — and a borrow reads through * db->rt->wal. At boot that pointer is not wired yet: main.c replays first * and assigns rt.wal afterwards, so the borrow found no log, the remove * failed, and replay reported a perfectly good tombstone as CORRUPTION. @@ -1171,8 +1296,9 @@ int64_t wo_wal_replay_ex(const char *path, wo_db *db, uint32_t *volatile_cid) { if (rc == -2 && volatile_cid) *volatile_cid = rec_cid; REPLAY_RETURN(rc == -2 ? -2 : -1); } - if ((rec_kind == WO_WAL_INSERT || rec_kind == WO_WAL_UPDATE) && rec_id && - wo_table_is_keys_resident(db, rec_cid)) + if ((rec_kind == WO_WAL_INSERT || rec_kind == WO_WAL_UPDATE || + rec_kind == WO_WAL_DELTA) && + rec_id && wo_table_is_keys_resident(db, rec_cid)) (void)wo_row_drop_payload(db, rec_cid, rec_id, off); off += 8u + len + 4u; applied++; diff --git a/runtime/test/test_wal.c b/runtime/test/test_wal.c index e620ebe..67aedfc 100644 --- a/runtime/test/test_wal.c +++ b/runtime/test/test_wal.c @@ -1445,6 +1445,253 @@ static void test_keys_resident_replay(void) { wo_rt_destroy(&rt); } +/* Task 5 (replay and compaction fold the same way): DELTA_CLASSES with a + * non-unique index on field 0 — the chain-of-deltas tests below need THREE + * touched fields (one delta each) but also an indexed column among them, + * which DELTA_CLASSES (no index) and KEYS_IDX_CLASSES (only two fields) + * don't together provide. */ +static const uint8_t delta_idx_kinds[] = {WO_K_SCALAR, WO_K_SCALAR, WO_K_SCALAR}; +static const uint32_t delta_idx_meta[] = {0 /*non-unique*/, 1, 0 /*col: field 0*/}; +static const wo_classdesc DELTA_IDX_CLASSES[] = { + {.name = 0, .flags = WO_CLASSF_RESIDENT_KEYS, .field_cnt = 3, .kinds = delta_idx_kinds, + .idx_cnt = 1, .idx_meta = delta_idx_meta}, +}; + +/* Task 5, step 1: a row with a chain of three deltas, replayed into a fresh + * database, must read exactly as it did before the restart — including + * through the secondary index on the column one of the deltas changed. + * Before this task apply_record had no DELTA arm: a DELTA record in the log + * made replay refuse the whole file as corruption. */ +static void test_delta_chain_replay(void) { + char path[128]; + snprintf(path, sizeof path, "%s/deltareplay.wal", g_dir); + wo_rt rt; + T_EQ(wo_rt_init(&rt, 1 << 20, DELTA_IDX_CLASSES, 1), 0); + wo_db db; + T_EQ(wo_db_init(&db, DELTA_IDX_CLASSES, 1, 0, 1), 0); + wo_wal w; + T_EQ(wo_wal_open(&w, path, 1 << 16), 0); + db.rt = &rt; rt.wal = &w; rt.db = &db; + const char *msg = ""; + + uint64_t vals[3] = {10, 20, 30}; + uint64_t id = wo_row_insert(&db, 0, vals, &msg, NULL); + T_CHECK(id != 0); + uint64_t base_off = wo_wal_next_offset(&w); + T_EQ(wo_wal_append_insert(&w, &db, 0, id), 0); + T_EQ(wo_wal_commit(&w), 0); + T_EQ(wo_row_drop_payload(&db, 0, id, base_off), 0); + + /* field 0 (indexed): 10 -> 111 */ + uint64_t d1_off = wo_wal_next_offset(&w); + T_EQ(wo_wal_append_delta(&w, &db, 0, id, 0, base_off, 111), 0); + T_EQ(wo_wal_commit(&w), 0); + T_EQ(wo_row_set_offset(&db, 0, id, d1_off), 0); + + /* field 1: 20 -> 222, chained off the first delta */ + uint64_t d2_off = wo_wal_next_offset(&w); + T_EQ(wo_wal_append_delta(&w, &db, 0, id, 1, d1_off, 222), 0); + T_EQ(wo_wal_commit(&w), 0); + T_EQ(wo_row_set_offset(&db, 0, id, d2_off), 0); + + /* field 2: 30 -> 333, chained off the second delta — three deltas total */ + uint64_t d3_off = wo_wal_next_offset(&w); + T_EQ(wo_wal_append_delta(&w, &db, 0, id, 2, d2_off, 333), 0); + T_EQ(wo_wal_commit(&w), 0); + T_EQ(wo_row_set_offset(&db, 0, id, d3_off), 0); + + /* what the row reads as BEFORE the restart */ + db_row *before = wo_row_borrow(&db, 0, id, &msg); + T_CHECK(before != NULL); + uint64_t want0 = before->slots[0], want1 = before->slots[1], want2 = before->slots[2]; + T_EQ(want0, 111); + T_EQ(want1, 222); + T_EQ(want2, 333); + wo_row_release(&db, 0, before); + + wo_wal_close(&w); + wo_db_destroy(&db); + + /* a fresh process would do exactly this: rt.wal is NULL (main.c wires + the real log in only AFTER replay) so a DELTA's fold, mid-replay, runs + through the lent view wo_wal_replay_ex sets up on its own — db.rt + itself must already be set, exactly as main.c sets DB.rt before + calling wo_wal_replay_ex. */ + wo_db db2; + T_EQ(wo_db_init(&db2, DELTA_IDX_CLASSES, 1, 0, 1), 0); + db2.rt = &rt; + rt.wal = NULL; + T_EQ(wo_wal_replay(path, &db2), 4); /* 1 insert + 3 deltas */ + wo_wal w2; + T_EQ(wo_wal_open(&w2, path, 1 << 16), 0); + rt.wal = &w2; rt.db = &db2; + + db_row *after = wo_row_borrow(&db2, 0, id, &msg); + T_CHECK(after != NULL); + T_EQ(after->slots[0], want0); + T_EQ(after->slots[1], want1); + T_EQ(after->slots[2], want2); + wo_row_release(&db2, 0, after); + + /* the index a delta changed: rebuilt at boot from the INSERT's value, + then folded forward by the delta that touched field 0 — a fold that + disagreed between reading and replaying would leave this probing the + stale value */ + uint64_t *ids; + uint32_t cnt; + T_EQ(wo_idx_probe(&db2, 0, 0, 111, NULL, 0, &ids, &cnt), 1); + T_CHECK(cnt == 1 && ids[0] == id); + free(ids); + T_EQ(wo_idx_probe(&db2, 0, 0, 10, NULL, 0, &ids, &cnt), 1); + T_CHECK(cnt == 0 && ids == NULL); /* the superseded value: no hits */ + + wo_wal_close(&w2); + wo_db_destroy(&db2); + wo_rt_destroy(&rt); +} + +/* Task 5, step 2: the same chain, then compacted. The row must read + * identically AND its chain must be length zero afterwards — the record its + * offset points at must be a full row, not a delta. That second assertion + * is the one the brief calls out as easy to skip: without it this test + * would still pass if compaction merely copied the chain byte-for-byte + * instead of flattening it, since a byte-for-byte copy still reads back + * correctly — it just never shortens the chain, which is the entire point + * of a checkpoint. */ +static void test_delta_chain_compact_flattens(void) { + char path[128]; + snprintf(path, sizeof path, "%s/deltacompact.wal", g_dir); + wo_rt rt; + T_EQ(wo_rt_init(&rt, 1 << 20, DELTA_IDX_CLASSES, 1), 0); + wo_db db; + T_EQ(wo_db_init(&db, DELTA_IDX_CLASSES, 1, 0, 1), 0); + wo_wal w; + T_EQ(wo_wal_open(&w, path, 1 << 16), 0); + db.rt = &rt; rt.wal = &w; rt.db = &db; + const char *msg = ""; + + uint64_t vals[3] = {10, 20, 30}; + uint64_t id = wo_row_insert(&db, 0, vals, &msg, NULL); + T_CHECK(id != 0); + uint64_t base_off = wo_wal_next_offset(&w); + T_EQ(wo_wal_append_insert(&w, &db, 0, id), 0); + T_EQ(wo_wal_commit(&w), 0); + T_EQ(wo_row_drop_payload(&db, 0, id, base_off), 0); + + uint64_t d1_off = wo_wal_next_offset(&w); + T_EQ(wo_wal_append_delta(&w, &db, 0, id, 0, base_off, 111), 0); + T_EQ(wo_wal_commit(&w), 0); + T_EQ(wo_row_set_offset(&db, 0, id, d1_off), 0); + + uint64_t d2_off = wo_wal_next_offset(&w); + T_EQ(wo_wal_append_delta(&w, &db, 0, id, 1, d1_off, 222), 0); + T_EQ(wo_wal_commit(&w), 0); + T_EQ(wo_row_set_offset(&db, 0, id, d2_off), 0); + + uint64_t d3_off = wo_wal_next_offset(&w); + T_EQ(wo_wal_append_delta(&w, &db, 0, id, 2, d2_off, 333), 0); + T_EQ(wo_wal_commit(&w), 0); + T_EQ(wo_row_set_offset(&db, 0, id, d3_off), 0); + + T_EQ(wo_wal_compact(&w, &db), 0); + + /* reads identically, through the compacted log */ + db_row *r = wo_row_borrow(&db, 0, id, &msg); + T_CHECK(r != NULL); + T_EQ(r->slots[0], 111); + T_EQ(r->slots[1], 222); + T_EQ(r->slots[2], 333); + wo_row_release(&db, 0, r); + + /* THE assertion the brief calls out: the offset now names a FULL ROW, + not a delta — chain length zero, not merely "still readable" */ + uint64_t o1 = wo_row_offset1(&db, 0, id); + T_CHECK(o1 != 0); + uint8_t kind_byte = 0xFF; + T_EQ((int)pread(w.fd, &kind_byte, 1, (off_t)(o1 - 1 + 8)), 1); + T_EQ(kind_byte, WO_WAL_INSERT); + + wo_wal_close(&w); + wo_db_destroy(&db); + + /* the compacted log replays to the same, flattened, state */ + wo_db db2; + T_EQ(wo_db_init(&db2, DELTA_IDX_CLASSES, 1, 0, 1), 0); + db2.rt = &rt; + rt.wal = NULL; /* see the replay comment above: set before replaying */ + T_EQ(wo_wal_replay(path, &db2), 1); /* one row, one INSERT — chain gone */ + wo_wal w2; + T_EQ(wo_wal_open(&w2, path, 1 << 16), 0); + rt.wal = &w2; rt.db = &db2; + db_row *r2 = wo_row_borrow(&db2, 0, id, &msg); + T_CHECK(r2 != NULL); + T_EQ(r2->slots[0], 111); + T_EQ(r2->slots[1], 222); + T_EQ(r2->slots[2], 333); + wo_row_release(&db2, 0, r2); + wo_wal_close(&w2); + wo_db_destroy(&db2); + wo_rt_destroy(&rt); +} + +/* Task 5, step 3: the crash window commit-before-re-point ordering exists + * for. Append a delta, commit it (durable), and deliberately do NOT + * re-point the map — exactly the state a crash between the barrier and the + * flush leaves behind (wo_db_flush_drops never got to run). Replaying the + * log into a FRESH database, which never sees this process's map at all, + * must still surface the update: the commit alone is what makes a delta + * recoverable, not the in-RAM re-point, and this is the test that would + * fail if that ordering were ever reversed. */ +static void test_delta_crash_window(void) { + char path[128]; + snprintf(path, sizeof path, "%s/deltacrash.wal", g_dir); + wo_rt rt; + T_EQ(wo_rt_init(&rt, 1 << 20, DELTA_CLASSES, 1), 0); + wo_db db; + T_EQ(wo_db_init(&db, DELTA_CLASSES, 1, 0, 1), 0); + wo_wal w; + T_EQ(wo_wal_open(&w, path, 1 << 16), 0); + db.rt = &rt; rt.wal = &w; rt.db = &db; + const char *msg = ""; + + uint64_t vals[3] = {1, 2, 3}; + uint64_t id = wo_row_insert(&db, 0, vals, &msg, NULL); + T_CHECK(id != 0); + uint64_t base_off = wo_wal_next_offset(&w); + T_EQ(wo_wal_append_insert(&w, &db, 0, id), 0); + T_EQ(wo_wal_commit(&w), 0); + T_EQ(wo_row_drop_payload(&db, 0, id, base_off), 0); + + /* field 0: 1 -> 999. Committed (durable) but NEVER re-pointed: nothing + after wo_wal_commit below runs — this IS the crash. */ + T_EQ(wo_wal_append_delta(&w, &db, 0, id, 0, base_off, 999), 0); + T_EQ(wo_wal_commit(&w), 0); + + wo_wal_close(&w); + wo_db_destroy(&db); + + /* a fresh process, with no memory of this one's (never re-pointed) map */ + wo_db db2; + T_EQ(wo_db_init(&db2, DELTA_CLASSES, 1, 0, 1), 0); + db2.rt = &rt; + rt.wal = NULL; /* see the replay comment in test_delta_chain_replay */ + T_EQ(wo_wal_replay(path, &db2), 2); /* insert + delta */ + wo_wal w2; + T_EQ(wo_wal_open(&w2, path, 1 << 16), 0); + rt.wal = &w2; rt.db = &db2; + + db_row *r = wo_row_borrow(&db2, 0, id, &msg); + T_CHECK(r != NULL); + T_EQ(r->slots[0], 999); /* the update IS present */ + T_EQ(r->slots[1], 2); + T_EQ(r->slots[2], 3); + wo_row_release(&db2, 0, r); + + wo_wal_close(&w2); + wo_db_destroy(&db2); + wo_rt_destroy(&rt); +} + static void test_torn_tail(void) { char path[128]; snprintf(path, sizeof path, "%s/torn.wal", g_dir); @@ -1877,6 +2124,9 @@ int main(void) { test_keys_resident_unique_clash_pending_repoint(); test_keys_resident_replay(); test_keys_resident_survives_compaction(); + test_delta_chain_replay(); + test_delta_chain_compact_flattens(); + test_delta_crash_window(); test_keys_resident_delete(); test_keys_resident_delete_then_replay(); test_stale_compact_temp_is_removed();