diff --git a/database/src/wal.c b/database/src/wal.c index caf682e..28efec8 100644 --- a/database/src/wal.c +++ b/database/src/wal.c @@ -439,6 +439,95 @@ void wo_wal_commit_fatal(wo_wal *w, uint32_t nrec) { wal_die(w, rc == WO_WAL_ERR_SYNC ? "fdatasync" : "pwrite", nrec); } +/* databasev2 3: how many records the dump stages before flushing. + * + * NOT unbounded: stage() grows the staging buffer by doubling and never + * shrinks it, so appending a whole store through one buffer would hold the + * entire store in RAM on top of the store itself — the unbounded growth + * databasev2 1 identified as how this engine dies. 256 records is a few tens + * of KiB per flush, which is large enough that the syscall cost is amortised + * and small enough that the buffer never matters. */ +#define WO_WAL_COMPACT_FLUSH 256u + +/* rename(2)'s atomicity is in-kernel: the new directory ENTRY is not durable + * until the parent directory is synced. Postgres does the same thing for the + * same reason. Best-effort — a filesystem that refuses to sync a directory + * still leaves a correct log, just one whose swap might not survive a power + * cut. */ +static void sync_parent_dir(const char *path) { + char dir[4096]; + size_t n = strlen(path); + if (n >= sizeof dir) return; + memcpy(dir, path, n + 1); + char *slash = strrchr(dir, '/'); + if (slash == dir) dir[1] = '\0'; + else if (slash) *slash = '\0'; + else memcpy(dir, ".", 2); + int fd = open(dir, O_RDONLY); + if (fd < 0) return; + (void)fsync(fd); + close(fd); +} + +int wo_wal_compact(wo_wal *w, wo_db *db) { + /* staged records would be written into a file about to be replaced */ + if (!w->path || w->len != 0) return -1; + + char tmp[4096]; + if ((size_t)snprintf(tmp, sizeof tmp, "%s%s", w->path, WO_WAL_TMP_SUFFIX) >= sizeof tmp) + return -1; + (void)unlink(tmp); /* a stale one would otherwise be appended to */ + + wo_wal nw; + if (wo_wal_open(&nw, tmp, 0) != 0) return -1; + + /* one INSERT per live row, in the existing grammar, through the existing + * append path — so replay needs no second decoder and ids are preserved + * exactly (wo_wal_append_insert takes the id and reads the row) */ + uint32_t pending = 0; + for (uint32_t cid = 0; cid < db->class_cnt; cid++) { + db_table *t = &db->tables[cid]; + if (!t->slabs) continue; /* tables are created lazily */ + uint32_t total = t->slab_cnt * DB_SLAB_ROWS; + for (uint32_t g = 0; g < total; g++) { + if (!(t->bitmap[g >> 6] & (1ull << (g & 63)))) continue; + db_row *r = (db_row *)(t->slabs[g / DB_SLAB_ROWS] + + (size_t)(g % DB_SLAB_ROWS) * t->row_size); + if (wo_wal_append_insert(&nw, db, cid, r->id) != 0) goto fail; + if (++pending >= WO_WAL_COMPACT_FLUSH) { + if (wo_wal_commit(&nw) != 0) goto fail; + pending = 0; + } + } + } + if (wo_wal_commit(&nw) != 0) goto fail; /* the tail batch */ + if (fsync(nw.fd) != 0) goto fail; /* commit fdatasyncs; this is for the size */ + + uint64_t new_bytes = nw.off; + wo_wal_close(&nw); + + /* THE SWITCH. Every crash point either side of this is safe. */ + if (rename(tmp, w->path) != 0) { + (void)unlink(tmp); + return -1; + } + sync_parent_dir(w->path); + + /* the old descriptor now refers to an unlinked inode */ + if (w->fd >= 0) close(w->fd); + w->fd = open(w->path, O_RDWR); + if (w->fd < 0) return -1; /* the log is correct on disk; this process cannot go on */ + w->off = new_bytes; + w->len = 0; + w->compacted_bytes = new_bytes; + return 0; + +fail: + wo_wal_close(&nw); + (void)unlink(tmp); + return -1; /* the live log is untouched and still usable */ +} + 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); diff --git a/database/src/wal.h b/database/src/wal.h index 6137e84..73c53bc 100644 --- a/database/src/wal.h +++ b/database/src/wal.h @@ -64,6 +64,10 @@ typedef struct wo_wal { uint64_t stat_records; /* records those commits carried */ uint64_t stat_peak_batch; /* most records in one barrier */ uint64_t stat_peak_staged; /* most bytes staged behind one barrier */ + /* databasev2 3: bytes the last compaction wrote. The trigger compares the + * log against THIS rather than an estimate of the live set — estimating + * would mean estimating Text, and the compactor knows the true number. */ + uint64_t compacted_bytes; } wo_wal; /* Open (create if missing) and preallocate [prealloc] bytes (best-effort; @@ -105,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 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). */ +#define WO_WAL_TMP_SUFFIX ".compact" + +/* databasev2 3: rewrite the log as one INSERT record per LIVE row, then swap + * it in with rename(2). + * + * Recovery is deliberately untouched: the result is an ordinary log in the + * ordinary grammar, replayed from byte 0. Crash safety comes from rename being + * atomic — before it the live log is intact and the temp file is not + * authoritative; after it the new log is complete. There is no window in which + * a reader sees a mixture, so this needs no recovery logic of its own. + * + * REFUSES if anything is staged (returns -1 without touching the log): those + * records would be written into a file about to be replaced. Callers must + * invoke this only where the staging buffer is empty — right after a barrier. + * + * A failure is a MISSED OPTIMISATION, not a durability event: the original log + * is left usable and the process keeps running. It must not take the fatal + * path wo_wal_commit_fatal takes. + * + * 0 ok, -1 on any failure. */ +int wo_wal_compact(wo_wal *w, wo_db *db); + /* 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. */ diff --git a/runtime/test/test_wal.c b/runtime/test/test_wal.c index 8d3bba2..bdc9be7 100644 --- a/runtime/test/test_wal.c +++ b/runtime/test/test_wal.c @@ -157,6 +157,83 @@ static void test_commit_failure_detected(void) { wo_rt_destroy(&rt); } +/* databasev2 3 Task 1: compaction rewrites the log as one record per LIVE row. + * Asserts BOTH halves on purpose: "the file got shorter" is also true of a + * truncating bug, so the replay comparison is what actually proves it. */ +static void test_compact_shortens_and_replays_equal(void) { + char path[128]; + snprintf(path, sizeof path, "%s/compact.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 = ""; + + uint64_t ids[3]; + for (int i = 0; i < 3; i++) { + wo_str *s = wo_str_new(&rt, "abc", 3); + uint64_t vals[2] = {(uint64_t)(i * 10), (uint64_t)(uintptr_t)s}; + ids[i] = wo_row_insert(&db, 0, vals, &msg, NULL); + T_CHECK(ids[i] != 0); + T_EQ(wo_wal_append_insert(&w, &db, 0, ids[i]), 0); + T_EQ(wo_wal_commit(&w), 0); + } + /* age it: the SAME row updated repeatedly, so HISTORY grows while the live + * set does not — the exact case checkpoint exists for */ + for (int k = 0; k < 40; k++) { + int ek = 0; + T_EQ(wo_row_update_field(&db, 0, ids[0], 0, (uint64_t)(500 + k), &msg, &ek), 0); + T_EQ(wo_wal_append_update(&w, &db, 0, ids[0]), 0); + T_EQ(wo_wal_commit(&w), 0); + } + uint64_t before_bytes = 0; + int64_t before_recs = wo_wal_check(path, &before_bytes); + T_CHECK(before_recs == 43); /* 3 inserts + 40 updates, all history */ + + T_EQ(wo_wal_compact(&w, &db), 0); + + uint64_t after_bytes = 0; + int64_t after_recs = wo_wal_check(path, &after_bytes); + T_CHECK(after_recs == 3); /* one record per LIVE row */ + T_CHECK(after_bytes < before_bytes); /* and the file really shrank */ + + /* the WAL stays usable: the descriptor was reopened and the offset reset, + * so a further write must land AFTER the compacted records, not over them */ + wo_str *s4 = wo_str_new(&rt, "xyz", 3); + uint64_t v4[2] = {99, (uint64_t)(uintptr_t)s4}; + uint64_t id4 = wo_row_insert(&db, 0, v4, &msg, NULL); + T_CHECK(id4 != 0); + T_EQ(wo_wal_append_insert(&w, &db, 0, id4), 0); + T_EQ(wo_wal_commit(&w), 0); + T_CHECK(wo_wal_check(path, NULL) == 4); + wo_wal_close(&w); + + /* the proof: a FRESH store replayed from the compacted log must hold the + * same rows, the same ids, and the LAST value each row had */ + 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_CHECK(out[0] == 539); /* the 40th update won, not the original 0 */ + wo_str_free(&rt, (wo_str *)(uintptr_t)out[1]); + T_EQ(wo_row_read(&db2, &rt, 0, ids[1], out, &msg), 0); + T_CHECK(out[0] == 10); + wo_str_free(&rt, (wo_str *)(uintptr_t)out[1]); + T_EQ(wo_row_read(&db2, &rt, 0, ids[2], out, &msg), 0); + T_CHECK(out[0] == 20); + wo_str_free(&rt, (wo_str *)(uintptr_t)out[1]); + T_EQ(wo_row_read(&db2, &rt, 0, id4, out, &msg), 0); + T_CHECK(out[0] == 99); + wo_str_free(&rt, (wo_str *)(uintptr_t)out[1]); + + wo_db_destroy(&db2); + 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); @@ -368,6 +445,7 @@ int main(void) { if (!mkdtemp(g_dir)) return 1; test_roundtrip_replay(); test_commit_failure_detected(); + test_compact_shortens_and_replays_equal(); test_torn_tail(); test_float_bytes_replay(); test_crash_battery();