databasev2 3, task 3. - wo_wal_should_compact is a PURE decision (used bytes, last compaction's measured output, floor, ratio) so it is testable without a store — which is the only way a policy like this gets tested at all. Denominator is the last compaction's real output, not an estimate of the live set: estimating would mean estimating Text - 8 boundary assertions incl. "exactly 3x is not MORE than 3x" and a zero ratio disabling the policy rather than dividing by nothing - MUTATION-TESTED instead of observing RED: implementation and test were written together, so removing the floor check was verified to fail exactly the two floor assertions. Equivalent evidence, stated plainly - WO_CHECKPOINT_BYTES / WO_CHECKPOINT_RATIO at boot beside WO_MAILBOX. The knobs are what make the policy testable — a gate sets a tiny floor and forces compaction in a few writes instead of megabytes - NO timer, per the spec: Postgres' CheckPointTimeout bounds loss from unflushed buffers; our records are durable at commit and an idle log does not grow - the ordering rule is now asserted, not trusted: a test stages a record, requests compaction, and requires REFUSAL with the log untouched and the staged record still committable afterwards FOUND AND FIXED a gap in my own wiring. The plan said to call the check "after the drain's barrier", and I did — but a statement running ON the owner shard never enters that drain, so WO_SHARDS=1 never compacted and its log grew forever: measured 536086 bytes where the multi-shard run held 446024. Now checked after the inline path's commit too (db.c maybe_compact), where the buffer is equally empty. WO_SHARDS=1 went 536086 -> 260657 bytes. For a checkpoint this mattered more than part A's equivalent gap: an unbounded log is an operational failure, not just lost throughput. Also corrected a measurement of my own: multi-shard logs looked unbounded (448KB -> 1013KB -> 1647KB across 8k/24k/48k updates). They are not. Instrumentation showed compaction ran 25 times with zero failures, each writing MORE than the last, because the live set genuinely grows — wmix's hist_dump and done-markers are themselves durable inserts. Final log 1631040 against a last compaction of 866432 is a ratio of 1.88, just under the 2x threshold: the policy holding exactly. Replies are released BEFORE compaction runs, deliberately: their records are already durable, and holding them across a stop-the-world rewrite would add its full duration to their latency for nothing. Verified: wovm-test 36 suites 0 fail, test_wal 360 pass; db-bench-quick crash.s1/crash.sN and both restart legs green, and part A still batches (sN mean 4.16, peak 30). Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
588 lines
22 KiB
C
588 lines
22 KiB
C
/* test_wal — iteration 9 Task 2: typed WAL + boot replay.
|
|
* Round-trip through a replay, torn-tail drop, reopen-overwrites-tear,
|
|
* and the commit-then-kill crash battery: a forked child inserts rows and
|
|
* acks each COMMITTED id over a pipe; SIGKILL lands mid-stream; the parent
|
|
* verifies with the offline oracle and a replay that every acked id is
|
|
* present with the right contents. */
|
|
#define _POSIX_C_SOURCE 200809L
|
|
|
|
#include <fcntl.h>
|
|
#include <signal.h>
|
|
#include <stdlib.h>
|
|
#include <string.h>
|
|
#include <sys/wait.h>
|
|
#include <time.h>
|
|
#include <unistd.h>
|
|
|
|
#include "gc.h"
|
|
#include "obj.h"
|
|
#include "t.h"
|
|
#include "table.h"
|
|
#include "wal.h"
|
|
|
|
/* class 0: Row { n: scalar, label: Text } */
|
|
static const uint8_t row_kinds[] = {WO_K_SCALAR, WO_K_TEXT};
|
|
static const wo_classdesc CLASSES[] = {
|
|
{.name = 0, .flags = 0, .field_cnt = 2, .kinds = row_kinds},
|
|
};
|
|
|
|
static char g_dir[64];
|
|
|
|
static void test_roundtrip_replay(void) {
|
|
char path[128];
|
|
snprintf(path, sizeof path, "%s/basic.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 = "";
|
|
|
|
/* three inserts and one remove, RAM first, WAL second, one commit */
|
|
uint64_t ids[3];
|
|
for (int i = 0; i < 3; i++) {
|
|
wo_str *s = wo_str_new(&rt, "abcXYZ" + i, 3); /* "abc","bcX","cXY" */
|
|
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);
|
|
wo_str_free(&rt, s);
|
|
}
|
|
T_EQ(wo_row_remove(&db, 0, ids[1]), 0);
|
|
T_EQ(wo_wal_append_remove(&w, 0, ids[1]), 0);
|
|
T_EQ(wo_wal_commit(&w), 0);
|
|
wo_wal_close(&w);
|
|
wo_db_destroy(&db);
|
|
|
|
/* boot: fresh engine, replay, deep-compare */
|
|
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_EQ(out[0], 0);
|
|
wo_str *s0 = (wo_str *)(uintptr_t)out[1];
|
|
T_CHECK(s0->len == 3 && memcmp(s0->data, "abc", 3) == 0);
|
|
wo_str_free(&rt, s0);
|
|
T_EQ(wo_row_read(&db2, &rt, 0, ids[1], out, &msg), -1); /* removed */
|
|
T_EQ(wo_row_read(&db2, &rt, 0, ids[2], out, &msg), 0);
|
|
T_EQ(out[0], 20);
|
|
wo_str_free(&rt, (wo_str *)(uintptr_t)out[1]);
|
|
/* next_id advanced past the replayed ids: a fresh insert never collides */
|
|
uint64_t vals[2] = {99, 0};
|
|
uint64_t fresh = wo_row_insert(&db2, 0, vals, &msg, NULL);
|
|
T_CHECK(fresh > ids[2]);
|
|
wo_db_destroy(&db2);
|
|
|
|
/* update record: re-log, replay replaces */
|
|
{
|
|
char upath[128];
|
|
snprintf(upath, sizeof upath, "%s/upd.wal", g_dir);
|
|
wo_db du;
|
|
T_EQ(wo_db_init(&du, CLASSES, 1, 0, 1), 0);
|
|
wo_wal wu;
|
|
T_EQ(wo_wal_open(&wu, upath, 0), 0);
|
|
wo_str *s1 = wo_str_new(&rt, "old", 3);
|
|
uint64_t uv[2] = {7, (uint64_t)(uintptr_t)s1};
|
|
uint64_t uid = wo_row_insert(&du, 0, uv, &msg, NULL);
|
|
T_EQ(wo_wal_append_insert(&wu, &du, 0, uid), 0);
|
|
int ek = 0;
|
|
wo_str *s2 = wo_str_new(&rt, "new!", 4);
|
|
T_EQ(wo_row_update_field(&du, 0, uid, 1, (uint64_t)(uintptr_t)s2, &msg, &ek), 0);
|
|
T_EQ(wo_row_update_field(&du, 0, uid, 0, 8, &msg, &ek), 0);
|
|
T_EQ(wo_wal_append_update(&wu, &du, 0, uid), 0);
|
|
T_EQ(wo_wal_commit(&wu), 0);
|
|
wo_wal_close(&wu);
|
|
wo_db_destroy(&du);
|
|
wo_db db4;
|
|
T_EQ(wo_db_init(&db4, CLASSES, 1, 0, 1), 0);
|
|
T_EQ(wo_wal_replay(upath, &db4), 2);
|
|
uint64_t uo[2];
|
|
T_EQ(wo_row_read(&db4, &rt, 0, uid, uo, &msg), 0);
|
|
T_EQ(uo[0], 8);
|
|
wo_str *us = (wo_str *)(uintptr_t)uo[1];
|
|
T_CHECK(us->len == 4 && memcmp(us->data, "new!", 4) == 0);
|
|
wo_str_free(&rt, us);
|
|
wo_drop_obj(&rt, (wo_hdr *)s1);
|
|
wo_drop_obj(&rt, (wo_hdr *)s2);
|
|
wo_db_destroy(&db4);
|
|
}
|
|
|
|
/* replay of a missing file is a fresh boot, not an error */
|
|
wo_db db3;
|
|
T_EQ(wo_db_init(&db3, CLASSES, 1, 0, 1), 0);
|
|
T_EQ(wo_wal_replay("/nonexistent/nope.wal", &db3), 0);
|
|
wo_db_destroy(&db3);
|
|
wo_rt_destroy(&rt);
|
|
}
|
|
|
|
/* databasev2 4 part A, Task 1: a failed barrier must be DETECTED, and the
|
|
* caller must be able to tell WHICH operation failed — a pwrite failure and
|
|
* an fdatasync failure are different operational problems and the diagnostic
|
|
* has to name the right one. This proves detection only; the fatal exit that
|
|
* follows it cannot be exercised in-process. */
|
|
static void test_commit_failure_detected(void) {
|
|
char path[128];
|
|
snprintf(path, sizeof path, "%s/commitfail.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 = "";
|
|
|
|
/* the WAL remembers where it lives — the abort diagnostic is worthless
|
|
* without it */
|
|
T_CHECK(w.path != NULL && strstr(w.path, "commitfail.wal") != NULL);
|
|
|
|
wo_str *s = wo_str_new(&rt, "abc", 3);
|
|
uint64_t vals[2] = {7, (uint64_t)(uintptr_t)s};
|
|
uint64_t id = wo_row_insert(&db, 0, vals, &msg, NULL);
|
|
T_CHECK(id != 0);
|
|
T_EQ(wo_wal_append_insert(&w, &db, 0, id), 0);
|
|
T_CHECK(w.len > 0); /* something really is staged */
|
|
|
|
/* an unusable descriptor: pwrite reports EBADF. -1 is used rather than
|
|
* closing the real fd so the close below cannot double-free it. */
|
|
int real = w.fd;
|
|
w.fd = -1;
|
|
T_EQ(wo_wal_commit(&w), WO_WAL_ERR_WRITE);
|
|
T_CHECK(w.len > 0); /* a failed commit consumes nothing */
|
|
w.fd = real;
|
|
|
|
wo_wal_close(&w);
|
|
wo_db_destroy(&db);
|
|
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);
|
|
}
|
|
|
|
/* databasev2 3 Task 2: a stale temp file is the one input that could be
|
|
* mistaken for data — a crash before the rename leaves one behind, full of
|
|
* well-formed records that are NOT yet authoritative. So the fixture uses
|
|
* plausible records (a byte copy of a real log), not garbage: garbage would be
|
|
* rejected by the CRC anyway and would prove nothing. */
|
|
static void test_stale_compact_temp_is_removed(void) {
|
|
char path[128], tmp[160];
|
|
snprintf(path, sizeof path, "%s/stale.wal", g_dir);
|
|
snprintf(tmp, sizeof tmp, "%s%s", path, WO_WAL_TMP_SUFFIX);
|
|
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 = "";
|
|
|
|
/* two live rows in the REAL log */
|
|
uint64_t ids[2];
|
|
for (int i = 0; i < 2; i++) {
|
|
wo_str *s = wo_str_new(&rt, "abc", 3);
|
|
uint64_t vals[2] = {(uint64_t)(i + 1), (uint64_t)(uintptr_t)s};
|
|
ids[i] = wo_row_insert(&db, 0, vals, &msg, NULL);
|
|
T_EQ(wo_wal_append_insert(&w, &db, 0, ids[i]), 0);
|
|
T_EQ(wo_wal_commit(&w), 0);
|
|
}
|
|
wo_wal_close(&w);
|
|
|
|
/* forge a plausible stale temp: a byte copy of the real log */
|
|
{
|
|
int src = open(path, O_RDONLY);
|
|
int dst = open(tmp, O_WRONLY | O_CREAT | O_TRUNC, 0644);
|
|
T_CHECK(src >= 0 && dst >= 0);
|
|
char buf[8192];
|
|
ssize_t n;
|
|
while ((n = read(src, buf, sizeof buf)) > 0) T_CHECK(write(dst, buf, (size_t)n) == n);
|
|
close(src);
|
|
close(dst);
|
|
T_EQ(access(tmp, F_OK), 0); /* it really is there before we open */
|
|
}
|
|
|
|
wo_wal w2;
|
|
T_EQ(wo_wal_open(&w2, path, 1 << 16), 0);
|
|
T_CHECK(access(tmp, F_OK) != 0); /* gone, and never consulted */
|
|
wo_wal_close(&w2);
|
|
|
|
/* and the live log still says exactly what it said */
|
|
wo_db db2;
|
|
T_EQ(wo_db_init(&db2, CLASSES, 1, 0, 1), 0);
|
|
T_EQ(wo_wal_replay(path, &db2), 2);
|
|
uint64_t out[2];
|
|
T_EQ(wo_row_read(&db2, &rt, 0, ids[0], out, &msg), 0);
|
|
T_CHECK(out[0] == 1);
|
|
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] == 2);
|
|
wo_str_free(&rt, (wo_str *)(uintptr_t)out[1]);
|
|
|
|
wo_db_destroy(&db2);
|
|
wo_db_destroy(&db);
|
|
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);
|
|
wo_db db;
|
|
T_EQ(wo_db_init(&db, CLASSES, 1, 0, 1), 0);
|
|
wo_wal w;
|
|
T_EQ(wo_wal_open(&w, path, 0), 0);
|
|
const char *msg = "";
|
|
for (int i = 0; i < 5; i++) {
|
|
uint64_t vals[2] = {(uint64_t)i, 0};
|
|
uint64_t id = wo_row_insert(&db, 0, vals, &msg, NULL);
|
|
T_EQ(wo_wal_append_insert(&w, &db, 0, id), 0);
|
|
T_EQ(wo_wal_commit(&w), 0);
|
|
}
|
|
uint64_t intact_end = w.off;
|
|
/* tear: append half a record's worth of a valid-looking header + junk */
|
|
uint32_t fake_len = 40, fake_crc = 0xDEAD;
|
|
uint8_t junk[20] = {7, 7, 7};
|
|
T_CHECK(pwrite(w.fd, &fake_len, 4, (off_t)intact_end) == 4);
|
|
T_CHECK(pwrite(w.fd, &fake_crc, 4, (off_t)(intact_end + 4)) == 4);
|
|
T_CHECK(pwrite(w.fd, junk, sizeof junk, (off_t)(intact_end + 8)) == (ssize_t)sizeof junk);
|
|
wo_wal_close(&w);
|
|
wo_db_destroy(&db);
|
|
|
|
/* the oracle sees exactly the intact prefix */
|
|
uint64_t at = 0;
|
|
T_EQ(wo_wal_check(path, &at), 5);
|
|
T_EQ(at, intact_end);
|
|
|
|
/* replay drops the tear whole */
|
|
wo_db db2;
|
|
T_EQ(wo_db_init(&db2, CLASSES, 1, 0, 1), 0);
|
|
T_EQ(wo_wal_replay(path, &db2), 5);
|
|
T_EQ(db2.tables[0].count, 5);
|
|
wo_db_destroy(&db2);
|
|
|
|
/* reopen positions AT the tear: the next commit overwrites it */
|
|
wo_db db3;
|
|
T_EQ(wo_db_init(&db3, CLASSES, 1, 0, 1), 0);
|
|
T_EQ(wo_wal_replay(path, &db3), 5);
|
|
wo_wal w2;
|
|
T_EQ(wo_wal_open(&w2, path, 0), 0);
|
|
T_EQ(w2.off, intact_end);
|
|
uint64_t vals[2] = {100, 0};
|
|
uint64_t id = wo_row_insert(&db3, 0, vals, &msg, NULL);
|
|
T_EQ(wo_wal_append_insert(&w2, &db3, 0, id), 0);
|
|
T_EQ(wo_wal_commit(&w2), 0);
|
|
wo_wal_close(&w2);
|
|
T_EQ(wo_wal_check(path, NULL), 6); /* tear gone, record in its place */
|
|
wo_db_destroy(&db3);
|
|
}
|
|
|
|
/* ---- the crash battery -------------------------------------------------- */
|
|
|
|
/* Child: insert forever — RAM, WAL, COMMIT, and only then ack the id down
|
|
* the pipe. Killed mid-stream by the parent. */
|
|
static void battery_child(const char *path, int ack_fd) {
|
|
wo_rt rt;
|
|
wo_db db;
|
|
wo_wal w;
|
|
if (wo_rt_init(&rt, 1 << 20, CLASSES, 1) != 0) _exit(9);
|
|
if (wo_db_init(&db, CLASSES, 1, 0, 1) != 0) _exit(9);
|
|
if (wo_wal_open(&w, path, 1 << 20) != 0) _exit(9);
|
|
const char *msg = "";
|
|
for (uint64_t i = 0;; i++) {
|
|
char label[32];
|
|
int n = snprintf(label, sizeof label, "row-%llu", (unsigned long long)i);
|
|
wo_str *s = wo_str_new(&rt, label, (uint32_t)n);
|
|
uint64_t vals[2] = {i * 3 + 1, (uint64_t)(uintptr_t)s};
|
|
uint64_t id = wo_row_insert(&db, 0, vals, &msg, NULL);
|
|
wo_str_free(&rt, s);
|
|
if (!id) _exit(9);
|
|
if (wo_wal_append_insert(&w, &db, 0, id) != 0) _exit(9);
|
|
if (wo_wal_commit(&w) != 0) _exit(9); /* durable BEFORE the ack */
|
|
ssize_t wr = write(ack_fd, &id, 8);
|
|
if (wr != 8) _exit(0); /* parent went away */
|
|
}
|
|
}
|
|
|
|
static void test_crash_battery(void) {
|
|
for (int round = 0; round < 5; round++) {
|
|
char path[128];
|
|
snprintf(path, sizeof path, "%s/crash-%d.wal", g_dir, round);
|
|
int pipefd[2];
|
|
T_EQ(pipe(pipefd), 0);
|
|
pid_t pid = fork();
|
|
T_CHECK(pid >= 0);
|
|
if (pid == 0) {
|
|
close(pipefd[0]);
|
|
battery_child(path, pipefd[1]);
|
|
_exit(0);
|
|
}
|
|
close(pipefd[1]);
|
|
/* collect acks for a few ms, then kill mid-stream — no sync with
|
|
the child's commit loop, which is the point */
|
|
struct timespec ts = {0, (20 + round * 13) * 1000000L};
|
|
while (nanosleep(&ts, &ts) != 0) {}
|
|
kill(pid, SIGKILL);
|
|
int status;
|
|
waitpid(pid, &status, 0);
|
|
/* drain every ack that made it into the pipe */
|
|
uint64_t acked[65536];
|
|
size_t n_acked = 0;
|
|
for (;;) {
|
|
uint64_t id;
|
|
ssize_t n = read(pipefd[0], &id, 8);
|
|
if (n != 8) break;
|
|
if (n_acked < 65536) acked[n_acked++] = id;
|
|
}
|
|
close(pipefd[0]);
|
|
T_CHECK(n_acked > 0); /* the child got at least one commit out */
|
|
|
|
/* offline oracle: the file's intact prefix covers every ack */
|
|
int64_t intact = wo_wal_check(path, NULL);
|
|
T_CHECK(intact >= (int64_t)n_acked);
|
|
|
|
/* replay and verify: every acked id present, contents exact */
|
|
wo_rt rt;
|
|
T_EQ(wo_rt_init(&rt, 1 << 22, CLASSES, 1), 0);
|
|
wo_db db;
|
|
T_EQ(wo_db_init(&db, CLASSES, 1, 0, 1), 0);
|
|
int64_t applied = wo_wal_replay(path, &db);
|
|
T_CHECK(applied >= (int64_t)n_acked);
|
|
const char *msg = "";
|
|
int bad = 0;
|
|
for (size_t i = 0; i < n_acked; i++) {
|
|
uint64_t out[2];
|
|
if (wo_row_read(&db, &rt, 0, acked[i], out, &msg) != 0) {
|
|
bad++;
|
|
continue;
|
|
}
|
|
/* id = i+1 (shard 0 of 1), field 0 = i*3+1, label = "row-i" */
|
|
char want[32];
|
|
int wl = snprintf(want, sizeof want, "row-%llu",
|
|
(unsigned long long)(acked[i] - 1));
|
|
wo_str *s = (wo_str *)(uintptr_t)out[1];
|
|
if (out[0] != (acked[i] - 1) * 3 + 1 || s->len != (uint32_t)wl ||
|
|
memcmp(s->data, want, (size_t)wl) != 0)
|
|
bad++;
|
|
wo_str_free(&rt, s);
|
|
}
|
|
T_EQ(bad, 0); /* zero acked-but-missing, zero acked-but-wrong */
|
|
wo_db_destroy(&db);
|
|
wo_rt_destroy(&rt);
|
|
}
|
|
}
|
|
|
|
/* iteration 19: a Float column and a Bytes column survive a WAL round trip
|
|
* BIT-EXACT. Bit-exact is the whole assertion — the durability path must not
|
|
* render a float as decimal anywhere, or NaN, the infinities and -0.0 would
|
|
* each come back as something else. Bytes goes through the same length-
|
|
* prefixed blob a Text does and must come back as a Bytes, not a Text. */
|
|
static const uint8_t fb_kinds[] = {WO_K_FLOAT, WO_K_BYTES};
|
|
static const wo_classdesc FB_CLASSES[] = {
|
|
{.name = 0, .flags = 0, .field_cnt = 2, .kinds = fb_kinds},
|
|
};
|
|
|
|
static void test_float_bytes_replay(void) {
|
|
char path[128];
|
|
snprintf(path, sizeof path, "%s/floatbytes.wal", g_dir);
|
|
wo_rt rt;
|
|
T_EQ(wo_rt_init(&rt, 1 << 20, FB_CLASSES, 1), 0);
|
|
wo_db db;
|
|
T_EQ(wo_db_init(&db, FB_CLASSES, 1, 0, 1), 0);
|
|
wo_wal w;
|
|
T_EQ(wo_wal_open(&w, path, 1 << 16), 0);
|
|
const char *msg = "";
|
|
|
|
/* the values a decimal round trip would destroy, plus a NUL-bearing blob
|
|
* that a NUL-terminated string path would truncate */
|
|
const double vals_f[] = {9.99, 0.0 / 0.0, 1.0 / 0.0, -1.0 / 0.0, -0.0, 1e308};
|
|
const char blob[] = {'a', '\0', 'b'};
|
|
enum { N = sizeof vals_f / sizeof vals_f[0] };
|
|
uint64_t ids[N];
|
|
for (int i = 0; i < N; i++) {
|
|
wo_str *b = wo_bytes_new(&rt, blob, sizeof blob);
|
|
T_CHECK(b != NULL);
|
|
uint64_t vals[2] = {wo_bits(vals_f[i]), (uint64_t)(uintptr_t)b};
|
|
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);
|
|
wo_str_free(&rt, b);
|
|
}
|
|
T_EQ(wo_wal_commit(&w), 0);
|
|
wo_wal_close(&w);
|
|
wo_db_destroy(&db);
|
|
|
|
wo_db db2;
|
|
T_EQ(wo_db_init(&db2, FB_CLASSES, 1, 0, 1), 0);
|
|
T_EQ(wo_wal_replay(path, &db2), N);
|
|
for (int i = 0; i < N; i++) {
|
|
uint64_t out[2];
|
|
T_EQ(wo_row_read(&db2, &rt, 0, ids[i], out, &msg), 0);
|
|
/* BITS, not value: NaN != NaN and -0.0 == 0.0, so a value comparison
|
|
* would pass while silently having lost the payload or the sign */
|
|
T_EQ(out[0], wo_bits(vals_f[i]));
|
|
wo_str *b = (wo_str *)(uintptr_t)out[1];
|
|
T_CHECK(b != NULL);
|
|
T_EQ(b->h.class_id, WO_CLS_BYTES); /* a Bytes column yields a Bytes */
|
|
T_CHECK(b->len == sizeof blob && memcmp(b->data, blob, sizeof blob) == 0);
|
|
wo_str_free(&rt, b);
|
|
}
|
|
wo_db_destroy(&db2);
|
|
wo_rt_destroy(&rt);
|
|
}
|
|
|
|
int main(void) {
|
|
snprintf(g_dir, sizeof g_dir, "/tmp/wo-wal-test-XXXXXX");
|
|
if (!mkdtemp(g_dir)) return 1;
|
|
test_roundtrip_replay();
|
|
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();
|
|
/* leave the dir for a failed run's forensics only */
|
|
if (!t_fail) {
|
|
char cmd[128];
|
|
snprintf(cmd, sizeof cmd, "rm -rf %s", g_dir);
|
|
if (system(cmd) != 0) fprintf(stderr, "cleanup failed, kept %s\n", g_dir);
|
|
} else {
|
|
fprintf(stderr, "kept %s\n", g_dir);
|
|
}
|
|
return t_report("test_wal");
|
|
}
|