Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
38 changes: 38 additions & 0 deletions docs/docs/storage/index.md
Original file line number Diff line number Diff line change
Expand Up @@ -64,6 +64,44 @@ 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.

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.

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.

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 |
Expand Down
40 changes: 31 additions & 9 deletions src/io/csv.c
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down Expand Up @@ -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);
Expand Down Expand Up @@ -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,
Expand Down
15 changes: 2 additions & 13 deletions src/ops/builtins.c
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand Down
5 changes: 0 additions & 5 deletions src/ops/system.c
Original file line number Diff line number Diff line change
Expand Up @@ -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;
}
Expand Down
17 changes: 13 additions & 4 deletions src/store/part.c
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down Expand Up @@ -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 */
Expand Down Expand Up @@ -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);
Expand Down
Loading
Loading