diff --git a/src/io/csv.c b/src/io/csv.c index 25302c1f..f4d583ef 100644 --- a/src/io/csv.c +++ b/src/io/csv.c @@ -151,6 +151,7 @@ static inline void scratch_free(ray_t* hdr) { typedef struct { const char* ptr; uint32_t len; + uint32_t hash; /* ray_hash_bytes of the field; filled for SYM columns */ } csv_strref_t; RAY_INLINE const char* scan_field(const char* p, const char* buf_end, @@ -952,9 +953,15 @@ static size_t csv_scan_split_at(const char* buf, size_t file_size, size_t s, /* Returns the row count (>=0), -1 if interrupted, or -2 when the parallel * path does not apply and the caller should run the serial scan. */ +/* force_quotes: the caller scans a byte window of a file and knows quotes + * occur somewhere in the file's data. The quote-aware state machine is + * then used even when this window holds none, so a window is split into + * rows exactly as the whole file would be (the quote-free fast path treats + * a lone '\r' and "\n\r" differently). */ static int64_t build_row_offsets_par(const char* buf, size_t buf_size, size_t data_offset, uint64_t prog_base, uint64_t prog_len, + bool force_quotes, int64_t** offsets_out, ray_t** hdr_out) { *offsets_out = NULL; *hdr_out = NULL; @@ -1023,7 +1030,7 @@ static int64_t build_row_offsets_par(const char* buf, size_t buf_size, total_q += quote_cnt[i]; total_term += term_cnt[i]; } - ctx.has_quotes = total_q != 0; + ctx.has_quotes = total_q != 0 || force_quotes; /* Now nudge the boundaries for the state machine pass 1 just chose. The * byte skipped is always '\n' or '\r', never '"', so the parities @@ -1201,7 +1208,7 @@ static int64_t build_row_offsets(const char* buf, size_t buf_size, * -2 means "not applicable" (small file, no pool, allocation refused) and * falls back to the serial scan, which defines the semantics. */ int64_t par = build_row_offsets_par(buf, buf_size, data_offset, - prog_base, prog_len, + prog_base, prog_len, false, offsets_out, hdr_out); if (par != -2) return par; @@ -1211,6 +1218,59 @@ static int64_t build_row_offsets(const char* buf, size_t buf_size, offsets_out, hdr_out); } +/* Row starts of the next `max_rows` rows from `data_offset`, found by the + * parallel scanner over a byte window instead of the serial walk: the + * window is sized from `avg_row` (bytes per row seen so far) with slack, and + * doubled when it holds fewer than max_rows + 1 row starts before the end + * of the file (the extra start proves the max_rows-th row is complete). + * Returns the row count, -1 if interrupted, or -2 when the parallel path + * does not apply (small window, no pool) — the caller then runs the serial + * limited scan. */ +static int64_t build_row_offsets_window(const char* buf, size_t buf_size, + size_t data_offset, int64_t max_rows, + size_t avg_row, bool data_has_quotes, + int64_t** offsets_out, ray_t** hdr_out, + size_t* next_offset_out) { + *offsets_out = NULL; *hdr_out = NULL; + if (next_offset_out) *next_offset_out = data_offset; + if (max_rows <= 0 || data_offset >= buf_size) return 0; + if (avg_row < 8) avg_row = 8; + size_t window = (size_t)max_rows * avg_row + (size_t)max_rows * avg_row / 4 + (64u << 10); + for (;;) { + size_t end = data_offset + window; + if (end > buf_size || end < data_offset) end = buf_size; + /* Readahead hint for the window about to be scanned: the scanner's + * tasks fault the pages in parallel, which a cold file serves best + * when the kernel already streams the range. */ + { + size_t ps = (size_t)sysconf(_SC_PAGESIZE); + size_t a = data_offset & ~(ps - 1); + madvise((void*)(buf + a), end - a, MADV_WILLNEED); + } + int64_t* offs = NULL; ray_t* hdr = NULL; + int64_t n = build_row_offsets_par(buf, end, data_offset, 0, 0, + data_has_quotes, &offs, &hdr); + if (n < 0) return n; /* -1 interrupted, -2 not applicable */ + if (n == 0) { scratch_free(hdr); return -2; } + if (n > max_rows) { + /* row max_rows - 1 ends before start max_rows: complete */ + if (next_offset_out) *next_offset_out = (size_t)offs[max_rows]; + *offsets_out = offs; *hdr_out = hdr; + return max_rows; + } + if (end == buf_size) { + /* every remaining row, the last one ended by the file */ + if (next_offset_out) *next_offset_out = buf_size; + *offsets_out = offs; *hdr_out = hdr; + return n; + } + /* the window held at most max_rows starts: widen and rescan */ + scratch_free(hdr); + if (window > SIZE_MAX / 2) return -2; + window *= 2; + } +} + static int64_t build_row_offsets_limited(const char* buf, size_t buf_size, size_t data_offset, int64_t max_rows, bool data_has_quotes, @@ -1381,6 +1441,7 @@ typedef struct { uint32_t len; const char* ptr; int64_t gid; /* global sym id, filled in step B */ + int64_t row; /* first row the string occurs on (domain batch order) */ } csv_dedup_ent_t; typedef struct { @@ -1438,8 +1499,32 @@ typedef struct { /* [n_cols] "this column wrote a canonical null", recorded where the value * is written. See csv_note_empty. */ bool* empties; + /* Target FILE domain (splayed save): the distinct strings are + * interned straight into the table's symfile domain, batched over + * every SYM column of the chunk, and the codes are positions in it. + * NULL: runtime symbol table, id-order-preserving serial walk. */ + struct ray_sym_domain_s* dom; + /* Hash partitions per SYM column (domain target only; 1 otherwise): + * dicts[i * n_part + p] holds the distinct strings of column cols[i] + * whose hash >> part_shift == p. Partitions of one column are + * deduplicated by independent tasks. */ + int n_part; + int part_shift; } csv_dedup_ctx_t; +static inline int csv_dedup_part(const csv_dedup_ctx_t* dd, uint32_t hash) { + return dd->n_part > 1 ? (int)(hash >> dd->part_shift) : 0; +} + +/* Every partition of column i deduplicated without overflow. */ +static inline bool csv_col_dict_ok(const csv_dedup_ctx_t* dd, int i) { + for (int p = 0; p < dd->n_part; p++) { + const csv_dedup_t* d = &dd->dicts[i * dd->n_part + p]; + if (!d->done || d->overflow) return false; + } + return true; +} + /* HAS_NULLS accounting for SYM/STR columns (step 9c). * * bfb5b380 made an empty SYM/STR cell a canonical null and, where the parse @@ -1460,10 +1545,14 @@ static void csv_dedup_task(void* arg, uint32_t worker_id, (void)worker_id; (void)end_i; csv_dedup_ctx_t* ctx = (csv_dedup_ctx_t*)arg; csv_dedup_t* d = &ctx->dicts[start]; - const csv_strref_t* refs = ctx->str_refs[ctx->cols[start]]; + int col_i = (int)(start / ctx->n_part); + int part = (int)(start % ctx->n_part); + const csv_strref_t* refs = ctx->str_refs[ctx->cols[col_i]]; /* Local codes land in the destination id array and are replaced in place - * by step C, so the dedupe needs no per-row scratch of its own. */ - uint32_t* codes = (uint32_t*)ctx->col_data[ctx->cols[start]]; + * by step C, so the dedupe needs no per-row scratch of its own. With + * partitions each task owns the rows whose hash falls in its partition; + * partition 0 also writes the null codes. */ + uint32_t* codes = (uint32_t*)ctx->col_data[ctx->cols[col_i]]; int64_t n_rows = ctx->n_rows; if (!csv_dedup_grow(d) || !csv_dedup_grow_ents(d)) { @@ -1475,8 +1564,9 @@ static void csv_dedup_task(void* arg, uint32_t worker_id, for (int64_t r = 0; r < n_rows; r++) { if (RAY_UNLIKELY((r & 1023) == 0 && ray_interrupted())) return; - if (refs[r].ptr == NULL) { codes[r] = 0; continue; } - uint32_t h = (uint32_t)ray_hash_bytes(refs[r].ptr, refs[r].len); + if (refs[r].ptr == NULL) { if (part == 0) codes[r] = 0; continue; } + uint32_t h = refs[r].hash; /* computed by the parse */ + if (csv_dedup_part(ctx, h) != part) continue; uint32_t mask = d->n_slots - 1; uint32_t j = h & mask; uint32_t found = 0; @@ -1507,6 +1597,7 @@ static void csv_dedup_task(void* arg, uint32_t worker_id, e->len = refs[r].len; e->ptr = refs[r].ptr; e->gid = 0; + e->row = r; d->n_ents++; codes[r] = d->n_ents; /* code = entry index + 1 */ d->slots[j] = d->n_ents; @@ -1523,8 +1614,9 @@ static void csv_dedup_map_task(void* arg, uint32_t worker_id, int64_t start, int64_t end_i) { (void)worker_id; (void)end_i; csv_dedup_ctx_t* ctx = (csv_dedup_ctx_t*)arg; - const csv_dedup_t* d = &ctx->dicts[start]; - if (d->overflow) return; /* step B already wrote real ids */ + if (!csv_col_dict_ok(ctx, (int)start)) return; /* step B already wrote real ids */ + const csv_dedup_t* dcol = &ctx->dicts[start * ctx->n_part]; + const csv_strref_t* refs = ctx->str_refs[ctx->cols[start]]; uint32_t* ids = (uint32_t*)ctx->col_data[ctx->cols[start]]; int64_t n_rows = ctx->n_rows; uint32_t empty = (uint32_t)ctx->empty_gid; @@ -1532,7 +1624,8 @@ static void csv_dedup_map_task(void* arg, uint32_t worker_id, for (int64_t r = 0; r < n_rows; r++) { if (RAY_UNLIKELY((r & 1023) == 0 && ray_interrupted())) return; uint32_t code = ids[r]; - uint32_t id = code ? (uint32_t)d->ents[code - 1].gid : empty; + uint32_t id = code ? (uint32_t)dcol[csv_dedup_part(ctx, refs[r].hash)].ents[code - 1].gid + : empty; ids[r] = id; saw_null |= (id == 0); } @@ -1551,9 +1644,135 @@ static void csv_dedup_map_task(void* arg, uint32_t worker_id, * empty string anyway, so collapsing them is the only deterministic answer the * parser can give. Empty fields are local code 0 out of the dedupe and are * mapped to that id by step C. */ +/* Step B for a FILE domain target: one batch for the whole chunk, laid out + * the way the cell-by-cell writer met the strings — columns in order, each + * column's strings by first occurrence — so the symfile gets the same + * positions whatever the worker count and however a column's dictionary + * was split into hash partitions. The domain probes the existing + * vocabulary in parallel and appends the new strings in batch order. A + * column whose dictionary overflowed contributes its rows directly (the + * batch dedupes them) and gets its ids written here. */ +/* Append column i's dictionary entries to out[] by first row. Each hash + * partition lists its strings in row order already, so this is a merge of + * n_part sorted runs (n_part is small). */ +static int64_t csv_col_ents_by_row(csv_dedup_ctx_t* dd, int i, csv_dedup_ent_t** out) { + int np = dd->n_part; + csv_dedup_t* d0 = &dd->dicts[i * np]; + if (np == 1) { + for (uint32_t e = 0; e < d0->n_ents; e++) out[e] = &d0->ents[e]; + return d0->n_ents; + } + uint32_t cur[64]; + for (int p = 0; p < np; p++) cur[p] = 0; + int64_t n = 0; + for (;;) { + int best = -1; + int64_t br = 0; + for (int p = 0; p < np; p++) { + if (cur[p] >= d0[p].n_ents) continue; + int64_t r = d0[p].ents[cur[p]].row; + if (best < 0 || r < br) { best = p; br = r; } + } + if (best < 0) break; + out[n++] = &d0[best].ents[cur[best]++]; + } + return n; +} + +static bool csv_intern_dicts_domain(csv_dedup_ctx_t* dd, int n_sym, + int64_t* col_max_ids, + uint64_t prog_base, uint64_t prog_len) { + struct ray_sym_domain_s* dom = dd->dom; + /* Position 0 of a symfile domain is "" (reserved on creation). */ + if (ray_sym_domain_intern(dom, "", 0) != 0) return false; + dd->empty_gid = 0; + + int64_t total = 0; + for (int i = 0; i < n_sym; i++) { + if (csv_col_dict_ok(dd, i)) { + for (int p = 0; p < dd->n_part; p++) total += dd->dicts[i * dd->n_part + p].n_ents; + } else { + total += dd->n_rows; + } + } + if (total == 0) { + if (prog_len) ray_progress_span_set(prog_base + prog_len); + return true; + } + + ray_t *hs = NULL, *hl = NULL, *hh = NULL, *hp = NULL, *ho = NULL; + const char** strs = (const char**)scratch_alloc(&hs, (size_t)total * sizeof(char*)); + size_t* lens = (size_t*)scratch_alloc(&hl, (size_t)total * sizeof(size_t)); + uint32_t* hashes = (uint32_t*)scratch_alloc(&hh, (size_t)total * sizeof(uint32_t)); + int64_t* pos = (int64_t*)scratch_alloc(&hp, (size_t)total * sizeof(int64_t)); + /* order[k]: the dictionary entry behind batch slot k (dict columns) */ + csv_dedup_ent_t** order = (csv_dedup_ent_t**)scratch_alloc(&ho, + (size_t)total * sizeof(csv_dedup_ent_t*)); + bool ok = strs && lens && hashes && pos && order; + + int64_t k = 0; + for (int i = 0; ok && i < n_sym; i++) { + if (csv_col_dict_ok(dd, i)) { + int64_t ne = csv_col_ents_by_row(dd, i, order + k); + for (int64_t e = 0; e < ne; e++, k++) { + strs[k] = order[k]->ptr; + lens[k] = order[k]->len; + hashes[k] = order[k]->hash; + } + } else { + const csv_strref_t* refs = dd->str_refs[dd->cols[i]]; + for (int64_t r = 0; r < dd->n_rows; r++) { + if (refs[r].ptr == NULL) continue; + strs[k] = refs[r].ptr; + lens[k] = refs[r].len; + hashes[k] = refs[r].hash; + k++; + } + } + } + if (ok) ok = ray_sym_domain_intern_batch(dom, k, strs, lens, hashes, pos); + + /* Hand the positions back in the same walk. */ + k = 0; + for (int i = 0; ok && i < n_sym; i++) { + int c = dd->cols[i]; + int64_t max_id = 0; + if (csv_col_dict_ok(dd, i)) { + int64_t ne = 0; + for (int p = 0; p < dd->n_part; p++) ne += dd->dicts[i * dd->n_part + p].n_ents; + for (int64_t e = 0; e < ne; e++, k++) { + int64_t id = pos[k]; + if (id < 0) { ok = false; id = 0; } + order[k]->gid = id; + if (id > max_id) max_id = id; + } + } else { + const csv_strref_t* refs = dd->str_refs[c]; + uint32_t* ids = (uint32_t*)dd->col_data[c]; + bool saw_null = false; + for (int64_t r = 0; r < dd->n_rows; r++) { + if (refs[r].ptr == NULL) { ids[r] = 0; saw_null = true; continue; } + int64_t id = pos[k++]; + if (id < 0) { ok = false; id = 0; } + ids[r] = (uint32_t)id; + saw_null |= (id == 0); + if (id > max_id) max_id = id; + } + csv_note_empty(dd->empties, c, saw_null); + } + if (col_max_ids) col_max_ids[c] = max_id; + } + + scratch_free(hs); scratch_free(hl); scratch_free(hh); scratch_free(hp); scratch_free(ho); + if (prog_len) ray_progress_span_set(prog_base + prog_len); + return ok; +} + static bool csv_intern_dicts(csv_dedup_ctx_t* dd, int n_sym, int64_t* col_max_ids, uint64_t prog_base, uint64_t prog_len) { + if (dd->dom) + return csv_intern_dicts_domain(dd, n_sym, col_max_ids, prog_base, prog_len); bool ok = true; /* Same first call, same reason, as the serial walk: sym 0 is reserved by @@ -1739,6 +1958,7 @@ typedef struct { int n_sym; bool* empties; /* [n_cols] SYM/STR column wrote a null */ bool intern_ok; + struct ray_sym_domain_s* dom; /* SYM target domain, NULL = runtime */ } csv_finalize_ctx_t; /* dispatch 1: [0, n_fill) fill a RAY_STR column, [n_fill, n_fill+n_sym) dedupe @@ -1768,7 +1988,7 @@ static void csv_finalize_task(void* arg, uint32_t worker_id, * prog_base/prog_len describe this phase's slice of the load's byte axis; pass * len 0 when no progress span is active (the streaming conversion path). */ static bool csv_finalize_run(csv_finalize_ctx_t* ctx, int* fill_cols, - bool* fill_ok, int* sym_cols, csv_dedup_t* dicts, + bool* fill_ok, int* sym_cols, bool* empties, uint64_t prog_base, uint64_t prog_len) { int n_fill = 0, n_sym = 0; @@ -1778,7 +1998,31 @@ static bool csv_finalize_run(csv_finalize_ctx_t* ctx, int* fill_cols, else sym_cols[n_sym++] = c; } for (int i = 0; i < n_fill; i++) fill_ok[i] = true; - memset(dicts, 0, (size_t)n_sym * sizeof(csv_dedup_t)); + + ray_pool_t* pool = ray_pool_get(); + bool par = pool && ray_pool_total_workers(pool) >= 2; + + /* Partitions per SYM column. The runtime path keeps one dictionary per + * column (its id order is the serial walk's); a domain target has no + * order to keep, so the heavy columns are split by hash until the + * dedupe tasks cover the pool about twice over. */ + int n_part = 1; + if (ctx->dom && par && n_sym > 0) { + int64_t want = (int64_t)ray_pool_total_workers(pool) * 2 / n_sym; + while (n_part * 2 <= want && n_part < 16) n_part *= 2; + while (n_part > 1 && (int64_t)n_fill + (int64_t)n_sym * n_part > (int64_t)RAY_POOL_INIT_TASKS) + n_part /= 2; + } + int part_shift = 32; + for (int q = n_part; q > 1; q >>= 1) part_shift--; + + ray_t* dicts_hdr = NULL; + csv_dedup_t* dicts = NULL; + if (n_sym > 0) { + dicts = (csv_dedup_t*)scratch_calloc(&dicts_hdr, + (size_t)n_sym * (size_t)n_part * sizeof(csv_dedup_t)); + if (!dicts) return false; + } ctx->fill_cols = fill_cols; ctx->n_fill = n_fill; @@ -1793,6 +2037,9 @@ static bool csv_finalize_run(csv_finalize_ctx_t* ctx, int* fill_cols, ctx->dd.n_rows = ctx->n_rows; ctx->dd.empty_gid = 0; ctx->dd.empties = ctx->empties; + ctx->dd.dom = ctx->dom; + ctx->dd.n_part = n_part; + ctx->dd.part_shift = part_shift; /* Slice the phase: dedupe+fill is the bulk, the serial intern touches only * distinct strings, the remap is one linear pass per SYM column. Measured @@ -1800,10 +2047,8 @@ static bool csv_finalize_run(csv_finalize_ctx_t* ctx, int* fill_cols, uint64_t w1 = prog_len * 6 / 10; uint64_t w2 = prog_len / 10; - int64_t n_tasks = (int64_t)n_fill + (int64_t)n_sym; - ray_pool_t* pool = ray_pool_get(); - bool par = pool && ray_pool_total_workers(pool) >= 2 && - n_tasks > 0 && n_tasks <= (int64_t)RAY_POOL_INIT_TASKS; + int64_t n_tasks = (int64_t)n_fill + (int64_t)n_sym * n_part; + par = par && n_tasks > 0 && n_tasks <= (int64_t)RAY_POOL_INIT_TASKS; if (prog_len) ray_progress_span_phase("finalize", prog_base, w1); if (par) ray_pool_dispatch_n(pool, csv_finalize_task, ctx, (uint32_t)n_tasks); @@ -1830,12 +2075,14 @@ static bool csv_finalize_run(csv_finalize_ctx_t* ctx, int* fill_cols, if (ray_interrupted()) goto fail; } - for (int i = 0; i < n_sym; i++) csv_dedup_release(&dicts[i]); + for (int i = 0; i < n_sym * n_part; i++) csv_dedup_release(&dicts[i]); + scratch_free(dicts_hdr); for (int i = 0; i < n_fill; i++) if (!fill_ok[i]) return false; return true; fail: - for (int i = 0; i < n_sym; i++) csv_dedup_release(&dicts[i]); + for (int i = 0; i < n_sym * n_part; i++) csv_dedup_release(&dicts[i]); + scratch_free(dicts_hdr); return false; } @@ -2021,6 +2268,10 @@ static void csv_parse_fn(void* arg, uint32_t worker_id, } ctx->str_refs[c][row].ptr = fld; ctx->str_refs[c][row].len = (uint32_t)flen; + /* SYM columns are deduplicated by hash right after + * the parse; hash here while the bytes are hot. */ + if (ctx->resolved_types[c] == RAY_SYM) + ctx->str_refs[c][row].hash = (uint32_t)ray_hash_bytes(fld, flen); } break; } @@ -2202,6 +2453,8 @@ static bool csv_parse_serial(const char* buf, size_t buf_size, } str_refs[c][row].ptr = fld; str_refs[c][row].len = (uint32_t)flen; + if (resolved_types[c] == RAY_SYM) + str_refs[c][row].hash = (uint32_t)ray_hash_bytes(fld, flen); } break; } @@ -2479,7 +2732,8 @@ static ray_t* csv_materialize_rows(const char* buf, size_t file_size, const int64_t* row_offsets, int64_t n_rows, int ncols, char delimiter, const int64_t* col_name_ids, - const int8_t* resolved_types) { + const int8_t* resolved_types, + struct ray_sym_domain_s* sym_dom) { /* Defensive guard: RAY_CSV_AUTO_TAG must be resolved to a concrete width * before reaching this point (by csv_resolve_auto_in_place or * csv_resolve_auto_streamed). If a marker slips through, the resolution @@ -2505,6 +2759,11 @@ static ray_t* csv_materialize_rows(const char* buf, size_t file_size, for (int j = 0; j < c; j++) ray_release(col_vecs[j]); return NULL; } + if (type == RAY_SYM && sym_dom) { + /* Cells are positions in the target symfile domain. */ + ray_sym_domain_retain(sym_dom); + col_vecs[c]->sym_domain = sym_dom; + } col_vecs[c]->len = n_rows; col_data[c] = ray_data(col_vecs[c]); } @@ -2645,14 +2904,14 @@ static ray_t* csv_materialize_rows(const char* buf, size_t file_size, .col_vecs = col_vecs, .n_rows = n_rows, .sym_max_ids = sym_max_ids, + .dom = sym_dom, }; int fill_cols[CSV_MAX_COLS]; int sym_cols[CSV_MAX_COLS]; bool fill_ok[CSV_MAX_COLS]; - csv_dedup_t dicts[CSV_MAX_COLS]; /* No progress span on the conversion path — pass a zero-length slice. */ bool fin_ok = csv_finalize_run(&fctx, fill_cols, fill_ok, - sym_cols, dicts, col_wrote_null, 0, 0); + sym_cols, col_wrote_null, 0, 0); if (!fin_ok || ray_interrupted()) { csv_free_escaped_strrefs(str_ref_bufs, ncols, parse_types, n_rows, buf, file_size, row_done, col_had_escaped); @@ -2693,6 +2952,7 @@ static ray_t* csv_materialize_rows(const char* buf, size_t file_size, if (new_w >= RAY_SYM_W32) continue; ray_t* narrow = ray_sym_vec_new(new_w, n_rows); if (!narrow || RAY_IS_ERR(narrow)) continue; + ray_sym_vec_adopt_domain(narrow, col_vecs[c]); narrow->len = n_rows; const uint32_t* src = (const uint32_t*)col_data[c]; void* dst = ray_data(narrow); @@ -3123,9 +3383,8 @@ static ray_t* csv_read_named_opts_inner(const char* path, char delimiter, bool h int fill_cols[CSV_MAX_COLS]; int sym_cols[CSV_MAX_COLS]; bool fill_ok[CSV_MAX_COLS]; - csv_dedup_t dicts[CSV_MAX_COLS]; bool fin_ok = csv_finalize_run(&fctx, fill_cols, fill_ok, - sym_cols, dicts, col_wrote_null, + sym_cols, col_wrote_null, prog_parse_end, (uint64_t)file_size - prog_parse_end); if (!fin_ok || ray_interrupted()) { @@ -3283,7 +3542,15 @@ typedef struct { * fixed at W32: a streaming writer can't know the final vocabulary * before the last chunk, and W32 covers any STRL count. */ struct ray_sym_domain_s* dom; + /* runtime id -> domain position, direct-mapped: the chunk vecs are + * runtime-domain and a column's values repeat across rows and chunks, + * so a value is interned into the symfile's domain (a locked probe) the + * first time it is met and looked up here after. */ + int64_t* lut_id; /* [CSV_SPLAYED_LUT] runtime id per slot, -1 empty */ + uint32_t* lut_pos; /* [CSV_SPLAYED_LUT] position per slot */ } csv_splayed_col_writer_t; +#define CSV_SPLAYED_LUT_BITS 19 +#define CSV_SPLAYED_LUT (1u << CSV_SPLAYED_LUT_BITS) static ray_err_t csv_splayed_writer_open(csv_splayed_col_writer_t* w, const char* dir, int64_t name_id, @@ -3314,9 +3581,25 @@ static ray_err_t csv_splayed_writer_open(csv_splayed_col_writer_t* w, if (!w->fp) return RAY_ERR_IO; ray_t zero = {0}; if (fwrite(&zero, 1, 32, w->fp) != 32) return RAY_ERR_IO; + if (type == RAY_SYM) { + /* best effort: without the cache every cell probes the domain */ + w->lut_id = (int64_t*)ray_alloc_raw((size_t)CSV_SPLAYED_LUT * sizeof(int64_t)); + w->lut_pos = (uint32_t*)ray_alloc_raw((size_t)CSV_SPLAYED_LUT * sizeof(uint32_t)); + if (!w->lut_id || !w->lut_pos) { + ray_free_raw(w->lut_id); ray_free_raw(w->lut_pos); + w->lut_id = NULL; w->lut_pos = NULL; + } else { + memset(w->lut_id, 0xff, (size_t)CSV_SPLAYED_LUT * sizeof(int64_t)); + } + } return RAY_OK; } +static void csv_splayed_writer_drop_lut(csv_splayed_col_writer_t* w) { + ray_free_raw(w->lut_id); ray_free_raw(w->lut_pos); + w->lut_id = NULL; w->lut_pos = NULL; +} + static ray_err_t csv_splayed_writer_append(csv_splayed_col_writer_t* w, ray_t* col) { if (!w->fp || !col || RAY_IS_ERR(col)) return RAY_ERR_TYPE; @@ -3328,18 +3611,40 @@ static ray_err_t csv_splayed_writer_append(csv_splayed_col_writer_t* w, * resolve each cell through the chunk vec's own domain and * find-or-append into the target (distinct work rides the * write). The domain is flushed before the column files are - * committed (close), preserving the sym-first crash ordering. */ + * committed (close), preserving the sym-first crash ordering. + * A runtime-domain chunk vec goes through the id -> position + * cache: only a value's first encounter pays the domain probe. */ + bool direct = ray_sym_vec_domain(col) == w->dom; + bool cached = ray_sym_vec_domain(col) == ray_sym_runtime_domain() && w->lut_id; + const void* cd = ray_data(col); uint32_t buf[8192]; for (int64_t off = 0; off < n; ) { int64_t cnt = n - off; if (cnt > (int64_t)(sizeof(buf) / sizeof(buf[0]))) cnt = (int64_t)(sizeof(buf) / sizeof(buf[0])); for (int64_t i = 0; i < cnt; i++) { - ray_t* s = ray_sym_vec_cell(col, off + i); - if (!s) return RAY_ERR_CORRUPT; - int64_t pos = ray_sym_domain_intern(w->dom, ray_str_ptr(s), - ray_str_len(s)); - if (pos < 0) return RAY_ERR_OOM; + int64_t pos; + if (direct) { + /* Already encoded over the target domain. */ + pos = ray_read_sym(cd, off + i, RAY_SYM, col->attrs); + } else if (cached) { + int64_t id = ray_read_sym(cd, off + i, RAY_SYM, col->attrs); + uint32_t slot = (uint32_t)(((uint64_t)id * 0x9E3779B97F4A7C15ull) >> (64 - CSV_SPLAYED_LUT_BITS)); + if (w->lut_id[slot] == id) { + pos = w->lut_pos[slot]; + } else { + ray_t* s = ray_sym_str(id); + if (!s) return RAY_ERR_CORRUPT; + pos = ray_sym_domain_intern(w->dom, ray_str_ptr(s), ray_str_len(s)); + if (pos < 0) return RAY_ERR_OOM; + w->lut_id[slot] = id; w->lut_pos[slot] = (uint32_t)pos; + } + } else { + ray_t* s = ray_sym_vec_cell(col, off + i); + if (!s) return RAY_ERR_CORRUPT; + pos = ray_sym_domain_intern(w->dom, ray_str_ptr(s), ray_str_len(s)); + if (pos < 0) return RAY_ERR_OOM; + } buf[i] = (uint32_t)pos; /* Position 0 of any symfile is the empty string (domain.c * enforces that reservation on open), so a re-encoded cell is @@ -3365,7 +3670,29 @@ static ray_err_t csv_splayed_writer_append(csv_splayed_col_writer_t* w, return RAY_OK; } +/* Append task: column `start` of the chunk table into its writer. The + * first failure is kept (a later task cannot clear it). */ +typedef struct { + csv_splayed_col_writer_t* writers; + ray_t* tbl; + int ncols; + _Atomic(ray_err_t) err; +} csv_splayed_append_ctx_t; + +static void csv_splayed_append_task(void* raw, uint32_t wid, int64_t start, int64_t end) { + (void)wid; (void)end; + csv_splayed_append_ctx_t* a = (csv_splayed_append_ctx_t*)raw; + if (atomic_load_explicit(&a->err, memory_order_relaxed) != RAY_OK) return; + ray_t* col = ray_table_get_col_idx(a->tbl, (int64_t)start); + ray_err_t e = csv_splayed_writer_append(&a->writers[start], col); + if (e != RAY_OK) { + ray_err_t ok = RAY_OK; + atomic_compare_exchange_strong_explicit(&a->err, &ok, e, memory_order_relaxed, memory_order_relaxed); + } +} + static ray_err_t csv_splayed_writer_close(csv_splayed_col_writer_t* w) { + csv_splayed_writer_drop_lut(w); if (!w->fp) return RAY_OK; ray_err_t err = RAY_OK; @@ -3397,6 +3724,7 @@ static ray_err_t csv_splayed_writer_close(csv_splayed_col_writer_t* w) { } static void csv_splayed_writer_abort(csv_splayed_col_writer_t* w) { + csv_splayed_writer_drop_lut(w); if (w->fp) fclose(w->fp); w->fp = NULL; remove(w->tmp_path); @@ -3630,16 +3958,27 @@ ray_err_t ray_csv_save_splayed_named_opts(const char* path, char delimiter, bool size_t chunk_offset = data_offset; bool wrote_any = false; + size_t avg_row_bytes = 64; /* refined from every chunk scanned */ while (chunk_offset < file_size || !wrote_any) { ray_t* row_offsets_hdr = NULL; int64_t* row_offsets = NULL; size_t next_offset = chunk_offset; int64_t cnt = 0; if (chunk_offset < file_size) { - cnt = build_row_offsets_limited(buf, file_size, chunk_offset, - rows_per_chunk, data_has_quotes, - &row_offsets, - &row_offsets_hdr, &next_offset); + /* Parallel scan over a byte window sized from the rows seen so + * far; the serial walk remains the fallback and the semantics. */ + cnt = build_row_offsets_window(buf, file_size, chunk_offset, + rows_per_chunk, avg_row_bytes, + data_has_quotes, + &row_offsets, &row_offsets_hdr, + &next_offset); + if (cnt == -2) + cnt = build_row_offsets_limited(buf, file_size, chunk_offset, + rows_per_chunk, data_has_quotes, + &row_offsets, + &row_offsets_hdr, &next_offset); + if (cnt > 0 && next_offset > chunk_offset) + avg_row_bytes = (next_offset - chunk_offset) / (size_t)cnt; if (cnt <= 0) { scratch_free(row_offsets_hdr); err = (cnt < 0) ? RAY_ERR_CANCEL : RAY_ERR_IO; @@ -3649,7 +3988,7 @@ ray_err_t ray_csv_save_splayed_named_opts(const char* path, char delimiter, bool ray_t* tbl = csv_materialize_rows(buf, file_size, row_offsets, cnt, ncols, delimiter, col_name_ids, - resolved_types); + resolved_types, sym_dom); scratch_free(row_offsets_hdr); if (!tbl || RAY_IS_ERR(tbl)) { err = (tbl && RAY_IS_ERR(tbl)) ? ray_err_from_obj(tbl) @@ -3658,14 +3997,31 @@ ray_err_t ray_csv_save_splayed_named_opts(const char* path, char delimiter, bool break; } - for (int c = 0; c < ncols; c++) { - ray_t* col = ray_table_get_col_idx(tbl, c); - err = csv_splayed_writer_append(&writers[c], col); - if (err != RAY_OK) break; + /* One task per column: each writer owns its file, its cache and + * its symfile domain (the domain probe takes the domain lock, the + * symbol table read its own), so the columns of a chunk are + * encoded and written side by side. */ + { + csv_splayed_append_ctx_t actx = { .writers = writers, .tbl = tbl, + .ncols = ncols, .err = RAY_OK }; + ray_pool_t* wpool = ray_pool_get(); + if (ray_pool_par_dispatch_ok(wpool, ncols, 2)) + ray_pool_dispatch_n(wpool, csv_splayed_append_task, &actx, (uint32_t)ncols); + else + for (int c = 0; c < ncols; c++) csv_splayed_append_task(&actx, 0, c, c + 1); + err = actx.err; } ray_release(tbl); if (err != RAY_OK) break; wrote_any = true; + /* The chunk's bytes are done with: drop them from the mapping so a + * long file does not pin its whole length in resident memory. */ + if (next_offset > chunk_offset) { + size_t ps = (size_t)sysconf(_SC_PAGESIZE); + size_t a = (chunk_offset + ps - 1) & ~(ps - 1); + size_t b = next_offset & ~(ps - 1); + if (b > a) madvise((void*)(buf + a), b - a, MADV_DONTNEED); + } if (cnt == 0) break; chunk_offset = next_offset; } @@ -3673,8 +4029,9 @@ ray_err_t ray_csv_save_splayed_named_opts(const char* path, char delimiter, bool /* Flush the symfile BEFORE committing column files (writer_close * renames tmp → final): columns must never reference positions the * symfile doesn't persist (sym-first crash ordering). */ - if (err == RAY_OK && sym_dom) + if (err == RAY_OK && sym_dom) { err = ray_sym_domain_flush(sym_dom, false); + } for (int c = 0; c < ncols; c++) { ray_err_t cerr = (err == RAY_OK) ? csv_splayed_writer_close(&writers[c]) @@ -3886,7 +4243,8 @@ ray_err_t ray_csv_save_parted_named_opts(const char* path, char delimiter, bool } ray_t* tbl = csv_materialize_rows(buf, file_size, row_offsets, - cnt, ncols, delimiter, col_name_ids, resolved_types); + cnt, ncols, delimiter, col_name_ids, resolved_types, + NULL); if (!tbl || RAY_IS_ERR(tbl)) { err = (tbl && RAY_IS_ERR(tbl)) ? ray_err_from_obj(tbl) : RAY_ERR_OOM; diff --git a/src/mem/arena.c b/src/mem/arena.c index df44caf8..cc2e6e0b 100644 --- a/src/mem/arena.c +++ b/src/mem/arena.c @@ -109,26 +109,54 @@ ray_t* ray_arena_alloc(ray_arena_t* arena, size_t nbytes) { return v; } -ray_t* ray_arena_str(ray_arena_t* arena, const char* s, size_t len) { +size_t ray_arena_str_bytes(size_t len) { + if (len < 7) return 32; + /* [U8 header (32) | data (len+1) | pad to 32 | STR header (32)] */ + return (((32 + len + 1) + 31) & ~(size_t)31) + 32; +} + +void* ray_arena_alloc_raw(ray_arena_t* arena, size_t nbytes) { + if (!arena) return NULL; + if (nbytes > SIZE_MAX - (ARENA_ALIGN - 1)) return NULL; + size_t block_size = ARENA_ALIGN_UP(nbytes); + ray_arena_chunk_t* c = arena->chunks; + if (c->used + block_size > c->cap) { + size_t new_cap = arena->chunk_size; + if (block_size > new_cap) new_cap = ARENA_ALIGN_UP(block_size); + ray_arena_chunk_t* nc = arena_new_chunk(new_cap); + if (!nc) return NULL; + nc->next = arena->chunks; + arena->chunks = nc; + c = nc; + } + void* p = chunk_data(c) + c->used; + c->used += block_size; + return p; +} + +ray_t* ray_arena_str_at(void* at, const char* s, size_t len) { if (len < 7) { - /* SSO: bytes inline in the header (ray_arena_alloc zeroes it and sets - * RAY_ATTR_ARENA + rc=1). */ - ray_t* v = ray_arena_alloc(arena, 0); - if (!v) return NULL; + /* SSO: bytes inline in the header. */ + ray_t* v = (ray_t*)at; + memset(v, 0, 32); + v->attrs = RAY_ATTR_ARENA; + ray_atomic_store(&v->rc, 1); v->type = -RAY_STR; v->slen = (uint8_t)len; if (len > 0) memcpy(v->sdata, s, len); v->sdata[len] = '\0'; return v; } - /* Long string: fused single allocation for the U8 data vec + the STR atom. + /* Long string: fused single block for the U8 data vec + the STR atom. * Layout: [U8 ray_t header (32) | data (len+1) | pad to 32 | STR header (32)]. - * One arena_alloc instead of two. 32-byte arena alignment keeps the atom's - * obj pointer low byte out of is_sso()'s 1..7 SSO range. */ + * 32-byte arena alignment keeps the atom's obj pointer low byte out of + * is_sso()'s 1..7 SSO range. */ size_t data_size = len + 1; size_t chars_block = ((32 + data_size) + 31) & ~(size_t)31; /* align up to 32 */ - ray_t* chars = ray_arena_alloc(arena, chars_block); - if (!chars) return NULL; + ray_t* chars = (ray_t*)at; + memset(chars, 0, 32); + chars->attrs = RAY_ATTR_ARENA; + ray_atomic_store(&chars->rc, 1); chars->type = RAY_U8; chars->len = (int64_t)len; memcpy(ray_data(chars), s, len); @@ -143,6 +171,12 @@ ray_t* ray_arena_str(ray_arena_t* arena, const char* s, size_t len) { return v; } +ray_t* ray_arena_str(ray_arena_t* arena, const char* s, size_t len) { + void* at = ray_arena_alloc_raw(arena, ray_arena_str_bytes(len)); + if (!at) return NULL; + return ray_arena_str_at(at, s, len); +} + bool ray_arena_reserve(ray_arena_t* arena, size_t bytes) { if (!arena) return false; if (bytes == 0) return true; diff --git a/src/mem/arena.h b/src/mem/arena.h index 1cce80eb..d7861fd5 100644 --- a/src/mem/arena.h +++ b/src/mem/arena.h @@ -45,6 +45,14 @@ ray_t* ray_arena_alloc(ray_arena_t* arena, size_t nbytes); * heap (used by the global sym table and FILE sym domains). NULL on OOM. */ ray_t* ray_arena_str(ray_arena_t* arena, const char* s, size_t len); +/* Bulk string construction: reserve one raw region for many atoms, then + * build each atom in place (ray_arena_str == alloc_raw + str_at). The + * region is arena memory: 32-byte aligned, released with the arena. Lets + * a caller carve one region per worker and build atoms in parallel. */ +size_t ray_arena_str_bytes(size_t len); /* bytes one atom needs */ +void* ray_arena_alloc_raw(ray_arena_t* arena, size_t nbytes); /* NULL on OOM */ +ray_t* ray_arena_str_at(void* at, const char* s, size_t len); /* at: ray_arena_str_bytes(len) */ + /* Ensure the arena can serve subsequent allocations totalling at least * `bytes` without the head chunk needing to grow. If the head chunk has * enough free space already, this is a no-op; otherwise a new chunk with diff --git a/src/ops/idxop.c b/src/ops/idxop.c index 6927a52d..b4ec5f78 100644 --- a/src/ops/idxop.c +++ b/src/ops/idxop.c @@ -32,6 +32,8 @@ #include "ops/ops.h" #include "ops/rowsel.h" #include "ops/hash.h" /* ray_hash_bytes: STR hash-index key word */ +#include "core/pool.h" /* parallel hash-index build */ +#include "mem/sys.h" /* ray_sys_alloc: build scratch off the buddy heap */ #include #include #include @@ -1084,6 +1086,324 @@ ray_t* ray_index_inline_map(uint8_t* region) { * hit yields the contiguous ascending slice rows[offs[gid]..offs[gid+1]). * -------------------------------------------------------------------------- */ +/* Parallel build of the CSR hash layout for numeric / SYM keys. The + * result is the serial walk's: groups numbered by first occurrence, rows + * ascending inside a group, nulls excluded — assembled in + * partition-parallel passes: + * + * A row ranges: key word per row, rows bucketed by hash partition + * (a partition's rows stay ascending: the ranges are in row order); + * B per partition: open-addressing dedupe into local groups, each + * group's first row marked in a bitmap; + * C a group's number is the rank of its first row among all marked + * rows (block popcounts + one prefix) — exactly the order the + * serial walk assigns; + * D per partition: keys and counts into the group's slots; + * E per partition: row scatter (of[] as cursor); then the key table, + * filled in group order on the calling thread so its bytes match the + * serial build. + * + * Every array a pass writes at random is faulted in beforehand from all + * workers in slices: a fresh mapping faulted at random from every worker + * serialises on the page-table locks. Returns false (nothing allocated) + * when the pool cannot be used; the caller then runs the serial walk. */ +typedef struct { + ray_t* v; + const uint8_t* base; + int64_t n; + int n_tasks; + int n_part; + int part_shift; /* partition = mix64(key) >> shift */ + uint64_t* kw; /* [n] key word (non-null rows) */ + int64_t* pr; /* [n_keys] row ids grouped by partition */ + int64_t* lg; /* [n_keys] local group of pr[j] */ + int64_t* cnt; /* [n_tasks * n_part] rows per (task, partition) */ + int64_t* part_off; /* [n_part + 1] */ + int64_t* ng_p; /* [n_part] local groups per partition */ + int64_t* gfirst; /* [n_keys] first row per local group (partition-relative) */ + int64_t* gcount; /* [n_keys] rows per local group, then its global number */ + uint64_t* bits; /* [n/64 + 1] first-row marks */ + int64_t* blk_rank; /* [n/HP_BLOCK + 1] exclusive prefix of block popcounts */ + int64_t* gk; /* [n_groups] keys */ + int64_t* of; /* [n_groups + 1] */ + int64_t* rw; /* [n_keys] */ + int64_t* tbl; /* [cap] */ + uint64_t tmask; + _Atomic(bool) oom; +} hash_par_t; + +#define HP_BLOCK 4096 + +static inline int64_t hp_task_lo(const hash_par_t* h, int64_t t) { return h->n * t / h->n_tasks; } + +static void hp_pass_a(void* raw, uint32_t wid, int64_t start, int64_t end) { + (void)wid; (void)end; + hash_par_t* h = (hash_par_t*)raw; + int64_t lo = hp_task_lo(h, start), hi = hp_task_lo(h, start + 1); + int64_t* cnt = h->cnt + start * h->n_part; + for (int64_t i = lo; i < hi; i++) { + if (ray_vec_is_null(h->v, i)) { h->kw[i] = 0; continue; } + uint64_t k = hash_row_key_word(h->v, h->base, i); + h->kw[i] = k; + cnt[mix64(k) >> h->part_shift]++; + } +} + +static void hp_pass_a2(void* raw, uint32_t wid, int64_t start, int64_t end) { + (void)wid; (void)end; + hash_par_t* h = (hash_par_t*)raw; + int64_t lo = hp_task_lo(h, start), hi = hp_task_lo(h, start + 1); + int64_t* cur = h->cnt + start * h->n_part; /* now the write cursors */ + for (int64_t i = lo; i < hi; i++) { + if (ray_vec_is_null(h->v, i)) continue; + h->pr[cur[mix64(h->kw[i]) >> h->part_shift]++] = i; + } +} + +static void hp_pass_b(void* raw, uint32_t wid, int64_t start, int64_t end) { + (void)wid; (void)end; + hash_par_t* h = (hash_par_t*)raw; + int64_t p = start; + int64_t lo = h->part_off[p], hi = h->part_off[p + 1]; + int64_t cnt = hi - lo; + h->ng_p[p] = 0; + if (cnt == 0) return; + uint64_t cap = next_pow2((uint64_t)cnt * 2 + 1); + if (cap < 16) cap = 16; + int64_t* tab = (int64_t*)ray_sys_alloc((size_t)cap * sizeof(int64_t)); + if (!tab) { atomic_store_explicit(&h->oom, true, memory_order_relaxed); return; } + memset(tab, 0, (size_t)cap * sizeof(int64_t)); + uint64_t mask = cap - 1; + int64_t* gfirst = h->gfirst + lo; + int64_t* gcount = h->gcount + lo; + int64_t ng = 0; + for (int64_t j = lo; j < hi; j++) { + int64_t i = h->pr[j]; + uint64_t k = h->kw[i]; + uint64_t slot = mix64(k) & mask; + for (;;) { + int64_t g1 = tab[slot]; + if (g1 == 0) { + tab[slot] = ng + 1; + gfirst[ng] = i; + gcount[ng] = 1; + h->lg[j] = ng; + /* the partition's rows are ascending, so i is this group's + * first row; the word is shared with other partitions */ + atomic_fetch_or_explicit((_Atomic(uint64_t)*)&h->bits[i >> 6], + (uint64_t)1 << (i & 63), memory_order_relaxed); + ng++; + break; + } + if (h->kw[gfirst[g1 - 1]] == k) { h->lg[j] = g1 - 1; gcount[g1 - 1]++; break; } + slot = (slot + 1) & mask; + } + } + h->ng_p[p] = ng; + ray_sys_free(tab); +} + +/* per-block popcount of the first-row marks (prefixed serially after) */ +static void hp_pass_c(void* raw, uint32_t wid, int64_t start, int64_t end) { + (void)wid; (void)end; + hash_par_t* h = (hash_par_t*)raw; + int64_t w0 = start * (HP_BLOCK / 64), w1 = w0 + HP_BLOCK / 64; + int64_t nw = (h->n + 63) / 64; + if (w1 > nw) w1 = nw; + int64_t c = 0; + for (int64_t w = w0; w < w1; w++) c += __builtin_popcountll(h->bits[w]); + h->blk_rank[start] = c; +} + +static inline int64_t hp_rank(const hash_par_t* h, int64_t i) { + int64_t blk = i / HP_BLOCK; + int64_t r = h->blk_rank[blk]; + int64_t w = i >> 6; + for (int64_t x = blk * (HP_BLOCK / 64); x < w; x++) r += __builtin_popcountll(h->bits[x]); + return r + __builtin_popcountll(h->bits[w] & (((uint64_t)1 << (i & 63)) - 1)); +} + +static void hp_pass_d(void* raw, uint32_t wid, int64_t start, int64_t end) { + (void)wid; (void)end; + hash_par_t* h = (hash_par_t*)raw; + int64_t p = start; + int64_t lo = h->part_off[p]; + const int64_t* gfirst = h->gfirst + lo; + int64_t* gcount = h->gcount + lo; + int64_t ng = h->ng_p[p]; + for (int64_t g = 0; g < ng; g++) { + int64_t G = hp_rank(h, gfirst[g]); + h->gk[G] = (int64_t)h->kw[gfirst[g]]; + h->of[G + 1] = gcount[g]; + gcount[g] = G; /* local -> global from here on */ + } +} + +static void hp_pass_e(void* raw, uint32_t wid, int64_t start, int64_t end) { + (void)wid; (void)end; + hash_par_t* h = (hash_par_t*)raw; + int64_t p = start; + int64_t lo = h->part_off[p], hi = h->part_off[p + 1]; + const int64_t* gmap = h->gcount + lo; + /* rows of a partition are ascending -> ascending inside each group; + * of[] doubles as the fill cursor (the caller shifts it back); the + * partition's groups are nobody else's, so the cursors are private */ + for (int64_t j = lo; j < hi; j++) + h->rw[h->of[gmap[h->lg[j]]]++] = h->pr[j]; +} + +/* Zero (and so fault in) up to 8 fresh regions in parallel slices: each + * task owns one contiguous slice per region, so the page faults spread + * over the workers without two of them ever meeting on a page. */ +typedef struct { void* p[8]; size_t bytes[8]; int n; int n_tasks; } hp_touch_t; +static void hp_touch_fn(void* raw, uint32_t wid, int64_t start, int64_t end) { + (void)wid; (void)end; + hp_touch_t* t = (hp_touch_t*)raw; + for (int r = 0; r < t->n; r++) { + size_t lo = t->bytes[r] * (size_t)start / (size_t)t->n_tasks; + size_t hi = t->bytes[r] * (size_t)(start + 1) / (size_t)t->n_tasks; + if (hi > lo) memset((char*)t->p[r] + lo, 0, hi - lo); + } +} +static void hp_touch(ray_pool_t* pool, int n_tasks, hp_touch_t* t) { + t->n_tasks = n_tasks; + ray_pool_dispatch_n(pool, hp_touch_fn, t, (uint32_t)n_tasks); +} + +static bool hash_build_par(ray_t* v, ray_t** gkeys_out, ray_t** offs_out, + ray_t** rows_out, ray_t** table_out, uint64_t* mask_out, + int64_t* n_keys_out, int64_t* n_groups_out) { + int64_t n = v->len; + ray_pool_t* pool = ray_pool_get(); + if (!ray_pool_par_dispatch_ok(pool, n, 1 << 16)) return false; + int workers = (int)ray_pool_total_workers(pool); + + hash_par_t h; + memset(&h, 0, sizeof(h)); + h.v = v; h.base = (const uint8_t*)ray_data(v); h.n = n; + h.n_tasks = workers * 4; + if (h.n_tasks > 512) h.n_tasks = 512; + h.n_part = 2; + while (h.n_part < workers * 4 && h.n_part < 512) h.n_part <<= 1; + h.part_shift = 64; + for (int q = h.n_part; q > 1; q >>= 1) h.part_shift--; + + int64_t nw = (n + 63) / 64; + int64_t nblk = (n + HP_BLOCK - 1) / HP_BLOCK; + size_t cnt_b = (size_t)h.n_tasks * (size_t)h.n_part * sizeof(int64_t); + size_t po_b = (size_t)(h.n_part + 1) * sizeof(int64_t); + size_t bits_b = (size_t)(nw + 1) * sizeof(uint64_t); + h.kw = (uint64_t*)ray_sys_alloc((size_t)n * sizeof(uint64_t)); + h.cnt = (int64_t*)ray_sys_alloc(cnt_b); + h.part_off = (int64_t*)ray_sys_alloc(po_b); + h.ng_p = (int64_t*)ray_sys_alloc(po_b); + h.bits = (uint64_t*)ray_sys_alloc(bits_b); + h.blk_rank = (int64_t*)ray_sys_alloc((size_t)(nblk + 1) * sizeof(int64_t)); + ray_t *gkeys = NULL, *offs = NULL, *rows = NULL, *table = NULL; + bool ok = h.kw && h.cnt && h.part_off && h.ng_p && h.bits && h.blk_rank; + if (!ok) goto done; + memset(h.cnt, 0, cnt_b); memset(h.part_off, 0, po_b); memset(h.ng_p, 0, po_b); + memset(h.blk_rank, 0, (size_t)(nblk + 1) * sizeof(int64_t)); + { + hp_touch_t t = { .p = { h.kw, h.bits }, .bytes = { (size_t)n * sizeof(uint64_t), bits_b }, .n = 2 }; + hp_touch(pool, h.n_tasks, &t); + } + + /* A: key words + partition counts, then the bucketed row ids */ + ray_pool_dispatch_n(pool, hp_pass_a, &h, (uint32_t)h.n_tasks); + if (ray_interrupted()) { ok = false; goto done; } + { + int64_t run = 0; + for (int p = 0; p < h.n_part; p++) { + h.part_off[p] = run; + for (int t = 0; t < h.n_tasks; t++) { + int64_t c = h.cnt[(int64_t)t * h.n_part + p]; + h.cnt[(int64_t)t * h.n_part + p] = run; + run += c; + } + } + h.part_off[h.n_part] = run; + } + int64_t n_keys = h.part_off[h.n_part]; + size_t kb = (size_t)(n_keys > 0 ? n_keys : 1) * sizeof(int64_t); + h.pr = (int64_t*)ray_sys_alloc(kb); + h.lg = (int64_t*)ray_sys_alloc(kb); + h.gfirst = (int64_t*)ray_sys_alloc(kb); + h.gcount = (int64_t*)ray_sys_alloc(kb); + if (!h.pr || !h.lg || !h.gfirst || !h.gcount) { ok = false; goto done; } + { + hp_touch_t t = { .p = { h.pr, h.lg, h.gfirst, h.gcount }, .bytes = { kb, kb, kb, kb }, .n = 4 }; + hp_touch(pool, h.n_tasks, &t); + } + ray_pool_dispatch_n(pool, hp_pass_a2, &h, (uint32_t)h.n_tasks); + if (ray_interrupted()) { ok = false; goto done; } + + /* B: per-partition dedupe */ + ray_pool_dispatch_n(pool, hp_pass_b, &h, (uint32_t)h.n_part); + if (ray_interrupted() || atomic_load_explicit(&h.oom, memory_order_relaxed)) { ok = false; goto done; } + int64_t n_groups = 0; + for (int p = 0; p < h.n_part; p++) n_groups += h.ng_p[p]; + + /* C: first-occurrence numbering = rank of the group's first row */ + ray_pool_dispatch_n(pool, hp_pass_c, &h, (uint32_t)nblk); + { + int64_t run = 0; + for (int64_t b = 0; b < nblk; b++) { int64_t c = h.blk_rank[b]; h.blk_rank[b] = run; run += c; } + h.blk_rank[nblk] = run; + } + + gkeys = ray_vec_new(RAY_I64, n_groups > 0 ? n_groups : 1); + offs = ray_vec_new(RAY_I64, n_groups + 1); + rows = ray_vec_new(RAY_I64, n_keys > 0 ? n_keys : 1); + uint64_t cap = next_pow2((uint64_t)(n_groups < 4 ? 8 : 2 * n_groups)); + if (cap < 8) cap = 8; + table = ray_vec_new(RAY_I64, (int64_t)cap); + if (!gkeys || RAY_IS_ERR(gkeys) || !offs || RAY_IS_ERR(offs) || + !rows || RAY_IS_ERR(rows) || !table || RAY_IS_ERR(table)) { ok = false; goto done; } + gkeys->len = n_groups; offs->len = n_groups + 1; rows->len = n_keys; table->len = (int64_t)cap; + h.gk = (int64_t*)ray_data(gkeys); h.of = (int64_t*)ray_data(offs); + h.rw = (int64_t*)ray_data(rows); h.tbl = (int64_t*)ray_data(table); + h.tmask = cap - 1; + { + hp_touch_t t = { .p = { h.gk, h.of, h.rw, h.tbl }, + .bytes = { (size_t)(n_groups > 0 ? n_groups : 1) * sizeof(int64_t), + (size_t)(n_groups + 1) * sizeof(int64_t), kb, + (size_t)cap * sizeof(int64_t) }, .n = 4 }; + hp_touch(pool, h.n_tasks, &t); + } + + /* D: keys and counts; E: rows and the key table */ + ray_pool_dispatch_n(pool, hp_pass_d, &h, (uint32_t)h.n_part); + for (int64_t g = 0; g < n_groups; g++) h.of[g + 1] += h.of[g]; + ray_pool_dispatch_n(pool, hp_pass_e, &h, (uint32_t)h.n_part); + for (int64_t g = n_groups; g > 0; g--) h.of[g] = h.of[g - 1]; + h.of[0] = 0; + /* Key table filled in group order, exactly as the serial build does: + * the persisted index is then the same bytes whatever the core count. */ + for (int64_t g = 0; g < n_groups; g++) { + uint64_t slot = mix64((uint64_t)h.gk[g]) & h.tmask; + while (h.tbl[slot] != 0) slot = (slot + 1) & h.tmask; + h.tbl[slot] = g + 1; + } + if (ray_interrupted()) { ok = false; goto done; } + + *gkeys_out = gkeys; *offs_out = offs; *rows_out = rows; *table_out = table; + *mask_out = h.tmask; *n_keys_out = n_keys; *n_groups_out = n_groups; + +done: + if (!ok) { + if (gkeys && !RAY_IS_ERR(gkeys)) ray_release(gkeys); + if (offs && !RAY_IS_ERR(offs)) ray_release(offs); + if (rows && !RAY_IS_ERR(rows)) ray_release(rows); + if (table && !RAY_IS_ERR(table)) ray_release(table); + } + ray_sys_free(h.kw); ray_sys_free(h.pr); ray_sys_free(h.lg); + ray_sys_free(h.gfirst); ray_sys_free(h.gcount); ray_sys_free(h.cnt); + ray_sys_free(h.part_off); ray_sys_free(h.ng_p); ray_sys_free(h.bits); + ray_sys_free(h.blk_rank); + return ok; +} + ray_t* ray_index_attach_hash(ray_t** vp) { /* allow_str: keyed on a byte hash with payload-verified compares; * allow_sym: RAY_SYM uses domain ids. */ @@ -1092,6 +1412,34 @@ ray_t* ray_index_attach_hash(ray_t** vp) { bool is_str = (v->type == RAY_STR); int64_t n = v->len; + ray_t* table = NULL; + uint64_t mask = 0; + { + ray_t *pg = NULL, *po = NULL, *pr = NULL, *pt = NULL; + int64_t pk = 0, pn = 0; + uint64_t pm = 0; + if (!is_str && hash_build_par(v, &pg, &po, &pr, &pt, &pm, &pk, &pn)) { + if (ray_interrupted()) { + ray_release(pg); ray_release(po); ray_release(pr); ray_release(pt); + return ray_error("cancel", "interrupted"); + } + ray_t* idx = ray_index_alloc(RAY_IDX_HASH, v->type, n); + if (!idx || RAY_IS_ERR(idx)) { + ray_release(pg); ray_release(po); ray_release(pr); ray_release(pt); + return idx ? idx : ray_error("oom", NULL); + } + ray_index_t* ix = ray_index_payload(idx); + ix->u.hash.table = pt; + ix->u.hash.gkeys = pg; + ix->u.hash.offs = po; + ix->u.hash.rows = pr; + ix->u.hash.mask = pm; + ix->u.hash.n_keys = pk; + ix->u.hash.n_groups = pn; + ix->u.hash.order_sym = -1; + return attach_finalize(v, idx); + } + } /* Build-time capacity: sized by rows for O(1) inserts. */ uint64_t bcap = next_pow2((uint64_t)(n < 4 ? 8 : 2 * n)); if (bcap < 8) bcap = 8; @@ -1195,8 +1543,8 @@ ray_t* ray_index_attach_hash(ray_t** vp) { /* Attached/persisted bucket table: sized by DISTINCT keys. */ uint64_t cap = next_pow2((uint64_t)(n_groups < 4 ? 8 : 2 * n_groups)); if (cap < 8) cap = 8; - uint64_t mask = cap - 1; - ray_t* table = ray_vec_new(RAY_I64, (int64_t)cap); + mask = cap - 1; + table = ray_vec_new(RAY_I64, (int64_t)cap); if (!table || RAY_IS_ERR(table)) { ray_release(gkeys); ray_release(offs); ray_release(rows); return table ? table : ray_error("oom", NULL); diff --git a/src/store/splay.c b/src/store/splay.c index d59bd0b2..4fe29bc1 100644 --- a/src/store/splay.c +++ b/src/store/splay.c @@ -23,6 +23,8 @@ #include "splay.h" #include "core/runtime.h" +#include "core/pool.h" +#include "mem/sys.h" #include "store/col.h" #include "store/fileio.h" #include "store/serde.h" @@ -512,12 +514,16 @@ ray_t* ray_splay_load(const char* dir, const char* sym_path) { * rewrite), so a later mmap load gets the same block-skip an in-memory build * has. Best-effort and idempotent-ish: ray_col_append_index refuses a file * that is not exactly payload-sized (already indexed), so re-runs are no-ops. */ -void ray_splay_build_indexes(const char* dir, ray_t* tbl) { - if (!dir || !tbl || RAY_IS_ERR(tbl) || tbl->type != RAY_TABLE) return; - int64_t nc = ray_table_ncols(tbl); - for (int64_t c = 0; c < nc; c++) { +/* Index one column of a just-written splayed table (see + * ray_splay_build_indexes). Columns are independent — each reads and + * appends to its own file — so the caller runs one task per column. */ +/* deferred: when non-NULL and the column qualifies for a hash index, the + * computed zone is stored there instead of persisted and nothing is written; + * splay_persist_hash_or_zone finishes the column. */ +static void splay_build_index_col(const char* dir, ray_t* tbl, int64_t c, ray_t** deferred) { + { ray_t* col = ray_table_get_col_idx(tbl, c); - if (!col || RAY_IS_ERR(col)) continue; + if (!col || RAY_IS_ERR(col)) return; /* Explicit SYM index: a SYM column carrying a grouped (hash) index in * memory gets a hash index persisted inline — regardless of length (the @@ -557,7 +563,7 @@ void ray_splay_build_indexes(const char* dir, ray_t* tbl) { } } } - continue; + return; } /* Explicit STR index: a grouped / unique hash on a STR column is keyed @@ -575,10 +581,10 @@ void ray_splay_build_indexes(const char* dir, ray_t* tbl) { (void)ray_col_append_index(path, ray_index_payload(col->index), col->len, RAY_STR); } - continue; + return; } - if (col->len < (1 << 16)) continue; + if (col->len < (1 << 16)) return; /* STR columns get a dictionary (group on int codes); numeric/temporal * get the per-chunk min/max for block-skip. @@ -593,11 +599,18 @@ void ray_splay_build_indexes(const char* dir, ray_t* tbl) { ray_t* idx = (col->type == RAY_STR) ? ray_index_dict_compute(col) : ray_index_chunk_zone_compute(col, 16); - if (!idx || RAY_IS_ERR(idx)) { if (idx) ray_error_free(idx); continue; } + if (!idx || RAY_IS_ERR(idx)) { if (idx) ray_error_free(idx); return; } if (col->type != RAY_STR && ray_csv_hash_upgrade_check(col->type, col->len, ray_index_payload(idx))) { + if (deferred) { + /* Inside a per-column task: the hash build has its own + * parallel path that needs the pool, so hand the zone + * back and let the caller build the hash afterwards. */ + *deferred = idx; + return; + } ray_t* hi = ray_idx_hash_fn(col); if (hi && !RAY_IS_ERR(hi) && (hi->attrs & RAY_ATTR_HAS_INDEX)) { ray_release(idx); /* zone sacrificed for the hash */ @@ -612,7 +625,7 @@ void ray_splay_build_indexes(const char* dir, ray_t* tbl) { ray_index_payload(hi->index), hi->len, hi->type); } ray_release(hi); - continue; + return; } if (hi) { if (RAY_IS_ERR(hi)) ray_error_free(hi); else ray_release(hi); } /* Hash build failed — fall through and persist the zone. */ @@ -631,6 +644,77 @@ void ray_splay_build_indexes(const char* dir, ray_t* tbl) { } } +/* Persist a deferred column: its hash index when the build succeeded + * (hashed[c]), else the zone it was computed with (deferred[c]). */ +typedef struct { const char* dir; ray_t* tbl; ray_t** deferred; ray_t** hashed; } splay_index_ctx_t; + +static void splay_persist_deferred(splay_index_ctx_t* x, int64_t c) { + ray_t* col = ray_table_get_col_idx(x->tbl, c); + ray_t* nstr = ray_sym_str(ray_table_col_name(x->tbl, c)); + char path[1100]; + int n = (nstr && !RAY_IS_ERR(nstr)) + ? snprintf(path, sizeof(path), "%s/%.*s", x->dir, (int)ray_str_len(nstr), ray_str_ptr(nstr)) + : -1; + bool have_path = n > 0 && n < (int)sizeof(path); + ray_t* hi = x->hashed[c]; + if (hi) { + if (have_path) + (void)ray_col_append_index(path, ray_index_payload(hi->index), hi->len, hi->type); + ray_release(hi); + } else if (have_path) { + (void)ray_col_append_index(path, ray_index_payload(x->deferred[c]), col->len, col->type); + } + ray_release(x->deferred[c]); + x->hashed[c] = NULL; x->deferred[c] = NULL; +} + +static void splay_persist_task(void* raw, uint32_t wid, int64_t start, int64_t end) { + (void)wid; (void)end; + splay_index_ctx_t* x = (splay_index_ctx_t*)raw; + if (x->deferred[start]) splay_persist_deferred(x, start); +} +static void splay_build_index_task(void* raw, uint32_t wid, int64_t start, int64_t end) { + (void)wid; (void)end; + splay_index_ctx_t* c = (splay_index_ctx_t*)raw; + splay_build_index_col(c->dir, c->tbl, start, &c->deferred[start]); +} + +void ray_splay_build_indexes(const char* dir, ray_t* tbl) { + + if (!dir || !tbl || RAY_IS_ERR(tbl) || tbl->type != RAY_TABLE) return; + int64_t nc = ray_table_ncols(tbl); + if (nc <= 0) return; + /* One task per column: the zone / dictionary / hash builds are per-row + * scans of each column and used to run one after another on the + * calling thread — the longest serial stretch of a CSV → splayed load. */ + ray_pool_t* pool = ray_pool_get(); + if (ray_pool_par_dispatch_ok(pool, nc, 2)) { + /* Zones / dictionaries per column in parallel; the columns that + * qualify for a hash index come back deferred and are built one + * after another on this thread, each hash build parallel inside. */ + ray_t** deferred = (ray_t**)ray_sys_alloc((size_t)nc * 2 * sizeof(ray_t*)); + if (!deferred) { + for (int64_t c = 0; c < nc; c++) splay_build_index_col(dir, tbl, c, NULL); + return; + } + memset(deferred, 0, (size_t)nc * 2 * sizeof(ray_t*)); + splay_index_ctx_t ctx = { .dir = dir, .tbl = tbl, .deferred = deferred, .hashed = deferred + nc }; + ray_pool_dispatch_n(pool, splay_build_index_task, &ctx, (uint32_t)nc); + /* Hash builds one after another (each parallel inside), then the + * writes of all deferred columns together. */ + for (int64_t c = 0; c < nc; c++) { + if (!deferred[c]) continue; + ray_t* hi = ray_idx_hash_fn(ray_table_get_col_idx(tbl, c)); + if (hi && !RAY_IS_ERR(hi) && (hi->attrs & RAY_ATTR_HAS_INDEX)) ctx.hashed[c] = hi; + else if (hi) { if (RAY_IS_ERR(hi)) ray_error_free(hi); else ray_release(hi); } + } + ray_pool_dispatch_n(pool, splay_persist_task, &ctx, (uint32_t)nc); + ray_sys_free(deferred); + } else { + for (int64_t c = 0; c < nc; c++) splay_build_index_col(dir, tbl, c, NULL); + } +} + ray_t* ray_read_splayed(const char* dir, const char* sym_path) { return splay_load_impl(dir, sym_path, true); } diff --git a/src/table/domain.c b/src/table/domain.c index 0e9f5869..a7d13f77 100644 --- a/src/table/domain.c +++ b/src/table/domain.c @@ -63,6 +63,8 @@ #include "mem/arena.h" /* ray_arena_t / ray_arena_str — domain atom storage */ #include "store/fileio.h" /* flock + tmp/rename protocol for flush */ #include "ops/hash.h" /* ray_hash_bytes (same hash family as g_sym) */ +#include "core/pool.h" /* batch intern: parallel read-only probe */ +#include "sym.h" /* ray_sym_intern_prehashed (runtime fallback) */ #include #include #include @@ -177,6 +179,10 @@ struct ray_sym_domain_s { * inserts incrementally; growth rebuilds). Guarded by g_dom_lock. */ uint64_t* buckets; uint64_t bucket_mask; /* cap - 1; 0 = not built yet */ + /* Batch interns currently probing a snapshot of `buckets` outside the + * lock (guarded by g_dom_lock). While non-zero a replaced table is + * retired instead of freed. */ + int32_t batch_inflight; dom_retired_t* retired; /* replaced atom arrays + old LUTs */ @@ -292,6 +298,15 @@ static bool dom_retire(ray_sym_domain_t* d, void* p) { d->retired = r; return true; } + +/* Drop a replaced reverse-index table: freed at once unless a batch intern + * is probing a snapshot outside the lock, then retired with the domain. */ +static bool dom_drop_buckets_locked(ray_sym_domain_t* d, uint64_t* old) { + if (!old) return true; + if (d->batch_inflight == 0) { ray_sys_free(old); return true; } + return dom_retire(d, old); +} + /* Publish (map, offsets, count) as the domain's raw snapshot; the previous * record is retired. Called with the file's current image in place. */ static bool dom_publish_raw_snap(ray_sym_domain_t* d) { @@ -516,7 +531,17 @@ static bool dom_extend_from_file_locked(ray_sym_domain_t* d, size_t st_size) { * the grown count (OOB reads for lock-free consumers) — the * same corner dom_append_locked hits; mirror its loud abort * (the count is already published, there is no clean undo). */ - if (d->buckets) { ray_sys_free(d->buckets); d->buckets = NULL; d->bucket_mask = 0; } + if (d->buckets) { + /* A batch intern may be probing a snapshot of this + * table outside the lock. */ + if (!dom_drop_buckets_locked(d, d->buckets)) { + fprintf(stderr, "rayforce: sym domain '%s': OOM retiring " + "reverse index after external extend\n", + d->path ? d->path : "?"); + abort(); + } + d->buckets = NULL; d->bucket_mask = 0; + } int64_t* lut = atomic_load_explicit(&d->runtime_lut, memory_order_relaxed); if (lut) { if (!dom_retire(d, lut)) { @@ -763,15 +788,33 @@ static bool dom_build_index_locked(ray_sym_domain_t* d, int64_t extra) { if (!buckets) return false; uint64_t mask = cap - 1; - for (int64_t i = 0; i < count; i++) { - ray_t* a = dom_atom_at_locked(d, i); - if (!a) { ray_sys_free(buckets); return false; } - uint32_t h = (uint32_t)ray_hash_bytes(ray_str_ptr(a), ray_str_len(a)); - uint64_t slot = h & mask; - while (buckets[slot] != 0) slot = (slot + 1) & mask; - buckets[slot] = ((uint64_t)h << 32) | ((uint64_t)(uint32_t)i + 1); + if (d->buckets) { + /* Growth: the old table covers every published entry and carries + * the hashes — re-slot its entries, no string hashing. */ + uint64_t old_cap = d->bucket_mask + 1; + for (uint64_t i = 0; i < old_cap; i++) { + uint64_t e = d->buckets[i]; + if (e == 0) continue; + uint64_t slot = (uint32_t)(e >> 32) & mask; + while (buckets[slot] != 0) slot = (slot + 1) & mask; + buckets[slot] = e; + } + } else { + for (int64_t i = 0; i < count; i++) { + ray_t* a = dom_atom_at_locked(d, i); + if (!a) { ray_sys_free(buckets); return false; } + uint32_t h = (uint32_t)ray_hash_bytes(ray_str_ptr(a), ray_str_len(a)); + uint64_t slot = h & mask; + while (buckets[slot] != 0) slot = (slot + 1) & mask; + buckets[slot] = ((uint64_t)h << 32) | ((uint64_t)(uint32_t)i + 1); + } + } + /* A batch intern may be probing a snapshot of the old table outside + * the lock (see ray_sym_domain_intern_batch). */ + if (d->buckets && !dom_drop_buckets_locked(d, d->buckets)) { + ray_sys_free(buckets); + return false; } - ray_sys_free(d->buckets); d->buckets = buckets; d->bucket_mask = mask; return true; @@ -919,6 +962,345 @@ int64_t ray_sym_domain_intern(ray_sym_domain_t* dom, const char* str, size_t len return pos; } +/* ---- batch intern ---------------------------------------------------------- */ + +/* Read-only probe over a snapshot of the reverse index taken under the + * lock. Everything the snapshot points at outlives it: replaced bucket + * tables and atom arrays are retired, not freed, and the file prefix is + * read through the pinned raw snapshot. Entries appended after the + * snapshot (pos >= count) are ignored here and resolved under the lock. */ +typedef struct { + const uint64_t* buckets; + uint64_t mask; + ray_t* const* atoms; + int64_t count; + ray_sym_domain_raw_t raw; /* raw.count == 0: no file prefix */ + const char* const* strs; + const size_t* lens; + const uint32_t* hashes; + int64_t* out_pos; + /* misses, grouped by hash partition (hash >> part_shift) */ + int64_t* miss; /* [n_miss] batch indices */ + int64_t* uniq; /* [n_miss] first occurrences, per partition segment */ + int64_t* part_off; /* [n_part + 1] */ + int64_t* uniq_n; /* [n_part] distinct misses per partition */ + int64_t* bytes_p; /* [n_part] arena bytes the partition's atoms need */ + void** region; /* [n_part] arena region per partition */ + ray_t** atoms_w; /* current atom array (fill target) */ + /* Set when a probe met an entry it could not compare (no atom, no + * raw bytes): the misses are then resolved under the lock instead. */ + _Atomic(bool) unsure; + int part_shift; + _Atomic(bool) oom; +} dom_batch_ctx_t; + +static void dom_batch_probe_fn(void* raw, uint32_t wid, int64_t start, int64_t end) { + (void)wid; + const dom_batch_ctx_t* b = (const dom_batch_ctx_t*)raw; + for (int64_t i = start; i < end; i++) { + uint32_t h = b->hashes[i]; + size_t len = b->lens[i]; + const char* s = b->strs[i]; + uint64_t slot = h & b->mask; + int64_t found = -1; + for (;;) { + uint64_t e = atomic_load_explicit((_Atomic(uint64_t)*)&b->buckets[slot], + memory_order_relaxed); + if (e == 0) break; + if ((uint32_t)(e >> 32) == h) { + int64_t pos = (int64_t)(uint32_t)e - 1; + if (pos < b->count) { + const char* p = NULL; + size_t l = 0; + ray_t* a = atomic_load_explicit((_Atomic(ray_t*)*)&b->atoms[pos], + memory_order_acquire); + if (a) { p = ray_str_ptr(a); l = ray_str_len(a); } + else if (pos < b->raw.count) p = ray_sym_domain_raw_str(&b->raw, pos, &l); + else atomic_store_explicit((_Atomic(bool)*)&b->unsure, true, memory_order_relaxed); + if (p && l == len && (len == 0 || memcmp(p, s, len) == 0)) { + found = pos; + break; + } + } + } + slot = (slot + 1) & b->mask; + } + b->out_pos[i] = found; + } +} + +/* Dedupe the misses of one hash partition among themselves: the first + * occurrence stays a miss (out_pos -1) and is listed in the partition's + * segment of `uniq`; a repeat records its representative as -(rep + 2). + * Partitions are disjoint by hash, so no two tasks ever see the same + * string. */ +static void dom_batch_dedupe_fn(void* raw, uint32_t wid, int64_t start, int64_t end) { + (void)wid; + dom_batch_ctx_t* b = (dom_batch_ctx_t*)raw; + for (int64_t p = start; p < end; p++) { + int64_t lo = b->part_off[p], hi = b->part_off[p + 1]; + int64_t cnt = hi - lo; + if (cnt == 0) { b->uniq_n[p] = 0; continue; } + uint64_t cap = 16; + while ((uint64_t)cnt * 2 > cap) cap <<= 1; + int64_t* tab = (int64_t*)ray_sys_alloc((size_t)cap * sizeof(int64_t)); + if (!tab) { atomic_store_explicit(&b->oom, true, memory_order_relaxed); b->uniq_n[p] = 0; continue; } + memset(tab, 0xff, (size_t)cap * sizeof(int64_t)); /* -1 = empty */ + uint64_t mask = cap - 1; + int64_t k = 0; + int64_t bytes = 0; + for (int64_t j = lo; j < hi; j++) { + int64_t i = b->miss[j]; + uint32_t h = b->hashes[i]; + uint64_t slot = h & mask; + int64_t rep = -1; + while (tab[slot] >= 0) { + int64_t r = tab[slot]; + if (b->hashes[r] == h && b->lens[r] == b->lens[i] && + (b->lens[i] == 0 || memcmp(b->strs[r], b->strs[i], b->lens[i]) == 0)) { + rep = r; + break; + } + slot = (slot + 1) & mask; + } + if (rep >= 0) { + b->out_pos[i] = -(rep + 2); + } else { + tab[slot] = i; + b->uniq[lo + k++] = i; + bytes += (int64_t)ray_arena_str_bytes(b->lens[i]); + } + } + b->uniq_n[p] = k; + b->bytes_p[p] = bytes; + ray_sys_free(tab); + } +} + +/* Append the partition's distinct strings: build the atoms in the + * partition's arena region at the positions the caller assigned (batch + * order), publish them in the reverse index (CAS on the empty slot keeps + * concurrent partitions from claiming one slot twice) and resolve the + * partition's repeats. Runs under the domain lock; the count is + * published by the caller once every partition is done. */ +static void dom_batch_insert_fn(void* raw, uint32_t wid, int64_t start, int64_t end) { + (void)wid; + dom_batch_ctx_t* b = (dom_batch_ctx_t*)raw; + for (int64_t p = start; p < end; p++) { + int64_t lo = b->part_off[p]; + int64_t hi = lo + b->uniq_n[p]; + char* at = (char*)b->region[p]; + for (int64_t j = lo; j < hi; j++) { + int64_t i = b->uniq[j]; + int64_t pos = b->out_pos[i]; /* assigned in batch order */ + ray_t* s = ray_arena_str_at(at, b->strs[i], b->lens[i]); + at += ray_arena_str_bytes(b->lens[i]); + b->atoms_w[pos] = s; + uint32_t h = b->hashes[i]; + uint64_t e = ((uint64_t)h << 32) | ((uint64_t)(uint32_t)pos + 1); + uint64_t slot = h & b->mask; + for (;;) { + uint64_t cur = 0; + if (atomic_compare_exchange_strong_explicit((_Atomic(uint64_t)*)&b->buckets[slot], + &cur, e, memory_order_relaxed, memory_order_relaxed)) + break; + slot = (slot + 1) & b->mask; + } + } + for (int64_t j = lo; j < b->part_off[p + 1]; j++) { + int64_t i = b->miss[j]; + int64_t v = b->out_pos[i]; + if (v < -1) b->out_pos[i] = b->out_pos[-(v + 2)]; + } + } +} + +/* Serial fallback for the misses: find-or-append one by one under the lock + * (the domain changed under us, or the parallel path ran out of memory). */ +static bool dom_batch_append_serial_locked(ray_sym_domain_t* dom, int64_t n, + const char* const* strs, const size_t* lens, + const uint32_t* hashes, int64_t* out_pos) { + for (int64_t i = 0; i < n; i++) { + if (out_pos[i] >= 0) continue; + int64_t pos = dom_probe_locked(dom, hashes[i], strs[i], lens[i]); + if (pos < 0) pos = dom_append_locked(dom, hashes[i], strs[i], lens[i]); + if (pos < 0) return false; + out_pos[i] = pos; + } + return true; +} + +bool ray_sym_domain_intern_batch(ray_sym_domain_t* dom, int64_t n, + const char* const* strs, const size_t* lens, + const uint32_t* hashes, int64_t* out_pos) { + if (!dom || n < 0) return false; + if (n == 0) return true; + if (dom->kind == DOM_RUNTIME) { + for (int64_t i = 0; i < n; i++) { + int64_t id = ray_sym_intern_prehashed(hashes[i], strs[i], lens[i]); + if (id < 0) return false; + out_pos[i] = id; + } + return true; + } + + dom_batch_ctx_t b; + memset(&b, 0, sizeof(b)); + + dom_lock(); + int64_t count = atomic_load_explicit(&dom->count, memory_order_relaxed); + /* Headroom for the whole batch up front so no rebuild happens while + * the batch is in flight. */ + if (!dom->buckets || + (double)(count + n + 1) > 0.7 * (double)(dom->bucket_mask + 1)) { + if (!dom_build_index_locked(dom, n + 1)) { dom_unlock(); return false; } + } + if (count == 0) { + uint32_t h0 = (uint32_t)ray_hash_bytes("", 0); + if (dom_append_locked(dom, h0, "", 0) != 0) { dom_unlock(); return false; } + } + b.buckets = dom->buckets; + b.mask = dom->bucket_mask; + b.atoms = atomic_load_explicit(&dom->atoms, memory_order_acquire); + b.count = atomic_load_explicit(&dom->count, memory_order_acquire); + dom->batch_inflight++; + dom_unlock(); + + if (!ray_sym_domain_raw_pin(dom, &b.raw)) b.raw.count = 0; + b.strs = strs; b.lens = lens; b.hashes = hashes; b.out_pos = out_pos; + + ray_pool_t* pool = ray_pool_get(); + bool par = ray_pool_par_dispatch_ok(pool, n, 4096); + if (par) ray_pool_dispatch(pool, dom_batch_probe_fn, &b, n); + else dom_batch_probe_fn(&b, 0, 0, n); + + /* Misses, grouped by hash partition. */ + int n_part = 1; + if (par) { + int64_t want = (int64_t)ray_pool_total_workers(pool) * 4; + while (n_part < want && n_part < 1024) n_part <<= 1; + } + b.part_shift = 31; + for (int p = n_part; p > 2; p >>= 1) b.part_shift--; + if (n_part == 1) n_part = 2; /* keep the shift below the type width */ + int64_t n_miss = 0; + for (int64_t i = 0; i < n; i++) n_miss += (out_pos[i] < 0); + if (n_miss == 0) { + dom_lock(); + dom->batch_inflight--; + dom_unlock(); + return true; + } + + b.miss = (int64_t*)ray_sys_alloc((size_t)n_miss * sizeof(int64_t)); + b.uniq = (int64_t*)ray_sys_alloc((size_t)n_miss * sizeof(int64_t)); + b.part_off = (int64_t*)ray_sys_alloc((size_t)(n_part + 1) * sizeof(int64_t)); + b.uniq_n = (int64_t*)ray_sys_alloc((size_t)n_part * sizeof(int64_t)); + b.bytes_p = (int64_t*)ray_sys_alloc((size_t)n_part * sizeof(int64_t)); + b.region = (void**)ray_sys_alloc((size_t)n_part * sizeof(void*)); + bool ok = b.miss && b.uniq && b.part_off && b.uniq_n && b.bytes_p && b.region; + if (ok) { + memset(b.part_off, 0, (size_t)(n_part + 1) * sizeof(int64_t)); + for (int64_t i = 0; i < n; i++) + if (out_pos[i] < 0) b.part_off[(hashes[i] >> b.part_shift) + 1]++; + for (int p = 0; p < n_part; p++) b.part_off[p + 1] += b.part_off[p]; + { + int64_t* fill = b.uniq_n; /* scratch cursor per partition */ + memcpy(fill, b.part_off, (size_t)n_part * sizeof(int64_t)); + for (int64_t i = 0; i < n; i++) + if (out_pos[i] < 0) b.miss[fill[hashes[i] >> b.part_shift]++] = i; + } + if (par) ray_pool_dispatch_n(pool, dom_batch_dedupe_fn, &b, (uint32_t)n_part); + else dom_batch_dedupe_fn(&b, 0, 0, n_part); + if (atomic_load_explicit(&b.oom, memory_order_relaxed)) { + /* Undo the repeat marks; the serial path resolves everything. */ + for (int64_t i = 0; i < n; i++) if (out_pos[i] < -1) out_pos[i] = -1; + ok = false; + } + } + + dom_lock(); + dom->batch_inflight--; /* no probe outside the lock past this point */ + bool unchanged = ok && dom->buckets == b.buckets && + atomic_load_explicit(&dom->count, memory_order_relaxed) == b.count && + !atomic_load_explicit(&b.unsure, memory_order_relaxed); + if (unchanged) { + int64_t total = 0; + for (int p = 0; p < n_part; p++) total += b.uniq_n[p]; + int64_t base = b.count; + ok = base + total < (int64_t)UINT32_MAX; + /* Atom array: grow once by replacement (lock-free readers may hold + * the old pointer). */ + if (ok && base + total > dom->atoms_cap) { + int64_t ncap = dom->atoms_cap < 8 ? 8 : dom->atoms_cap; + while (ncap < base + total) ncap *= 2; + ray_t** narr = (ray_t**)ray_sys_alloc((size_t)ncap * sizeof(ray_t*)); + if (!narr) ok = false; + else { + ray_t** old = atomic_load_explicit(&dom->atoms, memory_order_relaxed); + if (base > 0) memcpy(narr, old, (size_t)base * sizeof(ray_t*)); + if (!dom_retire(dom, old)) { ray_sys_free(narr); ok = false; } + else { + atomic_store_explicit(&dom->atoms, narr, memory_order_release); + dom->atoms_cap = ncap; + } + } + } + if (ok) { + /* One arena region per partition; positions in partition order. + * Nothing is published until every partition has built its + * atoms, so a failed reservation costs only arena space. */ + /* New strings take positions in batch order (first occurrence), + * so the symfile does not depend on how the batch was split. */ + int64_t pos = base; + for (int64_t i = 0; i < n; i++) + if (out_pos[i] == -1) out_pos[i] = pos++; + for (int p = 0; p < n_part; p++) { + b.region[p] = NULL; + if (b.bytes_p[p] > 0) { + b.region[p] = ray_arena_alloc_raw(dom->arena, (size_t)b.bytes_p[p]); + if (!b.region[p]) { ok = false; break; } + } + } + if (ok) { + b.atoms_w = atomic_load_explicit(&dom->atoms, memory_order_relaxed); + if (par) ray_pool_dispatch_n(pool, dom_batch_insert_fn, &b, (uint32_t)n_part); + else dom_batch_insert_fn(&b, 0, 0, n_part); + atomic_store_explicit(&dom->count, pos, memory_order_release); + /* Same invalidation as dom_append_locked: the runtime LUT + * no longer covers the vocabulary. */ + int64_t* lut = atomic_load_explicit(&dom->runtime_lut, memory_order_relaxed); + if (lut) { + if (!dom_retire(dom, lut)) { + fprintf(stderr, "rayforce: sym domain '%s': OOM retiring runtime " + "LUT after batch append\n", dom->path ? dom->path : "?"); + abort(); + } + atomic_store_explicit(&dom->runtime_lut, NULL, memory_order_release); + } + } else { + /* Nothing published: undo the assigned positions and the + * repeat marks; the serial path resolves them again. */ + for (int64_t i = 0; i < n; i++) + if (out_pos[i] < -1 || out_pos[i] >= base) out_pos[i] = -1; + } + } + } else { + for (int64_t i = 0; i < n; i++) if (out_pos[i] < -1) out_pos[i] = -1; + } + if (!unchanged || !ok) + ok = dom_batch_append_serial_locked(dom, n, strs, lens, hashes, out_pos); + dom_unlock(); + + ray_sys_free(b.miss); + ray_sys_free(b.uniq); + ray_sys_free(b.part_off); + ray_sys_free(b.uniq_n); + ray_sys_free(b.bytes_p); + ray_sys_free(b.region); + return ok; +} + int64_t ray_sym_domain_count(ray_sym_domain_t* dom) { if (!dom) return 0; if (dom->kind == DOM_RUNTIME) return (int64_t)ray_sym_count(); @@ -1013,18 +1395,37 @@ ray_err_t ray_sym_domain_flush(ray_sym_domain_t* dom, bool durable) { if (fwrite(&magic, 4, 1, f) != 1 || fwrite(&count, 8, 1, f) != 1) err = RAY_ERR_IO; written_size = 12; + /* Records are packed into a large buffer first: two stdio calls per + * entry dominated the flush of a big vocabulary. */ + enum { FLUSH_BUF = 1u << 20 }; + uint8_t* wb = (uint8_t*)ray_sys_alloc(FLUSH_BUF); + size_t wn = 0; + if (!wb) err = RAY_ERR_OOM; for (int64_t i = 0; err == RAY_OK && i < count; i++) { ray_t* s = atoms[i]; size_t slen = ray_str_len(s); if (slen > UINT32_MAX) { err = RAY_ERR_RANGE; break; } uint32_t len32 = (uint32_t)slen; - if (fwrite(&len32, 4, 1, f) != 1 || - (slen > 0 && fwrite(ray_str_ptr(s), 1, slen, f) != slen)) { - err = RAY_ERR_IO; - break; + if (wn + 4 + slen > FLUSH_BUF) { + if (wn && fwrite(wb, 1, wn, f) != wn) { err = RAY_ERR_IO; break; } + wn = 0; + } + if (4 + slen > FLUSH_BUF) { + /* Oversized record: straight through. */ + if (fwrite(&len32, 4, 1, f) != 1 || + fwrite(ray_str_ptr(s), 1, slen, f) != slen) { + err = RAY_ERR_IO; + break; + } + } else { + memcpy(wb + wn, &len32, 4); + if (slen) memcpy(wb + wn + 4, ray_str_ptr(s), slen); + wn += 4 + slen; } written_size += 4 + slen; } + if (err == RAY_OK && wn && fwrite(wb, 1, wn, f) != wn) err = RAY_ERR_IO; + ray_sys_free(wb); if (fclose(f) != 0 && err == RAY_OK) err = RAY_ERR_IO; } if (err != RAY_OK) goto fail_tmp; diff --git a/src/table/domain.h b/src/table/domain.h index e5c0f9e6..b9ce6610 100644 --- a/src/table/domain.h +++ b/src/table/domain.h @@ -164,6 +164,17 @@ const int64_t* ray_sym_domain_runtime_lut(ray_sym_domain_t* dom); * the shared object immediately; ray_sym_domain_flush persists them. */ int64_t ray_sym_domain_intern(ray_sym_domain_t* dom, const char* str, size_t len); +/* Batch find-or-append of n strings (hashes[i] = (uint32_t)ray_hash_bytes + * of strs[i]); out_pos[i] receives the position. Repeats inside the + * batch resolve to one position. FILE: the lookup of the existing + * vocabulary runs in parallel on the pool; new strings are appended in + * batch order (a string's first occurrence), whatever the worker count. RUNTIME: one ray_sym_intern per entry. + * Returns false on OOM (out_pos then holds -1 for the entries not + * resolved). */ +bool ray_sym_domain_intern_batch(ray_sym_domain_t* dom, int64_t n, + const char* const* strs, const size_t* lens, + const uint32_t* hashes, int64_t* out_pos); + /* Number of entries in the domain. */ int64_t ray_sym_domain_count(ray_sym_domain_t* dom); diff --git a/test/rfl/io/csv_splayed_dedupe_overflow.rfl b/test/rfl/io/csv_splayed_dedupe_overflow.rfl new file mode 100644 index 00000000..5364d72d --- /dev/null +++ b/test/rfl/io/csv_splayed_dedupe_overflow.rfl @@ -0,0 +1,27 @@ +;; Splayed CSV load of a SYM column whose per-partition dictionary +;; overflows (src/io/csv.c: csv_dedup_task caps a partition at +;; CSV_DEDUP_MAX_ENTS distinct strings — 8192 in debug builds — and the +;; column then interns row by row into the symfile domain, +;; csv_intern_dicts_domain's fallback). 140k distinct strings exceed the +;; cap in at least one partition however the column is split (at most 16 +;; partitions). Release builds keep the cap at 1M and take the batched +;; path; the answers are the same either way. + +(.sys.exec "rm -rf rf_test_splayed_ovf rf_test_splayed_ovf.csv") -- 0 +(.sys.exec "seq 0 139999 | awk '{print $1\",k\"$1}' > rf_test_splayed_ovf.csv") -- 0 +(set Tovf (.csv.splayed [id s] [I64 SYM] "rf_test_splayed_ovf.csv" "rf_test_splayed_ovf/")) +(count Tovf) -- 140000 +(sum (at Tovf 'id)) -- 9799930000 +(count (distinct (at Tovf 's))) -- 140000 +(at (at Tovf 's) 0) -- 'k0 +(at (at Tovf 's) 8192) -- 'k8192 +(at (at Tovf 's) 139999) -- 'k139999 +;; every row's symbol is the one written for its id +(count (select {from: Tovf where: (== s (as 'SYM "k77777"))})) -- 1 +(at (at (select {id: id from: Tovf where: (== s (as 'SYM "k77777"))}) 'id) 0) -- 77777 +;; reload from disk agrees +(set Rovf (.db.splayed.get "rf_test_splayed_ovf/")) +(count Rovf) -- 140000 +(count (distinct (at Rovf 's))) -- 140000 +(at (at Rovf 's) 139999) -- 'k139999 +(.sys.exec "rm -rf rf_test_splayed_ovf rf_test_splayed_ovf.csv") -- 0 diff --git a/test/test_domain.c b/test/test_domain.c index ce7186ac..6946c9a7 100644 --- a/test/test_domain.c +++ b/test/test_domain.c @@ -39,6 +39,9 @@ #include "table/domain.h" #include "store/col.h" #include "store/serde.h" +#include "core/pool.h" /* the batch intern probes on the pool */ +#include "ops/hash.h" /* ray_hash_bytes: batch intern takes prehashed entries */ +#include "mem/sys.h" #include "ops/ops.h" /* RAY_PARTED_BASE (parted-flatten adoption test) */ #include #include @@ -401,6 +404,108 @@ static test_result_t test_domain_open_basic(void) { PASS(); } +/* ray_sym_domain_intern_batch: one batch with repeats spread over the hash + * partitions, half the vocabulary already interned one by one. Every + * entry resolves to the position find() reports, repeats agree, the + * pre-interned half keeps its positions, the count grows by exactly the + * new distinct strings, a second identical batch changes nothing and "" + * is position 0. */ +static test_result_t test_domain_intern_batch(void) { + unlink(TMP_DOM_SYM_PATH); + unlink(TMP_DOM_SYM_PATH ".lk"); + TEST_ASSERT_NOT_NULL(ray_pool_get()); /* parallel probe path */ + + ray_sym_domain_t* dom = ray_sym_domain_open_or_create(TMP_DOM_SYM_PATH); + TEST_ASSERT_NOT_NULL(dom); + + enum { NV = 20000, NB = 60000, SL = 16 }; + char* vocab = (char*)ray_sys_alloc((size_t)NV * SL); + int64_t* pre = (int64_t*)ray_sys_alloc((size_t)NV * sizeof(int64_t)); + int64_t* seen = (int64_t*)ray_sys_alloc((size_t)NV * sizeof(int64_t)); + const char** strs = (const char**)ray_sys_alloc((size_t)NB * sizeof(char*)); + size_t* lens = (size_t*)ray_sys_alloc((size_t)NB * sizeof(size_t)); + uint32_t* hashes = (uint32_t*)ray_sys_alloc((size_t)NB * sizeof(uint32_t)); + int64_t* pos = (int64_t*)ray_sys_alloc((size_t)NB * sizeof(int64_t)); + int64_t* pos2 = (int64_t*)ray_sys_alloc((size_t)NB * sizeof(int64_t)); + TEST_ASSERT_NOT_NULL(vocab); TEST_ASSERT_NOT_NULL(pre); TEST_ASSERT_NOT_NULL(seen); + TEST_ASSERT_NOT_NULL(strs); TEST_ASSERT_NOT_NULL(lens); TEST_ASSERT_NOT_NULL(hashes); + TEST_ASSERT_NOT_NULL(pos); TEST_ASSERT_NOT_NULL(pos2); + + for (int v = 0; v < NV; v++) { + snprintf(vocab + (size_t)v * SL, SL, "v%05d_%c", v, 'a' + v % 26); + seen[v] = -1; + pre[v] = -1; + } + for (int v = 0; v < NV / 2; v++) { + const char* sv = vocab + (size_t)v * SL; + pre[v] = ray_sym_domain_intern(dom, sv, strlen(sv)); + TEST_ASSERT(pre[v] > 0, "pre-intern gets a position"); + } + int64_t count_before = ray_sym_domain_count(dom); + TEST_ASSERT_EQ_I(count_before, NV / 2 + 1); /* + reserved "" */ + + for (int j = 0; j < NB; j++) { + int v = (int)(((int64_t)j * 7919) % NV); /* every string ~3 times */ + const char* sv = vocab + (size_t)v * SL; + strs[j] = sv; + lens[j] = strlen(sv); + hashes[j] = (uint32_t)ray_hash_bytes(sv, lens[j]); + pos[j] = -7; + } + TEST_ASSERT_TRUE(ray_sym_domain_intern_batch(dom, NB, strs, lens, hashes, pos)); + TEST_ASSERT_EQ_I(ray_sym_domain_count(dom), NV + 1); + + for (int j = 0; j < NB; j++) { + int v = (int)(((int64_t)j * 7919) % NV); + TEST_ASSERT(pos[j] > 0 && pos[j] <= NV, "position in range, never 0"); + if (seen[v] < 0) seen[v] = pos[j]; + TEST_ASSERT_EQ_I(pos[j], seen[v]); /* repeats agree */ + if (pre[v] >= 0) TEST_ASSERT_EQ_I(pos[j], pre[v]); /* hits keep their position */ + TEST_ASSERT_EQ_I(ray_sym_domain_find(dom, strs[j], lens[j]), pos[j]); + ray_t* a = ray_sym_domain_str(dom, pos[j]); + TEST_ASSERT_NOT_NULL(a); + TEST_ASSERT_EQ_U(ray_str_len(a), lens[j]); + TEST_ASSERT_MEM_EQ(lens[j], ray_str_ptr(a), strs[j]); + } + /* distinct positions: every vocabulary entry got exactly one */ + for (int v = 0; v < NV; v++) TEST_ASSERT(seen[v] > 0, "every string appeared"); + + /* the same batch again: all hits, nothing appended */ + for (int j = 0; j < NB; j++) pos2[j] = -7; + TEST_ASSERT_TRUE(ray_sym_domain_intern_batch(dom, NB, strs, lens, hashes, pos2)); + TEST_ASSERT_EQ_I(ray_sym_domain_count(dom), NV + 1); + for (int j = 0; j < NB; j++) TEST_ASSERT_EQ_I(pos2[j], pos[j]); + + /* "" resolves to the reserved position 0; a fresh string still appends */ + { + const char* two[2] = { "", "brand_new_entry" }; + size_t tl[2] = { 0, strlen("brand_new_entry") }; + uint32_t th[2] = { (uint32_t)ray_hash_bytes("", 0), (uint32_t)ray_hash_bytes(two[1], tl[1]) }; + int64_t tp[2] = { -7, -7 }; + TEST_ASSERT_TRUE(ray_sym_domain_intern_batch(dom, 2, two, tl, th, tp)); + TEST_ASSERT_EQ_I(tp[0], 0); + TEST_ASSERT_EQ_I(tp[1], NV + 1); + TEST_ASSERT_EQ_I(ray_sym_domain_count(dom), NV + 2); + } + + /* the flushed file reopens with the same vocabulary */ + TEST_ASSERT_EQ_I(ray_sym_domain_flush(dom, false), RAY_OK); + ray_sym_domain_release(dom); + ray_sym_domain_t* re = ray_sym_domain_open(TMP_DOM_SYM_PATH); + TEST_ASSERT_NOT_NULL(re); + TEST_ASSERT_EQ_I(ray_sym_domain_count(re), NV + 2); + for (int j = 0; j < NB; j += 997) + TEST_ASSERT_EQ_I(ray_sym_domain_find(re, strs[j], lens[j]), pos[j]); + ray_sym_domain_release(re); + + ray_sys_free(vocab); ray_sys_free(pre); ray_sys_free(seen); + ray_sys_free(strs); ray_sys_free(lens); ray_sys_free(hashes); + ray_sys_free(pos); ray_sys_free(pos2); + unlink(TMP_DOM_SYM_PATH); + unlink(TMP_DOM_SYM_PATH ".lk"); + PASS(); +} + /* Task 7b: open_or_create on a missing file yields an empty writable * domain; "" is seeded at position 0 by the first intern; flush creates * the file; verify-base-unchanged makes a racing writer LOUD. */ @@ -1988,6 +2093,7 @@ const test_entry_t domain_entries[] = { { "domain/raw_pin", test_domain_raw_pin, domain_setup, domain_teardown }, { "domain/runtime_lut", test_domain_runtime_lut, domain_rt_setup, domain_rt_teardown }, { "domain/open_position0_validation", test_domain_open_position0_validation, domain_setup, domain_teardown }, + { "domain/intern_batch", test_domain_intern_batch, domain_setup, domain_teardown }, { "domain/dict_upsert_file_keys", test_domain_dict_upsert_file_keys, domain_rt_setup, domain_rt_teardown }, { NULL, NULL, NULL, NULL }, }; diff --git a/test/test_index.c b/test/test_index.c index af3a0213..d79a0e3d 100644 --- a/test/test_index.c +++ b/test/test_index.c @@ -26,6 +26,8 @@ #include "test.h" #include #include "mem/heap.h" +#include "mem/sys.h" +#include "core/pool.h" #include "mem/cow.h" #include "vec/vec.h" #include "table/sym.h" @@ -327,6 +329,83 @@ static test_result_t test_index_hash_with_nulls_preserved(void) { PASS(); } +/* Large column: the build runs partition-parallel above 64k rows and must + * produce the serial walk's layout — groups in first-occurrence order, rows + * ascending inside a group, nulls excluded — checked against a reference + * computed the obvious way. */ +static test_result_t test_index_hash_large_parallel(void) { + ray_heap_init(); + /* The parallel build needs the pool; create it before the attach so + * the test does not silently take the serial fallback. */ + ray_pool_t* pool = ray_pool_get(); + TEST_ASSERT_NOT_NULL(pool); + const int64_t n = 300000, kmax = 5003; + ray_t* v = ray_vec_new(RAY_I64, n); + TEST_ASSERT_NOT_NULL(v); + int64_t* xs = (int64_t*)ray_data(v); + for (int64_t i = 0; i < n; i++) + xs[i] = (int64_t)(((uint64_t)i * 2654435761ull) % (uint64_t)kmax) - 17; + v->len = n; + /* every 977th row null */ + for (int64_t i = 0; i < n; i += 977) + TEST_ASSERT_EQ_I(ray_vec_set_null_checked(v, i, true), RAY_OK); + + /* reference: first-occurrence group ids and counts */ + int64_t* gid_of_key = (int64_t*)ray_sys_alloc((size_t)kmax * sizeof(int64_t)); + int64_t* ref_key = (int64_t*)ray_sys_alloc((size_t)kmax * sizeof(int64_t)); + int64_t* ref_cnt = (int64_t*)ray_sys_alloc((size_t)kmax * sizeof(int64_t)); + TEST_ASSERT_NOT_NULL(gid_of_key); TEST_ASSERT_NOT_NULL(ref_key); TEST_ASSERT_NOT_NULL(ref_cnt); + for (int64_t k = 0; k < kmax; k++) { gid_of_key[k] = -1; ref_cnt[k] = 0; } + int64_t ref_groups = 0, ref_keys = 0; + for (int64_t i = 0; i < n; i++) { + if (ray_vec_is_null(v, i)) continue; + int64_t k = xs[i] + 17; + if (gid_of_key[k] < 0) { gid_of_key[k] = ref_groups; ref_key[ref_groups++] = xs[i]; } + ref_cnt[gid_of_key[k]]++; + ref_keys++; + } + + ray_t* w = v; + ray_t* r = ray_index_attach_hash(&w); + TEST_ASSERT_FALSE(RAY_IS_ERR(r)); + ray_index_t* ix = ray_index_payload(w->index); + TEST_ASSERT_EQ_I((int)ix->kind, RAY_IDX_HASH); + TEST_ASSERT_EQ_I(ix->u.hash.n_keys, ref_keys); + TEST_ASSERT_EQ_I(ix->u.hash.n_groups, ref_groups); + const int64_t* gk = (const int64_t*)ray_data(ix->u.hash.gkeys); + const int64_t* of = (const int64_t*)ray_data(ix->u.hash.offs); + const int64_t* rw = (const int64_t*)ray_data(ix->u.hash.rows); + TEST_ASSERT_EQ_I(of[0], 0); + TEST_ASSERT_EQ_I(of[ref_groups], ref_keys); + /* Serial layout: groups in first-occurrence order, each with its count, + * rows ascending and all storing the group's key. */ + for (int64_t g = 0; g < ref_groups; g++) { + TEST_ASSERT_EQ_I(gk[g], ref_key[g]); + TEST_ASSERT_EQ_I(of[g + 1] - of[g], ref_cnt[g]); + for (int64_t j = of[g]; j < of[g + 1]; j++) { + TEST_ASSERT_EQ_I(xs[rw[j]], ref_key[g]); + if (j > of[g]) TEST_ASSERT_TRUE(rw[j] > rw[j - 1]); + } + } + /* table probes: every key resolves to its group, an absent key misses */ + for (int64_t k = 0; k < kmax; k += 61) { + const int64_t* grows = NULL; + int64_t gn = 0; + TEST_ASSERT_EQ_I(ray_index_hash_group(w, k - 17, &grows, &gn), 1); + TEST_ASSERT_EQ_I(gn, ref_cnt[gid_of_key[k]]); + } + { + const int64_t* grows = NULL; + int64_t gn = 0; + TEST_ASSERT_EQ_I(ray_index_hash_group(w, kmax + 1000, &grows, &gn), 0); + } + + ray_sys_free(gid_of_key); ray_sys_free(ref_key); ray_sys_free(ref_cnt); + ray_release(w); + ray_heap_destroy(); + PASS(); +} + /* ─── Sort index ──────────────────────────────────────────────────── */ static test_result_t test_index_sort_attach_drop(void) { @@ -3697,6 +3776,7 @@ const test_entry_t index_entries[] = { { "index/unsupported_type", test_index_unsupported_type, NULL, NULL }, { "index/hash_attach_drop", test_index_hash_attach_drop, NULL, NULL }, { "index/hash_with_nulls_preserved", test_index_hash_with_nulls_preserved, NULL, NULL }, + { "index/hash_large_parallel", test_index_hash_large_parallel, NULL, NULL }, { "index/sort_attach_drop", test_index_sort_attach_drop, NULL, NULL }, { "index/bloom_attach_drop", test_index_bloom_attach_drop, NULL, NULL }, { "index/replace_cross_kind", test_index_replace_cross_kind, NULL, NULL }, diff --git a/test/test_splay.c b/test/test_splay.c index 34188b46..d037cf5e 100644 --- a/test/test_splay.c +++ b/test/test_splay.c @@ -40,6 +40,10 @@ #include "lang/internal.h" /* ray_set/get_splayed_fn (surface resolver) */ #include "mem/heap.h" #include "table/sym.h" +#include "table/domain.h" /* symfile positions of a streamed CSV load */ +#include "io/csv.h" /* ray_csv_save_splayed_named_opts */ +#include "mem/sys.h" +#include "core/pool.h" #include #include #include @@ -68,6 +72,139 @@ static void rm_rf(const char* path) { (void)ray_test_rm_rf(path); } +/* Streamed CSV -> splayed load, several chunks. The symfile positions of + * the SYM columns are the ones the cell-by-cell writer gives: per chunk, + * columns in order, each column's strings by first occurrence — whatever + * the worker count or how the dictionaries were split into partitions. + * Chunks of 20k rows with >4k new strings each take the parallel batch. */ +static test_result_t test_csv_splayed_symfile_order(void) { + TEST_ASSERT_NOT_NULL(ray_pool_get()); + const char* dir = TMP_SPLAY_BASE "/csvorder"; + const char* csv = TMP_SPLAY_BASE "/csvorder.csv"; + rm_rf(dir); + mkdir(TMP_SPLAY_BASE, 0755); + + enum { NROWS = 60000, CHUNK = 20000, NA = 13000, NB = 9000 }; + FILE* f = fopen(csv, "wb"); + TEST_ASSERT_NOT_NULL(f); + fputs("a,b,v\n", f); + for (int r = 0; r < NROWS; r++) { + int ai = (int)(((int64_t)r * 7919) % NA); + if (r % 5 == 0) fprintf(f, "a%d,a%d,%d\n", ai, (ai + 11) % NA, r); /* b reuses a's strings */ + else fprintf(f, "a%d,b%d,%d\n", ai, (int)(((int64_t)r * 104729) % NB), r); + } + fclose(f); + + int8_t types[] = { RAY_SYM, RAY_SYM, RAY_I64 }; + ray_err_t err = ray_csv_save_splayed_named_opts(csv, ',', true, types, 3, NULL, 0, dir, CHUNK); + TEST_ASSERT_EQ_I(err, RAY_OK); + + /* expected positions: walk as the writer does */ + int64_t* pa = (int64_t*)ray_sys_alloc(NA * sizeof(int64_t)); + int64_t* pb = (int64_t*)ray_sys_alloc(NB * sizeof(int64_t)); + TEST_ASSERT_NOT_NULL(pa); TEST_ASSERT_NOT_NULL(pb); + for (int i = 0; i < NA; i++) pa[i] = -1; + for (int i = 0; i < NB; i++) pb[i] = -1; + int64_t next = 1; /* 0 is "" */ + for (int c0 = 0; c0 < NROWS; c0 += CHUNK) { + for (int r = c0; r < c0 + CHUNK; r++) { /* column a */ + int ai = (int)(((int64_t)r * 7919) % NA); + if (pa[ai] < 0) pa[ai] = next++; + } + for (int r = c0; r < c0 + CHUNK; r++) { /* column b */ + int ai = (int)(((int64_t)r * 7919) % NA); + if (r % 5 == 0) { int x = (ai + 11) % NA; if (pa[x] < 0) pa[x] = next++; } + else { int bi = (int)(((int64_t)r * 104729) % NB); if (pb[bi] < 0) pb[bi] = next++; } + } + } + + char sym_path[256]; + snprintf(sym_path, sizeof(sym_path), "%s/.sym", dir); + ray_sym_domain_t* dom = ray_sym_domain_open(sym_path); + TEST_ASSERT_NOT_NULL(dom); + TEST_ASSERT_EQ_I(ray_sym_domain_count(dom), next); + char buf[32]; + for (int i = 0; i < NA; i++) { + if (pa[i] < 0) continue; + int n = snprintf(buf, sizeof(buf), "a%d", i); + TEST_ASSERT_EQ_I(ray_sym_domain_find(dom, buf, (size_t)n), pa[i]); + } + for (int i = 0; i < NB; i++) { + if (pb[i] < 0) continue; + int n = snprintf(buf, sizeof(buf), "b%d", i); + TEST_ASSERT_EQ_I(ray_sym_domain_find(dom, buf, (size_t)n), pb[i]); + } + ray_sym_domain_release(dom); + + /* and the cells read back the strings written */ + ray_t* t = ray_read_splayed(dir, sym_path); + TEST_ASSERT_FALSE(RAY_IS_ERR(t)); + TEST_ASSERT_EQ_I(ray_table_nrows(t), NROWS); + ray_t* ca = ray_table_get_col_idx(t, 0); + for (int r = 0; r < NROWS; r += 997) { + ray_t* cell = ray_sym_vec_cell(ca, r); + TEST_ASSERT_NOT_NULL(cell); + int n = snprintf(buf, sizeof(buf), "a%d", (int)(((int64_t)r * 7919) % NA)); + TEST_ASSERT_EQ_U(ray_str_len(cell), (size_t)n); + TEST_ASSERT_MEM_EQ((size_t)n, ray_str_ptr(cell), buf); + } + ray_release(t); + ray_sys_free(pa); ray_sys_free(pb); + rm_rf(dir); + unlink(csv); + PASS(); +} + +/* A chunk whose byte window holds no quote, in a file that has quotes in + * another chunk, is split into rows exactly as the whole file is: a lone + * '\r' ends a row. (The quote-free fast path of the parallel scanner did + * not treat it so, and merged two rows, losing a value.) */ +static test_result_t test_csv_splayed_quote_mode_per_file(void) { + TEST_ASSERT_NOT_NULL(ray_pool_get()); + const char* dir = TMP_SPLAY_BASE "/csvquote"; + const char* csv = TMP_SPLAY_BASE "/csvquote.csv"; + rm_rf(dir); + mkdir(TMP_SPLAY_BASE, 0755); + + enum { NROWS = 60000, CHUNK = 20000, LONE = 45007 }; + FILE* f = fopen(csv, "wb"); + TEST_ASSERT_NOT_NULL(f); + fputs("s,v\n", f); + fputs("\"q,1\",0\n", f); /* quotes only in chunk 0 */ + for (int r = 1; r < NROWS; r++) { + if (r == LONE) fprintf(f, "left\rright,%d\n", r); /* chunk 2 */ + else fprintf(f, "r%d,%d\n", r, r); + } + fclose(f); + + int8_t types[] = { RAY_SYM, RAY_I64 }; + ray_t* mem = ray_read_csv_named_opts(csv, ',', true, types, 2, NULL, 0); + TEST_ASSERT_FALSE(RAY_IS_ERR(mem)); + int64_t want = ray_table_nrows(mem); + TEST_ASSERT_EQ_I(want, NROWS + 1); /* the lone \r splits a row */ + + ray_err_t err = ray_csv_save_splayed_named_opts(csv, ',', true, types, 2, NULL, 0, dir, CHUNK); + TEST_ASSERT_EQ_I(err, RAY_OK); + char sym_path[256]; + snprintf(sym_path, sizeof(sym_path), "%s/.sym", dir); + ray_t* t = ray_read_splayed(dir, sym_path); + TEST_ASSERT_FALSE(RAY_IS_ERR(t)); + TEST_ASSERT_EQ_I(ray_table_nrows(t), want); + /* every row agrees with the in-memory read */ + ray_t* vm = ray_table_get_col_idx(mem, 1); + ray_t* vt = ray_table_get_col_idx(t, 1); + for (int64_t r = 0; r < want; r++) { + TEST_ASSERT_EQ_I(ray_vec_is_null(vt, r), ray_vec_is_null(vm, r)); + if (!ray_vec_is_null(vm, r)) + TEST_ASSERT_EQ_I(((int64_t*)ray_data(vt))[r], ((int64_t*)ray_data(vm))[r]); + } + ray_release(t); + ray_release(mem); + rm_rf(dir); + unlink(csv); + PASS(); +} + /* ========================================================================= * 1. ray_splay_save: NULL dir → RAY_ERR_IO * ========================================================================= */ @@ -2277,5 +2414,7 @@ const test_entry_t splay_entries[] = { { "splay/empty_sym_table_roundtrip", test_empty_sym_table_roundtrip, splay_setup, splay_teardown }, { "splay/resolution_order_independence", test_resolution_order_independence, splay_setup, splay_teardown }, { "splay/resolution_explicit_wins", test_resolution_explicit_wins, splay_setup, splay_teardown }, + { "splay/csv_symfile_order", test_csv_splayed_symfile_order, splay_setup, splay_teardown }, + { "splay/csv_quote_mode_per_file", test_csv_splayed_quote_mode_per_file, splay_setup, splay_teardown }, { NULL, NULL, NULL, NULL }, };