feat(db): one barrier per drain, replies held — T2

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) <noreply@anthropic.com>
This commit is contained in:
shoney.arickathil 2026-08-28 09:25:10 +02:00
parent ceea00e0b6
commit 76d80cc027
4 changed files with 63 additions and 28 deletions

View file

@ -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;
}

View file

@ -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) {

View file

@ -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,

View file

@ -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 <pthread.h>
#include <poll.h>
@ -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;
}