From 91a9dfba6285a0509906c9b52338bc2b91097f89 Mon Sep 17 00:00:00 2001 From: "shoney.arickathil" Date: Sat, 15 Aug 2026 11:01:41 +0200 Subject: [PATCH] feat(database): typed WAL + boot replay (iteration 9, Task 2) - database/src/wal.{c,h}: framed records len|crc32|payload|mark ("WOL1" written last -- no mark, no record), typed-row payloads walking the class-table kinds (nested records, containers, nil encodings), little-endian like the loader - commit order verbatim from the shipped phase-D pattern: RAM apply, stage, ONE pwrite + ONE fdatasync for the batch, ack after -- group commit is everything staged riding one sync - replay decodes straight into engine-owned values (no VM at boot) and re-enters rows through the choke-point row API, so Task 4's indexes will rebuild for free; next_id advances past replayed ids this shard owns (wo_row_create_raw) - torn tail = short/CRC-fail/no-mark/zero-len: intact prefix applies, tear dropped whole, wo_wal_open positions AT the tear so the next commit overwrites it; CRC-valid-but-undecodable = corruption, loud - wo_wal_check: offline oracle, no engine needed -- the crash battery's verifier - test_wal 90/0 ASan+UBSan incl. five crash-battery rounds (fork, insert/commit/ack-over-pipe, SIGKILL mid-stream: zero acked-but- missing, zero acked-but-wrong); all runtime suites green, oop-e2e 71/0; binding doc WAL section + CODE-LOGIC + plan Task 2 checked Co-Authored-By: Claude Opus 5 (1M context) --- database/src/CODE-LOGIC.md | 13 + database/src/table.c | 28 ++ database/src/table.h | 11 + database/src/wal.c | 445 ++++++++++++++++++ database/src/wal.h | 87 ++++ docs/plan/oop-vm/04-db-binding.md | 44 +- .../plans/2026-08-01-db-engine-binding.md | 18 +- runtime/test/test_wal.c | 247 ++++++++++ 8 files changed, 885 insertions(+), 8 deletions(-) create mode 100644 database/src/wal.c create mode 100644 database/src/wal.h create mode 100644 runtime/test/test_wal.c diff --git a/database/src/CODE-LOGIC.md b/database/src/CODE-LOGIC.md index 1836786..cab777f 100644 --- a/database/src/CODE-LOGIC.md +++ b/database/src/CODE-LOGIC.md @@ -30,6 +30,19 @@ VM values ──copy──▶ row slots (engine-owned malloc) ──copy── has no context parameter). One process, one class table; revisit at iteration 8 (shards share the same immutable table). +## wal.c — durability (iteration 9, Task 2) + +The commit order IS the module: RAM apply → stage → one pwrite + one +fdatasync → ack. `wo_wal_commit` returning 0 is the only thing "durable" +means. Replay never touches the VM heap — payloads decode straight into +engine-owned values and re-enter through the row API, so whatever hooks the +choke points (indexes, Task 4) applies to replayed rows identically. Torn +tails end the intact prefix and get overwritten by the next commit; +CRC-valid-but-undecodable records fail replay loudly (corruption is not a +tear). The crash battery in `runtime/test/test_wal.c` is the module's +meaning proven: acked-over-a-pipe after commit, SIGKILL mid-stream, replay, +zero acked-but-missing. + ## Verifying a change - `make -C runtime test` — `test_table` is this directory's suite (round diff --git a/database/src/table.c b/database/src/table.c index caef7eb..a9217b9 100644 --- a/database/src/table.c +++ b/database/src/table.c @@ -427,6 +427,34 @@ int wo_row_read(wo_db *db, wo_rt *rt, uint32_t class_id, uint64_t id, return 0; } +db_row *wo_row_create_raw(wo_db *db, uint32_t class_id, uint64_t id) { + db_table *t = table_of(db, class_id); + if (!t || !id) return NULL; + if (hget(t, id)) return NULL; /* duplicate id: corruption, not a tear */ + uint32_t g = slot_alloc(t); + if (g == UINT32_MAX) return NULL; + db_row *r = slot_row(t, g); + r->id = id; + r->class_id = class_id; + r->flags = 0; + memset(r->slots, 0, t->row_size - sizeof(db_row)); + if (hput(t, id, (uint64_t)g + 1) != 0) return NULL; + t->bitmap[g >> 6] |= 1ull << (g & 63); + t->count++; + /* keep the interleave: only ids this shard owns move its counter */ + if ((id - 1) % db->nshards == db->shard && id >= t->next_id) + t->next_id = id + db->nshards; + /* INDEX HOOK (Task 4): replayed rows re-index here, same as inserts — + the caller fills slots BEFORE indexes exist on them (Task 4 will move + the hook to a post-fill call, recorded in the binding doc). */ + return r; +} + +void wo_db_val_free(wo_db *db, uint8_t kind, uint64_t v) { + (void)db; + db_val_free(kind, v); +} + int wo_row_remove(wo_db *db, uint32_t class_id, uint64_t id) { if (class_id >= db->class_cnt) return -1; db_table *t = &db->tables[class_id]; diff --git a/database/src/table.h b/database/src/table.h index 2b98120..78588f2 100644 --- a/database/src/table.h +++ b/database/src/table.h @@ -129,4 +129,15 @@ int wo_row_remove(wo_db *db, uint32_t class_id, uint64_t id); * to the VM. */ db_row *wo_row_ptr(wo_db *db, uint32_t class_id, uint64_t id); +/* Engine-internal, for WAL replay only: create a row with a FIXED id, + * slots zeroed — the caller (wal.c) fills them with engine-encoded values + * it built while decoding. Advances the table's next_id past [id] when the + * id belongs to this shard, so post-replay inserts never collide. NULL = + * OOM or duplicate id (corruption beyond a torn tail). */ +db_row *wo_row_create_raw(wo_db *db, uint32_t class_id, uint64_t id); + +/* Engine-internal: free one engine-encoded slot value of [kind] (wal.c's + * decode error paths). */ +void wo_db_val_free(wo_db *db, uint8_t kind, uint64_t v); + #endif /* WO_TABLE_H */ diff --git a/database/src/wal.c b/database/src/wal.c new file mode 100644 index 0000000..3bea4df --- /dev/null +++ b/database/src/wal.c @@ -0,0 +1,445 @@ +/* pread/pwrite/fdatasync/posix_fallocate under -std=c11 */ +#define _POSIX_C_SOURCE 200809L + +#include "wal.h" + +#include +#include +#include +#include +#include + +/* ---- crc32 (poly 0xEDB88320) — ported from runtime/wo-rt.c ------------- */ + +static uint32_t crc_table[256]; +static int crc_ready; + +static void crc32_init(void) { + for (uint32_t i = 0; i < 256; i++) { + uint32_t c = i; + for (int k = 0; k < 8; k++) c = (c & 1) ? 0xEDB88320u ^ (c >> 1) : c >> 1; + crc_table[i] = c; + } + crc_ready = 1; +} + +static uint32_t crc32(const void *buf, size_t len) { + if (!crc_ready) crc32_init(); + const uint8_t *p = buf; + uint32_t c = 0xFFFFFFFFu; + while (len--) c = crc_table[(c ^ *p++) & 0xFF] ^ (c >> 8); + return c ^ 0xFFFFFFFFu; +} + +/* ---- byte buffer -------------------------------------------------------- */ + +typedef struct { + uint8_t *b; + size_t len, cap; + int oom; +} wbuf; + +static void wput(wbuf *w, const void *p, size_t n) { + if (w->oom) return; + if (w->len + n > w->cap) { + size_t nc = w->cap ? w->cap * 2 : 256; + while (nc < w->len + n) nc *= 2; + uint8_t *nb = realloc(w->b, nc); + if (!nb) { + w->oom = 1; + return; + } + w->b = nb; + w->cap = nc; + } + memcpy(w->b + w->len, p, n); + w->len += n; +} + +static void wput_u8(wbuf *w, uint8_t v) { wput(w, &v, 1); } +static void wput_u32(wbuf *w, uint32_t v) { wput(w, &v, 4); } +static void wput_u64(wbuf *w, uint64_t v) { wput(w, &v, 8); } + +/* bounds-checked reader */ +typedef struct { + const uint8_t *p, *end; + int bad; +} rbuf; + +static int rtake(rbuf *r, void *out, size_t n) { + if (r->bad || (size_t)(r->end - r->p) < n) { + r->bad = 1; + return -1; + } + memcpy(out, r->p, n); + r->p += n; + return 0; +} + +static uint8_t rd_u8(rbuf *r) { + uint8_t v = 0; + rtake(r, &v, 1); + return v; +} +static uint32_t rd_u32(rbuf *r) { + uint32_t v = 0; + rtake(r, &v, 4); + return v; +} +static uint64_t rd_u64(rbuf *r) { + uint64_t v = 0; + rtake(r, &v, 8); + return v; +} + +#define WAL_NIL_TEXT 0xFFFFFFFFu + +/* ---- engine-value <-> bytes (kind-driven, mirrors table.c's encoding) --- */ + +static void enc_val(wbuf *w, const wo_classdesc *classes, uint8_t kind, uint64_t v) { + switch (kind) { + case WO_K_SCALAR: wput_u64(w, v); return; + case WO_K_TEXT: { + if (!v) { + wput_u32(w, WAL_NIL_TEXT); + return; + } + const db_text *t = (const db_text *)(uintptr_t)v; + wput_u32(w, t->len); + wput(w, t->bytes, t->len); + return; + } + case WO_K_OWNED: { + if (!v) { + wput_u8(w, 0); + return; + } + const db_rec *r = (const db_rec *)(uintptr_t)v; + wput_u8(w, 1); + wput_u32(w, r->class_id); + const wo_classdesc *c = &classes[r->class_id]; + for (uint32_t i = 0; i < c->field_cnt; i++) + enc_val(w, classes, c->kinds[i], r->slots[i]); + return; + } + case WO_K_MULTI: { + if (!v) { + wput_u8(w, 0); + return; + } + const db_multi *m = (const db_multi *)(uintptr_t)v; + wput_u8(w, 1); + wput_u8(w, m->elem_kind); + wput_u32(w, m->len); + for (uint32_t i = 0; i < m->len; i++) enc_val(w, classes, m->elem_kind, m->items[i]); + return; + } + case WO_K_MAP: { + if (!v) { + wput_u8(w, 0); + return; + } + const db_map *m = (const db_map *)(uintptr_t)v; + wput_u8(w, 1); + wput_u8(w, m->key_kind); + wput_u8(w, m->val_kind); + wput_u32(w, m->len); + for (uint32_t i = 0; i < m->len; i++) { + enc_val(w, classes, m->key_kind, m->kv[2 * i]); + enc_val(w, classes, m->val_kind, m->kv[2 * i + 1]); + } + return; + } + default: return; /* GCREF never stored, so never logged */ + } +} + +/* Decode one value into an engine-owned allocation. Returns 0 on success + * with *out set (0 = genuine nil); -1 on truncation/corruption/OOM — the + * caller frees what it already built. */ +static int dec_val(rbuf *r, wo_db *db, uint8_t kind, uint64_t *out) { + *out = 0; + switch (kind) { + case WO_K_SCALAR: { + uint64_t v = rd_u64(r); + if (r->bad) return -1; + *out = v; + return 0; + } + case WO_K_TEXT: { + uint32_t len = rd_u32(r); + if (r->bad) return -1; + if (len == WAL_NIL_TEXT) return 0; + if ((size_t)(r->end - r->p) < len) return -1; + db_text *t = malloc(sizeof(db_text) + len); + if (!t) return -1; + t->len = len; + memcpy(t->bytes, r->p, len); + r->p += len; + *out = (uint64_t)(uintptr_t)t; + return 0; + } + case WO_K_OWNED: { + uint8_t tag = rd_u8(r); + if (r->bad) return -1; + if (!tag) return 0; + uint32_t cid = rd_u32(r); + if (r->bad || cid >= db->class_cnt) return -1; + const wo_classdesc *c = &db->classes[cid]; + db_rec *rec = malloc(sizeof(db_rec) + (size_t)c->field_cnt * 8u); + if (!rec) return -1; + rec->class_id = cid; + rec->_pad = 0; + for (uint32_t i = 0; i < c->field_cnt; i++) { + if (dec_val(r, db, c->kinds[i], &rec->slots[i]) != 0) { + for (uint32_t j = 0; j < i; j++) wo_db_val_free(db, c->kinds[j], rec->slots[j]); + free(rec); + return -1; + } + } + *out = (uint64_t)(uintptr_t)rec; + return 0; + } + case WO_K_MULTI: { + uint8_t tag = rd_u8(r); + if (r->bad) return -1; + if (!tag) return 0; + uint8_t ek = rd_u8(r); + uint32_t len = rd_u32(r); + if (r->bad || ek > WO_K_MAX) return -1; + if (len > (size_t)(r->end - r->p)) return -1; /* each elem >= 1 byte */ + db_multi *m = malloc(sizeof(db_multi) + (size_t)len * 8u); + if (!m) return -1; + m->elem_kind = ek; + m->len = len; + for (uint32_t i = 0; i < len; i++) { + if (dec_val(r, db, ek, &m->items[i]) != 0) { + for (uint32_t j = 0; j < i; j++) wo_db_val_free(db, ek, m->items[j]); + free(m); + return -1; + } + } + *out = (uint64_t)(uintptr_t)m; + return 0; + } + case WO_K_MAP: { + uint8_t tag = rd_u8(r); + if (r->bad) return -1; + if (!tag) return 0; + uint8_t kk = rd_u8(r), vk = rd_u8(r); + uint32_t len = rd_u32(r); + if (r->bad || kk > WO_K_MAX || vk > WO_K_MAX) return -1; + if (len > (size_t)(r->end - r->p)) return -1; + db_map *m = malloc(sizeof(db_map) + (size_t)len * 16u); + if (!m) return -1; + m->key_kind = kk; + m->val_kind = vk; + m->len = len; + for (uint32_t i = 0; i < len; i++) { + if (dec_val(r, db, kk, &m->kv[2 * i]) != 0 || + dec_val(r, db, vk, &m->kv[2 * i + 1]) != 0) { + m->len = i; /* free only the fully-built pairs plus a possible key */ + for (uint32_t j = 0; j < i; j++) { + wo_db_val_free(db, kk, m->kv[2 * j]); + wo_db_val_free(db, vk, m->kv[2 * j + 1]); + } + wo_db_val_free(db, kk, m->kv[2 * i]); /* 0 if the key failed */ + free(m); + return -1; + } + } + *out = (uint64_t)(uintptr_t)m; + return 0; + } + default: return -1; + } +} + +/* ---- record scan (shared by open, replay, check) ------------------------ */ + +/* Read the record at [off]. 0 = intact (*len_out = payload length, payload + * malloc'd into *payload_out if non-NULL); 1 = end of intact prefix (zero + * length, short read, bad crc, missing mark). */ +static int scan_record(int fd, uint64_t off, uint32_t *len_out, uint8_t **payload_out) { + uint8_t hdr[8]; + ssize_t n = pread(fd, hdr, 8, (off_t)off); + if (n != 8) return 1; + uint32_t len, crc; + memcpy(&len, hdr, 4); + memcpy(&crc, hdr + 4, 4); + if (len == 0 || len > (64u << 20)) return 1; /* preallocated tail or garbage */ + uint8_t *payload = malloc(len + 4); + if (!payload) return 1; + n = pread(fd, payload, len + 4, (off_t)(off + 8)); + if (n != (ssize_t)(len + 4)) { + free(payload); + return 1; + } + uint32_t mark; + memcpy(&mark, payload + len, 4); + if (mark != WO_WAL_MARK || crc32(payload, len) != crc) { + free(payload); + return 1; + } + *len_out = len; + if (payload_out) *payload_out = payload; + else free(payload); + return 0; +} + +/* ---- public API ---------------------------------------------------------- */ + +int wo_wal_open(wo_wal *w, const char *path, uint64_t prealloc) { + memset(w, 0, sizeof(*w)); + w->fd = open(path, O_RDWR | O_CREAT, 0644); + if (w->fd < 0) return -1; + if (prealloc) { + /* best-effort: a filesystem without fallocate still works */ + (void)posix_fallocate(w->fd, 0, (off_t)prealloc); + } + /* position after the intact prefix: a torn tail is OVERWRITTEN by the + * next append, never appended after */ + uint64_t off = 0; + uint32_t len; + while (scan_record(w->fd, off, &len, NULL) == 0) off += 8u + len + 4u; + w->off = off; + return 0; +} + +void wo_wal_close(wo_wal *w) { + if (w->fd >= 0) close(w->fd); + free(w->buf); + memset(w, 0, sizeof(*w)); + w->fd = -1; +} + +/* frame one payload into the staged batch */ +static int stage(wo_wal *w, const wbuf *payload) { + if (payload->oom) return -1; + wbuf rec = {0}; + wput_u32(&rec, (uint32_t)payload->len); + wput_u32(&rec, crc32(payload->b, payload->len)); + wput(&rec, payload->b, payload->len); + wput_u32(&rec, WO_WAL_MARK); + if (rec.oom) { + free(rec.b); + return -1; + } + if (w->len + rec.len > w->cap) { + size_t nc = w->cap ? w->cap * 2 : 4096; + while (nc < w->len + rec.len) nc *= 2; + uint8_t *nb = realloc(w->buf, nc); + if (!nb) { + free(rec.b); + return -1; + } + w->buf = nb; + w->cap = nc; + } + memcpy(w->buf + w->len, rec.b, rec.len); + w->len += rec.len; + free(rec.b); + return 0; +} + +int wo_wal_append_insert(wo_wal *w, wo_db *db, uint32_t class_id, uint64_t id) { + db_row *r = wo_row_ptr(db, class_id, id); + if (!r) return -1; /* commit order: RAM apply comes FIRST */ + wbuf p = {0}; + wput_u8(&p, WO_WAL_INSERT); + wput_u32(&p, class_id); + wput_u64(&p, id); + const wo_classdesc *c = &db->classes[class_id]; + for (uint32_t i = 0; i < c->field_cnt; i++) enc_val(&p, db->classes, c->kinds[i], r->slots[i]); + int rc = stage(w, &p); + free(p.b); + return rc; +} + +int wo_wal_append_remove(wo_wal *w, uint32_t class_id, uint64_t id) { + wbuf p = {0}; + wput_u8(&p, WO_WAL_REMOVE); + wput_u32(&p, class_id); + wput_u64(&p, id); + int rc = stage(w, &p); + free(p.b); + return rc; +} + +int wo_wal_commit(wo_wal *w) { + if (!w->len) return 0; + size_t at = 0; + while (at < w->len) { + ssize_t n = pwrite(w->fd, w->buf + at, w->len - at, (off_t)(w->off + at)); + if (n < 0) { + if (errno == EINTR) continue; + return -1; + } + at += (size_t)n; + } + if (fdatasync(w->fd) != 0) return -1; + w->off += w->len; + w->len = 0; /* acked: the batch is durable */ + 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); + uint32_t cid = rd_u32(&r); + uint64_t id = rd_u64(&r); + if (r.bad || cid >= db->class_cnt) return -1; + if (kind == WO_WAL_REMOVE) return wo_row_remove(db, cid, id); + if (kind != WO_WAL_INSERT) return -1; /* UPDATE lands with Task 5 */ + db_row *row = wo_row_create_raw(db, cid, id); + if (!row) return -1; + const wo_classdesc *c = &db->classes[cid]; + for (uint32_t i = 0; i < c->field_cnt; i++) { + if (dec_val(&r, db, c->kinds[i], &row->slots[i]) != 0) { + /* a record that CRC-passed but does not decode is corruption, + * not a tear: fail loudly (the row's built slots are freed by + * wo_row_remove, which also unregisters the id) */ + wo_row_remove(db, cid, id); + return -1; + } + } + return (size_t)(r.end - r.p) == 0 ? 0 : -1; /* trailing bytes = corrupt */ +} + +int64_t wo_wal_replay(const char *path, wo_db *db) { + int fd = open(path, O_RDONLY); + if (fd < 0) return errno == ENOENT ? 0 : -1; /* no WAL yet = fresh boot */ + uint64_t off = 0; + int64_t applied = 0; + for (;;) { + uint32_t len; + uint8_t *payload; + if (scan_record(fd, off, &len, &payload) != 0) break; /* intact prefix ends */ + int rc = apply_record(db, payload, len); + free(payload); + if (rc != 0) { + close(fd); + return -1; + } + off += 8u + len + 4u; + applied++; + } + close(fd); + return applied; +} + +int64_t wo_wal_check(const char *path, uint64_t *intact_bytes) { + int fd = open(path, O_RDONLY); + if (fd < 0) return -1; + uint64_t off = 0; + int64_t records = 0; + for (;;) { + uint32_t len; + if (scan_record(fd, off, &len, NULL) != 0) break; + off += 8u + len + 4u; + records++; + } + if (intact_bytes) *intact_bytes = off; + close(fd); + return records; +} diff --git a/database/src/wal.h b/database/src/wal.h new file mode 100644 index 0000000..0a05a47 --- /dev/null +++ b/database/src/wal.h @@ -0,0 +1,87 @@ +/* wal.h — typed-row write-ahead log + boot replay (iteration 9, Task 2). + * + * The c-runtime plan's shipped pattern (phases D/E), generalized to typed + * rows. The commit order is doctrine, verbatim: + * + * RAM apply → wal_append (staged) → wal_commit (write + fdatasync) + * → only then is the write ACKNOWLEDGED + * + * Record framing — replay-whole-or-not-at-all: + * + * record := len u32 | crc u32 | payload | mark u32 + * len = payload byte count (never 0; 0 = preallocated tail, stop) + * crc = CRC32 of payload + * mark = 0x574F4C31 "WOL1" — written LAST, so a record without its + * mark is torn by definition + * payload := kind u8 | class_id u32 | row_id u64 | body + * kind : 1 insert (body = the row's fields, engine encoding below) + * 2 remove (no body) + * 3 update (reserved for Task 5) + * + * Field encoding in a body walks the class table's kinds: + * SCALAR 8 bytes + * TEXT u32 len | bytes (0xFFFFFFFF = nil) + * OWNED u8 0 = nil, or u8 1 | u32 class_id | fields recursively + * MULTI u8 0 = nil, or u8 1 | u8 elem_kind | u32 len | elements + * MAP u8 0 = nil, or u8 1 | u8 kk | u8 vk | u32 len | k v pairs + * + * Replay decodes payloads STRAIGHT into engine-owned values — the VM heap + * is never involved (boot must not depend on a VM existing yet), and rows + * re-enter through the same choke-point row API, so Task 4's indexes are + * rebuilt for free. A torn tail (short record, bad CRC, missing mark) drops + * everything from the tear onward — never a partial record, never a record + * after a tear. Little-endian on-disk, matching the .wob loader's platform + * note. + * + * wo_wal_check is the offline oracle the crash battery verifies with: it + * walks a WAL file with no engine at all and reports how many records are + * intact and where the intact prefix ends. */ +#ifndef WO_WAL_H +#define WO_WAL_H + +#include "table.h" + +#define WO_WAL_MARK 0x574F4C31u /* "WOL1" LE */ + +enum { WO_WAL_INSERT = 1, WO_WAL_REMOVE = 2, WO_WAL_UPDATE = 3 }; + +typedef struct wo_wal { + int fd; + uint64_t off; /* next write offset (the intact tail) */ + /* staged batch: appended by wal_append_*, flushed by wal_commit */ + uint8_t *buf; + size_t len, cap; +} wo_wal; + +/* Open (create if missing) and preallocate [prealloc] bytes (best-effort; + * a filesystem without fallocate still works). Positions the write offset + * at the end of the INTACT record prefix — an existing file is scanned the + * same way replay scans it, so a torn tail is overwritten, not appended + * after. 0 ok, -1 errno-style failure. */ +int wo_wal_open(wo_wal *w, const char *path, uint64_t prealloc); +void wo_wal_close(wo_wal *w); + +/* Stage a record for the row that MUST already be applied to RAM (the + * commit-order doctrine). Insert/update read the row via wo_row_ptr. + * 0 ok, -1 OOM / no such row. */ +int wo_wal_append_insert(wo_wal *w, wo_db *db, uint32_t class_id, uint64_t id); +int wo_wal_append_remove(wo_wal *w, uint32_t class_id, uint64_t id); + +/* Write the staged batch and fdatasync — the ack line. Empty batch = ok, + * no syscall. 0 ok, -1 write/sync failure (the batch stays staged). */ +int wo_wal_commit(wo_wal *w); + +/* Boot replay: apply every intact record to [db] in order. Ids re-enter + * exactly as logged; each table's next_id advances past the replayed ids + * that belong to this shard. Returns the number of records applied, or -1 + * on open failure / a record naming an unknown class (corruption beyond + * what a torn tail explains). A torn tail is NOT an error: replay applies + * the intact prefix and reports it. */ +int64_t wo_wal_replay(const char *path, wo_db *db); + +/* Offline verification (no engine): scan [path], count intact records. + * *intact_bytes (optional) = where the intact prefix ends. -1 = open + * failure. */ +int64_t wo_wal_check(const char *path, uint64_t *intact_bytes); + +#endif /* WO_WAL_H */ diff --git a/docs/plan/oop-vm/04-db-binding.md b/docs/plan/oop-vm/04-db-binding.md index f038545..0dbcfbc 100644 --- a/docs/plan/oop-vm/04-db-binding.md +++ b/docs/plan/oop-vm/04-db-binding.md @@ -68,12 +68,48 @@ table. Task 4's secondary indexes hook exactly these two sites (marked beside the same calls. Anything else touching a slab is a defect by definition — the doctrine the Rust engine learned and this engine enforces. +## WAL (Task 2) — `database/src/wal.{c,h}` + +Record framing, replay-whole-or-not-at-all: + +``` +record := len u32 | crc u32 | payload | mark u32 +len = payload bytes (never 0: a zero length is the preallocated tail) +crc = CRC32 (poly 0xEDB88320) of the payload +mark = 0x574F4C31 "WOL1", the last bytes of the record — a record + without its mark is torn by definition +payload := kind u8 | class_id u32 | row_id u64 | body +kind = 1 insert (body = fields), 2 remove (no body), 3 update (Task 5) +``` + +Body fields walk the class table's kinds: `SCALAR` 8 bytes; `TEXT` u32 len + +bytes (`0xFFFFFFFF` = nil); `OWNED` presence u8 then class id + fields +recursively; `MULTI` presence + elem kind + len + elements; `MAP` presence + +both kinds + len + pairs. Little-endian, same platform note as the loader. + +**Commit order (doctrine, verbatim from the shipped phase-D pattern):** RAM +apply → stage record → `wo_wal_commit` (one pwrite of the batch + one +fdatasync) → only then acknowledge. Group commit = everything staged since +the last commit rides one sync. + +**Replay** decodes payloads straight into engine-owned values — no VM heap +involved, boot cannot depend on a VM existing — and rows re-enter through +the choke-point row API, so Task 4's indexes rebuild for free. A torn tail +(short record, bad CRC, missing mark, zero length) ends the intact prefix: +everything from the tear on is dropped whole, and `wo_wal_open` positions +its write offset AT the tear so the next commit overwrites it. A record that +CRC-passes but does not decode is corruption, not a tear — replay fails +loudly. A missing file is a fresh boot, not an error. After replay each +table's `next_id` sits past every replayed id this shard owns. + +**Oracle:** `wo_wal_check(path)` walks a file with no engine and reports the +intact record count and prefix end — the crash battery's verifier +(`runtime/test/test_wal.c`: five rounds of insert/commit/ack-over-pipe with +SIGKILL mid-stream; every acked row present and exact after replay). + ## Still to come in this document -- **Task 2**: WAL record framing (`length | crc | payload | commit-mark`), - payload encoding for typed rows, group-commit ordering, replay rules, - torn-tail handling, the `wal-check` oracle. - **Task 3**: the `insert` statement's builtin ids (appended to `00-wob-format.md`'s builtin table) and execution contract. - **Task 4**: secondary-index format, `@unique` trap code. -- **Task 5**: the select subset and its builtins. +- **Task 5**: the select subset, update record semantics, and its builtins. diff --git a/docs/superpowers/plans/2026-08-01-db-engine-binding.md b/docs/superpowers/plans/2026-08-01-db-engine-binding.md index 97acc6d..53b9eb4 100644 --- a/docs/superpowers/plans/2026-08-01-db-engine-binding.md +++ b/docs/superpowers/plans/2026-08-01-db-engine-binding.md @@ -1,6 +1,6 @@ # DB Engine Binding Implementation Plan -> **Status: 🔄 in progress — Task 1 done 2026-08-15** (story iteration 9) — class-shaped tables, typed WAL + recovery, `insert`/`select` execution. Story iteration 9b (`@table` relations + language-integrated query) follows it and needs a spec brainstormed first. Board: [00-status.md](../../00-status.md) +> **Status: 🔄 in progress — Tasks 1–2 done 2026-08-15** (story iteration 9) — class-shaped tables, typed WAL + recovery, `insert`/`select` execution. Story iteration 9b (`@table` relations + language-integrated query) follows it and needs a spec brainstormed first. Board: [00-status.md](../../00-status.md) > **For agentic workers:** REQUIRED SUB-SKILL: Use superpowers:subagent-driven-development (recommended) or superpowers:executing-plans to implement this plan task-by-task. Steps use checkbox (`- [ ]`) syntax for tracking. > @@ -67,9 +67,19 @@ sanitizers included. `database/` gets its own CODE-LOGIC.md as code lands. **Concept & reason:** port phase D/E to typed rows. Records frame `length | crc | payload | commit-mark` (replay-whole-or-not-at-all); payload = record kind (insert/remove/update), class id, row id, encoded fields. Per-shard WAL files, fallocate-preallocated, appended on commit AFTER the RAM apply, one fdatasync covering the tick's batch, ack after the sync completes — the shipped commit order, verbatim. Boot: per-shard parallel replay before listeners open; torn tails detected and dropped whole (phase E behavior). The offline `wal-check` verification mode ports too — it is the crash test's oracle. -- [ ] Failing tests: commit-then-kill crash battery (the phase-D test shape: concurrent writes, SIGKILL mid-stream, offline verification proves every acked write present and CRC-valid, zero acked-but-missing); torn-tail drop; parallel replay rebuilds identical RAM state (deep-compare against pre-crash snapshot dump). -- [ ] Implement; green. -- [ ] Record commit draft: `feat(runtime): typed-row WAL + replay — framed CRC records over class rows, group-commit ack-after-fsync, parallel boot replay with torn-tail drop, wal-check oracle; crash battery green.` +- [x] Tests first: replay round-trip with deep compare (Texts included) and + next_id advance; torn-tail drop with the oracle agreeing on the intact + prefix byte-for-byte; reopen positions AT the tear so the next commit + overwrites it; crash battery — five rounds of fork + insert/commit/ + ack-over-pipe + SIGKILL mid-stream, zero acked-but-missing, zero + acked-but-wrong, contents exact. `test_wal` 90/0 under ASan+UBSan. +- [x] Implemented in `database/src/wal.{c,h}`: framed CRC records + (`len|crc|payload|mark`), typed-row payloads decoded straight into + engine values (no VM at boot), group commit as one pwrite + one + fdatasync, replay through the choke-point row API (indexes rebuild for + free when Task 4 lands), `wo_wal_check` offline oracle. "Parallel" + replay is per-shard by construction and runs at N=1 until iteration 8. +- [x] Committed locally (2026-08-15). Binding doc's WAL section written. ### Task 3: `insert` executes diff --git a/runtime/test/test_wal.c b/runtime/test/test_wal.c new file mode 100644 index 0000000..701c46a --- /dev/null +++ b/runtime/test/test_wal.c @@ -0,0 +1,247 @@ +/* test_wal — iteration 9 Task 2: typed WAL + boot replay. + * Round-trip through a replay, torn-tail drop, reopen-overwrites-tear, + * and the commit-then-kill crash battery: a forked child inserts rows and + * acks each COMMITTED id over a pipe; SIGKILL lands mid-stream; the parent + * verifies with the offline oracle and a replay that every acked id is + * present with the right contents. */ +#define _POSIX_C_SOURCE 200809L + +#include +#include +#include +#include +#include +#include +#include + +#include "obj.h" +#include "t.h" +#include "table.h" +#include "wal.h" + +/* class 0: Row { n: scalar, label: Text } */ +static const uint8_t row_kinds[] = {WO_K_SCALAR, WO_K_TEXT}; +static const wo_classdesc CLASSES[] = { + {.name = 0, .flags = 0, .field_cnt = 2, .kinds = row_kinds}, +}; + +static char g_dir[64]; + +static void test_roundtrip_replay(void) { + char path[128]; + snprintf(path, sizeof path, "%s/basic.wal", g_dir); + wo_rt rt; + T_EQ(wo_rt_init(&rt, 1 << 20, CLASSES, 1), 0); + wo_db db; + T_EQ(wo_db_init(&db, CLASSES, 1, 0, 1), 0); + wo_wal w; + T_EQ(wo_wal_open(&w, path, 1 << 16), 0); + const char *msg = ""; + + /* three inserts and one remove, RAM first, WAL second, one commit */ + uint64_t ids[3]; + for (int i = 0; i < 3; i++) { + wo_str *s = wo_str_new(&rt, "abcXYZ" + i, 3); /* "abc","bcX","cXY" */ + uint64_t vals[2] = {(uint64_t)(i * 10), (uint64_t)(uintptr_t)s}; + ids[i] = wo_row_insert(&db, 0, vals, &msg); + T_CHECK(ids[i] != 0); + T_EQ(wo_wal_append_insert(&w, &db, 0, ids[i]), 0); + wo_str_free(&rt, s); + } + T_EQ(wo_row_remove(&db, 0, ids[1]), 0); + T_EQ(wo_wal_append_remove(&w, 0, ids[1]), 0); + T_EQ(wo_wal_commit(&w), 0); + wo_wal_close(&w); + wo_db_destroy(&db); + + /* boot: fresh engine, replay, deep-compare */ + wo_db db2; + T_EQ(wo_db_init(&db2, CLASSES, 1, 0, 1), 0); + T_EQ(wo_wal_replay(path, &db2), 4); + uint64_t out[2]; + T_EQ(wo_row_read(&db2, &rt, 0, ids[0], out, &msg), 0); + T_EQ(out[0], 0); + wo_str *s0 = (wo_str *)(uintptr_t)out[1]; + T_CHECK(s0->len == 3 && memcmp(s0->data, "abc", 3) == 0); + wo_str_free(&rt, s0); + T_EQ(wo_row_read(&db2, &rt, 0, ids[1], out, &msg), -1); /* removed */ + T_EQ(wo_row_read(&db2, &rt, 0, ids[2], out, &msg), 0); + T_EQ(out[0], 20); + wo_str_free(&rt, (wo_str *)(uintptr_t)out[1]); + /* next_id advanced past the replayed ids: a fresh insert never collides */ + uint64_t vals[2] = {99, 0}; + uint64_t fresh = wo_row_insert(&db2, 0, vals, &msg); + T_CHECK(fresh > ids[2]); + wo_db_destroy(&db2); + wo_rt_destroy(&rt); + + /* replay of a missing file is a fresh boot, not an error */ + wo_db db3; + T_EQ(wo_db_init(&db3, CLASSES, 1, 0, 1), 0); + T_EQ(wo_wal_replay("/nonexistent/nope.wal", &db3), 0); + wo_db_destroy(&db3); +} + +static void test_torn_tail(void) { + char path[128]; + snprintf(path, sizeof path, "%s/torn.wal", g_dir); + wo_db db; + T_EQ(wo_db_init(&db, CLASSES, 1, 0, 1), 0); + wo_wal w; + T_EQ(wo_wal_open(&w, path, 0), 0); + const char *msg = ""; + for (int i = 0; i < 5; i++) { + uint64_t vals[2] = {(uint64_t)i, 0}; + uint64_t id = wo_row_insert(&db, 0, vals, &msg); + T_EQ(wo_wal_append_insert(&w, &db, 0, id), 0); + T_EQ(wo_wal_commit(&w), 0); + } + uint64_t intact_end = w.off; + /* tear: append half a record's worth of a valid-looking header + junk */ + uint32_t fake_len = 40, fake_crc = 0xDEAD; + uint8_t junk[20] = {7, 7, 7}; + T_CHECK(pwrite(w.fd, &fake_len, 4, (off_t)intact_end) == 4); + T_CHECK(pwrite(w.fd, &fake_crc, 4, (off_t)(intact_end + 4)) == 4); + T_CHECK(pwrite(w.fd, junk, sizeof junk, (off_t)(intact_end + 8)) == (ssize_t)sizeof junk); + wo_wal_close(&w); + wo_db_destroy(&db); + + /* the oracle sees exactly the intact prefix */ + uint64_t at = 0; + T_EQ(wo_wal_check(path, &at), 5); + T_EQ(at, intact_end); + + /* replay drops the tear whole */ + wo_db db2; + T_EQ(wo_db_init(&db2, CLASSES, 1, 0, 1), 0); + T_EQ(wo_wal_replay(path, &db2), 5); + T_EQ(db2.tables[0].count, 5); + wo_db_destroy(&db2); + + /* reopen positions AT the tear: the next commit overwrites it */ + wo_db db3; + T_EQ(wo_db_init(&db3, CLASSES, 1, 0, 1), 0); + T_EQ(wo_wal_replay(path, &db3), 5); + wo_wal w2; + T_EQ(wo_wal_open(&w2, path, 0), 0); + T_EQ(w2.off, intact_end); + uint64_t vals[2] = {100, 0}; + uint64_t id = wo_row_insert(&db3, 0, vals, &msg); + T_EQ(wo_wal_append_insert(&w2, &db3, 0, id), 0); + T_EQ(wo_wal_commit(&w2), 0); + wo_wal_close(&w2); + T_EQ(wo_wal_check(path, NULL), 6); /* tear gone, record in its place */ + wo_db_destroy(&db3); +} + +/* ---- the crash battery -------------------------------------------------- */ + +/* Child: insert forever — RAM, WAL, COMMIT, and only then ack the id down + * the pipe. Killed mid-stream by the parent. */ +static void battery_child(const char *path, int ack_fd) { + wo_rt rt; + wo_db db; + wo_wal w; + if (wo_rt_init(&rt, 1 << 20, CLASSES, 1) != 0) _exit(9); + if (wo_db_init(&db, CLASSES, 1, 0, 1) != 0) _exit(9); + if (wo_wal_open(&w, path, 1 << 20) != 0) _exit(9); + const char *msg = ""; + for (uint64_t i = 0;; i++) { + char label[32]; + int n = snprintf(label, sizeof label, "row-%llu", (unsigned long long)i); + wo_str *s = wo_str_new(&rt, label, (uint32_t)n); + uint64_t vals[2] = {i * 3 + 1, (uint64_t)(uintptr_t)s}; + uint64_t id = wo_row_insert(&db, 0, vals, &msg); + wo_str_free(&rt, s); + if (!id) _exit(9); + if (wo_wal_append_insert(&w, &db, 0, id) != 0) _exit(9); + if (wo_wal_commit(&w) != 0) _exit(9); /* durable BEFORE the ack */ + ssize_t wr = write(ack_fd, &id, 8); + if (wr != 8) _exit(0); /* parent went away */ + } +} + +static void test_crash_battery(void) { + for (int round = 0; round < 5; round++) { + char path[128]; + snprintf(path, sizeof path, "%s/crash-%d.wal", g_dir, round); + int pipefd[2]; + T_EQ(pipe(pipefd), 0); + pid_t pid = fork(); + T_CHECK(pid >= 0); + if (pid == 0) { + close(pipefd[0]); + battery_child(path, pipefd[1]); + _exit(0); + } + close(pipefd[1]); + /* collect acks for a few ms, then kill mid-stream — no sync with + the child's commit loop, which is the point */ + struct timespec ts = {0, (20 + round * 13) * 1000000L}; + while (nanosleep(&ts, &ts) != 0) {} + kill(pid, SIGKILL); + int status; + waitpid(pid, &status, 0); + /* drain every ack that made it into the pipe */ + uint64_t acked[65536]; + size_t n_acked = 0; + for (;;) { + uint64_t id; + ssize_t n = read(pipefd[0], &id, 8); + if (n != 8) break; + if (n_acked < 65536) acked[n_acked++] = id; + } + close(pipefd[0]); + T_CHECK(n_acked > 0); /* the child got at least one commit out */ + + /* offline oracle: the file's intact prefix covers every ack */ + int64_t intact = wo_wal_check(path, NULL); + T_CHECK(intact >= (int64_t)n_acked); + + /* replay and verify: every acked id present, contents exact */ + wo_rt rt; + T_EQ(wo_rt_init(&rt, 1 << 22, CLASSES, 1), 0); + wo_db db; + T_EQ(wo_db_init(&db, CLASSES, 1, 0, 1), 0); + int64_t applied = wo_wal_replay(path, &db); + T_CHECK(applied >= (int64_t)n_acked); + const char *msg = ""; + int bad = 0; + for (size_t i = 0; i < n_acked; i++) { + uint64_t out[2]; + if (wo_row_read(&db, &rt, 0, acked[i], out, &msg) != 0) { + bad++; + continue; + } + /* id = i+1 (shard 0 of 1), field 0 = i*3+1, label = "row-i" */ + char want[32]; + int wl = snprintf(want, sizeof want, "row-%llu", + (unsigned long long)(acked[i] - 1)); + wo_str *s = (wo_str *)(uintptr_t)out[1]; + if (out[0] != (acked[i] - 1) * 3 + 1 || s->len != (uint32_t)wl || + memcmp(s->data, want, (size_t)wl) != 0) + bad++; + wo_str_free(&rt, s); + } + T_EQ(bad, 0); /* zero acked-but-missing, zero acked-but-wrong */ + wo_db_destroy(&db); + wo_rt_destroy(&rt); + } +} + +int main(void) { + snprintf(g_dir, sizeof g_dir, "/tmp/wo-wal-test-XXXXXX"); + if (!mkdtemp(g_dir)) return 1; + test_roundtrip_replay(); + test_torn_tail(); + test_crash_battery(); + /* leave the dir for a failed run's forensics only */ + if (!t_fail) { + char cmd[128]; + snprintf(cmd, sizeof cmd, "rm -rf %s", g_dir); + if (system(cmd) != 0) fprintf(stderr, "cleanup failed, kept %s\n", g_dir); + } else { + fprintf(stderr, "kept %s\n", g_dir); + } + return t_report("test_wal"); +}