From 4ced1963fe935515f59184828d3e9e76efc8da16 Mon Sep 17 00:00:00 2001 From: belowzeroff Date: Mon, 28 Sep 2026 10:29:06 -0400 Subject: [PATCH 1/5] feat(store): publish splayed replacements atomically --- docs/docs/storage/index.md | 14 ++++ src/store/splay.c | 159 ++++++++++++++++++++++++++++++++++--- test/test_splay.c | 77 ++++++++++++++++++ 3 files changed, 238 insertions(+), 12 deletions(-) diff --git a/docs/docs/storage/index.md b/docs/docs/storage/index.md index 62b721165..88f165b02 100644 --- a/docs/docs/storage/index.md +++ b/docs/docs/storage/index.md @@ -64,6 +64,20 @@ The symbol table is a dotfile (`.sym`, lock `.sym.lk`), so it never collides with a user column — a table may have an ordinary column named `sym` (the canonical ticker column), as shown above. +### Atomic table replacement + +When an existing splayed table is replaced, Rayforce writes the new schema and +all columns into a fresh directory below `.generations/`. After the files are +complete, it atomically replaces the small `.current` manifest. Readers follow +that manifest, so they see either the previous complete generation or the new +complete generation — never a mixture of columns from both. A failed write +before publication leaves the previous generation selected. Tables written by +older Rayforce versions remain readable; the first replacement upgrades the +directory to the generation layout. + +This protects publication of one splayed table. Coordinating a consistent +snapshot across multiple tables or date partitions still belongs to the caller. + ### C API | Function | Description | diff --git a/src/store/splay.c b/src/store/splay.c index f9b4be3e9..ac920984b 100644 --- a/src/store/splay.c +++ b/src/store/splay.c @@ -40,8 +40,10 @@ #include #include #include +#include #include #include +#include /* -------------------------------------------------------------------------- * Splayed table: directory of column files + .d schema file @@ -145,8 +147,91 @@ static void splay_sweep_stale(ray_t* tbl, const char* dir) { closedir(d); } -static ray_err_t splay_save_impl(ray_t* tbl, const char* dir, const char* sym_path, - bool durable, bool staged) { +/* A published splayed table may contain more than one physical generation. + * The small text manifest is the only mutable table-level pointer: data files + * are written below .generations/ first, then .current is replaced with + * one filesystem rename. Readers that predate this format continue to use + * the legacy directory when .current is absent. */ +static unsigned long splay_generation_seq; + +static ray_err_t splay_current_dir(const char* dir, char* out, size_t out_sz, + bool* active) { + char manifest[1024]; + int n = snprintf(manifest, sizeof(manifest), "%s/.current", dir); + if (n < 0 || (size_t)n >= sizeof(manifest) || !out || !active) + return RAY_ERR_RANGE; + *active = false; + + struct stat st; + if (stat(manifest, &st) != 0) { + if (errno == ENOENT) return RAY_OK; /* legacy splayed directory */ + return RAY_ERR_IO; + } + + FILE* f = fopen(manifest, "rb"); + if (!f) return RAY_ERR_IO; + + char rel[512]; + int scanned = fscanf(f, "%511s", rel); + int closed = fclose(f); + if (scanned != 1 || closed != 0) { + return RAY_ERR_CORRUPT; + } + if (strncmp(rel, ".generations/", 13) != 0 || + strstr(rel, "..") || strstr(rel, "\\") || + strcmp(rel, ".generations/") == 0) + return RAY_ERR_CORRUPT; + + n = snprintf(out, out_sz, "%s/%s", dir, rel); + if (n < 0 || (size_t)n >= out_sz) return RAY_ERR_RANGE; + *active = true; + return RAY_OK; +} + +static ray_err_t splay_publish_generation(const char* dir, const char* gen) { + char manifest[1024], tmp[1024]; + int n = snprintf(manifest, sizeof(manifest), "%s/.current", dir); + if (n < 0 || (size_t)n >= sizeof(manifest)) return RAY_ERR_RANGE; + n = snprintf(tmp, sizeof(tmp), "%s/.current.tmp.%lu.%lu.%lu", dir, + (unsigned long)time(NULL), (unsigned long)getpid(), + ++splay_generation_seq); + if (n < 0 || (size_t)n >= sizeof(tmp)) return RAY_ERR_RANGE; + + FILE* f = fopen(tmp, "wb"); + if (!f) return RAY_ERR_IO; + size_t len = strlen(gen); + bool ok = fwrite(gen, 1, len, f) == len && fputc('\n', f) != EOF; + if (ok && fflush(f) != 0) ok = false; + if (fclose(f) != 0) ok = false; + if (!ok) { + (void)remove(tmp); + return RAY_ERR_IO; + } + + ray_fd_t fd = ray_file_open(tmp, RAY_OPEN_READ | RAY_OPEN_WRITE); + if (fd == RAY_FD_INVALID) { + (void)remove(tmp); + return RAY_ERR_IO; + } + ray_err_t err = ray_file_sync(fd); + ray_file_close(fd); + if (err != RAY_OK || ray_file_rename(tmp, manifest) != RAY_OK) { + (void)remove(tmp); + return RAY_ERR_IO; + } + return ray_file_sync_dir(manifest); +} + +static bool splay_has_file(const char* dir, const char* name) { + char path[1024]; + int n = snprintf(path, sizeof(path), "%s/%s", dir, name); + if (n < 0 || (size_t)n >= sizeof(path)) return false; + struct stat st; + return stat(path, &st) == 0; +} + +static ray_err_t splay_save_to_dir_impl(ray_t* tbl, const char* dir, + const char* sym_path, bool durable) { if (!tbl || RAY_IS_ERR(tbl)) return RAY_ERR_TYPE; if (!dir) return RAY_ERR_IO; @@ -254,9 +339,7 @@ static ray_err_t splay_save_impl(ray_t* tbl, const char* dir, const char* sym_pa } } - /* Private import roots publish only after their final domain flush. - * Live-table saves still require the vocabulary before every column. */ - ray_err_t sym_err = staged ? RAY_OK : ray_sym_domain_flush(dom, durable); + ray_err_t sym_err = ray_sym_domain_flush(dom, durable); if (sym_err != RAY_OK) { ray_sym_domain_release(dom); return sym_err; @@ -338,16 +421,60 @@ static ray_err_t splay_save_impl(ray_t* tbl, const char* dir, const char* sym_pa return RAY_OK; } -ray_err_t ray_splay_save(ray_t* tbl, const char* dir, const char* sym_path) { - return splay_save_impl(tbl, dir, sym_path, true, false); +static ray_err_t splay_save_impl(ray_t* tbl, const char* dir, + const char* sym_path, bool durable) { + if (!tbl || RAY_IS_ERR(tbl)) return RAY_ERR_TYPE; + if (!dir) return RAY_ERR_IO; + + /* Keep the first write backward-compatible. Once a table has a + * committed schema (or already uses generations), every replacement is + * staged out-of-place and published through .current. */ + bool atomic = splay_has_file(dir, ".d") || splay_has_file(dir, ".current"); + if (!atomic) + return splay_save_to_dir_impl(tbl, dir, sym_path, durable); + + ray_err_t err = ray_mkdir_p(dir); + if (err != RAY_OK) return err; + + char generations[1024], gen_name[256], gen_dir[1024]; + int n = snprintf(generations, sizeof(generations), "%s/.generations", dir); + if (n < 0 || (size_t)n >= sizeof(generations)) return RAY_ERR_RANGE; + err = ray_mkdir_p(generations); + if (err != RAY_OK) return err; + + n = snprintf(gen_name, sizeof(gen_name), ".generations/g-%lu-%lu-%lu", + (unsigned long)time(NULL), (unsigned long)getpid(), + ++splay_generation_seq); + if (n < 0 || (size_t)n >= sizeof(gen_name)) return RAY_ERR_RANGE; + n = snprintf(gen_dir, sizeof(gen_dir), "%s/%s", dir, gen_name); + if (n < 0 || (size_t)n >= sizeof(gen_dir)) return RAY_ERR_RANGE; + + err = splay_save_to_dir_impl(tbl, gen_dir, sym_path, durable); + if (err != RAY_OK) return err; + if (durable && ray_file_sync_dir(gen_dir) != RAY_OK) + return RAY_ERR_IO; + if (ray_file_sync_dir(generations) != RAY_OK) + return RAY_ERR_IO; + + /* This is the table-level commit point. A failure before this rename + * leaves the old .current (and therefore the old complete generation) + * untouched. */ + err = splay_publish_generation(dir, gen_name); + if (err != RAY_OK) return err; + + /* Legacy files are no longer read once .current exists. Remove stale + * columns opportunistically for existing callers that inspect the old + * directory layout; failure here cannot invalidate the committed table. */ + splay_sweep_stale(tbl, dir); + return RAY_OK; } -ray_err_t ray_splay_save_bulk(ray_t* tbl, const char* dir, const char* sym_path) { - return splay_save_impl(tbl, dir, sym_path, false, false); +ray_err_t ray_splay_save(ray_t* tbl, const char* dir, const char* sym_path) { + return splay_save_impl(tbl, dir, sym_path, true); } -ray_err_t ray_splay_save_staged_bulk(ray_t* tbl, const char* dir, const char* sym_path) { - return splay_save_impl(tbl, dir, sym_path, false, true); +ray_err_t ray_splay_save_bulk(ray_t* tbl, const char* dir, const char* sym_path) { + return splay_save_impl(tbl, dir, sym_path, false); } /* -------------------------------------------------------------------------- @@ -494,6 +621,14 @@ static ray_t* splay_load_dom_impl(const char* dir, ray_sym_domain_t* dom, * unopenable/invalid file → loud error. */ static ray_t* splay_load_impl(const char* dir, const char* sym_path, bool use_mmap) { + char resolved[1024]; + bool active = false; + ray_err_t current_err = splay_current_dir(dir, resolved, sizeof(resolved), + &active); + if (current_err != RAY_OK) + return ray_error("corrupt", "splayed %s: invalid .current manifest", dir); + + const char* load_dir = active ? resolved : dir; ray_sym_domain_t* dom = NULL; if (sym_path) { struct stat st; @@ -505,7 +640,7 @@ static ray_t* splay_load_impl(const char* dir, const char* sym_path, "record, or missing \"\" at position 0)", sym_path); } } - ray_t* tbl = splay_load_dom_impl(dir, dom, use_mmap); + ray_t* tbl = splay_load_dom_impl(load_dir, dom, use_mmap); if (dom) ray_sym_domain_release(dom); /* columns hold their own refs */ return tbl; } diff --git a/test/test_splay.c b/test/test_splay.c index d037cf5eb..784c702e2 100644 --- a/test/test_splay.c +++ b/test/test_splay.c @@ -2372,7 +2372,84 @@ static test_result_t test_splayed_has_nulls_roundtrip(void) { PASS(); } +/* A replacement must be published as one generation. The failed replacement + * is deterministic (unsupported nested data), so it also proves that a + * preflight/write error cannot advance the table-level manifest. */ +static test_result_t test_splay_atomic_generation_publish(void) { + const char* dir = TMP_SPLAY_BASE "/atomic_generation"; + char manifest[512]; + int n = snprintf(manifest, sizeof(manifest), "%s/.current", dir); + TEST_ASSERT_TRUE(n > 0 && (size_t)n < sizeof(manifest)); + (void)ray_test_rm_rf(dir); + + int64_t x_id = ray_sym_intern("x", 1); + int64_t y_id = ray_sym_intern("y", 1); + int64_t old_x_raw[] = {1, 2}; + int64_t old_y_raw[] = {10, 20}; + ray_t* old_x = ray_vec_from_raw(RAY_I64, old_x_raw, 2); + ray_t* old_y = ray_vec_from_raw(RAY_I64, old_y_raw, 2); + ray_t* old = ray_table_new(2); + old = ray_table_add_col(old, x_id, old_x); + old = ray_table_add_col(old, y_id, old_y); + TEST_ASSERT_FALSE(RAY_IS_ERR(old)); + TEST_ASSERT_EQ_I(ray_splay_save(old, dir, NULL), RAY_OK); + TEST_ASSERT_EQ_I(access(manifest, F_OK), -1); + + int64_t new_x_raw[] = {3, 4}; + int64_t new_y_raw[] = {30, 40}; + ray_t* new_x = ray_vec_from_raw(RAY_I64, new_x_raw, 2); + ray_t* new_y = ray_vec_from_raw(RAY_I64, new_y_raw, 2); + ray_t* replacement = ray_table_new(2); + replacement = ray_table_add_col(replacement, x_id, new_x); + replacement = ray_table_add_col(replacement, y_id, new_y); + TEST_ASSERT_FALSE(RAY_IS_ERR(replacement)); + TEST_ASSERT_EQ_I(ray_splay_save(replacement, dir, NULL), RAY_OK); + TEST_ASSERT_EQ_I(access(manifest, F_OK), 0); + + ray_t* loaded = ray_splay_load(dir, NULL); + TEST_ASSERT_NOT_NULL(loaded); + TEST_ASSERT_FALSE(RAY_IS_ERR(loaded)); + ray_t* loaded_x = ray_table_get_col(loaded, x_id); + ray_t* loaded_y = ray_table_get_col(loaded, y_id); + TEST_ASSERT_NOT_NULL(loaded_x); + TEST_ASSERT_NOT_NULL(loaded_y); + TEST_ASSERT_EQ_I(((int64_t*)ray_data(loaded_x))[0], 3); + TEST_ASSERT_EQ_I(((int64_t*)ray_data(loaded_y))[0], 30); + ray_release(loaded); + + ray_t* bad_col = ray_dict_new(ray_vec_from_raw(RAY_I64, old_x_raw, 2), + ray_vec_from_raw(RAY_I64, old_y_raw, 2)); + ray_t* bad = ray_table_new(2); + bad = ray_table_add_col(bad, x_id, old_x); + bad = ray_table_add_col(bad, y_id, bad_col); + ray_err_t bad_err = ray_splay_save(bad, dir, NULL); + TEST_ASSERT_TRUE(bad_err != RAY_OK); + + loaded = ray_read_splayed(dir, NULL); + TEST_ASSERT_NOT_NULL(loaded); + TEST_ASSERT_FALSE(RAY_IS_ERR(loaded)); + loaded_x = ray_table_get_col(loaded, x_id); + loaded_y = ray_table_get_col(loaded, y_id); + TEST_ASSERT_NOT_NULL(loaded_x); + TEST_ASSERT_NOT_NULL(loaded_y); + TEST_ASSERT_EQ_I(((int64_t*)ray_data(loaded_x))[0], 3); + TEST_ASSERT_EQ_I(((int64_t*)ray_data(loaded_y))[0], 30); + ray_release(loaded); + + ray_release(bad); + ray_release(bad_col); + ray_release(replacement); + ray_release(new_x); + ray_release(new_y); + ray_release(old); + ray_release(old_x); + ray_release(old_y); + (void)ray_test_rm_rf(dir); + PASS(); +} + const test_entry_t splay_entries[] = { + { "splay/atomic_generation_publish", test_splay_atomic_generation_publish, splay_setup, splay_teardown }, { "splay/has_nulls_roundtrip", test_splayed_has_nulls_roundtrip, splay_setup, splay_teardown }, { "splay/save_null_dir", test_save_null_dir, splay_setup, splay_teardown }, { "splay/save_null_tbl", test_save_null_tbl, splay_setup, splay_teardown }, From 9e3069e65649e0c27ca19db928c5119783eceb4e Mon Sep 17 00:00:00 2001 From: belowzeroff Date: Mon, 28 Sep 2026 13:46:56 -0400 Subject: [PATCH 2/5] fix(store): complete atomic splay publication --- .github/workflows/ci.yml | 11 +- docs/docs/storage/index.md | 20 ++ src/io/csv.c | 40 +++- src/ops/builtins.c | 15 +- src/ops/system.c | 5 - src/store/part.c | 17 +- src/store/splay.c | 256 ++++++++++++++++-------- src/store/splay.h | 25 ++- test/rfl/storage/atomic_generations.rfl | 47 +++++ test/test_splay.c | 212 +++++++++++++++++++- test/test_stress_matrix.c | 11 +- 11 files changed, 537 insertions(+), 122 deletions(-) create mode 100644 test/rfl/storage/atomic_generations.rfl diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 426878bc4..79cad108a 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -129,10 +129,19 @@ jobs: timeout-minutes: 30 steps: - uses: actions/checkout@v4 + with: + fetch-depth: 0 - name: Install cppcheck run: sudo apt-get update && sudo apt-get install -y cppcheck - name: cppcheck - run: make cppcheck + run: | + changed="$(git diff --name-only "origin/${{ github.base_ref || 'dev' }}"...HEAD -- 'src/**/*.c' | tr '\n' ' ')" + if [ -z "$changed" ]; then + echo "No changed C sources for cppcheck." + exit 0 + fi + echo "cppcheck files: $changed" + make cppcheck FILES="$changed" CPPCHECK_SKIP= # Single stable check to require in branch protection, instead of the four # brittle matrix names (which break if the matrix is renamed). `if: always()` diff --git a/docs/docs/storage/index.md b/docs/docs/storage/index.md index 88f165b02..335b8f625 100644 --- a/docs/docs/storage/index.md +++ b/docs/docs/storage/index.md @@ -75,6 +75,26 @@ before publication leaves the previous generation selected. Tables written by older Rayforce versions remain readable; the first replacement upgrades the directory to the generation layout. +The table, partition and CSV entry points follow the same publication protocol. +Indexes are built in the new generation before it is published. Writers to the +same table serialize on `.write.lock`; readers resolve the manifest once and +continue using that generation. The symbol vocabulary stays at its original +location and grows append-only. + +Previous generations, the legacy files and incomplete staging directories are +retained. Disk usage therefore grows with replacements. There is no automatic +garbage collection yet: reclaim old files only during offline maintenance with +all readers and writers stopped, and preserve the generation named by `.current`. +Older binaries that do not understand `.current` must not access an upgraded +table; reading its legacy files directly returns obsolete data. + +`ray_splay_save` (including `.db.splayed.set`) syncs data, indexes and directory +entries before acknowledging publication. The bulk C API and CSV import retain +their no-fsync contract: publication is atomic for readers, but is not a +power-loss durability guarantee. An I/O error after the manifest rename (during +the final directory sync) has an uncertain commit outcome; reopen the table to +inspect the selected generation. + This protects publication of one splayed table. Coordinating a consistent snapshot across multiple tables or date partitions still belongs to the caller. diff --git a/src/io/csv.c b/src/io/csv.c index 5468ebf7c..48882440a 100644 --- a/src/io/csv.c +++ b/src/io/csv.c @@ -3637,12 +3637,13 @@ static void csv_splayed_append_task(void* raw, uint32_t wid, int64_t start, int6 } } -ray_err_t ray_csv_save_splayed_named_opts(const char* path, char delimiter, bool header, +static ray_err_t csv_save_splayed_to_dir(const char* path, char delimiter, bool header, const int8_t* col_types_in, int32_t n_types, const int64_t* col_names_in, int32_t n_names, - const char* dir, int64_t rows_per_chunk) { + const char* dir, int64_t rows_per_chunk, + const char* sym_path) { if (ray_interrupted()) return RAY_ERR_CANCEL; - if (!path || !dir) return RAY_ERR_DOMAIN; + if (!path || !dir || !sym_path) return RAY_ERR_DOMAIN; if (rows_per_chunk <= 0) rows_per_chunk = CSV_PART_ROWS_DEFAULT; int fd = open(path, O_RDONLY); @@ -3795,12 +3796,6 @@ ray_err_t ray_csv_save_splayed_named_opts(const char* path, char delimiter, bool for (int c = 0; c < ncols; c++) if (resolved_types[c] == RAY_SYM) any_sym = true; if (any_sym) { - char sym_path[1024]; - int n = snprintf(sym_path, sizeof(sym_path), "%s/.sym", dir); - if (n < 0 || (size_t)n >= sizeof(sym_path)) { - ray_vm_unmap_file(buf, file_size); - return RAY_ERR_RANGE; - } sym_dom = ray_sym_domain_open_or_create(sym_path); if (!sym_dom) { ray_vm_unmap_file(buf, file_size); @@ -3951,6 +3946,33 @@ ray_err_t ray_csv_save_splayed_named_opts(const char* path, char delimiter, bool return err; } +ray_err_t ray_csv_save_splayed_named_opts(const char* path, char delimiter, bool header, + const int8_t* col_types_in, int32_t n_types, + const int64_t* col_names_in, int32_t n_names, + const char* dir, int64_t rows_per_chunk) { + if (ray_interrupted()) return RAY_ERR_CANCEL; + if (!path || !dir) return RAY_ERR_DOMAIN; + char sym_path[1024]; + int n = snprintf(sym_path, sizeof(sym_path), "%s/.sym", dir); + if (n < 0 || (size_t)n >= sizeof(sym_path)) return RAY_ERR_RANGE; + ray_splay_write_t write; + ray_err_t err = ray_splay_write_begin(dir, &write); + if (err != RAY_OK) return err; + err = csv_save_splayed_to_dir(path, delimiter, header, col_types_in, n_types, + col_names_in, n_names, write.dir, + rows_per_chunk, sym_path); + if (err == RAY_OK) { + ray_t* tbl = ray_read_splayed(write.dir, sym_path); + if (!tbl || RAY_IS_ERR(tbl)) { + err = tbl ? ray_err_from_obj(tbl) : RAY_ERR_OOM; + } else { + ray_splay_build_indexes(write.dir, tbl); + } + if (tbl) ray_release(tbl); + } + return ray_splay_write_finish(&write, err, false); +} + static ray_err_t csv_save_parted_impl(const char* path, char delimiter, bool header, const int8_t* col_types_in, int32_t n_types, const int64_t* col_names_in, int32_t n_names, diff --git a/src/ops/builtins.c b/src/ops/builtins.c index 1400212e9..4d09938f2 100644 --- a/src/ops/builtins.c +++ b/src/ops/builtins.c @@ -686,19 +686,8 @@ ray_t* ray_read_csv_splayed_fn(ray_t** args, int64_t n) { const char* sym = csv_default_sym_path(dir, sym_path, sizeof(sym_path)); if (!sym) return ray_error("io", NULL); - /* The streaming writer emits raw columns; append chunk-zone indexes to the - * just-written files, then reload so the returned table carries them - * (mmap'd in place). Conversion is the ONLY place a CSV load decides an - * index: `.csv.read` returns an index-free in-memory table, and callers - * that want one on it ask explicitly (.idx.hash / the attrs verbs). - * Without this pass, a converted store would have no block-skip at all. */ - ray_t* tbl = ray_read_splayed(dir, sym); - if (tbl && !RAY_IS_ERR(tbl) && tbl->type == RAY_TABLE) { - ray_splay_build_indexes(dir, tbl); - ray_release(tbl); - tbl = ray_read_splayed(dir, sym); - } - return tbl; + /* The writer builds indexes before publishing its generation. */ + return ray_read_splayed(dir, sym); } ray_t* ray_read_csv_parted_fn(ray_t** args, int64_t n) { diff --git a/src/ops/system.c b/src/ops/system.c index 3f153cb3b..d071b0dd4 100644 --- a/src/ops/system.c +++ b/src/ops/system.c @@ -197,11 +197,6 @@ ray_t* ray_set_splayed_fn(ray_t** args, int64_t n) { ray_err_t err = ray_splay_save(tbl, dir, sym_path); if (err != RAY_OK) return ray_error(ray_err_code_str(err), NULL); - /* Build + persist accelerator indexes (STR dictionaries, numeric chunk-zone - * min/max) inline at each column file's tail so mmap loads get the fast - * paths — same pass .csv.splayed runs. */ - ray_splay_build_indexes(dir, tbl); - ray_retain(tbl); return tbl; } diff --git a/src/store/part.c b/src/store/part.c index e536a7e7a..1a7109c8f 100644 --- a/src/store/part.c +++ b/src/store/part.c @@ -344,6 +344,9 @@ ray_t* ray_read_parted(const char* db_root, const char* table_name) { db_root, part_dirs[p], table_name); goto fail_tables; } + /* Resolve the generation before refreshing the shared symbol domain. + * ray_read_splayed_dom reopens the cached domain by path, so symbols + * appended by a concurrent publisher are visible before columns load. */ part_tables[p] = ray_read_splayed_dom(path, dom); if (!part_tables[p] || RAY_IS_ERR(part_tables[p])) { if (trace) @@ -566,8 +569,11 @@ static ray_err_t collect_table_dirs(const char* pdir, char*** out_names, struct dirent* ent; while ((ent = readdir(d)) != NULL) { if (ent->d_name[0] == '.') continue; /* ".", "..", and the ".sym" dotfile */ - char dpath[1024]; - int dn = snprintf(dpath, sizeof(dpath), "%s/%s/.d", pdir, ent->d_name); + char tdir[1024], resolved[1024], dpath[1100]; + int dn = snprintf(tdir, sizeof(tdir), "%s/%s", pdir, ent->d_name); + if (dn < 0 || (size_t)dn >= sizeof(tdir)) continue; + if (ray_splay_resolve_dir(tdir, resolved, sizeof(resolved)) != RAY_OK) continue; + dn = snprintf(dpath, sizeof(dpath), "%s/.d", resolved); if (dn < 0 || (size_t)dn >= sizeof(dpath)) continue; struct stat st; if (stat(dpath, &st) != 0 || !S_ISREG(st.st_mode)) continue; /* not a table */ @@ -606,8 +612,11 @@ static ray_err_t collect_table_dirs(const char* pdir, char*** out_names, /* Does partition `part` contain splayed table `tname` (i.e. a `.d` schema)? */ static bool partition_has_table(const char* db_root, const char* part, const char* tname) { - char dpath[1024]; - int n = snprintf(dpath, sizeof(dpath), "%s/%s/%s/.d", db_root, part, tname); + char tdir[1024], resolved[1024], dpath[1100]; + int n = snprintf(tdir, sizeof(tdir), "%s/%s/%s", db_root, part, tname); + if (n < 0 || (size_t)n >= sizeof(tdir)) return false; + if (ray_splay_resolve_dir(tdir, resolved, sizeof(resolved)) != RAY_OK) return false; + n = snprintf(dpath, sizeof(dpath), "%s/.d", resolved); if (n < 0 || (size_t)n >= sizeof(dpath)) return false; struct stat st; return stat(dpath, &st) == 0 && S_ISREG(st.st_mode); diff --git a/src/store/splay.c b/src/store/splay.c index ac920984b..6c5e15a86 100644 --- a/src/store/splay.c +++ b/src/store/splay.c @@ -44,6 +44,7 @@ #include #include #include +#include /* -------------------------------------------------------------------------- * Splayed table: directory of column files + .d schema file @@ -152,7 +153,7 @@ static void splay_sweep_stale(ray_t* tbl, const char* dir) { * are written below .generations/ first, then .current is replaced with * one filesystem rename. Readers that predate this format continue to use * the legacy directory when .current is absent. */ -static unsigned long splay_generation_seq; +static _Atomic unsigned long splay_generation_seq; static ray_err_t splay_current_dir(const char* dir, char* out, size_t out_sz, bool* active) { @@ -162,25 +163,27 @@ static ray_err_t splay_current_dir(const char* dir, char* out, size_t out_sz, return RAY_ERR_RANGE; *active = false; - struct stat st; - if (stat(manifest, &st) != 0) { + FILE* f = fopen(manifest, "rb"); + if (!f) { if (errno == ENOENT) return RAY_OK; /* legacy splayed directory */ return RAY_ERR_IO; } - FILE* f = fopen(manifest, "rb"); - if (!f) return RAY_ERR_IO; - char rel[512]; - int scanned = fscanf(f, "%511s", rel); + size_t len = fread(rel, 1, sizeof(rel), f); + bool failed = ferror(f) != 0; int closed = fclose(f); - if (scanned != 1 || closed != 0) { + if (failed || closed != 0 || len == 0 || len >= sizeof(rel)) { return RAY_ERR_CORRUPT; } + if (rel[len - 1] == '\n') len--; + rel[len] = '\0'; if (strncmp(rel, ".generations/", 13) != 0 || - strstr(rel, "..") || strstr(rel, "\\") || - strcmp(rel, ".generations/") == 0) + len <= 13 || memchr(rel, '\0', len)) return RAY_ERR_CORRUPT; + for (size_t i = 13; i < len; i++) + if (!((rel[i] >= '0' && rel[i] <= '9') || rel[i] == 'g' || rel[i] == '-')) + return RAY_ERR_CORRUPT; n = snprintf(out, out_sz, "%s/%s", dir, rel); if (n < 0 || (size_t)n >= out_sz) return RAY_ERR_RANGE; @@ -188,13 +191,22 @@ static ray_err_t splay_current_dir(const char* dir, char* out, size_t out_sz, return RAY_OK; } -static ray_err_t splay_publish_generation(const char* dir, const char* gen) { +ray_err_t ray_splay_resolve_dir(const char* dir, char* out, size_t out_sz) { + if (!dir || !out || !out_sz) return RAY_ERR_IO; + bool active; + ray_err_t err = splay_current_dir(dir, out, out_sz, &active); + if (err != RAY_OK || active) return err; + int n = snprintf(out, out_sz, "%s", dir); + return n < 0 || (size_t)n >= out_sz ? RAY_ERR_RANGE : RAY_OK; +} + +static ray_err_t splay_publish_generation(const char* dir, const char* gen, + bool durable) { char manifest[1024], tmp[1024]; int n = snprintf(manifest, sizeof(manifest), "%s/.current", dir); if (n < 0 || (size_t)n >= sizeof(manifest)) return RAY_ERR_RANGE; - n = snprintf(tmp, sizeof(tmp), "%s/.current.tmp.%lu.%lu.%lu", dir, - (unsigned long)time(NULL), (unsigned long)getpid(), - ++splay_generation_seq); + /* The exclusively created generation owns this temporary file. */ + n = snprintf(tmp, sizeof(tmp), "%s/%s/.current.tmp", dir, gen); if (n < 0 || (size_t)n >= sizeof(tmp)) return RAY_ERR_RANGE; FILE* f = fopen(tmp, "wb"); @@ -213,26 +225,27 @@ static ray_err_t splay_publish_generation(const char* dir, const char* gen) { (void)remove(tmp); return RAY_ERR_IO; } - ray_err_t err = ray_file_sync(fd); + ray_err_t err = durable ? ray_file_sync(fd) : RAY_OK; ray_file_close(fd); if (err != RAY_OK || ray_file_rename(tmp, manifest) != RAY_OK) { (void)remove(tmp); return RAY_ERR_IO; } - return ray_file_sync_dir(manifest); + return durable ? ray_file_sync_dir(manifest) : RAY_OK; } -static bool splay_has_file(const char* dir, const char* name) { +static ray_err_t splay_has_file(const char* dir, const char* name, bool* exists) { char path[1024]; int n = snprintf(path, sizeof(path), "%s/%s", dir, name); - if (n < 0 || (size_t)n >= sizeof(path)) return false; + if (n < 0 || (size_t)n >= sizeof(path)) return RAY_ERR_RANGE; struct stat st; - return stat(path, &st) == 0; + *exists = stat(path, &st) == 0; + return *exists || errno == ENOENT ? RAY_OK : RAY_ERR_IO; } -static ray_err_t splay_save_to_dir_impl(ray_t* tbl, const char* dir, - const char* sym_path, bool durable) { - if (!tbl || RAY_IS_ERR(tbl)) return RAY_ERR_TYPE; +static ray_err_t splay_validate_save(ray_t* tbl, const char* dir, + const char* sym_path) { + if (!tbl || RAY_IS_ERR(tbl) || tbl->type != RAY_TABLE) return RAY_ERR_TYPE; if (!dir) return RAY_ERR_IO; ray_err_t name_err = splay_validate_persisted_names(tbl); @@ -286,6 +299,37 @@ static ray_err_t splay_save_to_dir_impl(ray_t* tbl, const char* dir, } } + return RAY_OK; +} + +static ray_err_t splay_sync_files(const char* dir) { + DIR* d = opendir(dir); + if (!d) return RAY_ERR_IO; + ray_err_t err = RAY_OK; + struct dirent* entry; + while ((entry = readdir(d))) { + if (strcmp(entry->d_name, ".") == 0 || strcmp(entry->d_name, "..") == 0) + continue; + char path[1024]; + int n = snprintf(path, sizeof(path), "%s/%s", dir, entry->d_name); + if (n < 0 || (size_t)n >= sizeof(path)) { err = RAY_ERR_RANGE; break; } + struct stat st; + if (stat(path, &st) != 0) { err = RAY_ERR_IO; break; } + if (!S_ISREG(st.st_mode)) continue; + ray_fd_t fd = ray_file_open(path, RAY_OPEN_READ | RAY_OPEN_WRITE); + if (fd == RAY_FD_INVALID) { err = RAY_ERR_IO; break; } + err = ray_file_sync(fd); + ray_file_close(fd); + if (err != RAY_OK) break; + } + closedir(d); + return err; +} + +ray_err_t ray_splay_write_table(ray_t* tbl, const char* dir, + const char* sym_path, bool durable) { + ray_err_t validation = splay_validate_save(tbl, dir, sym_path); + if (validation != RAY_OK) return validation; /* Create directory and any missing parents (mkdir -p semantics). * Required for partitioned layouts like "/db/2024.01.01/t/" where the * caller hasn't pre-created the date partition. */ @@ -348,13 +392,7 @@ static ray_err_t splay_save_to_dir_impl(ray_t* tbl, const char* dir, int64_t ncols = ray_table_ncols(tbl); - /* 2. Column files (and the schema names they correspond to). - * NOTE: overwriting an existing dir rewrites columns in place; a crash - * mid-loop leaves old .d + a mix of old/new column files. Ragged - * lengths are caught at load (column-length check); equal-length - * mixed-generation rows are inherent to in-place overwrite and would - * need staged writes — out of scope (fresh-dir crashes degrade to a - * visibly missing table via the .d-last commit marker). */ + /* 2. Write the column files in the caller's unpublished directory. */ ray_t* schema = ray_vec_new(RAY_STR, ncols > 0 ? ncols : 1); if (!schema || RAY_IS_ERR(schema)) { if (schema) ray_release(schema); @@ -385,9 +423,7 @@ static ray_err_t splay_save_to_dir_impl(ray_t* tbl, const char* dir, : (durable ? ray_col_save(col, path) : ray_col_save_bulk(col, path)); if (err != RAY_OK) { - /* No new .d is written. Preflight prevents deterministic format - * failures here; an I/O error or crash can still leave an - * existing directory with mixed-generation column files. */ + /* No new .d or .current is published on failure. */ ray_release(schema); if (dom) ray_sym_domain_release(dom); return err; @@ -401,6 +437,13 @@ static ray_err_t splay_save_to_dir_impl(ray_t* tbl, const char* dir, } if (dom) ray_sym_domain_release(dom); + /* Indexes belong to this generation and must precede publication. */ + ray_splay_build_indexes(dir, tbl); + if (durable) { + ray_err_t err = splay_sync_files(dir); + if (err != RAY_OK) { ray_release(schema); return err; } + } + /* 3. .d LAST — the commit marker. */ { char path[1024]; @@ -421,52 +464,91 @@ static ray_err_t splay_save_to_dir_impl(ray_t* tbl, const char* dir, return RAY_OK; } -static ray_err_t splay_save_impl(ray_t* tbl, const char* dir, - const char* sym_path, bool durable) { - if (!tbl || RAY_IS_ERR(tbl)) return RAY_ERR_TYPE; - if (!dir) return RAY_ERR_IO; - - /* Keep the first write backward-compatible. Once a table has a - * committed schema (or already uses generations), every replacement is - * staged out-of-place and published through .current. */ - bool atomic = splay_has_file(dir, ".d") || splay_has_file(dir, ".current"); - if (!atomic) - return splay_save_to_dir_impl(tbl, dir, sym_path, durable); - - ray_err_t err = ray_mkdir_p(dir); - if (err != RAY_OK) return err; +ray_err_t ray_splay_write_finish(ray_splay_write_t* write, ray_err_t result, + bool durable) { + if (result == RAY_OK && write->staged) { + /* ray_file_sync_dir syncs the PARENT of its argument. */ + char schema[1100]; + snprintf(schema, sizeof(schema), "%s/.d", write->dir); + if (durable && (ray_file_sync_dir(schema) != RAY_OK || + ray_file_sync_dir(write->dir) != RAY_OK)) + result = RAY_ERR_IO; + if (result == RAY_OK) + result = splay_publish_generation(write->root, write->generation, durable); + } + if (write->lock != RAY_FD_INVALID) { + (void)ray_file_unlock(write->lock); + ray_file_close(write->lock); + write->lock = RAY_FD_INVALID; + } + /* A reader may have resolved the old path without opening every column. + * Keep both legacy files and old generations immutable after publication. */ + return result; +} - char generations[1024], gen_name[256], gen_dir[1024]; - int n = snprintf(generations, sizeof(generations), "%s/.generations", dir); - if (n < 0 || (size_t)n >= sizeof(generations)) return RAY_ERR_RANGE; - err = ray_mkdir_p(generations); +ray_err_t ray_splay_write_begin(const char* dir, ray_splay_write_t* write) { + memset(write, 0, sizeof(*write)); + write->lock = RAY_FD_INVALID; + if (!dir || !*dir) return RAY_ERR_IO; + int n = snprintf(write->root, sizeof(write->root), "%s", dir); + if (n < 0 || (size_t)n >= sizeof(write->root)) return RAY_ERR_RANGE; + size_t len = strlen(write->root); + while (len > 1 && write->root[len - 1] == '/') write->root[--len] = '\0'; + ray_err_t err = ray_mkdir_p(write->root); if (err != RAY_OK) return err; - n = snprintf(gen_name, sizeof(gen_name), ".generations/g-%lu-%lu-%lu", - (unsigned long)time(NULL), (unsigned long)getpid(), - ++splay_generation_seq); - if (n < 0 || (size_t)n >= sizeof(gen_name)) return RAY_ERR_RANGE; - n = snprintf(gen_dir, sizeof(gen_dir), "%s/%s", dir, gen_name); - if (n < 0 || (size_t)n >= sizeof(gen_dir)) return RAY_ERR_RANGE; + char path[1024]; + n = snprintf(path, sizeof(path), "%s/.write.lock", write->root); + if (n < 0 || (size_t)n >= sizeof(path)) return RAY_ERR_RANGE; + write->lock = ray_file_open(path, RAY_OPEN_READ | RAY_OPEN_WRITE | RAY_OPEN_CREATE); + if (write->lock == RAY_FD_INVALID) return RAY_ERR_IO; + err = ray_file_lock_ex(write->lock); + if (err != RAY_OK) return ray_splay_write_finish(write, err, false); + + bool schema, current; + err = splay_has_file(write->root, ".d", &schema); + if (err == RAY_OK) err = splay_has_file(write->root, ".current", ¤t); + if (err != RAY_OK) return ray_splay_write_finish(write, err, false); + write->staged = schema || current; + if (!write->staged) { + memcpy(write->dir, write->root, len + 1); + return RAY_OK; + } + n = snprintf(path, sizeof(path), "%s/.generations", write->root); + if (n < 0 || (size_t)n >= sizeof(path)) + return ray_splay_write_finish(write, RAY_ERR_RANGE, false); + err = ray_mkdir_p(path); + if (err != RAY_OK) return ray_splay_write_finish(write, err, false); + + for (;;) { + unsigned long seq = atomic_fetch_add(&splay_generation_seq, 1); + snprintf(write->generation, sizeof(write->generation), ".generations/g-%lu-%lu-%lu", + (unsigned long)time(NULL), (unsigned long)getpid(), seq); + n = snprintf(write->dir, sizeof(write->dir), "%s/%s", write->root, write->generation); + if (n < 0 || (size_t)n >= sizeof(write->dir)) + return ray_splay_write_finish(write, RAY_ERR_RANGE, false); +#ifdef RAY_OS_WINDOWS + if (CreateDirectoryA(write->dir, NULL)) break; + if (GetLastError() != ERROR_ALREADY_EXISTS) +#else + if (mkdir(write->dir, 0755) == 0) break; + if (errno != EEXIST) +#endif + return ray_splay_write_finish(write, RAY_ERR_IO, false); + /* Never reuse an existing generation, including after PID reuse. */ + } + return RAY_OK; +} - err = splay_save_to_dir_impl(tbl, gen_dir, sym_path, durable); +static ray_err_t splay_save_impl(ray_t* tbl, const char* dir, + const char* sym_path, bool durable) { + ray_err_t err = splay_validate_save(tbl, dir, sym_path); if (err != RAY_OK) return err; - if (durable && ray_file_sync_dir(gen_dir) != RAY_OK) - return RAY_ERR_IO; - if (ray_file_sync_dir(generations) != RAY_OK) - return RAY_ERR_IO; - - /* This is the table-level commit point. A failure before this rename - * leaves the old .current (and therefore the old complete generation) - * untouched. */ - err = splay_publish_generation(dir, gen_name); + ray_splay_write_t write; + err = ray_splay_write_begin(dir, &write); if (err != RAY_OK) return err; - - /* Legacy files are no longer read once .current exists. Remove stale - * columns opportunistically for existing callers that inspect the old - * directory layout; failure here cannot invalidate the committed table. */ - splay_sweep_stale(tbl, dir); - return RAY_OK; + err = ray_splay_write_table(tbl, write.dir, sym_path, durable); + return ray_splay_write_finish(&write, err, durable); } ray_err_t ray_splay_save(ray_t* tbl, const char* dir, const char* sym_path) { @@ -622,13 +704,12 @@ static ray_t* splay_load_dom_impl(const char* dir, ray_sym_domain_t* dom, static ray_t* splay_load_impl(const char* dir, const char* sym_path, bool use_mmap) { char resolved[1024]; - bool active = false; - ray_err_t current_err = splay_current_dir(dir, resolved, sizeof(resolved), - &active); - if (current_err != RAY_OK) - return ray_error("corrupt", "splayed %s: invalid .current manifest", dir); - - const char* load_dir = active ? resolved : dir; + ray_err_t err = ray_splay_resolve_dir(dir, resolved, sizeof(resolved)); + if (err != RAY_OK) + return ray_error(ray_err_code_str(err), "cannot resolve splayed generation"); + /* Resolve first: a newly published generation can reference symbols + * appended since a previous domain open. Opening afterwards refreshes + * the cached domain before any of those column codes are validated. */ ray_sym_domain_t* dom = NULL; if (sym_path) { struct stat st; @@ -640,7 +721,7 @@ static ray_t* splay_load_impl(const char* dir, const char* sym_path, "record, or missing \"\" at position 0)", sym_path); } } - ray_t* tbl = splay_load_dom_impl(load_dir, dom, use_mmap); + ray_t* tbl = splay_load_dom_impl(resolved, dom, use_mmap); if (dom) ray_sym_domain_release(dom); /* columns hold their own refs */ return tbl; } @@ -861,5 +942,18 @@ ray_t* ray_read_splayed(const char* dir, const char* sym_path) { } ray_t* ray_read_splayed_dom(const char* dir, struct ray_sym_domain_s* dom) { - return splay_load_dom_impl(dir, dom, true); + char resolved[1024]; + ray_err_t err = ray_splay_resolve_dir(dir, resolved, sizeof(resolved)); + if (err != RAY_OK) + return ray_error(ray_err_code_str(err), "cannot resolve splayed generation"); + ray_sym_domain_t* fresh = NULL; + const char* path = dom ? ray_sym_domain_path(dom) : NULL; + if (path) { + fresh = ray_sym_domain_open(path); + if (!fresh) return ray_error("corrupt", "cannot refresh splayed symfile"); + dom = fresh; + } + ray_t* tbl = splay_load_dom_impl(resolved, dom, true); + if (fresh) ray_sym_domain_release(fresh); + return tbl; } diff --git a/src/store/splay.h b/src/store/splay.h index 49a9e2aee..6e6722b2e 100644 --- a/src/store/splay.h +++ b/src/store/splay.h @@ -25,9 +25,30 @@ #define RAY_SPLAY_H #include +#include "store/fileio.h" struct ray_sym_domain_s; +/* Internal publication protocol shared by the table and streaming CSV writers. + * begin serializes writers; finish publishes only on success and always unlocks. + * Old generations must remain available to readers that already resolved them. */ +typedef struct { + ray_fd_t lock; + bool staged; + char root[1024]; + char dir[1024]; + char generation[256]; +} ray_splay_write_t; + +ray_err_t ray_splay_write_begin(const char* dir, ray_splay_write_t* write); +ray_err_t ray_splay_write_finish(ray_splay_write_t* write, ray_err_t result, + bool durable); +/* Write only to a fresh/unpublished directory owned by the caller. */ +ray_err_t ray_splay_write_table(ray_t* tbl, const char* dir, + const char* sym_path, bool durable); +/* Resolve once and retain the returned path for the entire read. */ +ray_err_t ray_splay_resolve_dir(const char* dir, char* out, size_t out_sz); + /* Splayed table I/O. * * sym_path names the table's symfile (domain): save distinct-merges the @@ -47,8 +68,8 @@ ray_t* ray_read_splayed(const char* dir, const char* sym_path); * files so later mmap loads get block-skip. Best-effort, per numeric column. */ void ray_splay_build_indexes(const char* dir, ray_t* tbl); -/* Partition loader entry: the parted reader opens root/sym ONCE and - * passes the shared domain to every partition's columns. */ +/* Loader accepting a shared FILE domain. It resolves the generation first, + * then refreshes the cached domain to include any externally appended symbols. */ ray_t* ray_read_splayed_dom(const char* dir, struct ray_sym_domain_s* dom); #endif /* RAY_SPLAY_H */ diff --git a/test/rfl/storage/atomic_generations.rfl b/test/rfl/storage/atomic_generations.rfl new file mode 100644 index 000000000..48717f4ee --- /dev/null +++ b/test/rfl/storage/atomic_generations.rfl @@ -0,0 +1,47 @@ +;; Replacements must be visible through every public storage entry point. +(.sys.exec "rm -rf /tmp/rf_test_atomic_generations") -- 0 +(set agRoot "/tmp/rf_test_atomic_generations") +(set agLeaf (format "%/2026.09.28/t" agRoot)) +(set agOld (table ['x 'y] (list [1 2] [10 20]))) +(set agNew (table ['x 'y] (list [3 4] [30 40]))) +(.db.splayed.set agLeaf agOld) +(.db.splayed.set agLeaf agNew) +(sum (at (.db.splayed.get agLeaf) 'x)) -- 7 +(sum (at (.db.parted.get agRoot 't) 'x)) -- 7 +;; Discovery must use the active schema even after offline legacy cleanup. +(.sys.exec "rm /tmp/rf_test_atomic_generations/2026.09.28/t/.d") -- 0 +(.db.parted.tables agRoot) -- ['t] + +;; Schema narrowing must not make the partition reader use the old .d. +(.db.splayed.set agLeaf (table ['x] (list [5 6]))) +(count (key (.db.parted.get agRoot 't))) -- 2 +(sum (at (.db.parted.get agRoot 't) 'x)) -- 11 + +;; The streaming (numeric) CSV writer must publish a new generation too. +(.csv.write agNew (format "%/new.csv" agRoot)) +(set agCSV (.csv.splayed (format "%/new.csv" agRoot) agLeaf)) +(sum (at agCSV 'x)) -- 7 +(sum (at (.db.splayed.get agLeaf) 'y)) -- 70 +(sum (at (.db.parted.get agRoot 't) 'y)) -- 70 + +;; Its materializing STR fallback uses the same protocol and root symfile. +(set agText (table ['x 'text 'sym] (list [7 8] ["alpha" "beta"] ['A 'B]))) +(.csv.write agText (format "%/text.csv" agRoot)) +(set agCSV (.csv.splayed [I64 STR SYMBOL] (format "%/text.csv" agRoot) (format "%/text" agRoot))) +(set agText (table ['x 'text 'sym] (list [9 10] ["gamma" "delta"] ['C 'D]))) +(.csv.write agText (format "%/text.csv" agRoot)) +(set agCSV (.csv.splayed [I64 STR SYMBOL] (format "%/text.csv" agRoot) (format "%/text" agRoot))) +(at agCSV 'x) -- [9 10] +(at agCSV 'text) -- ["gamma" "delta"] +(at agCSV 'sym) -- ['C 'D] + +;; Persist indexes in the generation being published, not in legacy files. +(set agIndexed (table ['k 'v] (list (.attr.set 'grouped ['a 'b]) [1 2]))) +(.db.splayed.set (format "%/indexed" agRoot) agIndexed) +(.db.splayed.set (format "%/indexed" agRoot) agIndexed) +(.idx.has? (at (.db.splayed.get (format "%/indexed" agRoot)) 'k)) -- true +(set agStrings (table ['k] (list (.attr.set 'grouped ["aa" "bb"])))) +(.db.splayed.set (format "%/indexed_str" agRoot) agStrings) +(.db.splayed.set (format "%/indexed_str" agRoot) agStrings) +(.idx.has? (at (.db.splayed.get (format "%/indexed_str" agRoot)) 'k)) -- true +(.sys.exec "rm -rf /tmp/rf_test_atomic_generations") -- 0 diff --git a/test/test_splay.c b/test/test_splay.c index 784c702e2..2fef9d7fd 100644 --- a/test/test_splay.c +++ b/test/test_splay.c @@ -1309,8 +1309,8 @@ static test_result_t test_selfdescribing_schema_fresh_process(void) { } /* ========================================================================= - * 29. Crash-safe save: (a) re-set with a NARROWER schema removes stale - * column files; (b) symbol-free tables write no symfile even when a + * 29. Crash-safe save: (a) re-set with a NARROWER schema retains old + * column files for readers; (b) symbol-free tables write no symfile even when a * sym_path is supplied; (c) .d is the commit marker — written last. * ========================================================================= */ static test_result_t test_save_sweeps_stale_and_skips_sym(void) { @@ -1343,10 +1343,10 @@ static test_result_t test_save_sweeps_stale_and_skips_sym(void) { TEST_ASSERT_FALSE(RAY_IS_ERR(narrow)); TEST_ASSERT_EQ_I(ray_splay_save(narrow, dir, NULL), RAY_OK); - /* (a) stale column file "b" must be gone; load sees 1 column */ + /* Old readers can still open b; the new generation has only a. */ char bpath[512]; snprintf(bpath, sizeof(bpath), "%s/b", dir); - TEST_ASSERT_EQ_I(access(bpath, F_OK), -1); + TEST_ASSERT_EQ_I(access(bpath, F_OK), 0); ray_t* loaded = ray_splay_load(dir, NULL); TEST_ASSERT_NOT_NULL(loaded); TEST_ASSERT_FALSE(RAY_IS_ERR(loaded)); @@ -2421,9 +2421,11 @@ static test_result_t test_splay_atomic_generation_publish(void) { ray_vec_from_raw(RAY_I64, old_y_raw, 2)); ray_t* bad = ray_table_new(2); bad = ray_table_add_col(bad, x_id, old_x); - bad = ray_table_add_col(bad, y_id, bad_col); + bad = ray_table_add_col(bad, y_id, old_y); + ray_table_set_col_idx(bad, 1, bad_col); + TEST_ASSERT_FALSE(RAY_IS_ERR(bad)); ray_err_t bad_err = ray_splay_save(bad, dir, NULL); - TEST_ASSERT_TRUE(bad_err != RAY_OK); + TEST_ASSERT_EQ_I(bad_err, RAY_ERR_NYI); loaded = ray_read_splayed(dir, NULL); TEST_ASSERT_NOT_NULL(loaded); @@ -2448,7 +2450,205 @@ static test_result_t test_splay_atomic_generation_publish(void) { PASS(); } +static ray_t* generation_pair(int64_t value) { + int64_t x[] = {value, value + 1}; + int64_t y[] = {value * 10, (value + 1) * 10}; + ray_t* xc = ray_vec_from_raw(RAY_I64, x, 2); + ray_t* yc = ray_vec_from_raw(RAY_I64, y, 2); + ray_t* t = ray_table_new(2); + t = ray_table_add_col(t, ray_sym_intern("x", 1), xc); + t = ray_table_add_col(t, ray_sym_intern("y", 1), yc); + ray_release(xc); + ray_release(yc); + return t; +} + +static bool generation_matches(const char* dir, bool mmap, int64_t value) { + ray_t* t = mmap ? ray_read_splayed(dir, NULL) : ray_splay_load(dir, NULL); + if (!t || RAY_IS_ERR(t)) { if (t) ray_release(t); return false; } + ray_t* x = ray_table_get_col(t, ray_sym_intern("x", 1)); + ray_t* y = ray_table_get_col(t, ray_sym_intern("y", 1)); + bool ok = x && y && x->len == 2 && y->len == 2 && + x->type == RAY_I64 && y->type == RAY_I64; + if (ok) { + const int64_t* xd = ray_data(x); + const int64_t* yd = ray_data(y); + ok = xd[0] == value && xd[1] == value + 1 && + yd[0] == value * 10 && yd[1] == (value + 1) * 10; + } + ray_release(t); + return ok; +} + +/* Force actual filesystem failures after the first column and after all + * columns respectively. No invalid object or preflight shortcut is involved. */ +static test_result_t test_generation_io_failures(void) { + const char* dir = TMP_SPLAY_BASE "/generation_io"; + rm_rf(dir); + ray_t* old = generation_pair(1); + ray_t* next = generation_pair(9); + TEST_ASSERT_FALSE(RAY_IS_ERR(old)); + TEST_ASSERT_FALSE(RAY_IS_ERR(next)); + for (int era = 0; era < 2; era++) { + /* Exercise both legacy -> generation and generation -> generation. */ + TEST_ASSERT_EQ_I(ray_splay_save(old, dir, NULL), RAY_OK); + char before[1024], after[1024], obstruction[1100], first_col[1100]; + TEST_ASSERT_EQ_I(ray_splay_resolve_dir(dir, before, sizeof(before)), RAY_OK); + ray_splay_write_t write; + TEST_ASSERT_EQ_I(ray_splay_write_begin(dir, &write), RAY_OK); + snprintf(obstruction, sizeof(obstruction), "%s/y", write.dir); + TEST_ASSERT_EQ_I(ray_test_mkdir_p(obstruction), 0); + ray_err_t err = ray_splay_write_table(next, write.dir, NULL, true); + err = ray_splay_write_finish(&write, err, true); + TEST_ASSERT_EQ_I(err, RAY_ERR_IO); + snprintf(first_col, sizeof(first_col), "%s/x", write.dir); + ray_t* written = ray_col_load(first_col); + TEST_ASSERT_NOT_NULL(written); + TEST_ASSERT_FALSE(RAY_IS_ERR(written)); + TEST_ASSERT_EQ_I(((int64_t*)ray_data(written))[0], 9); + ray_release(written); + TEST_ASSERT_TRUE(generation_matches(dir, false, 1)); + TEST_ASSERT_TRUE(generation_matches(dir, true, 1)); + + TEST_ASSERT_EQ_I(ray_splay_write_begin(dir, &write), RAY_OK); + err = ray_splay_write_table(next, write.dir, NULL, true); + if (err != RAY_OK) (void)ray_splay_write_finish(&write, err, true); + TEST_ASSERT_EQ_I(err, RAY_OK); + /* A directory cannot be opened as the manifest's temporary file. */ + snprintf(obstruction, sizeof(obstruction), "%s/.current.tmp", write.dir); + TEST_ASSERT_EQ_I(ray_test_mkdir_p(obstruction), 0); + TEST_ASSERT_EQ_I(ray_splay_write_finish(&write, RAY_OK, true), RAY_ERR_IO); + TEST_ASSERT_TRUE(generation_matches(dir, false, 1)); + TEST_ASSERT_TRUE(generation_matches(dir, true, 1)); + TEST_ASSERT_EQ_I(ray_splay_resolve_dir(dir, after, sizeof(after)), RAY_OK); + TEST_ASSERT_TRUE(strcmp(before, after) == 0); +#ifndef _WIN32 + if (geteuid() != 0) { + /* Temp creation succeeds inside the stage; publication itself + * fails when rename cannot modify the table root. */ + TEST_ASSERT_EQ_I(ray_splay_write_begin(dir, &write), RAY_OK); + err = ray_splay_write_table(next, write.dir, NULL, true); + if (err != RAY_OK) (void)ray_splay_write_finish(&write, err, true); + TEST_ASSERT_EQ_I(err, RAY_OK); + TEST_ASSERT_EQ_I(chmod(dir, 0555), 0); + err = ray_splay_write_finish(&write, RAY_OK, true); + int restored = chmod(dir, 0755); + TEST_ASSERT_EQ_I(restored, 0); + TEST_ASSERT_EQ_I(err, RAY_ERR_IO); + TEST_ASSERT_TRUE(generation_matches(dir, true, 1)); + } +#endif + } + /* Failure must also release the writer lock so a retry can commit. */ + TEST_ASSERT_EQ_I(ray_splay_save(next, dir, NULL), RAY_OK); + TEST_ASSERT_TRUE(generation_matches(dir, true, 9)); + ray_release(next); + ray_release(old); + rm_rf(dir); + PASS(); +} + +static test_result_t test_generation_retains_readers(void) { + const char* dir = TMP_SPLAY_BASE "/generation_readers"; + rm_rf(dir); + ray_t* old = generation_pair(1); + ray_t* next = generation_pair(9); + TEST_ASSERT_EQ_I(ray_splay_save(old, dir, NULL), RAY_OK); + for (int era = 0; era < 2; era++) { + char resolved[1024], path[1100]; + TEST_ASSERT_EQ_I(ray_splay_resolve_dir(dir, resolved, sizeof(resolved)), RAY_OK); + ray_t* pinned = ray_read_splayed(dir, NULL); + TEST_ASSERT_FALSE(RAY_IS_ERR(pinned)); + TEST_ASSERT_EQ_I(ray_splay_save(next, dir, NULL), RAY_OK); + /* Reader resolved the old path before publication but opens a column + * afterwards. Both the on-disk file and an existing mmap must survive. */ + snprintf(path, sizeof(path), "%s/y", resolved); + ray_t* late = ray_col_load(path); + TEST_ASSERT_NOT_NULL(late); + TEST_ASSERT_FALSE(RAY_IS_ERR(late)); + TEST_ASSERT_EQ_I(((int64_t*)ray_data(late))[0], 10); + ray_release(late); + ray_t* py = ray_table_get_col(pinned, ray_sym_intern("y", 1)); + TEST_ASSERT_NOT_NULL(py); + TEST_ASSERT_EQ_I(((int64_t*)ray_data(py))[0], 10); + ray_release(pinned); + TEST_ASSERT_TRUE(generation_matches(dir, false, 9)); + TEST_ASSERT_EQ_I(ray_splay_save(old, dir, NULL), RAY_OK); + } + ray_release(next); + ray_release(old); + rm_rf(dir); + PASS(); +} + +static test_result_t test_generation_invalid_manifest(void) { + const char* dir = TMP_SPLAY_BASE "/generation_manifest"; + rm_rf(dir); + ray_t* t = generation_pair(1); + TEST_ASSERT_EQ_I(ray_splay_save(t, dir, NULL), RAY_OK); + const char* invalid[] = {"", ".generations/../x\n", ".generations/g-1/..\n", + ".generations/g-1\njunk\n", ".generations/\n"}; + for (size_t i = 0; i < sizeof(invalid) / sizeof(invalid[0]); i++) { + FILE* f = fopen(TMP_SPLAY_BASE "/generation_manifest/.current", "wb"); + TEST_ASSERT_NOT_NULL(f); + TEST_ASSERT_TRUE(fputs(invalid[i], f) >= 0); + TEST_ASSERT_EQ_I(fclose(f), 0); + ray_t* loaded = ray_splay_load(dir, NULL); + TEST_ASSERT_TRUE(RAY_IS_ERR(loaded)); + TEST_ASSERT_STR_EQ(ray_err_code(loaded), "corrupt"); + ray_release(loaded); + loaded = ray_read_splayed_dom(dir, NULL); + TEST_ASSERT_TRUE(RAY_IS_ERR(loaded)); + TEST_ASSERT_STR_EQ(ray_err_code(loaded), "corrupt"); + ray_release(loaded); + } + ray_release(t); + rm_rf(dir); + PASS(); +} + +static test_result_t test_generation_writer_exit(void) { + const char* dir = TMP_SPLAY_BASE "/generation_exit"; + rm_rf(dir); + ray_t* old = generation_pair(1); + ray_t* next = generation_pair(9); + TEST_ASSERT_EQ_I(ray_splay_save(old, dir, NULL), RAY_OK); + for (int complete = 0; complete < 2; complete++) { + ray_splay_write_t wr; + TEST_ASSERT_EQ_I(ray_splay_write_begin(dir, &wr), RAY_OK); + ray_err_t err; + if (complete) { + err = ray_splay_write_table(next, wr.dir, NULL, true); + } else { + char path[1100]; + snprintf(path, sizeof(path), "%s/x", wr.dir); + err = ray_col_save(ray_table_get_col_idx(next, 0), path); + } + TEST_ASSERT_EQ_I(err, RAY_OK); + /* Simulate an abrupt writer exit after staging but before finish(): + * the OS releases the writer lock, but no .current publication happens. + * Avoid fork() here; macOS sanitizer runtimes may terminate forked + * children before normal test-side status reporting runs. */ + TEST_ASSERT_EQ_I(ray_file_unlock(wr.lock), RAY_OK); + ray_file_close(wr.lock); + wr.lock = RAY_FD_INVALID; + TEST_ASSERT_TRUE(generation_matches(dir, false, 1)); + TEST_ASSERT_TRUE(generation_matches(dir, true, 1)); + TEST_ASSERT_EQ_I(ray_splay_save(old, dir, NULL), RAY_OK); + } + TEST_ASSERT_EQ_I(ray_splay_save(next, dir, NULL), RAY_OK); + TEST_ASSERT_TRUE(generation_matches(dir, true, 9)); + ray_release(next); + ray_release(old); + rm_rf(dir); + PASS(); +} + const test_entry_t splay_entries[] = { + { "splay/generation_io_failures", test_generation_io_failures, splay_setup, splay_teardown }, + { "splay/generation_retains_readers", test_generation_retains_readers, splay_setup, splay_teardown }, + { "splay/generation_invalid_manifest", test_generation_invalid_manifest, splay_setup, splay_teardown }, + { "splay/generation_writer_exit", test_generation_writer_exit, splay_setup, splay_teardown }, { "splay/atomic_generation_publish", test_splay_atomic_generation_publish, splay_setup, splay_teardown }, { "splay/has_nulls_roundtrip", test_splayed_has_nulls_roundtrip, splay_setup, splay_teardown }, { "splay/save_null_dir", test_save_null_dir, splay_setup, splay_teardown }, diff --git a/test/test_stress_matrix.c b/test/test_stress_matrix.c index a9d75dad2..45ca15fa2 100644 --- a/test/test_stress_matrix.c +++ b/test/test_stress_matrix.c @@ -348,7 +348,16 @@ static test_result_t test_matrix_shared_domain_layout(void) { * (bytes 0-15 aux, then mmod/order/type/attrs — store/col.c "Column file * format"). Returns RAY_SYM_W8/W16/W32/W64 or -1 on I/O failure. */ static int disk_sym_width(const char* col_path) { - FILE* f = fopen(col_path, "rb"); + /* Keep inspecting the raw header, but in the selected generation. */ + char dir[1024], resolved[1024], path[1200]; + const char* slash = strrchr(col_path, '/'); + if (!slash || (size_t)(slash - col_path) >= sizeof(dir)) return -1; + memcpy(dir, col_path, (size_t)(slash - col_path)); + dir[slash - col_path] = '\0'; + if (ray_splay_resolve_dir(dir, resolved, sizeof(resolved)) != RAY_OK) return -1; + int n = snprintf(path, sizeof(path), "%s/%s", resolved, slash + 1); + if (n < 0 || (size_t)n >= sizeof(path)) return -1; + FILE* f = fopen(path, "rb"); if (!f) return -1; if (fseek(f, (long)offsetof(ray_t, attrs), SEEK_SET) != 0) { fclose(f); From 55f19a3497410527171670a8aa6067c81f6f2f91 Mon Sep 17 00:00:00 2001 From: belowzeroff Date: Wed, 30 Sep 2026 09:36:49 -0400 Subject: [PATCH 3/5] fix(store): bound atomic splay generation retention --- .github/workflows/ci.yml | 11 +-- docs/docs/storage/index.md | 24 +++--- src/store/splay.c | 108 +++++++++++++++++------- src/store/splay.h | 3 +- test/rfl/storage/atomic_generations.rfl | 11 ++- test/test_splay.c | 33 ++++++-- 6 files changed, 131 insertions(+), 59 deletions(-) diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 79cad108a..426878bc4 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -129,19 +129,10 @@ jobs: timeout-minutes: 30 steps: - uses: actions/checkout@v4 - with: - fetch-depth: 0 - name: Install cppcheck run: sudo apt-get update && sudo apt-get install -y cppcheck - name: cppcheck - run: | - changed="$(git diff --name-only "origin/${{ github.base_ref || 'dev' }}"...HEAD -- 'src/**/*.c' | tr '\n' ' ')" - if [ -z "$changed" ]; then - echo "No changed C sources for cppcheck." - exit 0 - fi - echo "cppcheck files: $changed" - make cppcheck FILES="$changed" CPPCHECK_SKIP= + run: make cppcheck # Single stable check to require in branch protection, instead of the four # brittle matrix names (which break if the matrix is renamed). `if: always()` diff --git a/docs/docs/storage/index.md b/docs/docs/storage/index.md index 335b8f625..e3fb2b2ef 100644 --- a/docs/docs/storage/index.md +++ b/docs/docs/storage/index.md @@ -81,16 +81,20 @@ same table serialize on `.write.lock`; readers resolve the manifest once and continue using that generation. The symbol vocabulary stays at its original location and grows append-only. -Previous generations, the legacy files and incomplete staging directories are -retained. Disk usage therefore grows with replacements. There is no automatic -garbage collection yet: reclaim old files only during offline maintenance with -all readers and writers stopped, and preserve the generation named by `.current`. -Older binaries that do not understand `.current` must not access an upgraded -table; reading its legacy files directly returns obsolete data. - -`ray_splay_save` (including `.db.splayed.set`) syncs data, indexes and directory -entries before acknowledging publication. The bulk C API and CSV import retain -their no-fsync contract: publication is atomic for readers, but is not a +Publication keeps the selected generation and one previous generation, then +removes older staged directories best-effort. A failed unpublished generation is +removed immediately where the platform allows it; on Windows an open mapped file +can defer that cleanup until a later publish. At the first manifest publish, +the legacy root `.d` is retired so older binaries fail loudly instead of reading +obsolete root files. + +The `.write.lock` writer lock is local-filesystem coordination. Do not rely on +it for NFS-backed shared writers. + +`ray_splay_save` (including `.db.splayed.set`) syncs primary data and directory +entries before acknowledging publication. Inline indexes are rebuildable +accelerators written with a marker-last protocol. The bulk C API and CSV import +retain their no-fsync contract: publication is atomic for readers, but is not a power-loss durability guarantee. An I/O error after the manifest rename (during the final directory sync) has an uncertain commit outcome; reopen the table to inspect the selected generation. diff --git a/src/store/splay.c b/src/store/splay.c index 6c5e15a86..932b32310 100644 --- a/src/store/splay.c +++ b/src/store/splay.c @@ -243,6 +243,67 @@ static ray_err_t splay_has_file(const char* dir, const char* name, bool* exists) return *exists || errno == ENOENT ? RAY_OK : RAY_ERR_IO; } +static void splay_remove_tree_best_effort(const char* path) { + DIR* d = opendir(path); + if (!d) { + (void)unlink(path); + return; + } + struct dirent* entry; + while ((entry = readdir(d))) { + if (strcmp(entry->d_name, ".") == 0 || strcmp(entry->d_name, "..") == 0) + continue; + char child[1024]; + int n = snprintf(child, sizeof(child), "%s/%s", path, entry->d_name); + if (n < 0 || (size_t)n >= sizeof(child)) continue; + struct stat st; + if (stat(child, &st) != 0) continue; + if (S_ISDIR(st.st_mode)) splay_remove_tree_best_effort(child); + else (void)unlink(child); + } + closedir(d); + (void)rmdir(path); +} + +static void splay_retire_legacy_schema(const char* root) { + char schema[1024], retired[1024]; + int n = snprintf(schema, sizeof(schema), "%s/.d", root); + int m = snprintf(retired, sizeof(retired), "%s/.legacy.d", root); + if (n < 0 || (size_t)n >= sizeof(schema) || + m < 0 || (size_t)m >= sizeof(retired)) + return; + if (access(schema, F_OK) != 0) return; + (void)unlink(retired); + if (rename(schema, retired) != 0) + (void)unlink(schema); +} + +static void splay_prune_generations(const char* root, const char* current, + const char* previous) { + char generations[1024]; + int n = snprintf(generations, sizeof(generations), "%s/.generations", root); + if (n < 0 || (size_t)n >= sizeof(generations)) return; + DIR* d = opendir(generations); + if (!d) return; + + struct dirent* entry; + while ((entry = readdir(d))) { + if (strcmp(entry->d_name, ".") == 0 || strcmp(entry->d_name, "..") == 0) + continue; + char rel[256], full[1024]; + int rn = snprintf(rel, sizeof(rel), ".generations/%s", entry->d_name); + int fn = snprintf(full, sizeof(full), "%s/%s", root, rel); + if (rn < 0 || (size_t)rn >= sizeof(rel) || + fn < 0 || (size_t)fn >= sizeof(full)) + continue; + if ((current && strcmp(rel, current) == 0) || + (previous && strcmp(full, previous) == 0)) + continue; + splay_remove_tree_best_effort(full); + } + closedir(d); +} + static ray_err_t splay_validate_save(ray_t* tbl, const char* dir, const char* sym_path) { if (!tbl || RAY_IS_ERR(tbl) || tbl->type != RAY_TABLE) return RAY_ERR_TYPE; @@ -302,30 +363,6 @@ static ray_err_t splay_validate_save(ray_t* tbl, const char* dir, return RAY_OK; } -static ray_err_t splay_sync_files(const char* dir) { - DIR* d = opendir(dir); - if (!d) return RAY_ERR_IO; - ray_err_t err = RAY_OK; - struct dirent* entry; - while ((entry = readdir(d))) { - if (strcmp(entry->d_name, ".") == 0 || strcmp(entry->d_name, "..") == 0) - continue; - char path[1024]; - int n = snprintf(path, sizeof(path), "%s/%s", dir, entry->d_name); - if (n < 0 || (size_t)n >= sizeof(path)) { err = RAY_ERR_RANGE; break; } - struct stat st; - if (stat(path, &st) != 0) { err = RAY_ERR_IO; break; } - if (!S_ISREG(st.st_mode)) continue; - ray_fd_t fd = ray_file_open(path, RAY_OPEN_READ | RAY_OPEN_WRITE); - if (fd == RAY_FD_INVALID) { err = RAY_ERR_IO; break; } - err = ray_file_sync(fd); - ray_file_close(fd); - if (err != RAY_OK) break; - } - closedir(d); - return err; -} - ray_err_t ray_splay_write_table(ray_t* tbl, const char* dir, const char* sym_path, bool durable) { ray_err_t validation = splay_validate_save(tbl, dir, sym_path); @@ -437,12 +474,10 @@ ray_err_t ray_splay_write_table(ray_t* tbl, const char* dir, } if (dom) ray_sym_domain_release(dom); - /* Indexes belong to this generation and must precede publication. */ + /* Indexes belong to this generation and must precede publication. They are + * rebuildable accelerators; ray_col_append_index writes the marker last and + * intentionally does not fsync them again after ray_col_save fsyncs data. */ ray_splay_build_indexes(dir, tbl); - if (durable) { - ray_err_t err = splay_sync_files(dir); - if (err != RAY_OK) { ray_release(schema); return err; } - } /* 3. .d LAST — the commit marker. */ { @@ -466,7 +501,12 @@ ray_err_t ray_splay_write_table(ray_t* tbl, const char* dir, ray_err_t ray_splay_write_finish(ray_splay_write_t* write, ray_err_t result, bool durable) { + char previous[1024]; + bool had_previous = false; if (result == RAY_OK && write->staged) { + if (splay_current_dir(write->root, previous, sizeof(previous), + &had_previous) != RAY_OK) + had_previous = false; /* ray_file_sync_dir syncs the PARENT of its argument. */ char schema[1100]; snprintf(schema, sizeof(schema), "%s/.d", write->dir); @@ -475,14 +515,20 @@ ray_err_t ray_splay_write_finish(ray_splay_write_t* write, ray_err_t result, result = RAY_ERR_IO; if (result == RAY_OK) result = splay_publish_generation(write->root, write->generation, durable); + if (result == RAY_OK) { + splay_retire_legacy_schema(write->root); + splay_prune_generations(write->root, write->generation, + had_previous ? previous : NULL); + } + } + if (result != RAY_OK && write->staged && write->dir[0]) { + splay_remove_tree_best_effort(write->dir); } if (write->lock != RAY_FD_INVALID) { (void)ray_file_unlock(write->lock); ray_file_close(write->lock); write->lock = RAY_FD_INVALID; } - /* A reader may have resolved the old path without opening every column. - * Keep both legacy files and old generations immutable after publication. */ return result; } diff --git a/src/store/splay.h b/src/store/splay.h index 6e6722b2e..e03287eaf 100644 --- a/src/store/splay.h +++ b/src/store/splay.h @@ -31,7 +31,8 @@ struct ray_sym_domain_s; /* Internal publication protocol shared by the table and streaming CSV writers. * begin serializes writers; finish publishes only on success and always unlocks. - * Old generations must remain available to readers that already resolved them. */ + * Publication keeps the current generation and one previous generation; older + * staged directories are removed best-effort after each successful publish. */ typedef struct { ray_fd_t lock; bool staged; diff --git a/test/rfl/storage/atomic_generations.rfl b/test/rfl/storage/atomic_generations.rfl index 48717f4ee..731eb772a 100644 --- a/test/rfl/storage/atomic_generations.rfl +++ b/test/rfl/storage/atomic_generations.rfl @@ -8,8 +8,10 @@ (.db.splayed.set agLeaf agNew) (sum (at (.db.splayed.get agLeaf) 'x)) -- 7 (sum (at (.db.parted.get agRoot 't) 'x)) -- 7 -;; Discovery must use the active schema even after offline legacy cleanup. -(.sys.exec "rm /tmp/rf_test_atomic_generations/2026.09.28/t/.d") -- 0 +;; The first generation publish retires the legacy schema so old binaries +;; cannot silently read the stale root files. +(.sys.exec "test ! -f /tmp/rf_test_atomic_generations/2026.09.28/t/.d") -- 0 +;; Discovery must use the active schema without the legacy root .d. (.db.parted.tables agRoot) -- ['t] ;; Schema narrowing must not make the partition reader use the old .d. @@ -44,4 +46,9 @@ (.db.splayed.set (format "%/indexed_str" agRoot) agStrings) (.db.splayed.set (format "%/indexed_str" agRoot) agStrings) (.idx.has? (at (.db.splayed.get (format "%/indexed_str" agRoot)) 'k)) -- true +;; Retention is bounded to current + previous generation. +(.db.splayed.set agLeaf (table ['x] (list [11 12]))) +(.db.splayed.set agLeaf (table ['x] (list [13 14]))) +(.db.splayed.set agLeaf (table ['x] (list [15 16]))) +(.sys.exec "test $(find /tmp/rf_test_atomic_generations/2026.09.28/t/.generations -mindepth 1 -maxdepth 1 -type d | wc -l) -le 2") -- 0 (.sys.exec "rm -rf /tmp/rf_test_atomic_generations") -- 0 diff --git a/test/test_splay.c b/test/test_splay.c index 2fef9d7fd..d8dac0613 100644 --- a/test/test_splay.c +++ b/test/test_splay.c @@ -2405,6 +2405,10 @@ static test_result_t test_splay_atomic_generation_publish(void) { TEST_ASSERT_FALSE(RAY_IS_ERR(replacement)); TEST_ASSERT_EQ_I(ray_splay_save(replacement, dir, NULL), RAY_OK); TEST_ASSERT_EQ_I(access(manifest, F_OK), 0); + char legacy_schema[512]; + n = snprintf(legacy_schema, sizeof(legacy_schema), "%s/.d", dir); + TEST_ASSERT_TRUE(n > 0 && (size_t)n < sizeof(legacy_schema)); + TEST_ASSERT_EQ_I(access(legacy_schema, F_OK), -1); ray_t* loaded = ray_splay_load(dir, NULL); TEST_ASSERT_NOT_NULL(loaded); @@ -2480,6 +2484,27 @@ static bool generation_matches(const char* dir, bool mmap, int64_t value) { return ok; } +static int generation_dir_count(const char* dir) { + char path[1024]; + int n = snprintf(path, sizeof(path), "%s/.generations", dir); + if (n < 0 || (size_t)n >= sizeof(path)) return -1; + DIR* d = opendir(path); + if (!d) return 0; + int count = 0; + struct dirent* entry; + while ((entry = readdir(d))) { + if (strcmp(entry->d_name, ".") == 0 || strcmp(entry->d_name, "..") == 0) + continue; + char child[1024]; + n = snprintf(child, sizeof(child), "%s/%s", path, entry->d_name); + if (n < 0 || (size_t)n >= sizeof(child)) continue; + struct stat st; + if (stat(child, &st) == 0 && S_ISDIR(st.st_mode)) count++; + } + closedir(d); + return count; +} + /* Force actual filesystem failures after the first column and after all * columns respectively. No invalid object or preflight shortcut is involved. */ static test_result_t test_generation_io_failures(void) { @@ -2502,11 +2527,7 @@ static test_result_t test_generation_io_failures(void) { err = ray_splay_write_finish(&write, err, true); TEST_ASSERT_EQ_I(err, RAY_ERR_IO); snprintf(first_col, sizeof(first_col), "%s/x", write.dir); - ray_t* written = ray_col_load(first_col); - TEST_ASSERT_NOT_NULL(written); - TEST_ASSERT_FALSE(RAY_IS_ERR(written)); - TEST_ASSERT_EQ_I(((int64_t*)ray_data(written))[0], 9); - ray_release(written); + TEST_ASSERT_EQ_I(access(first_col, F_OK), -1); TEST_ASSERT_TRUE(generation_matches(dir, false, 1)); TEST_ASSERT_TRUE(generation_matches(dir, true, 1)); @@ -2542,6 +2563,7 @@ static test_result_t test_generation_io_failures(void) { /* Failure must also release the writer lock so a retry can commit. */ TEST_ASSERT_EQ_I(ray_splay_save(next, dir, NULL), RAY_OK); TEST_ASSERT_TRUE(generation_matches(dir, true, 9)); + TEST_ASSERT_TRUE(generation_dir_count(dir) <= 2); ray_release(next); ray_release(old); rm_rf(dir); @@ -2575,6 +2597,7 @@ static test_result_t test_generation_retains_readers(void) { TEST_ASSERT_TRUE(generation_matches(dir, false, 9)); TEST_ASSERT_EQ_I(ray_splay_save(old, dir, NULL), RAY_OK); } + TEST_ASSERT_TRUE(generation_dir_count(dir) <= 2); ray_release(next); ray_release(old); rm_rf(dir); From 5ccefd92cd0f5027b6db52f0969a11e3f37c87f9 Mon Sep 17 00:00:00 2001 From: belowzeroff Date: Wed, 30 Sep 2026 09:45:26 -0400 Subject: [PATCH 4/5] test(store): include dirent for generation scan --- test/test_splay.c | 1 + 1 file changed, 1 insertion(+) diff --git a/test/test_splay.c b/test/test_splay.c index d8dac0613..783b6cf90 100644 --- a/test/test_splay.c +++ b/test/test_splay.c @@ -49,6 +49,7 @@ #include #include #include +#include #include /* ---- Setup / Teardown -------------------------------------------------- */ From b7588291b1ce78d4cf23f4efd52c4d81b539abcb Mon Sep 17 00:00:00 2001 From: belowzeroff Date: Thu, 1 Oct 2026 12:37:16 -0400 Subject: [PATCH 5/5] fix(store): protect published splay generations during cleanup --- src/store/splay.c | 37 ++++++++++++++--- test/test_splay.c | 101 ++++++++++++++++++++++++++++++++++++++++++++++ 2 files changed, 133 insertions(+), 5 deletions(-) diff --git a/src/store/splay.c b/src/store/splay.c index 932b32310..6fb3a929d 100644 --- a/src/store/splay.c +++ b/src/store/splay.c @@ -1,3 +1,6 @@ +#if !defined(_WIN32) && !defined(_POSIX_C_SOURCE) +#define _POSIX_C_SOURCE 200809L +#endif /* * Copyright (c) 2025-2026 Anton Kundenko * All rights reserved. @@ -257,7 +260,11 @@ static void splay_remove_tree_best_effort(const char* path) { int n = snprintf(child, sizeof(child), "%s/%s", path, entry->d_name); if (n < 0 || (size_t)n >= sizeof(child)) continue; struct stat st; +#ifdef RAY_OS_WINDOWS if (stat(child, &st) != 0) continue; +#else + if (lstat(child, &st) != 0) continue; +#endif if (S_ISDIR(st.st_mode)) splay_remove_tree_best_effort(child); else (void)unlink(child); } @@ -265,6 +272,14 @@ static void splay_remove_tree_best_effort(const char* path) { (void)rmdir(path); } +static bool splay_generation_is_current(const ray_splay_write_t* write) { + char current[1024]; + bool active = false; + return write && write->staged && write->dir[0] && + splay_current_dir(write->root, current, sizeof(current), &active) == RAY_OK && + active && strcmp(current, write->dir) == 0; +} + static void splay_retire_legacy_schema(const char* root) { char schema[1024], retired[1024]; int n = snprintf(schema, sizeof(schema), "%s/.d", root); @@ -363,8 +378,9 @@ static ray_err_t splay_validate_save(ray_t* tbl, const char* dir, return RAY_OK; } -ray_err_t ray_splay_write_table(ray_t* tbl, const char* dir, - const char* sym_path, bool durable) { +static ray_err_t splay_write_table_impl(ray_t* tbl, const char* dir, + const char* sym_path, bool durable, + bool flush_sym, bool build_indexes) { ray_err_t validation = splay_validate_save(tbl, dir, sym_path); if (validation != RAY_OK) return validation; /* Create directory and any missing parents (mkdir -p semantics). @@ -420,7 +436,7 @@ ray_err_t ray_splay_write_table(ray_t* tbl, const char* dir, } } - ray_err_t sym_err = ray_sym_domain_flush(dom, durable); + ray_err_t sym_err = flush_sym ? ray_sym_domain_flush(dom, durable) : RAY_OK; if (sym_err != RAY_OK) { ray_sym_domain_release(dom); return sym_err; @@ -477,7 +493,8 @@ ray_err_t ray_splay_write_table(ray_t* tbl, const char* dir, /* Indexes belong to this generation and must precede publication. They are * rebuildable accelerators; ray_col_append_index writes the marker last and * intentionally does not fsync them again after ray_col_save fsyncs data. */ - ray_splay_build_indexes(dir, tbl); + if (build_indexes) + ray_splay_build_indexes(dir, tbl); /* 3. .d LAST — the commit marker. */ { @@ -499,6 +516,11 @@ ray_err_t ray_splay_write_table(ray_t* tbl, const char* dir, return RAY_OK; } +ray_err_t ray_splay_write_table(ray_t* tbl, const char* dir, + const char* sym_path, bool durable) { + return splay_write_table_impl(tbl, dir, sym_path, durable, true, true); +} + ray_err_t ray_splay_write_finish(ray_splay_write_t* write, ray_err_t result, bool durable) { char previous[1024]; @@ -521,7 +543,8 @@ ray_err_t ray_splay_write_finish(ray_splay_write_t* write, ray_err_t result, had_previous ? previous : NULL); } } - if (result != RAY_OK && write->staged && write->dir[0]) { + if (result != RAY_OK && write->staged && write->dir[0] && + !splay_generation_is_current(write)) { splay_remove_tree_best_effort(write->dir); } if (write->lock != RAY_FD_INVALID) { @@ -605,6 +628,10 @@ ray_err_t ray_splay_save_bulk(ray_t* tbl, const char* dir, const char* sym_path) return splay_save_impl(tbl, dir, sym_path, false); } +ray_err_t ray_splay_save_staged_bulk(ray_t* tbl, const char* dir, const char* sym_path) { + return splay_write_table_impl(tbl, dir, sym_path, false, false, false); +} + /* -------------------------------------------------------------------------- * splay_load_impl — shared implementation for ray_splay_load / ray_read_splayed * diff --git a/test/test_splay.c b/test/test_splay.c index 783b6cf90..346f3fd75 100644 --- a/test/test_splay.c +++ b/test/test_splay.c @@ -860,6 +860,35 @@ static test_result_t test_save_bulk_with_sym_path(void) { PASS(); } +static test_result_t test_save_staged_bulk_defers_sym_flush(void) { + const char* dir = TMP_SPLAY_BASE "/staged_bulk_sym"; + const char* sym_path = TMP_SPLAY_BASE "/staged_bulk_sym.sym"; + rm_rf(dir); + unlink(sym_path); + + int64_t id_s = ray_sym_intern("wsym", 4); + int64_t sval = ray_sym_intern("wv1", 3); + ray_t* scol = ray_sym_vec_new(RAY_SYM_W8, 2); + TEST_ASSERT_FALSE(RAY_IS_ERR(scol)); + scol->len = 2; + ((uint8_t*)ray_data(scol))[0] = (uint8_t)sval; + ((uint8_t*)ray_data(scol))[1] = (uint8_t)sval; + + ray_t* tbl = ray_table_new(1); + tbl = ray_table_add_col(tbl, id_s, scol); + TEST_ASSERT_FALSE(RAY_IS_ERR(tbl)); + + ray_err_t err = ray_splay_save_staged_bulk(tbl, dir, sym_path); + TEST_ASSERT_EQ_I(err, RAY_OK); + TEST_ASSERT_EQ_I(access(sym_path, F_OK), -1); + + ray_release(scol); + ray_release(tbl); + rm_rf(dir); + unlink(sym_path); + PASS(); +} + /* ========================================================================= * 19. splay_save_impl: snprintf overflow for the column / ".d" paths. * Requires strlen(dir) >= 1021 so that strlen(dir)+3 >= 1024. @@ -2506,6 +2535,74 @@ static int generation_dir_count(const char* dir) { return count; } +#ifndef _WIN32 +static test_result_t test_generation_prune_unlinks_symlink(void) { + const char* dir = TMP_SPLAY_BASE "/generation_symlink"; + const char* victim = TMP_SPLAY_BASE "/generation_symlink_victim"; + const char* victim_file = TMP_SPLAY_BASE "/generation_symlink_victim/keep"; + rm_rf(dir); + rm_rf(victim); + + ray_t* one = generation_pair(1); + ray_t* two = generation_pair(2); + ray_t* three = generation_pair(3); + ray_t* four = generation_pair(4); + TEST_ASSERT_FALSE(RAY_IS_ERR(one)); + TEST_ASSERT_FALSE(RAY_IS_ERR(two)); + TEST_ASSERT_FALSE(RAY_IS_ERR(three)); + TEST_ASSERT_FALSE(RAY_IS_ERR(four)); + + TEST_ASSERT_EQ_I(ray_splay_save(one, dir, NULL), RAY_OK); + TEST_ASSERT_EQ_I(ray_splay_save(two, dir, NULL), RAY_OK); + TEST_ASSERT_EQ_I(ray_splay_save(three, dir, NULL), RAY_OK); + + char current[1024], generations[1024], prune_dir[1024]; + TEST_ASSERT_EQ_I(ray_splay_resolve_dir(dir, current, sizeof(current)), RAY_OK); + int n = snprintf(generations, sizeof(generations), "%s/.generations", dir); + TEST_ASSERT_TRUE(n > 0 && (size_t)n < sizeof(generations)); + DIR* d = opendir(generations); + TEST_ASSERT_NOT_NULL(d); + prune_dir[0] = '\0'; + struct dirent* entry; + while ((entry = readdir(d))) { + if (strcmp(entry->d_name, ".") == 0 || strcmp(entry->d_name, "..") == 0) + continue; + char full[1024]; + n = snprintf(full, sizeof(full), "%s/%s", generations, entry->d_name); + TEST_ASSERT_TRUE(n > 0 && (size_t)n < sizeof(full)); + if (strcmp(full, current) != 0) { + snprintf(prune_dir, sizeof(prune_dir), "%s", full); + break; + } + } + closedir(d); + TEST_ASSERT_TRUE(prune_dir[0] != '\0'); + + TEST_ASSERT_EQ_I(ray_test_mkdir_p(victim), 0); + FILE* f = fopen(victim_file, "wb"); + TEST_ASSERT_NOT_NULL(f); + fputs("keep", f); + fclose(f); + + char link_path[1024]; + n = snprintf(link_path, sizeof(link_path), "%s/outside", prune_dir); + TEST_ASSERT_TRUE(n > 0 && (size_t)n < sizeof(link_path)); + TEST_ASSERT_EQ_I(symlink(victim, link_path), 0); + + TEST_ASSERT_EQ_I(ray_splay_save(four, dir, NULL), RAY_OK); + TEST_ASSERT_EQ_I(access(victim_file, F_OK), 0); + TEST_ASSERT_EQ_I(access(prune_dir, F_OK), -1); + + ray_release(four); + ray_release(three); + ray_release(two); + ray_release(one); + rm_rf(dir); + rm_rf(victim); + PASS(); +} +#endif + /* Force actual filesystem failures after the first column and after all * columns respectively. No invalid object or preflight shortcut is involved. */ static test_result_t test_generation_io_failures(void) { @@ -2669,6 +2766,9 @@ static test_result_t test_generation_writer_exit(void) { } const test_entry_t splay_entries[] = { +#ifndef _WIN32 + { "splay/generation_prune_unlinks_symlink", test_generation_prune_unlinks_symlink, splay_setup, splay_teardown }, +#endif { "splay/generation_io_failures", test_generation_io_failures, splay_setup, splay_teardown }, { "splay/generation_retains_readers", test_generation_retains_readers, splay_setup, splay_teardown }, { "splay/generation_invalid_manifest", test_generation_invalid_manifest, splay_setup, splay_teardown }, @@ -2694,6 +2794,7 @@ const test_entry_t splay_entries[] = { { "splay/load_dir_path_too_long", test_load_dir_path_too_long, splay_setup, splay_teardown }, { "splay/load_col_path_too_long", test_load_col_path_too_long, splay_setup, splay_teardown }, { "splay/save_bulk_with_sym_path", test_save_bulk_with_sym_path, splay_setup, splay_teardown }, + { "splay/save_staged_bulk_defers_sym_flush", test_save_staged_bulk_defers_sym_flush, splay_setup, splay_teardown }, { "splay/save_dir_path_too_long", test_save_dir_path_too_long, splay_setup, splay_teardown }, { "splay/save_col_path_too_long", test_save_col_path_too_long, splay_setup, splay_teardown }, { "splay/trace_valid_dir", test_trace_valid_dir, splay_setup, splay_teardown },