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