diff --git a/database/src/db.c b/database/src/db.c index 4a2863e..c75397c 100644 --- a/database/src/db.c +++ b/database/src/db.c @@ -7,6 +7,23 @@ #include "table.h" #include "wal.h" +/* databasev2 3: the inline path's compaction check. + * + * The drain has its own (vm.c, after the barrier). This one exists because a + * statement running ON the owner shard never enters that drain, so without it + * a single-shard durable program's log grows FOREVER — measured: WO_SHARDS=1 + * reached 536 KB where the multi-shard run held 446 KB, because the check was + * only wired into the drain. + * + * Safe here for the same reason it is safe there: the commit above just + * emptied the staging buffer. The result is ignored because a failed + * compaction is a missed optimisation, not a durability event. */ +static void maybe_compact(wo_db *db, wo_wal *w) { + if (wo_wal_should_compact(w->off, w->compacted_bytes, wo_wal_ckpt_floor, + wo_wal_ckpt_ratio)) + (void)wo_wal_compact(w, db); +} + int wo_builtin_db(wo_vm *vm, uint64_t *R, uint32_t ins, const char **msg) { uint32_t A = wo_ins_a(ins), B = wo_ins_b(ins), C = wo_ins_c(ins); wo_db *db = (wo_db *)vm->rt.db; @@ -44,6 +61,7 @@ int wo_builtin_db(wo_vm *vm, uint64_t *R, uint32_t ins, const char **msg) { * Failure is fatal, not a trap: the row is already in RAM. */ if (wo_wal_append_insert(w, db, cid, id) != 0) wo_wal_stage_fatal(w); wo_wal_commit_fatal(w, 1); + maybe_compact(db, w); } R[A] = id; return 0; @@ -61,6 +79,7 @@ int wo_builtin_db(wo_vm *vm, uint64_t *R, uint32_t ins, const char **msg) { * admitted. Now fatal — see the insert arm. */ if (wo_wal_append_update(w, db, cid, id) != 0) wo_wal_stage_fatal(w); wo_wal_commit_fatal(w, 1); + maybe_compact(db, w); } R[A] = 0; return 0; @@ -82,6 +101,7 @@ int wo_builtin_db(wo_vm *vm, uint64_t *R, uint32_t ins, const char **msg) { if (w) { if (wo_wal_append_remove(w, cid, id) != 0) wo_wal_stage_fatal(w); wo_wal_commit_fatal(w, 1); + maybe_compact(db, w); } R[A] = 0; return 0; diff --git a/database/src/wal.c b/database/src/wal.c index c37dbcf..19e797f 100644 --- a/database/src/wal.c +++ b/database/src/wal.c @@ -450,6 +450,18 @@ void wo_wal_commit_fatal(wo_wal *w, uint32_t nrec) { wal_die(w, rc == WO_WAL_ERR_SYNC ? "fdatasync" : "pwrite", nrec); } +uint64_t wo_wal_ckpt_floor = 4u << 20; /* 4 MiB: below this there is nothing worth reclaiming */ +uint32_t wo_wal_ckpt_ratio = 3u; /* 3x the live-set's own size is enough history */ + +int wo_wal_should_compact(uint64_t used, uint64_t last, uint64_t floor, uint32_t ratio) { + if (used < floor) return 0; /* a small log has nothing to reclaim */ + if (last == 0) return 1; /* past the floor and never compacted: do it once + * to establish the denominator */ + if (ratio == 0) return 0; /* a zero ratio disables the policy rather than + * dividing by nothing */ + return used > last * (uint64_t)ratio; +} + /* databasev2 3: how many records the dump stages before flushing. * * NOT unbounded: stage() grows the staging buffer by doubling and never diff --git a/database/src/wal.h b/database/src/wal.h index 73c53bc..e219a14 100644 --- a/database/src/wal.h +++ b/database/src/wal.h @@ -109,6 +109,31 @@ 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 3: the checkpoint trigger, as a PURE decision so it can be tested + * without a store — which is the only way a policy like this gets tested at all. + * + * [used] the log's used bytes; [last] what the LAST compaction wrote (0 if it + * has never run); [floor] the size below which compacting is not worth it; + * [ratio] the multiple of [last] that counts as too much history. + * + * The denominator is the last compaction's MEASURED output rather than an + * estimate of the live set: estimating would mean estimating Text, and the + * compactor already knows the true number. + * + * There is deliberately NO TIME component. Postgres' CheckPointTimeout exists + * to bound data loss from unflushed buffers; our records are durable at commit, + * so a checkpoint only reclaims space and shortens boot. An idle log does not + * grow, so a timer would fire with nothing to do. + * + * 1 = compact now, 0 = leave it. */ +int wo_wal_should_compact(uint64_t used, uint64_t last, uint64_t floor, uint32_t ratio); + +/* Defaults, overridable at boot by WO_CHECKPOINT_BYTES / WO_CHECKPOINT_RATIO. + * The knobs are what make the policy testable: a test sets a tiny floor and + * forces compaction in a few writes instead of waiting for megabytes. */ +extern uint64_t wo_wal_ckpt_floor; +extern uint32_t wo_wal_ckpt_ratio; + /* databasev2 3: the temporary file compaction writes before the swap. Named * next to the log so it lands on the same filesystem — rename(2) is only * atomic within one. Boot removes a stale one (a crash before the rename). */ diff --git a/runtime/src/main.c b/runtime/src/main.c index f8f6191..3cb0f1c 100644 --- a/runtime/src/main.c +++ b/runtime/src/main.c @@ -240,6 +240,25 @@ int main(int argc, char **argv) { if (v >= 1 && v <= 0x7FFFFFFFul) wo_mailbox_cap = (uint32_t)v; } } + /* databasev2 3: the checkpoint policy. WO_CHECKPOINT_BYTES is the floor + * below which a log is too small to bother compacting; WO_CHECKPOINT_RATIO + * is how many times the live set's own size counts as too much history. + * Both exist mainly so the policy is TESTABLE — a gate sets a tiny floor + * and forces compaction in a few writes rather than waiting for megabytes. + * There is no time-based trigger, by design: our records are durable at + * commit, so an idle log does not grow. */ + { + const char *cb = getenv("WO_CHECKPOINT_BYTES"); + if (cb && cb[0]) { + unsigned long long v = strtoull(cb, NULL, 10); + if (v > 0) wo_wal_ckpt_floor = (uint64_t)v; + } + const char *cr = getenv("WO_CHECKPOINT_RATIO"); + if (cr && cr[0]) { + unsigned long v = strtoul(cr, NULL, 10); + if (v <= 0xFFFFFFFFul) wo_wal_ckpt_ratio = (uint32_t)v; + } + } /* the arc's stage 2: all cores by default (the brave landing), one * pinned worker vm per extra core; WO_SHARDS caps or forces it */ { diff --git a/runtime/src/vm.c b/runtime/src/vm.c index 169c4c6..4eda1e1 100644 --- a/runtime/src/vm.c +++ b/runtime/src/vm.c @@ -236,6 +236,25 @@ static int wo_vm_adopt(wo_vm *vm) { inbox_push_to(rq->from_shard, rhead); rhead = rn; } + /* databasev2 3: the ONE point where compaction is safe — the barrier above + * just ran, so the staging buffer is empty. Anywhere else, a staged record + * would be written into a file about to be replaced. This is a correctness + * requirement, not a scheduling preference; wo_wal_compact also refuses a + * non-empty buffer as a backstop. + * + * Replies are released FIRST, deliberately: their records are already + * durable, and holding them across a stop-the-world rewrite would add the + * rewrite's full duration to their latency for no benefit. + * + * The result is ignored because a failed compaction is a missed + * optimisation, not a durability event — the original log is left intact + * and the process carries on. */ + if (staged) { + wo_wal *cw = (wo_wal *)vm->rt.wal; + if (cw && wo_wal_should_compact(cw->off, cw->compacted_bytes, + wo_wal_ckpt_floor, wo_wal_ckpt_ratio)) + (void)wo_wal_compact(cw, (wo_db *)vm->rt.db); + } return n; } diff --git a/runtime/test/test_wal.c b/runtime/test/test_wal.c index d94a148..1be2996 100644 --- a/runtime/test/test_wal.c +++ b/runtime/test/test_wal.c @@ -297,6 +297,67 @@ static void test_stale_compact_temp_is_removed(void) { wo_rt_destroy(&rt); } +/* databasev2 3 Task 3: the trigger, tested as a pure decision. Kept pure + * precisely so it CAN be tested — a policy only observable by writing megabytes + * and waiting is a policy nobody checks. */ +static void test_should_compact_policy(void) { + /* below the floor, nothing fires however bad the ratio looks */ + T_EQ(wo_wal_should_compact(1000, 10, 4096, 3), 0); + T_EQ(wo_wal_should_compact(4095, 1, 4096, 3), 0); + /* past the floor with no prior compaction: run once to learn the size */ + T_EQ(wo_wal_should_compact(4096, 0, 4096, 3), 1); + /* with a known denominator it is a straight ratio test */ + T_EQ(wo_wal_should_compact(30000, 10000, 4096, 3), 0); /* exactly 3x is not MORE than 3x */ + T_EQ(wo_wal_should_compact(30001, 10000, 4096, 3), 1); + T_EQ(wo_wal_should_compact(19999, 10000, 4096, 2), 0); + T_EQ(wo_wal_should_compact(20001, 10000, 4096, 2), 1); + /* a zero ratio disables the policy rather than dividing by nothing */ + T_EQ(wo_wal_should_compact(1u << 30, 10, 4096, 0), 0); +} + +/* databasev2 3 Task 3: the ordering rule, asserted rather than trusted. + * Compaction with records staged would write them into a file about to be + * replaced, so it must be REFUSED — and refused without touching the log. */ +static void test_compact_refuses_with_staged_records(void) { + char path[128]; + snprintf(path, sizeof path, "%s/staged.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 = ""; + + wo_str *s1 = wo_str_new(&rt, "abc", 3); + uint64_t v1[2] = {7, (uint64_t)(uintptr_t)s1}; + uint64_t id1 = wo_row_insert(&db, 0, v1, &msg, NULL); + T_EQ(wo_wal_append_insert(&w, &db, 0, id1), 0); + T_EQ(wo_wal_commit(&w), 0); /* durable, buffer empty */ + + /* now stage WITHOUT committing */ + wo_str *s2 = wo_str_new(&rt, "xyz", 3); + uint64_t v2[2] = {8, (uint64_t)(uintptr_t)s2}; + uint64_t id2 = wo_row_insert(&db, 0, v2, &msg, NULL); + T_EQ(wo_wal_append_insert(&w, &db, 0, id2), 0); + T_CHECK(w.len > 0); + + uint64_t before = 0; + int64_t recs = wo_wal_check(path, &before); + T_EQ(wo_wal_compact(&w, &db), -1); /* refused */ + T_CHECK(w.len > 0); /* and the staged record is still there */ + uint64_t after = 0; + T_CHECK(wo_wal_check(path, &after) == recs && after == before); /* log untouched */ + + /* the staged record still commits normally afterwards */ + T_EQ(wo_wal_commit(&w), 0); + T_CHECK(wo_wal_check(path, NULL) == recs + 1); + + wo_wal_close(&w); + wo_db_destroy(&db); + wo_rt_destroy(&rt); +} + static void test_torn_tail(void) { char path[128]; snprintf(path, sizeof path, "%s/torn.wal", g_dir); @@ -510,6 +571,8 @@ int main(void) { test_commit_failure_detected(); test_compact_shortens_and_replays_equal(); test_stale_compact_temp_is_removed(); + test_should_compact_policy(); + test_compact_refuses_with_staged_records(); test_torn_tail(); test_float_bytes_replay(); test_crash_battery();