feat(db2-delta): replay and compaction fold delta chains
- 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) <noreply@anthropic.com> (cherry picked from commit 7e4ae70d7beb2a2e396a002184bc06f41253decb)
This commit is contained in:
parent
1f402ae4bb
commit
f138ac0abe
2 changed files with 382 additions and 6 deletions
|
|
@ -705,6 +705,45 @@ static int copy_record(wo_wal *nw, wo_wal *ow, uint64_t off, uint32_t want_cid,
|
||||||
return rc;
|
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) {
|
int wo_wal_compact(wo_wal *w, wo_db *db) {
|
||||||
/* staged records would be written into a file about to be replaced */
|
/* staged records would be written into a file about to be replaced */
|
||||||
if (!w->path || w->len != 0) return -1;
|
if (!w->path || w->len != 0) return -1;
|
||||||
|
|
@ -757,7 +796,25 @@ int wo_wal_compact(wo_wal *w, wo_db *db) {
|
||||||
goto fail;
|
goto fail;
|
||||||
}
|
}
|
||||||
uint64_t at = nw.off + (uint64_t)nw.len;
|
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 */
|
/* value-only update: cannot rehash, so `cur` stays valid */
|
||||||
if (wo_row_set_offset(db, cid, id, at) != 0) {
|
if (wo_row_set_offset(db, cid, id, at) != 0) {
|
||||||
why = "row vanished from the id map mid-compaction";
|
why = "row vanished from the id map mid-compaction";
|
||||||
|
|
@ -822,6 +879,72 @@ fail:
|
||||||
return -1; /* the live log is untouched and still usable */
|
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) {
|
static int apply_record(wo_db *db, const uint8_t *payload, uint32_t len) {
|
||||||
rbuf r = {payload, payload + len, 0};
|
rbuf r = {payload, payload + len, 0};
|
||||||
uint8_t kind = rd_u8(&r);
|
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. */
|
* any. -2 so the caller can say which of the two it is. */
|
||||||
if (db->classes[cid].flags & WO_CLASSF_VOLATILE) return -2;
|
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_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_INSERT && kind != WO_WAL_UPDATE) return -1;
|
||||||
if (kind == WO_WAL_UPDATE) {
|
if (kind == WO_WAL_UPDATE) {
|
||||||
/* replace: the row must exist (its insert precedes its update in a
|
/* 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;
|
int64_t applied = 0;
|
||||||
|
|
||||||
/* databasev2 2: replay must be able to READ rows back, not only write them.
|
/* 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
|
* A REMOVE (or an UPDATE or a DELTA, both of which replay as
|
||||||
* keys-resident table reaches wo_row_remove, whose keys arm borrows the row
|
* remove-then-recreate — see apply_delta, Task 5) on a keys-resident
|
||||||
* out of the log to find its index entries — and a borrow reads through
|
* 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
|
* 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
|
* and assigns rt.wal afterwards, so the borrow found no log, the remove
|
||||||
* failed, and replay reported a perfectly good tombstone as CORRUPTION.
|
* 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;
|
if (rc == -2 && volatile_cid) *volatile_cid = rec_cid;
|
||||||
REPLAY_RETURN(rc == -2 ? -2 : -1);
|
REPLAY_RETURN(rc == -2 ? -2 : -1);
|
||||||
}
|
}
|
||||||
if ((rec_kind == WO_WAL_INSERT || rec_kind == WO_WAL_UPDATE) && rec_id &&
|
if ((rec_kind == WO_WAL_INSERT || rec_kind == WO_WAL_UPDATE ||
|
||||||
wo_table_is_keys_resident(db, rec_cid))
|
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);
|
(void)wo_row_drop_payload(db, rec_cid, rec_id, off);
|
||||||
off += 8u + len + 4u;
|
off += 8u + len + 4u;
|
||||||
applied++;
|
applied++;
|
||||||
|
|
|
||||||
|
|
@ -1445,6 +1445,253 @@ static void test_keys_resident_replay(void) {
|
||||||
wo_rt_destroy(&rt);
|
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) {
|
static void test_torn_tail(void) {
|
||||||
char path[128];
|
char path[128];
|
||||||
snprintf(path, sizeof path, "%s/torn.wal", g_dir);
|
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_unique_clash_pending_repoint();
|
||||||
test_keys_resident_replay();
|
test_keys_resident_replay();
|
||||||
test_keys_resident_survives_compaction();
|
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();
|
||||||
test_keys_resident_delete_then_replay();
|
test_keys_resident_delete_then_replay();
|
||||||
test_stale_compact_temp_is_removed();
|
test_stale_compact_temp_is_removed();
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue