From 76d80cc027e70753fc37a6846e4e8529ce815455 Mon Sep 17 00:00:00 2001 From: "shoney.arickathil" Date: Fri, 28 Aug 2026 09:25:10 +0200 Subject: [PATCH] =?UTF-8?q?feat(db):=20one=20barrier=20per=20drain,=20repl?= =?UTF-8?q?ies=20held=20=E2=80=94=20T2?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit databasev2 4 part A, task 2. The core change, and mostly deletion. - the REQUEST path (wo_db_exec_req) no longer commits after each append. Applying to RAM and staging stay exactly where they were - wo_vm_adopt holds each DB reply envelope in a local FIFO instead of pushing it as the statement finishes. Pushing there would unpark the requester before its record is durable — the ack contract this iteration exists to make literally true rather than true by accident of every batch having one member - at the end of the drain: ONE wo_wal_commit_fatal for everything staged, then every held reply. Locals rather than per-shard state: nothing needs to outlive the batch it describes - "did this statement stage anything" is asked of the buffer, not guessed from the opcode, and that count is what the failure diagnostic reports - the drain commits unconditionally when anything is staged, because the inline path relies on finding the buffer empty (task 3 documents that) - staging failure on the request path is now FATAL via wo_wal_stage_fatal: the row is already in RAM and of the three verbs only insert could undo itself, so continuing means RAM ahead of disk. One rule - wal_die is now shared by both fatal points Verified — the ack contract is the thing that could break, so it is what was tested: - just wovm-test: 36 suites (18 x both dispatch flavors) 0 fail, cli_smoke OK - just db-bench-quick: 85 checks, 0 failures. The legs that matter: crash.sN.0 — 612 acked rows all present after kill -9, which is the BATCHING path (multi-shard requests, held replies, one barrier); crash.s1.0 — 800 acked rows; restart.s1 and restart.sN replay byte-true Co-Authored-By: Claude Opus 5 (1M context) --- database/src/db.c | 23 +++++++---------------- database/src/wal.c | 27 ++++++++++++++++----------- database/src/wal.h | 5 +++++ runtime/src/vm.c | 36 +++++++++++++++++++++++++++++++++++- 4 files changed, 63 insertions(+), 28 deletions(-) diff --git a/database/src/db.c b/database/src/db.c index 58273a3..24a5fc9 100644 --- a/database/src/db.c +++ b/database/src/db.c @@ -216,12 +216,11 @@ void wo_db_exec_req(wo_vm *vm, wo_db_req *q) { break; } if (w) { - if (wo_wal_append_insert(w, db, q->cid, id) != 0 || wo_wal_commit(w) != 0) { - wo_row_remove(db, q->cid, id); - q->status = WO_T_IO; - q->msg = "wal commit failed"; - break; - } + /* databasev2 4: staging failure is FATAL, not a trap. The row is + * already in RAM; of the three verbs only insert could undo + * itself, so continuing means RAM ahead of disk. One rule: once a + * statement has mutated RAM, the outcomes are durable or death. */ + if (wo_wal_append_insert(w, db, q->cid, id) != 0) wo_wal_stage_fatal(w); } q->result = id; break; @@ -234,11 +233,7 @@ void wo_db_exec_req(wo_vm *vm, wo_db_req *q) { break; } if (w) { - if (wo_wal_append_update(w, db, q->cid, q->id) != 0 || wo_wal_commit(w) != 0) { - q->status = WO_T_IO; - q->msg = "wal commit failed"; - break; - } + if (wo_wal_append_update(w, db, q->cid, q->id) != 0) wo_wal_stage_fatal(w); } break; } @@ -254,11 +249,7 @@ void wo_db_exec_req(wo_vm *vm, wo_db_req *q) { break; } if (w) { - if (wo_wal_append_remove(w, q->cid, q->id) != 0 || wo_wal_commit(w) != 0) { - q->status = WO_T_IO; - q->msg = "wal commit failed"; - break; - } + if (wo_wal_append_remove(w, q->cid, q->id) != 0) wo_wal_stage_fatal(w); } break; } diff --git a/database/src/wal.c b/database/src/wal.c index b49fd93..600f2b2 100644 --- a/database/src/wal.c +++ b/database/src/wal.c @@ -409,20 +409,25 @@ int wo_wal_commit(wo_wal *w) { return 0; } +/* Nothing at either fatal point is recoverable: RAM holds changes the log + * does not, and this process can no longer serve reads that would survive a + * restart. Name what failed precisely enough to act on, then stop. */ +static void wal_die(const wo_wal *w, const char *op, uint32_t nrec) { + fprintf(stderr, + "writeonce: DURABILITY FAILURE — %s failed on %s: %s\n" + " %u record(s) were NOT made durable and are not acknowledged.\n" + " The process is stopping: replay restores the last durable state.\n", + op, w->path ? w->path : "(the write-ahead log)", strerror(errno), + nrec); + exit(WO_EXIT_DURABILITY); +} + +void wo_wal_stage_fatal(const wo_wal *w) { wal_die(w, "staging a record", 1); } + void wo_wal_commit_fatal(wo_wal *w, uint32_t nrec) { int rc = wo_wal_commit(w); if (rc == 0) return; - /* Nothing here is recoverable: RAM holds changes the log does not, and - * this process can no longer serve reads that would survive a restart. - * Name what failed precisely enough to act on, then stop. */ - fprintf(stderr, - "writeonce: DURABILITY FAILURE — %s failed on %s: %s\n" - " %u record(s) in the batch were NOT made durable and are not acknowledged.\n" - " The process is stopping: replay restores the last durable state.\n", - rc == WO_WAL_ERR_SYNC ? "fdatasync" : "pwrite", - w->path ? w->path : "(the write-ahead log)", strerror(errno), - nrec); - exit(WO_EXIT_DURABILITY); + wal_die(w, rc == WO_WAL_ERR_SYNC ? "fdatasync" : "pwrite", nrec); } static int apply_record(wo_db *db, const uint8_t *payload, uint32_t len) { diff --git a/database/src/wal.h b/database/src/wal.h index 7ff8a83..aeedc53 100644 --- a/database/src/wal.h +++ b/database/src/wal.h @@ -90,6 +90,11 @@ int wo_wal_append_update(wo_wal *w, wo_db *db, uint32_t class_id, uint64_t id); * batch stays staged: a failed commit consumes nothing). */ int wo_wal_commit(wo_wal *w); +/* databasev2 4: a record could not even be STAGED (the row is already in + * RAM, so this is the same unrecoverable position as a failed barrier — see + * wo_wal_commit_fatal). Never returns. */ +void wo_wal_stage_fatal(const wo_wal *w); + /* databasev2 4: commit, or END THE PROCESS. * * The one rule this iteration introduces: once a statement has mutated RAM, diff --git a/runtime/src/vm.c b/runtime/src/vm.c index 31c0b9e..0c63dbe 100644 --- a/runtime/src/vm.c +++ b/runtime/src/vm.c @@ -16,6 +16,7 @@ #include "db.h" /* arc stage 3: the transparent DB RPC (wo_db_req) */ #include "table.h" /* slot encode/decode for the RPC marshaling */ +#include "wal.h" /* databasev2 4: the drain issues the barrier */ #include #include @@ -91,6 +92,11 @@ static int wo_vm_adopt(wo_vm *vm) { ib->head = ib->tail = NULL; pthread_mutex_unlock(&ib->mu); int n = 0; + /* databasev2 4 (group commit): DB replies are HELD until one barrier has + * covered the whole drain. Locals, not per-shard state: nothing here needs + * to outlive the batch it describes. */ + wo_envelope *rhead = NULL, *rtail = NULL; + uint32_t staged = 0; while (e) { wo_envelope *nx = e->next; switch (e->kind) { @@ -171,13 +177,25 @@ static int wo_vm_adopt(wo_vm *vm) { * the same request back as the reply. */ wo_db_req *q = (wo_db_req *)(uintptr_t)e->payload; assert(vm->is_primary && "DB requests route to shard 0 only"); + wo_wal *dw = (wo_wal *)vm->rt.wal; + size_t before = dw ? dw->len : 0; wo_db_exec_req(vm, q); q->done = 1; + /* did this statement actually stage a record? Asking the buffer + * beats guessing from the opcode, and the count is what the + * failure diagnostic reports. */ + if (dw && dw->len > before) staged++; wo_envelope *re = calloc(1, sizeof *re); if (re) { re->kind = 4; re->payload = e->payload; - inbox_push_to(q->from_shard, re); + /* HELD, not pushed: pushing here would unpark the requester + * before its record is durable, which is the ack contract + * this iteration exists to make literally true. FIFO so the + * first waiter is released first. */ + re->next = NULL; + if (rtail) rtail->next = re; else rhead = re; + rtail = re; } /* OOM: the requester stays parked until stop — leak, not UB */ break; } @@ -192,6 +210,22 @@ static int wo_vm_adopt(wo_vm *vm) { n++; e = nx; } + /* databasev2 4: ONE barrier for everything this drain staged, then every + * held reply. Each requester therefore unparks having been acknowledged + * after the barrier that carried ITS record. Commit unconditionally when + * anything is staged — the inline path relies on finding the buffer empty + * (see db.c), so a drain must never leave a record behind. */ + if (staged) { + wo_wal *cw = (wo_wal *)vm->rt.wal; + if (cw) wo_wal_commit_fatal(cw, staged); + } + while (rhead) { + wo_envelope *rn = rhead->next; + wo_db_req *rq = (wo_db_req *)(uintptr_t)rhead->payload; + rhead->next = NULL; + inbox_push_to(rq->from_shard, rhead); + rhead = rn; + } return n; }