#if defined(__linux__)
#define _GNU_SOURCE
#endif
#include "csv.h"
#include "mem/heap.h"
#include "mem/sys.h"
#include "core/numparse.h"
#include "core/pool.h"
#include "core/profile.h"
#include "core/platform.h"
#include "lang/format.h"
#include "ops/hash.h"
#include "ops/idxop.h"
#include "store/col.h"
#include "store/fileio.h"
#include "store/splay.h"
#include "store/stream.h"
#include "table/sym.h"
#include "table/domain.h"
#include "vec/str.h"
#include "vec/vec.h"
#include <inttypes.h>
#include <math.h>
#include <stdarg.h>
#include <string.h>
#include <stdio.h>
#include <stdlib.h>
#include <sys/stat.h>
#include <fcntl.h>
#ifndef RAY_OS_WINDOWS
#include <unistd.h>
#include <sys/mman.h>
#endif
#define CSV_MAX_COLS 256
#define CSV_SAMPLE_ROWS 4096
#define CSV_STR_DISTINCT_MIN 256
#define CSV_PART_ROWS_DEFAULT 1000000
#define CSV_PROG_PCT_SCAN 27
#define CSV_PROG_PCT_PARSE 41
static inline uint64_t csv_prog_at(size_t file_size, unsigned pct) {
return (uint64_t)((double)file_size * (double)pct / 100.0);
}
#ifndef RAY_OS_WINDOWS
#define MMAP_FLAGS MAP_PRIVATE
#endif
static inline void* scratch_alloc(ray_t** hdr_out, size_t nbytes) {
ray_t* h = ray_alloc(nbytes);
if (!h) { *hdr_out = NULL; return NULL; }
*hdr_out = h;
return ray_data(h);
}
static inline void* scratch_realloc(ray_t** hdr_out, size_t old_bytes, size_t new_bytes) {
ray_t* old_h = *hdr_out;
ray_t* new_h = ray_alloc(new_bytes);
if (!new_h) return NULL;
void* new_p = ray_data(new_h);
if (old_h) {
memcpy(new_p, ray_data(old_h), old_bytes < new_bytes ? old_bytes : new_bytes);
ray_free(old_h);
}
*hdr_out = new_h;
return new_p;
}
static inline void* scratch_calloc(ray_t** hdr_out, size_t nbytes) {
void* p = scratch_alloc(hdr_out, nbytes);
if (p) memset(p, 0, nbytes);
return p;
}
static inline void scratch_free(ray_t* hdr) {
if (hdr) ray_free(hdr);
}
typedef struct {
const char* ptr;
uint32_t len;
uint32_t hash;
} csv_strref_t;
RAY_INLINE const char* scan_field(const char* p, const char* buf_end,
char delimiter,
const char** out, size_t* out_len,
char* esc_buf, char** dyn_esc);
csv_type_t csv_resolve_int_width(int64_t min, int64_t max, bool has_null) {
if (min > max) return CSV_TYPE_I64;
if (has_null) {
if (min > NULL_I16 && max <= INT16_MAX) return CSV_TYPE_I16;
if (min > NULL_I32 && max <= INT32_MAX) return CSV_TYPE_I32;
return CSV_TYPE_I64;
}
if (min >= INT16_MIN && max <= INT16_MAX) return CSV_TYPE_I16;
if (min >= INT32_MIN && max <= INT32_MAX) return CSV_TYPE_I32;
return CSV_TYPE_I64;
}
RAY_INLINE int32_t fast_date(const char* p, size_t len, bool* is_null);
RAY_INLINE int32_t fast_time(const char* p, size_t len, bool* is_null);
RAY_INLINE int64_t fast_timestamp(const char* p, size_t len, bool* is_null);
static csv_type_t detect_type(const char* f, size_t len) {
if (len == 0) return CSV_TYPE_UNKNOWN;
if ((len == 3 && (memcmp(f, "N/A", 3) == 0 || memcmp(f, "n/a", 3) == 0)) ||
(len == 2 && (memcmp(f, "NA", 2) == 0 || memcmp(f, "na", 2) == 0)) ||
(len == 4 && (memcmp(f, "null", 4) == 0 || memcmp(f, "NULL", 4) == 0 ||
memcmp(f, "None", 4) == 0 || memcmp(f, "none", 4) == 0)) ||
(len == 1 && f[0] == '.'))
return CSV_TYPE_UNKNOWN;
if (len == 3) {
if ((f[0]=='n'||f[0]=='N') && (f[1]=='a'||f[1]=='A') && (f[2]=='n'||f[2]=='N'))
return CSV_TYPE_F64;
if ((f[0]=='i'||f[0]=='I') && (f[1]=='n'||f[1]=='N') && (f[2]=='f'||f[2]=='F'))
return CSV_TYPE_F64;
}
if ((len == 4 && (f[0]=='+' || f[0]=='-')) &&
(f[1]=='i'||f[1]=='I') && (f[2]=='n'||f[2]=='N') && (f[3]=='f'||f[3]=='F'))
return CSV_TYPE_F64;
if ((len == 4 && memcmp(f, "true", 4) == 0) ||
(len == 5 && memcmp(f, "false", 5) == 0) ||
(len == 4 && memcmp(f, "TRUE", 4) == 0) ||
(len == 5 && memcmp(f, "FALSE", 5) == 0))
return CSV_TYPE_BOOL;
const char* p = f;
const char* end = f + len;
if (*p == '-' || *p == '+') p++;
bool has_dot = false, has_e = false, has_digit = false;
while (p < end) {
unsigned char c = (unsigned char)*p;
if (c >= '0' && c <= '9') { has_digit = true; p++; continue; }
if (c == '.' && !has_dot) { has_dot = true; p++; continue; }
if ((c == 'e' || c == 'E') && !has_e) {
has_e = true; p++;
if (p < end && (*p == '-' || *p == '+')) p++;
continue;
}
break;
}
if (p == end && has_digit) {
if (!has_dot && !has_e) return CSV_TYPE_I64;
return CSV_TYPE_F64;
}
bool temporal_null = true;
if (len >= 19) {
(void)fast_timestamp(f, len, &temporal_null);
if (!temporal_null) return CSV_TYPE_TIMESTAMP;
}
(void)fast_time(f, len, &temporal_null);
if (!temporal_null) return CSV_TYPE_TIME;
if (len == 10) {
(void)fast_date(f, len, &temporal_null);
if (!temporal_null) return CSV_TYPE_DATE;
}
return CSV_TYPE_STR;
}
static csv_type_t promote_csv_type(csv_type_t cur, csv_type_t obs) {
if (cur == CSV_TYPE_UNKNOWN) return obs;
if (obs == CSV_TYPE_UNKNOWN) return cur;
if (cur == obs) return cur;
if (cur == CSV_TYPE_STR || obs == CSV_TYPE_STR) return CSV_TYPE_STR;
if ((cur == CSV_TYPE_DATE && obs == CSV_TYPE_TIMESTAMP) ||
(cur == CSV_TYPE_TIMESTAMP && obs == CSV_TYPE_DATE))
return CSV_TYPE_TIMESTAMP;
if (cur <= CSV_TYPE_F64 && obs <= CSV_TYPE_F64) {
if (cur == CSV_TYPE_F64 || obs == CSV_TYPE_F64) return CSV_TYPE_F64;
if (cur == CSV_TYPE_I64 || obs == CSV_TYPE_I64) return CSV_TYPE_I64;
return cur;
}
return CSV_TYPE_STR;
}
static void csv_cardinality_note(uint32_t* hashes, uint16_t* lens,
uint16_t* distinct, uint16_t* non_null,
int col, const char* fld, size_t flen) {
if (flen == 0) return;
size_t base = (size_t)col * CSV_SAMPLE_ROWS;
uint32_t h = (uint32_t)ray_hash_bytes(fld, flen);
uint16_t l = flen > UINT16_MAX ? UINT16_MAX : (uint16_t)flen;
for (uint16_t i = 0; i < distinct[col]; i++) {
if (hashes[base + i] == h && lens[base + i] == l) {
non_null[col]++;
return;
}
}
if (distinct[col] < CSV_SAMPLE_ROWS) {
hashes[base + distinct[col]] = h;
lens[base + distinct[col]] = l;
distinct[col]++;
}
non_null[col]++;
}
static int8_t csv_resolve_inferred_type(csv_type_t t,
uint16_t distinct,
uint16_t non_null) {
switch (t) {
case CSV_TYPE_BOOL: return RAY_BOOL;
case CSV_TYPE_I64: return RAY_I64;
case CSV_TYPE_F64: return RAY_F64;
case CSV_TYPE_DATE: return RAY_DATE;
case CSV_TYPE_TIME: return RAY_TIME;
case CSV_TYPE_TIMESTAMP: return RAY_TIMESTAMP;
case CSV_TYPE_GUID: return RAY_GUID;
case CSV_TYPE_STR:
return (distinct >= CSV_STR_DISTINCT_MIN ||
(non_null >= 64 &&
(uint32_t)distinct * 100u >= (uint32_t)non_null * 80u))
? RAY_STR : RAY_SYM;
default:
return RAY_SYM;
}
}
static bool csv_infer_types_from_offsets(const char* buf, const char* buf_end,
const int64_t* row_offsets,
int64_t n_rows, int ncols,
char delimiter, char* esc_buf,
int8_t* resolved_types) {
csv_type_t col_types[CSV_MAX_COLS];
memset(col_types, 0, (size_t)ncols * sizeof(csv_type_t));
ray_t *hash_hdr = NULL, *lens_hdr = NULL;
uint32_t* text_hashes = (uint32_t*)scratch_calloc(&hash_hdr,
(size_t)ncols * CSV_SAMPLE_ROWS * sizeof(uint32_t));
uint16_t* text_lens = (uint16_t*)scratch_calloc(&lens_hdr,
(size_t)ncols * CSV_SAMPLE_ROWS * sizeof(uint16_t));
uint16_t text_distinct[CSV_MAX_COLS] = {0};
uint16_t text_non_null[CSV_MAX_COLS] = {0};
if (!text_hashes || !text_lens) {
scratch_free(hash_hdr);
scratch_free(lens_hdr);
return false;
}
int64_t sample_n = n_rows < CSV_SAMPLE_ROWS ? n_rows : CSV_SAMPLE_ROWS;
for (int64_t si = 0; si < sample_n; si++) {
int64_t r = si;
if (sample_n > 1 && sample_n < n_rows)
r = (si * (n_rows - 1)) / (sample_n - 1);
const char* rp = buf + row_offsets[r];
for (int c = 0; c < ncols; c++) {
const char* fld;
size_t flen;
char* dyn_esc = NULL;
rp = scan_field(rp, buf_end, delimiter, &fld, &flen, esc_buf, &dyn_esc);
csv_type_t t = detect_type(fld, flen);
if (t == CSV_TYPE_STR)
csv_cardinality_note(text_hashes, text_lens,
text_distinct, text_non_null,
c, fld, flen);
if (dyn_esc) ray_sys_free(dyn_esc);
col_types[c] = promote_csv_type(col_types[c], t);
}
}
for (int c = 0; c < ncols; c++)
resolved_types[c] = csv_resolve_inferred_type(
col_types[c], text_distinct[c], text_non_null[c]);
scratch_free(hash_hdr);
scratch_free(lens_hdr);
return true;
}
static const char* scan_field_quoted(const char* p, const char* buf_end,
char delim,
const char** out, size_t* out_len,
char* esc_buf, char** dyn_esc) {
p++;
const char* fld_start = p;
bool has_escape = false;
while (p < buf_end) {
if (*p == '"') {
if (p + 1 < buf_end && *(p + 1) == '"') {
has_escape = true;
p += 2;
} else {
break;
}
} else {
p++;
}
}
size_t raw_len = (size_t)(p - fld_start);
if (p < buf_end && *p == '"') p++;
if (has_escape) {
char* dest = esc_buf;
if (RAY_UNLIKELY(raw_len > 8192)) {
dest = (char*)ray_sys_alloc(raw_len);
if (!dest) {
*out = fld_start;
*out_len = raw_len;
goto advance;
}
*dyn_esc = dest;
}
size_t olen = 0;
for (const char* s = fld_start; s < fld_start + raw_len; s++) {
if (*s == '"' && s + 1 < fld_start + raw_len && *(s + 1) == '"') {
dest[olen++] = '"';
s++;
} else {
dest[olen++] = *s;
}
}
*out = dest;
*out_len = olen;
} else {
*out = fld_start;
*out_len = raw_len;
}
advance:
if (p < buf_end && *p == delim) p++;
return p;
}
RAY_INLINE const char* scan_field(const char* p, const char* buf_end,
char delim,
const char** out, size_t* out_len,
char* esc_buf, char** dyn_esc) {
if (RAY_UNLIKELY(p >= buf_end)) {
*out = p;
*out_len = 0;
return p;
}
if (RAY_LIKELY(*p != '"')) {
const char* s = p;
while (p < buf_end && *p != delim && *p != '\n' && *p != '\r') p++;
*out = s;
*out_len = (size_t)(p - s);
if (p < buf_end && *p == delim) return p + 1;
return p;
}
return scan_field_quoted(p, buf_end, delim, out, out_len, esc_buf, dyn_esc);
}
RAY_INLINE int64_t fast_i64(const char* p, size_t len, bool* is_null) {
int64_t v = 0;
size_t n = ray_parse_i64(p, len, &v);
*is_null = (n == 0 || n != len);
return *is_null ? 0 : v;
}
RAY_INLINE double fast_f64(const char* p, size_t len, bool* is_null) {
double v = 0.0;
size_t n = ray_parse_f64(p, len, &v);
*is_null = (n == 0 || n != len || v != v);
return *is_null ? NULL_F64 : v;
}
RAY_INLINE int32_t civil_to_days(int y, int m, int d) {
if (m <= 2) { y--; m += 9; } else { m -= 3; }
int era = (y >= 0 ? y : y - 399) / 400;
int yoe = y - era * 400;
int doy = (153 * m + 2) / 5 + d - 1;
int doe = yoe * 365 + yoe / 4 - yoe / 100 + doy;
return (int32_t)(era * 146097 + doe - 719468 - 10957);
}
RAY_INLINE int32_t fast_date(const char* p, size_t len, bool* is_null) {
bool sep_ok = len == 10 && (p[4] == '-' || p[4] == '.') && p[7] == p[4];
if (RAY_UNLIKELY(!sep_ok)) { *is_null = true; return 0; }
if (RAY_UNLIKELY((unsigned)(p[0]-'0') > 9u || (unsigned)(p[1]-'0') > 9u ||
(unsigned)(p[2]-'0') > 9u || (unsigned)(p[3]-'0') > 9u ||
(unsigned)(p[5]-'0') > 9u || (unsigned)(p[6]-'0') > 9u ||
(unsigned)(p[8]-'0') > 9u || (unsigned)(p[9]-'0') > 9u)) {
*is_null = true; return 0;
}
*is_null = false;
int y = (p[0]-'0')*1000 + (p[1]-'0')*100 + (p[2]-'0')*10 + (p[3]-'0');
int m = (p[5]-'0')*10 + (p[6]-'0');
int d = (p[8]-'0')*10 + (p[9]-'0');
if (RAY_UNLIKELY(m < 1 || m > 12 || d < 1)) { *is_null = true; return 0; }
static const int md[] = {0,31,28,31,30,31,30,31,31,30,31,30,31};
int leap = (y % 4 == 0 && (y % 100 != 0 || y % 400 == 0));
if (RAY_UNLIKELY(d > md[m] + (m == 2 && leap ? 1 : 0))) { *is_null = true; return 0; }
return civil_to_days(y, m, d);
}
RAY_INLINE int32_t fast_time(const char* p, size_t len, bool* is_null) {
*is_null = false;
size_t o = 0;
int64_t sign = 1;
if (len > 0 && p[0] == '-') { sign = -1; o = 1; }
int64_t h = 0;
size_t hd = o;
while (hd < len && (unsigned)(p[hd] - '0') <= 9u) {
h = h * 10 + (p[hd] - '0');
if (RAY_UNLIKELY(++hd - o > 7)) { *is_null = true; return 0; }
}
size_t w = hd - o;
if (RAY_UNLIKELY(w == 0 || o + w + 6 > len)) { *is_null = true; return 0; }
if (RAY_UNLIKELY(p[o+w] != ':' || p[o+w+3] != ':')) { *is_null = true; return 0; }
if (RAY_UNLIKELY((unsigned)(p[o+w+1]-'0') > 9u || (unsigned)(p[o+w+2]-'0') > 9u ||
(unsigned)(p[o+w+4]-'0') > 9u || (unsigned)(p[o+w+5]-'0') > 9u)) {
*is_null = true; return 0;
}
int mi = (p[o+w+1]-'0')*10 + (p[o+w+2]-'0');
int s = (p[o+w+4]-'0')*10 + (p[o+w+5]-'0');
if (RAY_UNLIKELY(mi > 59 || s > 59)) { *is_null = true; return 0; }
int64_t ms = h * 3600000 + (int64_t)mi * 60000 + (int64_t)s * 1000;
size_t i = o + w + 6;
if (i < len) {
if (RAY_UNLIKELY(p[i] != '.')) { *is_null = true; return 0; }
i++;
int frac = 0, digits = 0;
for (; i < len && (unsigned)(p[i]-'0') <= 9u; i++) {
if (digits < 3) { frac = frac * 10 + (p[i] - '0'); digits++; }
}
if (RAY_UNLIKELY(digits == 0 || i != len)) { *is_null = true; return 0; }
while (digits < 3) { frac *= 10; digits++; }
ms += frac;
}
ms *= sign;
if (RAY_UNLIKELY(ms > INT32_MAX || ms < INT32_MIN)) { *is_null = true; return 0; }
return (int32_t)ms;
}
RAY_INLINE int64_t fast_time_ns(const char* p, size_t len, bool* is_null, size_t* consumed) {
if (RAY_UNLIKELY(len < 8 || p[2] != ':' || p[5] != ':')) { *is_null = true; return 0; }
if (RAY_UNLIKELY((unsigned)(p[0]-'0') > 9u || (unsigned)(p[1]-'0') > 9u ||
(unsigned)(p[3]-'0') > 9u || (unsigned)(p[4]-'0') > 9u ||
(unsigned)(p[6]-'0') > 9u || (unsigned)(p[7]-'0') > 9u)) {
*is_null = true; return 0;
}
*is_null = false;
int h = (p[0]-'0')*10 + (p[1]-'0');
int mi = (p[3]-'0')*10 + (p[4]-'0');
int s = (p[6]-'0')*10 + (p[7]-'0');
if (RAY_UNLIKELY(h > 23 || mi > 59 || s > 59)) { *is_null = true; return 0; }
int64_t ns = (int64_t)h * 3600000000000LL + (int64_t)mi * 60000000000LL +
(int64_t)s * 1000000000LL;
size_t used = 8;
if (len > 8 && p[8] == '.') {
int64_t frac = 0;
int digits = 0;
size_t i = 9;
for (; i < len && (unsigned)(p[i]-'0') <= 9u; i++) {
if (digits < 9) { frac = frac * 10 + (int64_t)(p[i] - '0'); digits++; }
}
if (RAY_UNLIKELY(digits == 0)) { *is_null = true; return 0; }
while (digits < 9) { frac *= 10; digits++; }
ns += frac;
used = i;
}
if (consumed) *consumed = used;
return ns;
}
RAY_INLINE bool parse_tz_offset(const char* p, size_t len, int64_t* out_ns, size_t* consumed) {
if (len == 0) return false;
if (p[0] == 'Z' || p[0] == 'z') { *out_ns = 0; if (consumed) *consumed = 1; return true; }
int sign;
if (p[0] == '+') sign = 1;
else if (p[0] == '-') sign = -1;
else return false;
if (len < 3 || p[1] < '0' || p[1] > '9' || p[2] < '0' || p[2] > '9') return false;
int hh = (p[1]-'0')*10 + (p[2]-'0');
int mm = 0;
size_t i = 3;
size_t msep = (len > 3 && p[3] == ':') ? 1 : 0;
if (len >= 3 + msep + 2 &&
p[3+msep] >= '0' && p[3+msep] <= '9' && p[4+msep] >= '0' && p[4+msep] <= '9') {
mm = (p[3+msep]-'0')*10 + (p[4+msep]-'0');
i = 3 + msep + 2;
}
if (RAY_UNLIKELY(hh > 23 || mm > 59)) return false;
*out_ns = (int64_t)sign * ((int64_t)hh * 3600 + (int64_t)mm * 60) * 1000000000LL;
if (consumed) *consumed = i;
return true;
}
RAY_INLINE int64_t fast_timestamp(const char* p, size_t len, bool* is_null) {
if (RAY_UNLIKELY(len < 19)) { *is_null = true; return 0; }
bool rayfall_sep = (p[10] == 'D');
bool dt_sep_ok = p[10] == 'T' || p[10] == 't' || p[10] == ' ' ||
(rayfall_sep && p[4] == '.');
if (RAY_UNLIKELY(!dt_sep_ok)) { *is_null = true; return 0; }
*is_null = false;
int32_t days = fast_date(p, 10, is_null);
if (*is_null) return 0;
bool time_null = false;
size_t time_used = 8;
int64_t time_ns = fast_time_ns(p + 11, len - 11, &time_null, &time_used);
if (time_null) { *is_null = true; return 0; }
const int64_t NS_PER_DAY = 86400000000000LL;
int64_t result = (int64_t)days * NS_PER_DAY + time_ns;
size_t off = 11 + time_used;
if (RAY_UNLIKELY(off < len)) {
int64_t adj;
size_t tz_used = 0;
if (RAY_UNLIKELY(!parse_tz_offset(p + off, len - off, &adj, &tz_used) ||
off + tz_used != len)) {
*is_null = true; return 0;
}
result -= adj;
}
return result;
}
RAY_INLINE uint8_t fast_bool(const char* s, size_t len, bool* is_null) {
if (len == 0) { *is_null = true; return 0; }
*is_null = false;
if ((len == 4 && (memcmp(s, "true", 4) == 0 || memcmp(s, "TRUE", 4) == 0)) ||
(len == 1 && s[0] == '1'))
return 1;
if ((len == 5 && (memcmp(s, "false", 5) == 0 || memcmp(s, "FALSE", 5) == 0)) ||
(len == 1 && s[0] == '0'))
return 0;
*is_null = true;
return 0;
}
RAY_INLINE int hex_nibble(unsigned char c) {
if (c >= '0' && c <= '9') return c - '0';
if (c >= 'a' && c <= 'f') return c - 'a' + 10;
if (c >= 'A' && c <= 'F') return c - 'A' + 10;
return -1;
}
RAY_INLINE void fast_guid(const char* p, size_t len, uint8_t* dst, bool* is_null) {
if (RAY_UNLIKELY(len != 36 ||
p[8] != '-' || p[13] != '-' ||
p[18] != '-' || p[23] != '-')) {
*is_null = true;
return;
}
static const uint8_t pos[16] = { 0,2,4,6, 9,11, 14,16, 19,21, 24,26,28,30,32,34 };
for (int i = 0; i < 16; i++) {
int hi = hex_nibble((unsigned char)p[pos[i]]);
int lo = hex_nibble((unsigned char)p[pos[i] + 1]);
if (RAY_UNLIKELY((hi | lo) < 0)) { *is_null = true; return; }
dst[i] = (uint8_t)((hi << 4) | lo);
}
*is_null = false;
}
#if defined(DEBUG) || defined(RAY_HARDENED)
#define CSV_SCAN_PAR_MIN_BYTES (8u << 10)
#define CSV_SCAN_CHUNK_MIN 64u
#define CSV_SCAN_CHUNKS(pool) ((int64_t)RAY_POOL_INIT_TASKS)
#else
#define CSV_SCAN_PAR_MIN_BYTES (4u << 20)
#define CSV_SCAN_CHUNK_MIN (256u << 10)
#define CSV_SCAN_CHUNKS(pool) ((int64_t)ray_pool_total_workers(pool) * 4)
#endif
typedef struct {
const char* buf;
size_t file_size;
const size_t* bound;
uint64_t* quote_cnt;
uint64_t* term_cnt;
const uint8_t* start_quoted;
const int64_t* slot;
int64_t* n_out;
int64_t* offs;
bool has_quotes;
} csv_scan_ctx_t;
static void csv_scan_count_fn(void* arg, uint32_t worker_id,
int64_t start, int64_t end_i) {
(void)worker_id;
csv_scan_ctx_t* ctx = (csv_scan_ctx_t*)arg;
for (int64_t i = start; i < end_i; i++) {
const char* p = ctx->buf + ctx->bound[i];
const char* e = ctx->buf + ctx->bound[i + 1];
uint64_t q = 0, t = 0;
while (p < e) {
const char* stop = p + (1 << 20);
if (stop > e) stop = e;
for (; p < stop; p++) {
char c = *p;
q += (c == '"');
t += (c == '\n' || c == '\r');
}
if (RAY_UNLIKELY(ray_interrupted())) break;
}
ctx->quote_cnt[i] = q;
ctx->term_cnt[i] = t;
}
}
static void csv_scan_rows_fn(void* arg, uint32_t worker_id,
int64_t start, int64_t end_i) {
(void)worker_id;
csv_scan_ctx_t* ctx = (csv_scan_ctx_t*)arg;
const char* buf = ctx->buf;
const char* end = buf + ctx->file_size;
for (int64_t i = start; i < end_i; i++) {
const char* p = buf + ctx->bound[i];
const char* cend = buf + ctx->bound[i + 1];
int64_t* out = ctx->offs + ctx->slot[i];
int64_t n = 0;
if (RAY_LIKELY(!ctx->has_quotes)) {
size_t checked = 0;
while (p < cend) {
if (RAY_UNLIKELY((++checked & 0xFFFF) == 0 && ray_interrupted())) break;
const char* nl = (const char*)memchr(p, '\n', (size_t)(cend - p));
if (!nl) break;
p = nl + 1;
if (p < end && *p == '\r') p++;
if (p >= end) break;
out[n++] = (int64_t)(p - buf);
}
} else {
bool in_quote = ctx->start_quoted[i] != 0;
size_t checked = 0;
while (p < cend) {
if (in_quote) {
if (RAY_UNLIKELY((++checked & 0xFFF) == 0 && ray_interrupted())) break;
const char* q = (const char*)memchr(p, '"', (size_t)(cend - p));
if (!q) { p = cend; break; }
p = q + 1;
in_quote = false;
continue;
}
if (RAY_UNLIKELY((++checked & 0xFFFF) == 0 && ray_interrupted())) break;
char c = *p;
if (c == '"') {
in_quote = true;
p++;
} else if (c == '\n' || c == '\r') {
if (c == '\r' && p + 1 < end && *(p + 1) == '\n') p++;
p++;
if (p < end) out[n++] = (int64_t)(p - buf);
} else {
p++;
}
}
}
ctx->n_out[i] = n;
}
}
static size_t csv_scan_split_at(const char* buf, size_t file_size, size_t s,
bool has_quotes) {
if (s == 0 || s >= file_size) return s;
char prev = buf[s - 1], cur = buf[s];
if (has_quotes) { if (prev == '\r' && cur == '\n') s++; }
else { if (prev == '\n' && cur == '\r') s++; }
return s;
}
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;
size_t remaining = buf_size - data_offset;
ray_pool_t* pool = ray_pool_get();
if (!ray_pool_par_dispatch_ok(pool, (int64_t)remaining,
(int64_t)CSV_SCAN_PAR_MIN_BYTES))
return -2;
int64_t n_chunks = CSV_SCAN_CHUNKS(pool);
if (n_chunks > (int64_t)RAY_POOL_INIT_TASKS) n_chunks = RAY_POOL_INIT_TASKS;
size_t span = remaining / (size_t)n_chunks;
if (span < CSV_SCAN_CHUNK_MIN) {
span = CSV_SCAN_CHUNK_MIN;
n_chunks = (int64_t)((remaining + span - 1) / span);
}
if (n_chunks < 2) return -2;
size_t hdr_bytes = (size_t)(n_chunks + 1) * sizeof(size_t)
+ (size_t)n_chunks * (2 * sizeof(uint64_t)
+ 2 * sizeof(int64_t)
+ sizeof(uint8_t));
void* blk = ray_sys_alloc(hdr_bytes);
if (!blk) return -2;
size_t* bound = (size_t*)blk;
uint64_t* quote_cnt = (uint64_t*)(bound + n_chunks + 1);
uint64_t* term_cnt = quote_cnt + n_chunks;
int64_t* slot = (int64_t*)(term_cnt + n_chunks);
int64_t* n_out = slot + n_chunks;
uint8_t* start_q = (uint8_t*)(n_out + n_chunks);
csv_scan_ctx_t ctx = {
.buf = buf, .file_size = buf_size, .bound = bound,
.quote_cnt = quote_cnt, .term_cnt = term_cnt,
.start_quoted = start_q, .slot = slot, .n_out = n_out,
.offs = NULL, .has_quotes = false,
};
bound[0] = data_offset;
for (int64_t i = 1; i < n_chunks; i++) {
size_t s = data_offset + (size_t)i * span;
if (s > buf_size) s = buf_size;
if (s < bound[i - 1]) s = bound[i - 1];
bound[i] = s;
}
bound[n_chunks] = buf_size;
ray_progress_span_phase("scan: quotes", prog_base, prog_len / 5);
ray_pool_dispatch_n(pool, csv_scan_count_fn, &ctx, (uint32_t)n_chunks);
if (ray_interrupted()) { ray_sys_free(blk); return -1; }
uint64_t total_q = 0, total_term = 0;
for (int64_t i = 0; i < n_chunks; i++) {
start_q[i] = (uint8_t)(total_q & 1u);
total_q += quote_cnt[i];
total_term += term_cnt[i];
}
ctx.has_quotes = total_q != 0 || force_quotes;
for (int64_t i = 1; i < n_chunks; i++) {
size_t s = csv_scan_split_at(buf, buf_size, bound[i], ctx.has_quotes);
if (s < bound[i - 1]) s = bound[i - 1];
bound[i] = s;
}
uint64_t slack = 2u * (uint64_t)n_chunks + 1u;
if (total_term > (uint64_t)INT64_MAX / (uint64_t)sizeof(int64_t) - slack) {
ray_sys_free(blk);
return -2;
}
int64_t cap = (int64_t)(total_term + slack);
{
int64_t base = 1;
for (int64_t i = 0; i < n_chunks; i++) {
slot[i] = base;
base += (int64_t)term_cnt[i] + 2;
}
}
ray_t* hdr = NULL;
int64_t* offs = (int64_t*)scratch_alloc(&hdr, (size_t)cap * sizeof(int64_t));
if (!offs) { ray_sys_free(blk); return -2; }
offs[0] = (int64_t)data_offset;
ctx.offs = offs;
ray_progress_span_phase("scan: rows", prog_base + prog_len / 5,
prog_len - prog_len / 5);
ray_pool_dispatch_n(pool, csv_scan_rows_fn, &ctx, (uint32_t)n_chunks);
if (ray_interrupted()) {
scratch_free(hdr);
ray_sys_free(blk);
return -1;
}
int64_t n = 1;
for (int64_t i = 0; i < n_chunks; i++) {
int64_t cnt = n_out[i];
if (cnt <= 0) continue;
if (n != slot[i])
memmove(offs + n, offs + slot[i], (size_t)cnt * sizeof(int64_t));
n += cnt;
}
ray_sys_free(blk);
*offsets_out = offs;
*hdr_out = hdr;
return n;
}
static int64_t build_row_offsets_serial(const char* buf, size_t buf_size,
size_t data_offset,
uint64_t prog_base, uint64_t prog_len,
int64_t** offsets_out, ray_t** hdr_out) {
const char* p = buf + data_offset;
const char* end = buf + buf_size;
if (p >= end) { *offsets_out = NULL; *hdr_out = NULL; return 0; }
size_t remaining = (size_t)(end - p);
int64_t est = (int64_t)(remaining / 40) + 16;
ray_t* hdr = NULL;
int64_t* offs = (int64_t*)scratch_alloc(&hdr, (size_t)est * sizeof(int64_t));
if (!offs) { *offsets_out = NULL; *hdr_out = NULL; return 0; }
int64_t n = 0;
offs[n++] = (int64_t)(p - buf);
bool has_quotes = (memchr(p, '"', remaining) != NULL);
if (RAY_LIKELY(!has_quotes)) {
for (;;) {
if (RAY_UNLIKELY((n & 0xFFFF) == 0)) {
if (ray_interrupted()) {
scratch_free(hdr);
*offsets_out = NULL;
*hdr_out = NULL;
return -1;
}
ray_progress_span_set(prog_base +
(uint64_t)((double)prog_len * (double)(size_t)(p - buf - (ptrdiff_t)data_offset)
/ (double)remaining));
}
const char* nl = (const char*)memchr(p, '\n', (size_t)(end - p));
if (!nl) break;
p = nl + 1;
if (p < end && *p == '\r') p++;
if (p >= end) break;
if (n >= est) {
est *= 2;
offs = (int64_t*)scratch_realloc(&hdr,
(size_t)n * sizeof(int64_t),
(size_t)est * sizeof(int64_t));
if (!offs) { scratch_free(hdr); *offsets_out = NULL; *hdr_out = NULL; return 0; }
}
offs[n++] = (int64_t)(p - buf);
}
} else {
bool in_quote = false;
size_t checked = 0;
while (p < end) {
if (RAY_UNLIKELY(++checked == 65536)) {
checked = 0;
if (ray_interrupted()) {
scratch_free(hdr);
*offsets_out = NULL;
*hdr_out = NULL;
return -1;
}
ray_progress_span_set(prog_base +
(uint64_t)((double)prog_len * (double)(size_t)(p - buf - (ptrdiff_t)data_offset)
/ (double)remaining));
}
char c = *p;
if (c == '"') {
in_quote = !in_quote;
p++;
} else if (!in_quote && (c == '\n' || c == '\r')) {
if (c == '\r' && p + 1 < end && *(p + 1) == '\n') p++;
p++;
if (p < end) {
if (n >= est) {
est *= 2;
offs = (int64_t*)scratch_realloc(&hdr,
(size_t)n * sizeof(int64_t),
(size_t)est * sizeof(int64_t));
if (!offs) { scratch_free(hdr); *offsets_out = NULL; *hdr_out = NULL; return 0; }
}
offs[n++] = (int64_t)(p - buf);
}
} else {
p++;
}
}
}
*offsets_out = offs;
*hdr_out = hdr;
return n;
}
static int64_t build_row_offsets(const char* buf, size_t buf_size,
size_t data_offset,
uint64_t prog_base, uint64_t prog_len,
int64_t** offsets_out, ray_t** hdr_out) {
if (data_offset >= buf_size) { *offsets_out = NULL; *hdr_out = NULL; return 0; }
int64_t par = build_row_offsets_par(buf, buf_size, data_offset,
prog_base, prog_len, false,
offsets_out, hdr_out);
if (par != -2) return par;
ray_progress_span_phase("scan", prog_base, prog_len);
return build_row_offsets_serial(buf, buf_size, data_offset,
prog_base, prog_len,
offsets_out, hdr_out);
}
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;
{
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;
if (n == 0) { scratch_free(hdr); return -2; }
if (n > max_rows) {
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) {
if (next_offset_out) *next_offset_out = buf_size;
*offsets_out = offs; *hdr_out = hdr;
return n;
}
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,
int64_t** offsets_out, ray_t** hdr_out,
size_t* next_offset_out) {
const char* p = buf + data_offset;
const char* end = buf + buf_size;
*offsets_out = NULL;
*hdr_out = NULL;
if (next_offset_out) *next_offset_out = data_offset;
if (max_rows <= 0 || p >= end) return 0;
size_t remaining = (size_t)(end - p);
int64_t est = (int64_t)(remaining / 40) + 16;
if (est < 1) est = 1;
if (est > max_rows) est = max_rows;
ray_t* hdr = NULL;
int64_t* offs = (int64_t*)scratch_alloc(&hdr, (size_t)est * sizeof(int64_t));
if (!offs) return 0;
int64_t n = 0;
offs[n++] = (int64_t)(p - buf);
if (RAY_LIKELY(!data_has_quotes)) {
for (;;) {
if (RAY_UNLIKELY((n & 0xFFFF) == 0 && ray_interrupted())) {
scratch_free(hdr);
return -1;
}
const char* nl = (const char*)memchr(p, '\n', (size_t)(end - p));
if (!nl) {
p = end;
break;
}
p = nl + 1;
if (p < end && *p == '\r') p++;
if (p >= end) break;
if (n >= max_rows) break;
if (n >= est) {
int64_t new_est = est * 2;
if (new_est > max_rows) new_est = max_rows;
offs = (int64_t*)scratch_realloc(&hdr,
(size_t)n * sizeof(int64_t),
(size_t)new_est * sizeof(int64_t));
if (!offs) {
scratch_free(hdr);
return 0;
}
est = new_est;
}
offs[n++] = (int64_t)(p - buf);
}
} else {
bool in_quote = false;
size_t checked = 0;
while (p < end) {
if (RAY_UNLIKELY(++checked == 65536)) {
checked = 0;
if (ray_interrupted()) {
scratch_free(hdr);
return -1;
}
}
char c = *p;
if (c == '"') {
in_quote = !in_quote;
p++;
} else if (!in_quote && (c == '\n' || c == '\r')) {
if (c == '\r' && p + 1 < end && *(p + 1) == '\n') p++;
p++;
if (p >= end) break;
if (n >= max_rows) break;
if (n >= est) {
int64_t new_est = est * 2;
if (new_est > max_rows) new_est = max_rows;
offs = (int64_t*)scratch_realloc(&hdr,
(size_t)n * sizeof(int64_t),
(size_t)new_est * sizeof(int64_t));
if (!offs) {
scratch_free(hdr);
return 0;
}
est = new_est;
}
offs[n++] = (int64_t)(p - buf);
} else {
p++;
}
}
}
*offsets_out = offs;
*hdr_out = hdr;
if (next_offset_out) *next_offset_out = (size_t)(p - buf);
return n;
}
static int64_t csv_scan_rows_chunk(const char* buf, size_t file_size,
size_t off, int64_t max_rows,
size_t* avg_row, bool data_has_quotes,
int64_t** offs, ray_t** hdr, size_t* next) {
int64_t cnt = build_row_offsets_window(buf, file_size, off, max_rows,
*avg_row, data_has_quotes,
offs, hdr, next);
if (cnt == -2)
cnt = build_row_offsets_limited(buf, file_size, off, max_rows,
data_has_quotes, offs, hdr, next);
if (cnt > 0 && *next > off) *avg_row = (*next - off) / (size_t)cnt;
return cnt;
}
#define CSV_SAMPLE_SCAN_ROWS (1 << 20)
static int64_t csv_streaming_sample(const char* buf, size_t file_size,
size_t data_offset, bool data_has_quotes,
int64_t** offsets_out, ray_t** hdr_out) {
*offsets_out = NULL;
*hdr_out = NULL;
if (data_offset >= file_size) return 0;
int64_t total = 0;
size_t avg_row = 64;
for (size_t off = data_offset; off < file_size; ) {
int64_t* o = NULL; ray_t* h = NULL; size_t next = off;
int64_t cnt = csv_scan_rows_chunk(buf, file_size, off,
CSV_SAMPLE_SCAN_ROWS, &avg_row,
data_has_quotes, &o, &h, &next);
scratch_free(h);
if (cnt < 0) return -1;
if (cnt == 0 || next <= off) break;
total += cnt;
off = next;
}
if (total <= CSV_SAMPLE_ROWS)
return build_row_offsets_limited(buf, file_size, data_offset,
CSV_SAMPLE_ROWS, data_has_quotes,
offsets_out, hdr_out, NULL);
ray_t* hdr = NULL;
int64_t* picked = (int64_t*)scratch_alloc(&hdr,
(size_t)CSV_SAMPLE_ROWS * sizeof(int64_t));
if (!picked) return 0;
int64_t si = 0, base = 0;
int64_t target = 0;
for (size_t off = data_offset; off < file_size && si < CSV_SAMPLE_ROWS; ) {
int64_t* o = NULL; ray_t* h = NULL; size_t next = off;
int64_t cnt = csv_scan_rows_chunk(buf, file_size, off,
CSV_SAMPLE_SCAN_ROWS, &avg_row,
data_has_quotes, &o, &h, &next);
if (cnt < 0) { scratch_free(h); scratch_free(hdr); return -1; }
if (cnt == 0 || next <= off) { scratch_free(h); break; }
while (si < CSV_SAMPLE_ROWS && target < base + cnt) {
picked[si++] = o[target - base];
if (si < CSV_SAMPLE_ROWS)
target = si * (total - 1) / (CSV_SAMPLE_ROWS - 1);
}
scratch_free(h);
base += cnt;
off = next;
}
if (si != CSV_SAMPLE_ROWS) {
scratch_free(hdr);
return 0;
}
*offsets_out = picked;
*hdr_out = hdr;
return si;
}
#if defined(DEBUG) || defined(RAY_HARDENED)
#define CSV_DEDUP_MAX_ENTS 8192u
#else
#define CSV_DEDUP_MAX_ENTS (1u << 20)
#endif
#define CSV_DEDUP_MIN_SLOTS 64u
typedef struct {
uint32_t hash;
uint32_t len;
const char* ptr;
int64_t gid;
int64_t row;
} csv_dedup_ent_t;
typedef struct {
uint32_t* slots;
uint32_t n_slots;
csv_dedup_ent_t* ents;
uint32_t n_ents;
uint32_t ents_cap;
bool done;
bool overflow;
} csv_dedup_t;
static void csv_dedup_release(csv_dedup_t* d) {
if (d->slots) { ray_sys_free(d->slots); d->slots = NULL; }
if (d->ents) { ray_sys_free(d->ents); d->ents = NULL; }
d->n_slots = 0; d->n_ents = 0; d->ents_cap = 0;
}
static bool csv_dedup_grow(csv_dedup_t* d) {
uint32_t new_slots = d->n_slots ? d->n_slots * 2 : CSV_DEDUP_MIN_SLOTS;
uint32_t* s = (uint32_t*)ray_sys_alloc((size_t)new_slots * sizeof(uint32_t));
if (!s) return false;
memset(s, 0, (size_t)new_slots * sizeof(uint32_t));
uint32_t mask = new_slots - 1;
for (uint32_t i = 0; i < d->n_ents; i++) {
uint32_t j = d->ents[i].hash & mask;
while (s[j]) j = (j + 1) & mask;
s[j] = i + 1;
}
if (d->slots) ray_sys_free(d->slots);
d->slots = s;
d->n_slots = new_slots;
return true;
}
static bool csv_dedup_grow_ents(csv_dedup_t* d) {
uint32_t cap = d->ents_cap ? d->ents_cap * 2 : CSV_DEDUP_MIN_SLOTS;
csv_dedup_ent_t* e = (csv_dedup_ent_t*)ray_sys_realloc(
d->ents, (size_t)cap * sizeof(csv_dedup_ent_t));
if (!e) return false;
d->ents = e;
d->ents_cap = cap;
return true;
}
typedef struct {
csv_strref_t** str_refs;
void** col_data;
const int* cols;
csv_dedup_t* dicts;
int64_t n_rows;
int64_t empty_gid;
bool* empties;
struct ray_sym_domain_s* dom;
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;
}
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;
}
static inline void csv_note_empty(bool* empties, int col, bool is_null) {
if (is_null && empties) empties[col] = true;
}
static void csv_dedup_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;
csv_dedup_t* d = &ctx->dicts[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]];
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)) {
d->overflow = true;
d->done = true;
csv_dedup_release(d);
return;
}
for (int64_t r = 0; r < n_rows; r++) {
if (RAY_UNLIKELY((r & 1023) == 0 && ray_interrupted())) return;
if (refs[r].ptr == NULL) { if (part == 0) codes[r] = 0; continue; }
uint32_t h = refs[r].hash;
if (csv_dedup_part(ctx, h) != part) continue;
uint32_t mask = d->n_slots - 1;
uint32_t j = h & mask;
uint32_t found = 0;
for (;;) {
uint32_t slot = d->slots[j];
if (!slot) break;
const csv_dedup_ent_t* e = &d->ents[slot - 1];
if (e->hash == h && e->len == refs[r].len &&
memcmp(e->ptr, refs[r].ptr, refs[r].len) == 0) {
found = slot;
break;
}
j = (j + 1) & mask;
}
if (found) { codes[r] = found; continue; }
if (d->n_ents >= CSV_DEDUP_MAX_ENTS) {
d->overflow = true;
d->done = true;
csv_dedup_release(d);
return;
}
if (d->n_ents == d->ents_cap && !csv_dedup_grow_ents(d)) {
d->overflow = true; d->done = true; csv_dedup_release(d); return;
}
csv_dedup_ent_t* e = &d->ents[d->n_ents];
e->hash = h;
e->len = refs[r].len;
e->ptr = refs[r].ptr;
e->gid = 0;
e->row = r;
d->n_ents++;
codes[r] = d->n_ents;
d->slots[j] = d->n_ents;
if (d->n_ents * 4u >= d->n_slots * 3u && !csv_dedup_grow(d)) {
d->overflow = true; d->done = true; csv_dedup_release(d); return;
}
}
d->done = true;
}
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;
if (!csv_col_dict_ok(ctx, (int)start)) return;
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;
bool saw_null = false;
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)dcol[csv_dedup_part(ctx, refs[r].hash)].ents[code - 1].gid
: empty;
ids[r] = id;
saw_null |= (id == 0);
}
csv_note_empty(ctx->empties, ctx->cols[start], saw_null);
}
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;
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));
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);
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;
int64_t empty_sym_id = ray_sym_intern_prehashed(
(uint32_t)ray_hash_bytes("", 0), "", 0);
if (empty_sym_id < 0) empty_sym_id = 0;
dd->empty_gid = empty_sym_id;
for (int i = 0; i < n_sym; i++) {
int c = dd->cols[i];
csv_dedup_t* d = &dd->dicts[i];
int64_t max_id = empty_sym_id;
uint32_t current = ray_sym_count();
if (!ray_sym_ensure_cap(current +
(uint32_t)(dd->n_rows < UINT32_MAX ? dd->n_rows : UINT32_MAX)))
return false;
if (d->done && !d->overflow) {
for (uint32_t e = 0; e < d->n_ents; e++) {
if (RAY_UNLIKELY((e & 1023) == 0 && ray_interrupted())) return false;
int64_t id = ray_sym_intern_no_split_unlocked(d->ents[e].ptr,
d->ents[e].len);
if (id < 0) { ok = false; id = 0; }
d->ents[e].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 (RAY_UNLIKELY((r & 1023) == 0 && ray_interrupted())) return false;
if (refs[r].ptr == NULL) {
ids[r] = (uint32_t)empty_sym_id;
saw_null |= (empty_sym_id == 0);
continue;
}
int64_t id = ray_sym_intern_no_split_unlocked(refs[r].ptr, refs[r].len);
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;
if (prog_len)
ray_progress_span_set(prog_base + (uint64_t)((double)prog_len *
(double)(i + 1) / (double)n_sym));
}
return ok;
}
static void csv_free_escaped_strrefs(csv_strref_t** str_refs, int n_cols,
const csv_type_t* col_types,
int64_t n_rows,
const char* buf, size_t buf_size,
const uint8_t* row_done,
const bool* col_had_escaped) {
const char* buf_end = buf + buf_size;
for (int c = 0; c < n_cols; c++) {
if (col_types[c] != CSV_TYPE_STR || !str_refs[c]) continue;
if (col_had_escaped && !col_had_escaped[c]) continue;
for (int64_t r = 0; r < n_rows; r++) {
if (row_done && !row_done[r]) continue;
const char* p = str_refs[c][r].ptr;
if (p && (p < buf || p >= buf_end))
ray_sys_free((void*)p);
}
}
}
static bool csv_fill_str_col(csv_strref_t* refs, ray_t* vec, int64_t n_rows,
bool* saw_null_out) {
{
ray_str_t* dst = (ray_str_t*)ray_data(vec);
uint64_t pool_bytes = 0;
for (int64_t r = 0; r < n_rows; r++) {
if (RAY_UNLIKELY((r & 1023) == 0 && ray_interrupted()))
return false;
if (refs[r].ptr == NULL) continue;
uint32_t l = refs[r].len;
if (l > RAY_STR_INLINE_MAX) pool_bytes += l;
}
if (pool_bytes > UINT32_MAX) return false;
if (pool_bytes > 0) {
ray_t* pool = ray_alloc((size_t)pool_bytes);
if (!pool || RAY_IS_ERR(pool)) return false;
pool->type = RAY_U8;
pool->len = 0;
vec->str_pool = pool;
}
char* pool_base = vec->str_pool ? (char*)ray_data(vec->str_pool) : NULL;
uint32_t pool_off = 0;
bool saw_null = false;
for (int64_t r = 0; r < n_rows; r++) {
if (RAY_UNLIKELY((r & 1023) == 0 && ray_interrupted()))
return false;
memset(&dst[r], 0, sizeof(ray_str_t));
if (refs[r].ptr == NULL || refs[r].len == 0) { saw_null = true; }
if (refs[r].ptr == NULL) continue;
const char* p = refs[r].ptr;
uint32_t l = refs[r].len;
dst[r].len = l;
if (l <= RAY_STR_INLINE_MAX) {
if (l > 0) memcpy(dst[r].data, p, l);
} else {
memcpy(dst[r].prefix, p, 4);
dst[r].pool_off = pool_off;
memcpy(pool_base + pool_off, p, l);
ray_str_t_cache_hash(&dst[r], pool_base);
pool_off += l;
}
}
if (vec->str_pool) vec->str_pool->len = (int64_t)pool_off;
if (saw_null_out) *saw_null_out = saw_null;
}
return true;
}
typedef struct {
csv_strref_t** str_refs;
int n_cols;
const csv_type_t* parse_types;
const int8_t* resolved_types;
void** col_data;
ray_t** col_vecs;
int64_t n_rows;
int64_t* sym_max_ids;
const int* fill_cols;
int n_fill;
bool* fill_ok;
csv_dedup_ctx_t dd;
int n_sym;
bool* empties;
bool intern_ok;
struct ray_sym_domain_s* dom;
} csv_finalize_ctx_t;
static void csv_finalize_task(void* arg, uint32_t worker_id,
int64_t start, int64_t end_idx) {
csv_finalize_ctx_t* ctx = (csv_finalize_ctx_t*)arg;
if (start < (int64_t)ctx->n_fill) {
int c = ctx->fill_cols[start];
bool saw_null = false;
ctx->fill_ok[start] = csv_fill_str_col(ctx->str_refs[c],
ctx->col_vecs[c], ctx->n_rows,
&saw_null);
csv_note_empty(ctx->empties, c, saw_null);
} else {
csv_dedup_task(&ctx->dd, worker_id, start - (int64_t)ctx->n_fill, end_idx);
}
}
static bool csv_finalize_run(csv_finalize_ctx_t* ctx, int* fill_cols,
bool* fill_ok, int* sym_cols,
bool* empties,
uint64_t prog_base, uint64_t prog_len) {
int n_fill = 0, n_sym = 0;
for (int c = 0; c < ctx->n_cols; c++) {
if (ctx->parse_types[c] != CSV_TYPE_STR || !ctx->str_refs[c]) continue;
if (ctx->resolved_types[c] == RAY_STR) fill_cols[n_fill++] = c;
else sym_cols[n_sym++] = c;
}
for (int i = 0; i < n_fill; i++) fill_ok[i] = true;
ray_pool_t* pool = ray_pool_get();
bool par = pool && ray_pool_total_workers(pool) >= 2;
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;
ctx->fill_ok = fill_ok;
ctx->n_sym = n_sym;
ctx->empties = empties;
ctx->intern_ok = true;
ctx->dd.str_refs = ctx->str_refs;
ctx->dd.col_data = ctx->col_data;
ctx->dd.cols = sym_cols;
ctx->dd.dicts = dicts;
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;
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 * 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);
else for (int64_t i = 0; i < n_tasks; i++) csv_finalize_task(ctx, 0, i, i + 1);
if (ray_interrupted()) goto fail;
if (prog_len) ray_progress_span_phase("intern", prog_base + w1, w2);
ctx->intern_ok = csv_intern_dicts(&ctx->dd, n_sym, ctx->sym_max_ids,
prog_base + w1, w2);
if (!ctx->intern_ok || ray_interrupted()) goto fail;
if (n_sym > 0) {
if (prog_len)
ray_progress_span_phase("sym ids", prog_base + w1 + w2,
prog_len - w1 - w2);
if (par) ray_pool_dispatch_n(pool, csv_dedup_map_task, &ctx->dd,
(uint32_t)n_sym);
else for (int64_t i = 0; i < n_sym; i++)
csv_dedup_map_task(&ctx->dd, 0, i, i + 1);
if (ray_interrupted()) goto fail;
}
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 * n_part; i++) csv_dedup_release(&dicts[i]);
scratch_free(dicts_hdr);
return false;
}
typedef struct {
const char* buf;
size_t buf_size;
const int64_t* row_offsets;
int64_t n_rows;
int n_cols;
char delim;
const csv_type_t* col_types;
const int8_t* resolved_types;
void** col_data;
csv_strref_t** str_refs;
bool* worker_had_null;
bool* worker_had_escaped;
uint8_t* row_done;
} csv_par_ctx_t;
static void csv_parse_fn(void* arg, uint32_t worker_id,
int64_t start, int64_t end_row) {
csv_par_ctx_t* ctx = (csv_par_ctx_t*)arg;
char esc_buf[8192];
const char* buf_end = ctx->buf + ctx->buf_size;
bool* my_had_null = &ctx->worker_had_null[(size_t)worker_id * (size_t)ctx->n_cols];
bool* my_had_escaped = &ctx->worker_had_escaped[(size_t)worker_id * (size_t)ctx->n_cols];
for (int64_t row = start; row < end_row; row++) {
if (RAY_UNLIKELY(((row - start) & 1023) == 0 && ray_interrupted()))
return;
const char* p = ctx->buf + ctx->row_offsets[row];
const char* row_end = (row + 1 < ctx->n_rows)
? ctx->buf + ctx->row_offsets[row + 1]
: buf_end;
for (int c = 0; c < ctx->n_cols; c++) {
if (p >= row_end) {
for (; c < ctx->n_cols; c++) {
switch (ctx->col_types[c]) {
case CSV_TYPE_BOOL: ((uint8_t*)ctx->col_data[c])[row] = 0; break;
case CSV_TYPE_U8: ((uint8_t*)ctx->col_data[c])[row] = 0; break;
case CSV_TYPE_I16: ((int16_t*)ctx->col_data[c])[row] = NULL_I16; break;
case CSV_TYPE_I32: ((int32_t*)ctx->col_data[c])[row] = NULL_I32; break;
case CSV_TYPE_I64: ((int64_t*)ctx->col_data[c])[row] = NULL_I64; break;
case CSV_TYPE_F32: ((float*)ctx->col_data[c])[row] = NULL_F32; break;
case CSV_TYPE_F64: ((double*)ctx->col_data[c])[row] = NULL_F64; break;
case CSV_TYPE_DATE: ((int32_t*)ctx->col_data[c])[row] = NULL_I32; break;
case CSV_TYPE_TIME: ((int32_t*)ctx->col_data[c])[row] = NULL_I32; break;
case CSV_TYPE_TIMESTAMP:
((int64_t*)ctx->col_data[c])[row] = NULL_I64; break;
case CSV_TYPE_GUID:
memset((uint8_t*)ctx->col_data[c] + (size_t)row * 16, 0, 16);
break;
case CSV_TYPE_STR:
ctx->str_refs[c][row].ptr = NULL;
ctx->str_refs[c][row].len = 0;
break;
default: break;
}
if (ctx->col_types[c] != CSV_TYPE_BOOL &&
ctx->col_types[c] != CSV_TYPE_U8) {
my_had_null[c] = true;
}
}
break;
}
const char* fld;
size_t flen;
char* dyn_esc = NULL;
p = scan_field(p, buf_end, ctx->delim, &fld, &flen, esc_buf, &dyn_esc);
if (c == ctx->n_cols - 1 && flen > 0 && fld[flen - 1] == '\r')
flen--;
switch (ctx->col_types[c]) {
case CSV_TYPE_BOOL: {
bool is_null;
uint8_t v = fast_bool(fld, flen, &is_null);
((uint8_t*)ctx->col_data[c])[row] = v;
break;
}
case CSV_TYPE_I64: {
bool is_null;
int64_t v = fast_i64(fld, flen, &is_null);
((int64_t*)ctx->col_data[c])[row] = is_null ? NULL_I64 : v;
if (is_null) my_had_null[c] = true;
break;
}
case CSV_TYPE_U8: {
bool is_null;
int64_t v = fast_i64(fld, flen, &is_null);
((uint8_t*)ctx->col_data[c])[row] = (uint8_t)v;
break;
}
case CSV_TYPE_I16: {
bool is_null;
int64_t v = fast_i64(fld, flen, &is_null);
((int16_t*)ctx->col_data[c])[row] = is_null ? NULL_I16 : (int16_t)v;
if (is_null) my_had_null[c] = true;
break;
}
case CSV_TYPE_I32: {
bool is_null;
int64_t v = fast_i64(fld, flen, &is_null);
((int32_t*)ctx->col_data[c])[row] = is_null ? NULL_I32 : (int32_t)v;
if (is_null) my_had_null[c] = true;
break;
}
case CSV_TYPE_F64: {
bool is_null;
double v = fast_f64(fld, flen, &is_null);
((double*)ctx->col_data[c])[row] = is_null ? NULL_F64 : v;
if (is_null) my_had_null[c] = true;
break;
}
case CSV_TYPE_F32: {
bool is_null;
double v = fast_f64(fld, flen, &is_null);
((float*)ctx->col_data[c])[row] = is_null ? NULL_F32 : (float)v;
if (is_null) my_had_null[c] = true;
break;
}
case CSV_TYPE_DATE: {
bool is_null;
int32_t v = fast_date(fld, flen, &is_null);
((int32_t*)ctx->col_data[c])[row] = is_null ? NULL_I32 : v;
if (is_null) my_had_null[c] = true;
break;
}
case CSV_TYPE_TIME: {
bool is_null;
int32_t v = fast_time(fld, flen, &is_null);
((int32_t*)ctx->col_data[c])[row] = is_null ? NULL_I32 : v;
if (is_null) my_had_null[c] = true;
break;
}
case CSV_TYPE_TIMESTAMP: {
bool is_null;
int64_t v = fast_timestamp(fld, flen, &is_null);
((int64_t*)ctx->col_data[c])[row] = is_null ? NULL_I64 : v;
if (is_null) my_had_null[c] = true;
break;
}
case CSV_TYPE_GUID: {
bool is_null;
uint8_t* slot = (uint8_t*)ctx->col_data[c] + (size_t)row * 16;
fast_guid(fld, flen, slot, &is_null);
if (is_null) {
memset(slot, 0, 16);
my_had_null[c] = true;
}
break;
}
case CSV_TYPE_STR: {
if (flen == 0) {
ctx->str_refs[c][row].ptr = NULL;
ctx->str_refs[c][row].len = 0;
my_had_null[c] = true;
} else {
if (fld < ctx->buf || fld >= buf_end) {
my_had_escaped[c] = true;
if (dyn_esc && fld == dyn_esc) {
dyn_esc = NULL;
} else {
char* cp = (char*)ray_sys_alloc(flen);
if (cp) { memcpy(cp, fld, flen); fld = cp; }
}
}
ctx->str_refs[c][row].ptr = fld;
ctx->str_refs[c][row].len = (uint32_t)flen;
if (ctx->resolved_types[c] == RAY_SYM)
ctx->str_refs[c][row].hash = (uint32_t)ray_hash_bytes(fld, flen);
}
break;
}
default:
break;
}
if (RAY_UNLIKELY(dyn_esc != NULL)) ray_sys_free(dyn_esc);
}
if (ctx->row_done) ctx->row_done[row] = 1;
}
}
static bool csv_parse_serial(const char* buf, size_t buf_size,
const int64_t* row_offsets, int64_t n_rows,
int n_cols, char delim,
const csv_type_t* col_types,
const int8_t* resolved_types,
void** col_data,
csv_strref_t** str_refs,
bool* col_had_null,
bool* col_had_escaped,
uint8_t* row_done) {
char esc_buf[8192];
const char* buf_end = buf + buf_size;
for (int64_t row = 0; row < n_rows; row++) {
if (RAY_UNLIKELY((row & 1023) == 0 && ray_interrupted()))
return false;
const char* p = buf + row_offsets[row];
const char* row_end = (row + 1 < n_rows)
? buf + row_offsets[row + 1]
: buf_end;
for (int c = 0; c < n_cols; c++) {
if (p >= row_end) {
for (; c < n_cols; c++) {
switch (col_types[c]) {
case CSV_TYPE_BOOL: ((uint8_t*)col_data[c])[row] = 0; break;
case CSV_TYPE_U8: ((uint8_t*)col_data[c])[row] = 0; break;
case CSV_TYPE_I16: ((int16_t*)col_data[c])[row] = NULL_I16; break;
case CSV_TYPE_I32: ((int32_t*)col_data[c])[row] = NULL_I32; break;
case CSV_TYPE_I64: ((int64_t*)col_data[c])[row] = NULL_I64; break;
case CSV_TYPE_F32: ((float*)col_data[c])[row] = NULL_F32; break;
case CSV_TYPE_F64: ((double*)col_data[c])[row] = NULL_F64; break;
case CSV_TYPE_DATE: ((int32_t*)col_data[c])[row] = NULL_I32; break;
case CSV_TYPE_TIME: ((int32_t*)col_data[c])[row] = NULL_I32; break;
case CSV_TYPE_TIMESTAMP:
((int64_t*)col_data[c])[row] = NULL_I64; break;
case CSV_TYPE_GUID:
memset((uint8_t*)col_data[c] + (size_t)row * 16, 0, 16);
break;
case CSV_TYPE_STR:
str_refs[c][row].ptr = NULL;
str_refs[c][row].len = 0;
break;
default: break;
}
if (col_types[c] != CSV_TYPE_BOOL &&
col_types[c] != CSV_TYPE_U8) {
col_had_null[c] = true;
}
}
break;
}
const char* fld;
size_t flen;
char* dyn_esc = NULL;
p = scan_field(p, buf_end, delim, &fld, &flen, esc_buf, &dyn_esc);
if (c == n_cols - 1 && flen > 0 && fld[flen - 1] == '\r')
flen--;
switch (col_types[c]) {
case CSV_TYPE_BOOL: {
bool is_null;
uint8_t v = fast_bool(fld, flen, &is_null);
((uint8_t*)col_data[c])[row] = v;
break;
}
case CSV_TYPE_I64: {
bool is_null;
int64_t v = fast_i64(fld, flen, &is_null);
((int64_t*)col_data[c])[row] = is_null ? NULL_I64 : v;
if (is_null) col_had_null[c] = true;
break;
}
case CSV_TYPE_U8: {
bool is_null;
int64_t v = fast_i64(fld, flen, &is_null);
((uint8_t*)col_data[c])[row] = (uint8_t)v;
break;
}
case CSV_TYPE_I16: {
bool is_null;
int64_t v = fast_i64(fld, flen, &is_null);
((int16_t*)col_data[c])[row] = is_null ? NULL_I16 : (int16_t)v;
if (is_null) col_had_null[c] = true;
break;
}
case CSV_TYPE_I32: {
bool is_null;
int64_t v = fast_i64(fld, flen, &is_null);
((int32_t*)col_data[c])[row] = is_null ? NULL_I32 : (int32_t)v;
if (is_null) col_had_null[c] = true;
break;
}
case CSV_TYPE_F64: {
bool is_null;
double v = fast_f64(fld, flen, &is_null);
((double*)col_data[c])[row] = is_null ? NULL_F64 : v;
if (is_null) col_had_null[c] = true;
break;
}
case CSV_TYPE_F32: {
bool is_null;
double v = fast_f64(fld, flen, &is_null);
((float*)col_data[c])[row] = is_null ? NULL_F32 : (float)v;
if (is_null) col_had_null[c] = true;
break;
}
case CSV_TYPE_DATE: {
bool is_null;
int32_t v = fast_date(fld, flen, &is_null);
((int32_t*)col_data[c])[row] = is_null ? NULL_I32 : v;
if (is_null) col_had_null[c] = true;
break;
}
case CSV_TYPE_TIME: {
bool is_null;
int32_t v = fast_time(fld, flen, &is_null);
((int32_t*)col_data[c])[row] = is_null ? NULL_I32 : v;
if (is_null) col_had_null[c] = true;
break;
}
case CSV_TYPE_TIMESTAMP: {
bool is_null;
int64_t v = fast_timestamp(fld, flen, &is_null);
((int64_t*)col_data[c])[row] = is_null ? NULL_I64 : v;
if (is_null) col_had_null[c] = true;
break;
}
case CSV_TYPE_GUID: {
bool is_null;
uint8_t* slot = (uint8_t*)col_data[c] + (size_t)row * 16;
fast_guid(fld, flen, slot, &is_null);
if (is_null) {
memset(slot, 0, 16);
col_had_null[c] = true;
}
break;
}
case CSV_TYPE_STR: {
if (flen == 0) {
str_refs[c][row].ptr = NULL;
str_refs[c][row].len = 0;
col_had_null[c] = true;
} else {
if (fld < buf || fld >= buf_end) {
col_had_escaped[c] = true;
if (dyn_esc && fld == dyn_esc) {
dyn_esc = NULL;
} else {
char* cp = (char*)ray_sys_alloc(flen);
if (cp) { memcpy(cp, fld, flen); fld = cp; }
}
}
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;
}
default:
break;
}
if (RAY_UNLIKELY(dyn_esc != NULL)) ray_sys_free(dyn_esc);
}
if (row_done) row_done[row] = 1;
}
return true;
}
static int csv_hash_elem_size(int8_t t) {
switch (t) {
case RAY_BOOL: case RAY_U8: return 1;
case RAY_I16: return 2;
case RAY_I32: case RAY_DATE: return 4;
case RAY_I64: case RAY_TIME: case RAY_TIMESTAMP: return 8;
default: return 0;
}
}
int ray_csv_hash_upgrade_check(int8_t type, int64_t len,
const void* index_payload) {
const ray_index_t* ix = (const ray_index_t*)index_payload;
int esz = csv_hash_elem_size(type);
if (esz == 0) return 0;
if (!ix || ix->kind != RAY_IDX_CHUNK_ZONE || ix->u.chunk_zone.is_f64)
return 0;
uint32_t n_chunks = ix->u.chunk_zone.n_chunks;
if (n_chunks < 4) return 0;
const int64_t* mins = (const int64_t*)ray_data(ix->u.chunk_zone.mins);
const int64_t* maxs = (const int64_t*)ray_data(ix->u.chunk_zone.maxs);
int64_t gmin = INT64_MAX, gmax = INT64_MIN;
for (uint32_t g = 0; g < n_chunks; g++) {
if (mins[g] > maxs[g]) continue;
if (mins[g] < gmin) gmin = mins[g];
if (maxs[g] > gmax) gmax = maxs[g];
}
if (gmin == INT64_MAX || gmax == INT64_MIN) return 0;
uint64_t global_range = (uint64_t)gmax - (uint64_t)gmin;
if (global_range == 0) return 0;
double dgr = (double)global_range;
double span_sum = 0.0;
uint32_t n_eff = 0;
for (uint32_t g = 0; g < n_chunks; g++) {
if (mins[g] > maxs[g]) continue;
uint64_t span = (uint64_t)maxs[g] - (uint64_t)mins[g];
span_sum += (double)span;
n_eff++;
}
if (n_eff < 4) return 0;
double mean_ratio = (span_sum / (double)n_eff) / dgr;
if (mean_ratio <= 0.5) return 0;
int64_t n = len;
if (n <= 0) return 0;
uint64_t cap = 8;
uint64_t want = (uint64_t)(2 * n);
while (cap < want) cap <<= 1;
uint64_t aux_bytes = cap * 8u + (uint64_t)n * 8u;
uint64_t data_bytes = (uint64_t)n * (uint64_t)esz;
if (aux_bytes > 5u * data_bytes) return 0;
return 1;
}
static int8_t csv_auto_width_to_ray(csv_type_t w) {
switch (w) {
case CSV_TYPE_I16: return RAY_I16;
case CSV_TYPE_I32: return RAY_I32;
default: return RAY_I64;
}
}
static bool csv_auto_scan_rows(const char* buf, const char* buf_end,
const int64_t* row_offsets, int64_t n_rows,
int ncols, char delim,
const int8_t* resolved_types,
int64_t* col_min, int64_t* col_max,
bool* col_had_null) {
char esc_buf[8192];
for (int64_t row = 0; row < n_rows; row++) {
if (RAY_UNLIKELY((row & 1023) == 0 && ray_interrupted()))
return false;
const char* p = buf + row_offsets[row];
const char* row_end = (row + 1 < n_rows)
? buf + row_offsets[row + 1] : buf_end;
for (int c = 0; c < ncols; c++) {
if (p >= row_end) {
for (; c < ncols; c++)
if (resolved_types[c] == RAY_CSV_AUTO_TAG)
col_had_null[c] = true;
break;
}
const char* fld;
size_t flen;
char* dyn_esc = NULL;
p = scan_field(p, buf_end, delim, &fld, &flen, esc_buf, &dyn_esc);
if (c == ncols - 1 && flen > 0 && fld[flen - 1] == '\r') flen--;
if (resolved_types[c] == RAY_CSV_AUTO_TAG) {
bool is_null;
int64_t v = fast_i64(fld, flen, &is_null);
if (is_null) {
col_had_null[c] = true;
} else {
if (v < col_min[c]) col_min[c] = v;
if (v > col_max[c]) col_max[c] = v;
}
}
if (RAY_UNLIKELY(dyn_esc != NULL)) ray_sys_free(dyn_esc);
}
}
return true;
}
static bool csv_resolve_auto_in_place(const char* buf, size_t file_size,
const int64_t* row_offsets, int64_t n_rows,
int ncols, char delim,
int8_t* resolved_types) {
bool any = false;
for (int c = 0; c < ncols; c++)
if (resolved_types[c] == RAY_CSV_AUTO_TAG) any = true;
if (!any) return true;
int64_t col_min[CSV_MAX_COLS], col_max[CSV_MAX_COLS];
bool col_had_null[CSV_MAX_COLS];
for (int c = 0; c < ncols; c++) {
col_min[c] = INT64_MAX; col_max[c] = INT64_MIN; col_had_null[c] = false;
}
if (!csv_auto_scan_rows(buf, buf + file_size, row_offsets, n_rows, ncols,
delim, resolved_types, col_min, col_max,
col_had_null))
return false;
for (int c = 0; c < ncols; c++)
if (resolved_types[c] == RAY_CSV_AUTO_TAG)
resolved_types[c] = csv_auto_width_to_ray(
csv_resolve_int_width(col_min[c], col_max[c], col_had_null[c]));
return true;
}
static bool csv_resolve_auto_streamed(const char* buf, size_t file_size,
size_t data_offset, int ncols, char delim,
bool data_has_quotes,
int8_t* resolved_types) {
bool any = false;
for (int c = 0; c < ncols; c++)
if (resolved_types[c] == RAY_CSV_AUTO_TAG) any = true;
if (!any) return true;
int64_t col_min[CSV_MAX_COLS], col_max[CSV_MAX_COLS];
bool col_had_null[CSV_MAX_COLS];
for (int c = 0; c < ncols; c++) {
col_min[c] = INT64_MAX; col_max[c] = INT64_MIN; col_had_null[c] = false;
}
size_t off = data_offset;
while (off < file_size) {
if (ray_interrupted()) return false;
ray_t* hdr = NULL;
int64_t* roff = NULL;
size_t next = off;
int64_t cnt = build_row_offsets_limited(buf, file_size, off,
CSV_PART_ROWS_DEFAULT,
data_has_quotes,
&roff, &hdr, &next);
if (cnt < 0) { scratch_free(hdr); return false; }
if (cnt == 0) { scratch_free(hdr); break; }
bool scan_ok = csv_auto_scan_rows(buf, buf + file_size, roff, cnt,
ncols, delim, resolved_types,
col_min, col_max, col_had_null);
scratch_free(hdr);
if (!scan_ok) return false;
if (next <= off) break;
off = next;
}
for (int c = 0; c < ncols; c++)
if (resolved_types[c] == RAY_CSV_AUTO_TAG)
resolved_types[c] = csv_auto_width_to_ray(
csv_resolve_int_width(col_min[c], col_max[c], col_had_null[c]));
return true;
}
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,
struct ray_sym_domain_s* sym_dom) {
for (int c = 0; c < ncols; c++) {
if (resolved_types[c] == RAY_CSV_AUTO_TAG) {
fprintf(stderr,
"csv: BUG: unresolved INT (RAY_CSV_AUTO_TAG=%d) reached "
"csv_materialize_rows for column %d\n",
(int)RAY_CSV_AUTO_TAG, c);
return NULL;
}
}
ray_t* col_vecs[CSV_MAX_COLS];
void* col_data[CSV_MAX_COLS];
for (int c = 0; c < ncols; c++) {
int8_t type = resolved_types[c];
col_vecs[c] = (type == RAY_SYM) ? ray_sym_vec_new(RAY_SYM_W32, n_rows)
: ray_vec_new(type, n_rows);
if (!col_vecs[c] || RAY_IS_ERR(col_vecs[c])) {
for (int j = 0; j < c; j++) ray_release(col_vecs[j]);
return NULL;
}
if (type == RAY_SYM && sym_dom) {
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]);
}
bool col_had_null[CSV_MAX_COLS];
if (ncols > 0) memset(col_had_null, 0, (size_t)ncols * sizeof(bool));
bool col_wrote_null[CSV_MAX_COLS];
if (ncols > 0) memset(col_wrote_null, 0, (size_t)ncols * sizeof(bool));
bool col_had_escaped[CSV_MAX_COLS];
memset(col_had_escaped, 0, sizeof(col_had_escaped));
csv_type_t parse_types[CSV_MAX_COLS];
for (int c = 0; c < ncols; c++) {
switch (resolved_types[c]) {
case RAY_BOOL: parse_types[c] = CSV_TYPE_BOOL; break;
case RAY_U8: parse_types[c] = CSV_TYPE_U8; break;
case RAY_I16: parse_types[c] = CSV_TYPE_I16; break;
case RAY_I32: parse_types[c] = CSV_TYPE_I32; break;
case RAY_I64: parse_types[c] = CSV_TYPE_I64; break;
case RAY_F32: parse_types[c] = CSV_TYPE_F32; break;
case RAY_F64: parse_types[c] = CSV_TYPE_F64; break;
case RAY_DATE: parse_types[c] = CSV_TYPE_DATE; break;
case RAY_TIME: parse_types[c] = CSV_TYPE_TIME; break;
case RAY_TIMESTAMP: parse_types[c] = CSV_TYPE_TIMESTAMP; break;
case RAY_GUID: parse_types[c] = CSV_TYPE_GUID; break;
default: parse_types[c] = CSV_TYPE_STR; break;
}
}
int64_t sym_max_ids[CSV_MAX_COLS];
memset(sym_max_ids, 0, (size_t)ncols * sizeof(int64_t));
int has_text_cols = 0;
for (int c = 0; c < ncols; c++) {
if (parse_types[c] == CSV_TYPE_STR) {
has_text_cols = 1;
break;
}
}
csv_strref_t* str_ref_bufs[CSV_MAX_COLS];
ray_t* str_ref_hdrs[CSV_MAX_COLS];
memset(str_ref_bufs, 0, sizeof(str_ref_bufs));
memset(str_ref_hdrs, 0, sizeof(str_ref_hdrs));
ray_t* row_done_hdr = NULL;
uint8_t* row_done = has_text_cols
? (uint8_t*)scratch_calloc(&row_done_hdr, (size_t)n_rows)
: NULL;
if (has_text_cols && !row_done) {
for (int c = 0; c < ncols; c++) ray_release(col_vecs[c]);
return NULL;
}
for (int c = 0; c < ncols; c++) {
if (parse_types[c] == CSV_TYPE_STR) {
size_t sz = (size_t)n_rows * sizeof(csv_strref_t);
str_ref_bufs[c] = (csv_strref_t*)scratch_alloc(&str_ref_hdrs[c], sz);
if (!str_ref_bufs[c]) {
for (int j = 0; j < ncols; j++) ray_release(col_vecs[j]);
for (int j = 0; j < c; j++) scratch_free(str_ref_hdrs[j]);
scratch_free(row_done_hdr);
return NULL;
}
}
}
{
ray_pool_t* pool = ray_pool_get();
bool use_parallel = pool && n_rows > 8192;
if (use_parallel) {
uint32_t n_workers = ray_pool_total_workers(pool);
size_t worker_flags_sz = (size_t)n_workers * (size_t)ncols * sizeof(bool);
bool* worker_flags = (bool*)ray_sys_alloc(worker_flags_sz * 2);
if (!worker_flags) {
use_parallel = false;
} else {
memset(worker_flags, 0, worker_flags_sz * 2);
bool* worker_had_null_buf = worker_flags;
bool* worker_had_escaped_buf = worker_flags +
(size_t)n_workers * (size_t)ncols;
csv_par_ctx_t ctx = {
.buf = buf,
.buf_size = file_size,
.row_offsets = row_offsets,
.n_rows = n_rows,
.n_cols = ncols,
.delim = delimiter,
.col_types = parse_types,
.resolved_types = resolved_types,
.col_data = col_data,
.str_refs = str_ref_bufs,
.worker_had_null = worker_had_null_buf,
.worker_had_escaped = worker_had_escaped_buf,
.row_done = row_done,
};
ray_pool_dispatch(pool, csv_parse_fn, &ctx, n_rows);
for (uint32_t w = 0; w < n_workers; w++) {
for (int c = 0; c < ncols; c++) {
if (worker_had_null_buf[(size_t)w * (size_t)ncols + (size_t)c])
col_had_null[c] = true;
if (worker_had_escaped_buf[(size_t)w * (size_t)ncols + (size_t)c])
col_had_escaped[c] = true;
}
}
ray_sys_free(worker_flags);
}
}
if (!use_parallel) {
(void)csv_parse_serial(buf, file_size, row_offsets, n_rows,
ncols, delimiter, parse_types, resolved_types,
col_data, str_ref_bufs, col_had_null,
col_had_escaped, row_done);
}
}
if (ray_interrupted()) {
csv_free_escaped_strrefs(str_ref_bufs, ncols, parse_types, n_rows,
buf, file_size, row_done, col_had_escaped);
for (int c = 0; c < ncols; c++) scratch_free(str_ref_hdrs[c]);
scratch_free(row_done_hdr);
for (int c = 0; c < ncols; c++) ray_release(col_vecs[c]);
return ray_error("cancel", "interrupted");
}
if (has_text_cols) {
csv_finalize_ctx_t fctx = {
.str_refs = str_ref_bufs,
.n_cols = ncols,
.parse_types = parse_types,
.resolved_types = resolved_types,
.col_data = col_data,
.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];
bool fin_ok = csv_finalize_run(&fctx, fill_cols, fill_ok,
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);
for (int c = 0; c < ncols; c++) scratch_free(str_ref_hdrs[c]);
scratch_free(row_done_hdr);
for (int c = 0; c < ncols; c++) ray_release(col_vecs[c]);
return ray_interrupted() ? ray_error("cancel", "interrupted") : NULL;
}
}
csv_free_escaped_strrefs(str_ref_bufs, ncols, parse_types, n_rows,
buf, file_size, NULL, col_had_escaped);
for (int c = 0; c < ncols; c++) scratch_free(str_ref_hdrs[c]);
scratch_free(row_done_hdr);
if (ray_interrupted()) {
for (int c = 0; c < ncols; c++) ray_release(col_vecs[c]);
return ray_error("cancel", "interrupted");
}
for (int c = 0; c < ncols; c++) {
ray_t* vec = col_vecs[c];
if (col_had_null[c] || col_wrote_null[c])
vec->attrs |= RAY_ATTR_HAS_NULLS;
else
vec->attrs &= (uint8_t)~RAY_ATTR_HAS_NULLS;
}
for (int c = 0; c < ncols; c++) {
if (resolved_types[c] != RAY_SYM) continue;
uint8_t new_w = ray_sym_dict_width(sym_max_ids[c]);
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);
if (new_w == RAY_SYM_W8) {
uint8_t* d = (uint8_t*)dst;
for (int64_t r = 0; r < n_rows; r++) {
if (RAY_UNLIKELY((r & 1023) == 0 && ray_interrupted())) {
ray_release(narrow);
for (int j = 0; j < ncols; j++) ray_release(col_vecs[j]);
return ray_error("cancel", "interrupted");
}
d[r] = (uint8_t)src[r];
}
} else {
uint16_t* d = (uint16_t*)dst;
for (int64_t r = 0; r < n_rows; r++) {
if (RAY_UNLIKELY((r & 1023) == 0 && ray_interrupted())) {
ray_release(narrow);
for (int j = 0; j < ncols; j++) ray_release(col_vecs[j]);
return ray_error("cancel", "interrupted");
}
d[r] = (uint16_t)src[r];
}
}
narrow->attrs |= (col_vecs[c]->attrs & RAY_ATTR_HAS_NULLS);
ray_release(col_vecs[c]);
col_vecs[c] = narrow;
col_data[c] = dst;
}
ray_t* tbl = ray_table_new(ncols);
if (!tbl || RAY_IS_ERR(tbl)) {
for (int c = 0; c < ncols; c++) ray_release(col_vecs[c]);
return NULL;
}
for (int c = 0; c < ncols; c++) {
tbl = ray_table_add_col(tbl, col_name_ids[c], col_vecs[c]);
ray_release(col_vecs[c]);
}
return tbl;
}
static const char* csv_skip_matching_header(const char* p, const char* buf_end,
char delimiter, int ncols,
const int64_t* names, char* esc_buf) {
bool is_header = true;
for (int c = 0; c < ncols; c++) {
const char* fld; size_t flen; char* dyn = NULL;
p = scan_field(p, buf_end, delimiter, &fld, &flen, esc_buf, &dyn);
if (is_header) {
ray_t* nm = ray_sym_str(names[c]);
const char* ns = nm ? ray_str_ptr(nm) : NULL;
size_t nl = nm ? ray_str_len(nm) : 0;
if (!ns || flen != nl || memcmp(fld, ns, flen) != 0) is_header = false;
}
if (dyn) ray_sys_free(dyn);
}
if (!is_header) return NULL;
if (p < buf_end && *p == '\r') p++;
if (p < buf_end && *p == '\n') p++;
return p;
}
static ray_t* csv_read_named_opts_inner(const char* path, char delimiter, bool header,
const int8_t* col_types_in, int32_t n_types,
const int64_t* col_names_in, int32_t n_names) {
if (ray_interrupted()) return ray_error("cancel", "interrupted");
int fd = open(path, O_RDONLY);
if (fd < 0) return ray_error("io", NULL);
struct stat st;
if (fstat(fd, &st) != 0 || st.st_size <= 0) {
close(fd);
return ray_error("io", NULL);
}
size_t file_size = (size_t)st.st_size;
char* buf = (char*)ray_vm_map_fd_ro(fd, file_size);
close(fd);
if (!buf) return ray_error("io", NULL);
#ifdef __APPLE__
madvise(buf, file_size, MADV_SEQUENTIAL);
#endif
const char* buf_end = buf + file_size;
ray_t* result = NULL;
ray_progress_span_begin(NULL, (uint64_t)file_size);
if (delimiter == 0) {
int commas = 0, tabs = 0;
for (const char* p = buf; p < buf_end && *p != '\n'; p++) {
if (*p == ',') commas++;
if (*p == '\t') tabs++;
}
delimiter = (tabs > commas) ? '\t' : ',';
}
int ncols = 1;
{
const char* p = buf;
bool in_quote = false;
while (p < buf_end && (in_quote || (*p != '\n' && *p != '\r'))) {
if (*p == '"') in_quote = !in_quote;
else if (!in_quote && *p == delimiter) ncols++;
p++;
}
}
if (ncols > CSV_MAX_COLS) {
ray_vm_unmap_file(buf, file_size);
return ray_error("range", "csv read: header has too many columns, got %lld (max %d)",
(long long)ncols, CSV_MAX_COLS);
}
if (col_types_in && n_types != ncols) {
ray_vm_unmap_file(buf, file_size);
return ray_error("length",
"csv read: schema has %lld types but file has %d columns",
(long long)n_types, ncols);
}
if (col_names_in && n_names != ncols) {
ray_vm_unmap_file(buf, file_size);
return ray_error("length",
"csv read: schema has %lld names but file has %d columns",
(long long)n_names, ncols);
}
const char* p = buf;
char esc_buf[8192];
int64_t col_name_ids[CSV_MAX_COLS];
if (header) {
for (int c = 0; c < ncols; c++) {
const char* fld;
size_t flen;
char* dyn_esc = NULL;
p = scan_field(p, buf_end, delimiter, &fld, &flen, esc_buf, &dyn_esc);
col_name_ids[c] = ray_sym_intern(fld, flen);
if (dyn_esc) ray_sys_free(dyn_esc);
}
if (p < buf_end && *p == '\r') p++;
if (p < buf_end && *p == '\n') p++;
} else if (col_names_in) {
for (int c = 0; c < ncols; c++)
col_name_ids[c] = col_names_in[c];
const char* after = csv_skip_matching_header(buf, buf_end, delimiter,
ncols, col_names_in, esc_buf);
if (after) p = after;
} else {
for (int c = 0; c < ncols; c++) {
char name[32];
snprintf(name, sizeof(name), "V%d", c + 1);
col_name_ids[c] = ray_sym_intern(name, strlen(name));
}
}
size_t data_offset = (size_t)(p - buf);
ray_t* row_offsets_hdr = NULL;
int64_t* row_offsets = NULL;
uint64_t prog_scan_end = csv_prog_at(file_size, CSV_PROG_PCT_SCAN);
uint64_t prog_parse_end = csv_prog_at(file_size,
CSV_PROG_PCT_SCAN + CSV_PROG_PCT_PARSE);
int64_t n_rows = build_row_offsets(buf, file_size, data_offset,
0, prog_scan_end,
&row_offsets, &row_offsets_hdr);
if (n_rows < 0) {
ray_vm_unmap_file(buf, file_size);
return ray_error("cancel", "interrupted");
}
if (n_rows == 0) {
ray_t* tbl = ray_table_new(ncols);
if (!tbl || RAY_IS_ERR(tbl)) goto fail_unmap;
for (int c = 0; c < ncols; c++) {
ray_t* empty_vec = ray_vec_new(RAY_F64, 0);
if (empty_vec && !RAY_IS_ERR(empty_vec)) {
tbl = ray_table_add_col(tbl, col_name_ids[c], empty_vec);
ray_release(empty_vec);
}
}
ray_vm_unmap_file(buf, file_size);
return tbl;
}
int8_t resolved_types[CSV_MAX_COLS];
if (col_types_in) {
for (int c = 0; c < ncols; c++) {
int8_t t = col_types_in[c];
if (t < RAY_BOOL ||
(t >= RAY_TYPE_COUNT && t != RAY_CSV_AUTO_TAG) ||
t == RAY_TABLE) {
scratch_free(row_offsets_hdr);
ray_vm_unmap_file(buf, file_size);
return ray_error("type",
"csv read: unsupported type code %d in column %d",
(int)t, c);
}
resolved_types[c] = t;
}
} else if (!col_types_in) {
if (!csv_infer_types_from_offsets(buf, buf_end, row_offsets, n_rows,
ncols, delimiter, esc_buf,
resolved_types))
goto fail_offsets;
}
if (!csv_resolve_auto_in_place(buf, file_size, row_offsets, n_rows,
ncols, delimiter, resolved_types))
goto fail_offsets_cancel;
ray_t* col_vecs[CSV_MAX_COLS];
void* col_data[CSV_MAX_COLS];
for (int c = 0; c < ncols; c++) {
int8_t type = resolved_types[c];
col_vecs[c] = (type == RAY_SYM) ? ray_sym_vec_new(RAY_SYM_W32, n_rows)
: ray_vec_new(type, n_rows);
if (!col_vecs[c] || RAY_IS_ERR(col_vecs[c])) {
for (int j = 0; j < c; j++) ray_release(col_vecs[j]);
goto fail_offsets;
}
col_vecs[c]->len = n_rows;
col_data[c] = ray_data(col_vecs[c]);
}
bool col_had_null[CSV_MAX_COLS];
if (ncols > 0) memset(col_had_null, 0, (size_t)ncols * sizeof(bool));
bool col_wrote_null[CSV_MAX_COLS];
if (ncols > 0) memset(col_wrote_null, 0, (size_t)ncols * sizeof(bool));
bool col_had_escaped[CSV_MAX_COLS];
memset(col_had_escaped, 0, sizeof(col_had_escaped));
csv_type_t parse_types[CSV_MAX_COLS];
for (int c = 0; c < ncols; c++) {
switch (resolved_types[c]) {
case RAY_BOOL: parse_types[c] = CSV_TYPE_BOOL; break;
case RAY_U8: parse_types[c] = CSV_TYPE_U8; break;
case RAY_I16: parse_types[c] = CSV_TYPE_I16; break;
case RAY_I32: parse_types[c] = CSV_TYPE_I32; break;
case RAY_I64: parse_types[c] = CSV_TYPE_I64; break;
case RAY_F32: parse_types[c] = CSV_TYPE_F32; break;
case RAY_F64: parse_types[c] = CSV_TYPE_F64; break;
case RAY_DATE: parse_types[c] = CSV_TYPE_DATE; break;
case RAY_TIME: parse_types[c] = CSV_TYPE_TIME; break;
case RAY_TIMESTAMP: parse_types[c] = CSV_TYPE_TIMESTAMP; break;
case RAY_GUID: parse_types[c] = CSV_TYPE_GUID; break;
default: parse_types[c] = CSV_TYPE_STR; break;
}
}
int64_t sym_max_ids[CSV_MAX_COLS];
memset(sym_max_ids, 0, (size_t)ncols * sizeof(int64_t));
int has_text_cols = 0;
for (int c = 0; c < ncols; c++) {
if (parse_types[c] == CSV_TYPE_STR) {
has_text_cols = 1;
break;
}
}
csv_strref_t* str_ref_bufs[CSV_MAX_COLS];
ray_t* str_ref_hdrs[CSV_MAX_COLS];
memset(str_ref_bufs, 0, sizeof(str_ref_bufs));
memset(str_ref_hdrs, 0, sizeof(str_ref_hdrs));
ray_t* row_done_hdr = NULL;
uint8_t* row_done = has_text_cols
? (uint8_t*)scratch_calloc(&row_done_hdr, (size_t)n_rows)
: NULL;
if (has_text_cols && !row_done) {
for (int c = 0; c < ncols; c++) ray_release(col_vecs[c]);
goto fail_offsets;
}
for (int c = 0; c < ncols; c++) {
if (parse_types[c] == CSV_TYPE_STR) {
size_t sz = (size_t)n_rows * sizeof(csv_strref_t);
str_ref_bufs[c] = (csv_strref_t*)scratch_alloc(&str_ref_hdrs[c], sz);
if (!str_ref_bufs[c]) {
for (int j = 0; j < ncols; j++) ray_release(col_vecs[j]);
for (int j = 0; j < c; j++) scratch_free(str_ref_hdrs[j]);
scratch_free(row_done_hdr);
goto fail_offsets;
}
}
}
ray_progress_span_phase("parse", prog_scan_end, prog_parse_end - prog_scan_end);
{
ray_pool_t* pool = ray_pool_get();
bool use_parallel = pool && n_rows > 8192;
if (use_parallel) {
uint32_t n_workers = ray_pool_total_workers(pool);
size_t worker_flags_sz = (size_t)n_workers * (size_t)ncols * sizeof(bool);
bool* worker_flags = (bool*)ray_sys_alloc(worker_flags_sz * 2);
if (!worker_flags) {
use_parallel = false;
} else {
memset(worker_flags, 0, worker_flags_sz * 2);
bool* worker_had_null_buf = worker_flags;
bool* worker_had_escaped_buf = worker_flags +
(size_t)n_workers * (size_t)ncols;
csv_par_ctx_t ctx = {
.buf = buf,
.buf_size = file_size,
.row_offsets = row_offsets,
.n_rows = n_rows,
.n_cols = ncols,
.delim = delimiter,
.col_types = parse_types,
.resolved_types = resolved_types,
.col_data = col_data,
.str_refs = str_ref_bufs,
.worker_had_null = worker_had_null_buf,
.worker_had_escaped = worker_had_escaped_buf,
.row_done = row_done,
};
ray_pool_dispatch(pool, csv_parse_fn, &ctx, n_rows);
for (uint32_t w = 0; w < n_workers; w++) {
for (int c = 0; c < ncols; c++) {
if (worker_had_null_buf[(size_t)w * (size_t)ncols + (size_t)c])
col_had_null[c] = true;
if (worker_had_escaped_buf[(size_t)w * (size_t)ncols + (size_t)c])
col_had_escaped[c] = true;
}
}
ray_sys_free(worker_flags);
}
}
if (!use_parallel) {
(void)csv_parse_serial(buf, file_size, row_offsets, n_rows,
ncols, delimiter, parse_types, resolved_types,
col_data, str_ref_bufs, col_had_null,
col_had_escaped, row_done);
}
}
if (ray_interrupted()) goto fail_parsed_cancel;
if (has_text_cols) {
csv_finalize_ctx_t fctx = {
.str_refs = str_ref_bufs,
.n_cols = ncols,
.parse_types = parse_types,
.resolved_types = resolved_types,
.col_data = col_data,
.col_vecs = col_vecs,
.n_rows = n_rows,
.sym_max_ids = sym_max_ids,
};
int fill_cols[CSV_MAX_COLS];
int sym_cols[CSV_MAX_COLS];
bool fill_ok[CSV_MAX_COLS];
bool fin_ok = csv_finalize_run(&fctx, fill_cols, fill_ok,
sym_cols, col_wrote_null,
prog_parse_end,
(uint64_t)file_size - prog_parse_end);
if (!fin_ok || ray_interrupted()) {
if (ray_interrupted()) goto fail_parsed_cancel;
csv_free_escaped_strrefs(str_ref_bufs, ncols, parse_types, n_rows,
buf, file_size, NULL, col_had_escaped);
for (int c = 0; c < ncols; c++) scratch_free(str_ref_hdrs[c]);
scratch_free(row_done_hdr);
for (int c = 0; c < ncols; c++) ray_release(col_vecs[c]);
goto fail_offsets;
}
}
csv_free_escaped_strrefs(str_ref_bufs, ncols, parse_types, n_rows,
buf, file_size, NULL, col_had_escaped);
for (int c = 0; c < ncols; c++) scratch_free(str_ref_hdrs[c]);
scratch_free(row_done_hdr);
for (int c = 0; c < ncols; c++) {
ray_t* vec = col_vecs[c];
bool has_null = col_had_null[c] || col_wrote_null[c];
if (has_null)
vec->attrs |= RAY_ATTR_HAS_NULLS;
else
vec->attrs &= (uint8_t)~RAY_ATTR_HAS_NULLS;
}
for (int c = 0; c < ncols; c++) {
if (resolved_types[c] != RAY_SYM) continue;
uint8_t new_w = ray_sym_dict_width(sym_max_ids[c]);
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;
narrow->len = n_rows;
const uint32_t* src = (const uint32_t*)col_data[c];
void* dst = ray_data(narrow);
if (new_w == RAY_SYM_W8) {
uint8_t* d = (uint8_t*)dst;
for (int64_t r = 0; r < n_rows; r++) {
if (RAY_UNLIKELY((r & 1023) == 0 && ray_interrupted())) {
ray_release(narrow);
goto fail_cols_cancel;
}
d[r] = (uint8_t)src[r];
}
} else {
uint16_t* d = (uint16_t*)dst;
for (int64_t r = 0; r < n_rows; r++) {
if (RAY_UNLIKELY((r & 1023) == 0 && ray_interrupted())) {
ray_release(narrow);
goto fail_cols_cancel;
}
d[r] = (uint16_t)src[r];
}
}
narrow->attrs |= (col_vecs[c]->attrs & RAY_ATTR_HAS_NULLS);
ray_release(col_vecs[c]);
col_vecs[c] = narrow;
col_data[c] = dst;
}
{
ray_t* tbl = ray_table_new(ncols);
if (!tbl || RAY_IS_ERR(tbl)) {
for (int c = 0; c < ncols; c++) ray_release(col_vecs[c]);
goto fail_offsets;
}
for (int c = 0; c < ncols; c++) {
tbl = ray_table_add_col(tbl, col_name_ids[c], col_vecs[c]);
ray_release(col_vecs[c]);
}
result = tbl;
}
scratch_free(row_offsets_hdr);
ray_vm_unmap_file(buf, file_size);
return result;
fail_parsed_cancel:
csv_free_escaped_strrefs(str_ref_bufs, ncols, parse_types, n_rows,
buf, file_size, row_done, col_had_escaped);
for (int c = 0; c < ncols; c++) scratch_free(str_ref_hdrs[c]);
scratch_free(row_done_hdr);
fail_cols_cancel:
for (int c = 0; c < ncols; c++) ray_release(col_vecs[c]);
fail_offsets_cancel:
scratch_free(row_offsets_hdr);
ray_vm_unmap_file(buf, file_size);
return ray_error("cancel", "interrupted");
fail_offsets:
scratch_free(row_offsets_hdr);
fail_unmap:
ray_vm_unmap_file(buf, file_size);
return ray_error("oom", NULL);
}
ray_t* ray_read_csv_named_opts(const char* path, char delimiter, bool header,
const int8_t* col_types_in, int32_t n_types,
const int64_t* col_names_in, int32_t n_names) {
ray_t* r = csv_read_named_opts_inner(path, delimiter, header,
col_types_in, n_types,
col_names_in, n_names);
ray_progress_span_end();
return r;
}
ray_t* ray_read_csv_opts(const char* path, char delimiter, bool header,
const int8_t* col_types_in, int32_t n_types) {
return ray_read_csv_named_opts(path, delimiter, header,
col_types_in, n_types, NULL, 0);
}
typedef struct {
ray_col_stream_t* writers;
ray_t* tbl;
int ncols;
_Atomic(ray_err_t) err;
int64_t* col_ns;
} 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);
int64_t tt0 = a->col_ns ? ray_profile_now_ns() : 0;
ray_err_t e = ray_col_stream_append(&a->writers[start], col);
if (a->col_ns) a->col_ns[start] += ray_profile_now_ns() - tt0;
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_save_splayed_to_dir(const char* path, char delimiter, bool header,
const int8_t* col_types_in, int32_t n_types,
const int64_t* col_names_in, int32_t n_names,
const char* dir, int64_t rows_per_chunk,
const char* sym_path) {
if (ray_interrupted()) return RAY_ERR_CANCEL;
if (!path || !dir || !sym_path) return RAY_ERR_DOMAIN;
if (rows_per_chunk <= 0) rows_per_chunk = CSV_PART_ROWS_DEFAULT;
bool trace = getenv("RAY_CSV_TRACE") != NULL;
int64_t tr_t0 = ray_profile_now_ns();
int64_t tr_scan = 0, tr_parse = 0, tr_append = 0, tr_chunks = 0;
int64_t tr_col_ns[CSV_MAX_COLS];
memset(tr_col_ns, 0, sizeof(tr_col_ns));
#define TR_MS(ns) ((double)(ns) / 1e6)
int fd = open(path, O_RDONLY);
if (fd < 0) return RAY_ERR_IO;
struct stat st;
if (fstat(fd, &st) != 0 || st.st_size <= 0) {
close(fd);
return RAY_ERR_IO;
}
size_t file_size = (size_t)st.st_size;
char* buf = (char*)ray_vm_map_fd_ro(fd, file_size);
close(fd);
if (!buf) return RAY_ERR_IO;
#ifdef __APPLE__
madvise(buf, file_size, MADV_SEQUENTIAL);
#endif
const char* buf_end = buf + file_size;
ray_err_t err = RAY_OK;
if (delimiter == 0) {
int commas = 0, tabs = 0;
for (const char* q = buf; q < buf_end && *q != '\n'; q++) {
if (*q == ',') commas++;
if (*q == '\t') tabs++;
}
delimiter = (tabs > commas) ? '\t' : ',';
}
int ncols = 1;
{
const char* q = buf;
bool in_quote = false;
while (q < buf_end && (in_quote || (*q != '\n' && *q != '\r'))) {
if (*q == '"') in_quote = !in_quote;
else if (!in_quote && *q == delimiter) ncols++;
q++;
}
}
if (ncols > CSV_MAX_COLS) {
ray_vm_unmap_file(buf, file_size);
return RAY_ERR_RANGE;
}
if ((col_types_in && n_types != ncols) ||
(col_names_in && n_names != ncols)) {
ray_vm_unmap_file(buf, file_size);
return RAY_ERR_LENGTH;
}
const char* p = buf;
char esc_buf[8192];
int64_t col_name_ids[CSV_MAX_COLS];
if (header) {
for (int c = 0; c < ncols; c++) {
const char* fld;
size_t flen;
char* dyn_esc = NULL;
p = scan_field(p, buf_end, delimiter, &fld, &flen, esc_buf, &dyn_esc);
col_name_ids[c] = ray_sym_intern(fld, flen);
if (dyn_esc) ray_sys_free(dyn_esc);
}
if (p < buf_end && *p == '\r') p++;
if (p < buf_end && *p == '\n') p++;
} else if (col_names_in) {
for (int c = 0; c < ncols; c++)
col_name_ids[c] = col_names_in[c];
const char* after = csv_skip_matching_header(buf, buf_end, delimiter,
ncols, col_names_in, esc_buf);
if (after) p = after;
} else {
for (int c = 0; c < ncols; c++) {
char name[32];
snprintf(name, sizeof(name), "V%d", c + 1);
col_name_ids[c] = ray_sym_intern(name, strlen(name));
}
}
size_t data_offset = (size_t)(p - buf);
bool data_has_quotes = memchr(buf + data_offset, '"', file_size - data_offset) != NULL;
int8_t resolved_types[CSV_MAX_COLS];
if (col_types_in) {
for (int c = 0; c < ncols; c++) {
int8_t t = col_types_in[c];
if (t < RAY_BOOL ||
(t >= RAY_TYPE_COUNT && t != RAY_CSV_AUTO_TAG) ||
t == RAY_TABLE) {
ray_vm_unmap_file(buf, file_size);
return RAY_ERR_TYPE;
}
resolved_types[c] = t;
}
} else if (!col_types_in) {
ray_t* sample_offsets_hdr = NULL;
int64_t* sample_offsets = NULL;
int64_t sample_n = csv_streaming_sample(buf, file_size, data_offset,
data_has_quotes,
&sample_offsets,
&sample_offsets_hdr);
if (sample_n < 0) {
ray_vm_unmap_file(buf, file_size);
return RAY_ERR_CANCEL;
}
bool infer_ok = csv_infer_types_from_offsets(
buf, buf_end, sample_offsets, sample_n, ncols, delimiter,
esc_buf, resolved_types);
scratch_free(sample_offsets_hdr);
if (!infer_ok) {
ray_vm_unmap_file(buf, file_size);
return RAY_ERR_OOM;
}
}
if (!csv_resolve_auto_streamed(buf, file_size, data_offset, ncols,
delimiter, data_has_quotes, resolved_types)) {
ray_vm_unmap_file(buf, file_size);
return RAY_ERR_CANCEL;
}
err = ray_mkdir_p(dir);
if (err != RAY_OK) {
ray_vm_unmap_file(buf, file_size);
return err;
}
char schema_path[1024];
int sn = snprintf(schema_path, sizeof(schema_path), "%s/.d", dir);
if (sn < 0 || (size_t)sn >= sizeof(schema_path)) {
ray_vm_unmap_file(buf, file_size);
return RAY_ERR_RANGE;
}
remove(schema_path);
struct ray_sym_domain_s* sym_dom = NULL;
{
bool any_sym = false;
for (int c = 0; c < ncols; c++)
if (resolved_types[c] == RAY_SYM) any_sym = true;
if (any_sym) {
sym_dom = ray_sym_domain_open_or_create(sym_path);
if (!sym_dom) {
ray_vm_unmap_file(buf, file_size);
return RAY_ERR_IO;
}
if (ray_sym_domain_count(sym_dom) == 0 &&
ray_sym_domain_intern(sym_dom, "", 0) != 0) {
ray_sym_domain_release(sym_dom);
ray_vm_unmap_file(buf, file_size);
return RAY_ERR_OOM;
}
}
}
ray_col_stream_t writers[CSV_MAX_COLS];
memset(writers, 0, sizeof(writers));
for (int c = 0; c < ncols; c++) {
err = ray_col_stream_open(&writers[c], dir, col_name_ids[c],
resolved_types[c], sym_dom);
if (err == RAY_OK) err = ray_col_stream_index_begin(&writers[c], 0);
if (err != RAY_OK) {
for (int j = 0; j <= c; j++) ray_col_stream_abort(&writers[j]);
if (sym_dom) ray_sym_domain_release(sym_dom);
ray_vm_unmap_file(buf, file_size);
return err;
}
}
size_t chunk_offset = data_offset;
bool wrote_any = false;
size_t avg_row_bytes = 64;
int64_t tr_setup = ray_profile_now_ns() - tr_t0;
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;
int64_t tr_c0 = ray_profile_now_ns();
if (chunk_offset < file_size) {
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;
break;
}
}
int64_t tr_c1 = ray_profile_now_ns();
ray_t* tbl = csv_materialize_rows(buf, file_size, row_offsets,
cnt, ncols, delimiter, col_name_ids,
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)
: RAY_ERR_OOM;
if (tbl) ray_release(tbl);
break;
}
int64_t tr_c2 = ray_profile_now_ns();
{
csv_splayed_append_ctx_t actx = { .writers = writers, .tbl = tbl,
.ncols = ncols, .err = RAY_OK,
.col_ns = trace ? tr_col_ns : NULL };
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);
int64_t tr_c3 = ray_profile_now_ns();
tr_scan += tr_c1 - tr_c0; tr_parse += tr_c2 - tr_c1; tr_append += tr_c3 - tr_c2; tr_chunks++;
if (trace)
fprintf(stderr, "csv.splayed: chunk=%" PRId64 " rows=%" PRId64 " scan=%.1fms parse=%.1fms append=%.1fms\n",
tr_chunks, cnt, TR_MS(tr_c1 - tr_c0), TR_MS(tr_c2 - tr_c1), TR_MS(tr_c3 - tr_c2));
if (err != RAY_OK) break;
wrote_any = true;
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;
}
int64_t tr_loop_end = ray_profile_now_ns();
if (trace) {
fprintf(stderr, "csv.splayed: file=%s ncols=%d chunks=%" PRId64 " setup=%.1fms scan=%.1fms parse=%.1fms append=%.1fms loop=%.1fms\n",
path, ncols, tr_chunks, TR_MS(tr_setup), TR_MS(tr_scan), TR_MS(tr_parse), TR_MS(tr_append),
TR_MS(tr_loop_end - tr_t0 - tr_setup));
for (int k = 0; k < 3 && k < ncols; k++) {
int best = -1;
for (int c = 0; c < ncols; c++)
if (tr_col_ns[c] > 0 && (best < 0 || tr_col_ns[c] > tr_col_ns[best])) best = c;
if (best < 0) break;
ray_t* na = ray_sym_str(col_name_ids[best]);
fprintf(stderr, "csv.splayed: append_top%d col=%s type=%d total=%.1fms\n", k + 1,
na ? ray_str_ptr(na) : "?", (int)resolved_types[best], TR_MS(tr_col_ns[best]));
tr_col_ns[best] = -tr_col_ns[best];
}
}
if (err == RAY_OK && sym_dom) {
err = ray_sym_domain_flush(sym_dom, false);
}
int64_t tr_flush_end = ray_profile_now_ns();
if (trace) fprintf(stderr, "csv.splayed: symfile_flush=%.1fms\n", TR_MS(tr_flush_end - tr_loop_end));
memset(tr_col_ns, 0, sizeof(tr_col_ns));
if (err == RAY_OK) err = ray_col_stream_close_all(writers, ncols, false, trace ? tr_col_ns : NULL);
else for (int c = 0; c < ncols; c++) ray_col_stream_abort(&writers[c]);
int64_t tr_close_end = ray_profile_now_ns();
if (trace) {
int64_t mx = 0; int mxc = -1;
for (int c = 0; c < ncols; c++) if (tr_col_ns[c] > mx) { mx = tr_col_ns[c]; mxc = c; }
ray_t* na = mxc >= 0 ? ray_sym_str(col_name_ids[mxc]) : NULL;
fprintf(stderr, "csv.splayed: finish+publish=%.1fms longest_finish=%s %.1fms\n",
TR_MS(tr_close_end - tr_flush_end),
na ? ray_str_ptr(na) : "-", TR_MS(mx));
}
if (sym_dom) { ray_sym_domain_release(sym_dom); sym_dom = NULL; }
int64_t tr_rel_end = ray_profile_now_ns();
if (trace) fprintf(stderr, "csv.splayed: domain_release=%.1fms\n", TR_MS(tr_rel_end - tr_close_end));
memset(tr_col_ns, 0, sizeof(tr_col_ns));
if (err == RAY_OK) ray_col_stream_hash_all(writers, ncols, trace ? tr_col_ns : NULL);
for (int c = 0; c < ncols; c++)
if (writers[c].index) { ray_release(writers[c].index); writers[c].index = NULL; }
int64_t tr_hash_end = ray_profile_now_ns();
if (trace) {
for (int c = 0; c < ncols; c++) {
if (!tr_col_ns[c]) continue;
ray_t* na = ray_sym_str(col_name_ids[c]);
fprintf(stderr, "csv.splayed: hash col=%s type=%d task=%.1fms\n",
na ? ray_str_ptr(na) : "?", (int)resolved_types[c], TR_MS(tr_col_ns[c]));
}
fprintf(stderr, "csv.splayed: hash_phase=%.1fms\n", TR_MS(tr_hash_end - tr_rel_end));
}
if (err == RAY_OK) {
ray_t* schema = ray_vec_new(RAY_STR, ncols > 0 ? ncols : 1);
if (!schema || RAY_IS_ERR(schema)) {
if (schema) ray_release(schema);
schema = NULL;
err = RAY_ERR_OOM;
}
for (int c = 0; schema && c < ncols; c++) {
ray_t* na = ray_sym_str(col_name_ids[c]);
if (na)
schema = ray_str_vec_append(schema, ray_str_ptr(na),
ray_str_len(na));
if (!na || !schema || RAY_IS_ERR(schema)) {
if (schema && RAY_IS_ERR(schema)) ray_release(schema);
schema = NULL;
err = RAY_ERR_OOM;
}
}
if (schema) {
err = ray_col_save_bulk(schema, schema_path);
ray_release(schema);
}
}
if (trace) fprintf(stderr, "csv.splayed: schema=%.1fms total=%.1fms\n",
TR_MS(ray_profile_now_ns() - tr_hash_end), TR_MS(ray_profile_now_ns() - tr_t0));
#undef TR_MS
if (sym_dom) ray_sym_domain_release(sym_dom);
ray_vm_unmap_file(buf, file_size);
return err;
}
ray_err_t ray_csv_save_splayed_named_opts(const char* path, char delimiter, bool header,
const int8_t* col_types_in, int32_t n_types,
const int64_t* col_names_in, int32_t n_names,
const char* dir, int64_t rows_per_chunk) {
if (ray_interrupted()) return RAY_ERR_CANCEL;
if (!path || !dir) return RAY_ERR_DOMAIN;
char sym_path[1024];
int n = snprintf(sym_path, sizeof(sym_path), "%s/.sym", dir);
if (n < 0 || (size_t)n >= sizeof(sym_path)) return RAY_ERR_RANGE;
ray_splay_write_t write;
ray_err_t err = ray_splay_write_begin(dir, &write);
if (err != RAY_OK) return err;
err = csv_save_splayed_to_dir(path, delimiter, header, col_types_in, n_types,
col_names_in, n_names, write.dir,
rows_per_chunk, sym_path);
return ray_splay_write_finish(&write, err, false);
}
static ray_err_t csv_save_parted_impl(const char* path, char delimiter, bool header,
const int8_t* col_types_in, int32_t n_types,
const int64_t* col_names_in, int32_t n_names,
const char* root, const char* table_name,
int64_t rows_per_part, bool staged) {
if (ray_interrupted()) return RAY_ERR_CANCEL;
if (!path || !root || !table_name) return RAY_ERR_DOMAIN;
if (rows_per_part <= 0) rows_per_part = CSV_PART_ROWS_DEFAULT;
bool trace = getenv("RAY_CSV_TRACE") != NULL;
int fd = open(path, O_RDONLY);
if (fd < 0) return RAY_ERR_IO;
struct stat st;
if (fstat(fd, &st) != 0 || st.st_size <= 0) {
close(fd);
return RAY_ERR_IO;
}
size_t file_size = (size_t)st.st_size;
char* buf = (char*)ray_vm_map_fd_ro(fd, file_size);
close(fd);
if (!buf) return RAY_ERR_IO;
#ifdef __APPLE__
madvise(buf, file_size, MADV_SEQUENTIAL);
#endif
const char* buf_end = buf + file_size;
ray_err_t err = RAY_OK;
if (delimiter == 0) {
int commas = 0, tabs = 0;
for (const char* q = buf; q < buf_end && *q != '\n'; q++) {
if (*q == ',') commas++;
if (*q == '\t') tabs++;
}
delimiter = (tabs > commas) ? '\t' : ',';
}
int ncols = 1;
{
const char* q = buf;
bool in_quote = false;
while (q < buf_end && (in_quote || (*q != '\n' && *q != '\r'))) {
if (*q == '"') in_quote = !in_quote;
else if (!in_quote && *q == delimiter) ncols++;
q++;
}
}
if (ncols > CSV_MAX_COLS) {
ray_vm_unmap_file(buf, file_size);
return RAY_ERR_RANGE;
}
if ((col_types_in && n_types != ncols) ||
(col_names_in && n_names != ncols)) {
ray_vm_unmap_file(buf, file_size);
return RAY_ERR_LENGTH;
}
const char* p = buf;
char esc_buf[8192];
int64_t col_name_ids[CSV_MAX_COLS];
if (header) {
for (int c = 0; c < ncols; c++) {
const char* fld;
size_t flen;
char* dyn_esc = NULL;
p = scan_field(p, buf_end, delimiter, &fld, &flen, esc_buf, &dyn_esc);
col_name_ids[c] = ray_sym_intern(fld, flen);
if (dyn_esc) ray_sys_free(dyn_esc);
}
if (p < buf_end && *p == '\r') p++;
if (p < buf_end && *p == '\n') p++;
} else if (col_names_in) {
for (int c = 0; c < ncols; c++)
col_name_ids[c] = col_names_in[c];
const char* after = csv_skip_matching_header(buf, buf_end, delimiter,
ncols, col_names_in, esc_buf);
if (after) p = after;
} else {
for (int c = 0; c < ncols; c++) {
char name[32];
snprintf(name, sizeof(name), "V%d", c + 1);
col_name_ids[c] = ray_sym_intern(name, strlen(name));
}
}
size_t data_offset = (size_t)(p - buf);
bool data_has_quotes = memchr(buf + data_offset, '"', file_size - data_offset) != NULL;
if (trace) {
fprintf(stderr,
"csv.parted: file=%s size=%zu ncols=%d data_offset=%zu rows_per_part=%" PRId64 " root=%s table=%s\n",
path, file_size, ncols, data_offset, rows_per_part, root, table_name);
}
int8_t resolved_types[CSV_MAX_COLS];
if (col_types_in) {
for (int c = 0; c < ncols; c++) {
int8_t t = col_types_in[c];
if (t < RAY_BOOL ||
(t >= RAY_TYPE_COUNT && t != RAY_CSV_AUTO_TAG) ||
t == RAY_TABLE) {
ray_vm_unmap_file(buf, file_size);
return RAY_ERR_TYPE;
}
resolved_types[c] = t;
}
} else if (!col_types_in) {
ray_t* sample_offsets_hdr = NULL;
int64_t* sample_offsets = NULL;
int64_t sample_n = csv_streaming_sample(buf, file_size, data_offset,
data_has_quotes,
&sample_offsets,
&sample_offsets_hdr);
if (sample_n < 0) {
ray_vm_unmap_file(buf, file_size);
return RAY_ERR_CANCEL;
}
bool infer_ok = csv_infer_types_from_offsets(
buf, buf_end, sample_offsets, sample_n, ncols, delimiter,
esc_buf, resolved_types);
scratch_free(sample_offsets_hdr);
if (!infer_ok) {
ray_vm_unmap_file(buf, file_size);
return RAY_ERR_OOM;
}
}
if (!csv_resolve_auto_streamed(buf, file_size, data_offset, ncols,
delimiter, data_has_quotes, resolved_types)) {
ray_vm_unmap_file(buf, file_size);
return RAY_ERR_CANCEL;
}
err = ray_mkdir_p(root);
if (err != RAY_OK) {
ray_vm_unmap_file(buf, file_size);
return err;
}
int64_t part = 0;
size_t chunk_offset = data_offset;
bool wrote_any = false;
ray_sym_domain_t* import_domain = NULL;
for (int c = 0; c < ncols; c++) if (resolved_types[c] == RAY_SYM) {
char sym_path[1024];
int n = snprintf(sym_path,sizeof(sym_path),"%s/.sym",root);
if (n < 0 || (size_t)n >= sizeof(sym_path)) err = RAY_ERR_RANGE;
else if (!(import_domain = ray_sym_domain_open_or_create(sym_path))) err = RAY_ERR_IO;
break;
}
if (err != RAY_OK) { ray_vm_unmap_file(buf,file_size); return err; }
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_part, data_has_quotes,
&row_offsets,
&row_offsets_hdr, &next_offset);
if (cnt <= 0) {
if (trace)
fprintf(stderr, "csv.parted: row-offset failure part=%" PRId64 " offset=%zu\n",
part, chunk_offset);
scratch_free(row_offsets_hdr);
err = (cnt < 0) ? RAY_ERR_CANCEL : RAY_ERR_IO;
break;
}
}
ray_t* tbl = csv_materialize_rows(buf, file_size, row_offsets,
cnt, ncols, delimiter, col_name_ids, resolved_types,
import_domain);
if (!tbl || RAY_IS_ERR(tbl)) {
err = (tbl && RAY_IS_ERR(tbl)) ? ray_err_from_obj(tbl)
: RAY_ERR_OOM;
if (tbl) ray_release(tbl);
scratch_free(row_offsets_hdr);
if (trace)
fprintf(stderr, "csv.parted: materialize failure part=%" PRId64 " rows=%" PRId64 "\n",
part, cnt);
break;
}
for (int64_t c = 0; c < ray_table_ncols(tbl); c++) {
ray_t* col = ray_table_get_col_idx(tbl,c);
if (col->type == RAY_STR || col->type == RAY_SYM || col->type == RAY_F32 || col->type == RAY_GUID) continue;
ray_retain(col);
ray_t* indexed = col->len >= 65536 ? ray_index_attach_chunk_zone(&col,16) : ray_index_attach_zone(&col);
if (indexed && RAY_IS_ERR(indexed)) {
err = ray_err_from_obj(indexed); ray_release(indexed); ray_release(col);
break;
}
ray_table_set_col_idx(tbl,c,col); ray_release(col);
}
if (err != RAY_OK) { ray_release(tbl); scratch_free(row_offsets_hdr); break; }
char leaf[1024];
int n = snprintf(leaf, sizeof(leaf), "%s/%" PRId64 "/%s", root, part, table_name);
if (n < 0 || (size_t)n >= sizeof(leaf)) {
ray_release(tbl);
scratch_free(row_offsets_hdr);
err = RAY_ERR_RANGE;
break;
}
if (trace)
fprintf(stderr, "csv.parted: save part=%" PRId64 " rows=%" PRId64 " leaf=%s\n",
part, cnt, leaf);
char root_sym[1024];
int sn = snprintf(root_sym, sizeof(root_sym), "%s/.sym", root);
if (sn < 0 || (size_t)sn >= sizeof(root_sym)) {
ray_release(tbl);
scratch_free(row_offsets_hdr);
err = RAY_ERR_RANGE;
break;
}
err = staged ? ray_splay_save_staged_bulk(tbl, leaf, root_sym)
: ray_splay_save_bulk(tbl, leaf, root_sym);
ray_release(tbl);
scratch_free(row_offsets_hdr);
if (err != RAY_OK) {
if (trace)
fprintf(stderr, "csv.parted: save failure part=%" PRId64 " err=%s\n",
part, ray_err_code_str(err));
break;
}
wrote_any = true;
if (cnt == 0) break;
chunk_offset = next_offset;
part++;
}
if (err == RAY_OK && staged && import_domain)
err = ray_sym_domain_flush(import_domain, false);
if (import_domain) ray_sym_domain_release(import_domain);
ray_vm_unmap_file(buf, file_size);
if (trace)
fprintf(stderr, "csv.parted: done err=%s\n", ray_err_code_str(err));
return err;
}
ray_err_t ray_csv_parted_paths(const char* root, char* dest, size_t dest_size,
char* staging, size_t staging_size) {
size_t len = strlen(root);
while (len > 1 && (root[len-1] == '/' || root[len-1] == '\\')) len--;
if (len >= dest_size) return RAY_ERR_RANGE;
memcpy(dest,root,len); dest[len] = 0;
int n = snprintf(staging,staging_size,"%s.csv-partial",dest);
if (n < 0 || (size_t)n >= staging_size) return RAY_ERR_RANGE;
return RAY_OK;
}
ray_err_t ray_csv_save_parted_named_opts(const char* path, char delimiter, bool header,
const int8_t* col_types, int32_t n_types,
const int64_t* col_names, int32_t n_names,
const char* root, const char* table_name,
int64_t rows_per_part) {
if (!path || !root || !*root || !table_name || !*table_name ||
table_name[0] == '.' || strchr(table_name,'/') || strchr(table_name,'\\'))
return RAY_ERR_DOMAIN;
if (ray_interrupted()) return RAY_ERR_CANCEL;
char dest[1024], staging[1100];
ray_err_t perr = ray_csv_parted_paths(root,dest,sizeof(dest),staging,sizeof(staging));
if (perr != RAY_OK) return perr;
struct stat st;
if (stat(dest,&st) == 0)
return csv_save_parted_impl(path,delimiter,header,col_types,n_types,
col_names,n_names,dest,table_name,rows_per_part,false);
FILE* in = fopen(path,"rb");
if (!in) return RAY_ERR_IO;
fclose(in);
char* slash = strrchr(dest,'/');
if (slash && slash != dest) {
*slash = 0;
ray_err_t e = ray_mkdir_p(dest);
*slash = '/';
if (e != RAY_OK) return e;
}
#ifdef RAY_OS_WINDOWS
if (!CreateDirectoryA(staging,NULL)) return RAY_ERR_IO;
#else
if (mkdir(staging,0755) != 0) return RAY_ERR_IO;
#endif
ray_err_t err = csv_save_parted_impl(path,delimiter,header,col_types,n_types,
col_names,n_names,staging,table_name,rows_per_part,true);
if (err != RAY_OK) {
char probe[1200];
snprintf(probe,sizeof(probe),"%s/0",staging);
if (stat(probe,&st) != 0) {
snprintf(probe,sizeof(probe),"%s/.sym",staging);
remove(probe);
#ifdef RAY_OS_WINDOWS
RemoveDirectoryA(staging);
#else
rmdir(staging);
#endif
}
return err;
}
if (ray_interrupted()) return RAY_ERR_CANCEL;
if (stat(dest,&st) == 0) return RAY_ERR_IO;
return ray_file_rename_new(staging,dest);
}
ray_t* ray_read_csv(const char* path) {
return ray_read_csv_opts(path, 0, true, NULL, 0);
}
typedef struct csv_writer_t {
FILE* fp;
int err;
} csv_writer_t;
static inline void cw_putc(csv_writer_t* w, int c) {
if (w->err) return;
if (fputc(c, w->fp) == EOF) w->err = 1;
}
static inline void cw_write(csv_writer_t* w, const char* s, size_t len) {
if (w->err || len == 0) return;
if (fwrite(s, 1, len, w->fp) != len) w->err = 1;
}
static inline void cw_puts(csv_writer_t* w, const char* s) {
if (!s) return;
cw_write(w, s, strlen(s));
}
static void cw_printf(csv_writer_t* w, const char* fmt, ...) {
if (w->err) return;
char buf[64];
va_list ap;
va_start(ap, fmt);
int n = vsnprintf(buf, sizeof(buf), fmt, ap);
va_end(ap);
if (n < 0) { w->err = 1; return; }
if ((size_t)n >= sizeof(buf)) { w->err = 1; return; }
cw_write(w, buf, (size_t)n);
}
static void csv_write_str(csv_writer_t* w, const char* s, size_t len) {
int need_quote = 0;
for (size_t i = 0; i < len; i++) {
if (s[i] == ',' || s[i] == '"' || s[i] == '\n' || s[i] == '\r') {
need_quote = 1;
break;
}
}
if (need_quote) {
cw_putc(w, '"');
size_t start = 0;
for (size_t i = 0; i < len; i++) {
if (s[i] == '"') {
cw_write(w, s + start, i - start);
cw_putc(w, '"');
start = i;
}
}
cw_write(w, s + start, len - start);
cw_putc(w, '"');
} else {
cw_write(w, s, len);
}
}
static void csv_write_date(csv_writer_t* w, int32_t v) {
int32_t z = v + 10957 + 719468;
int32_t era = (z >= 0 ? z : z - 146096) / 146097;
uint32_t doe = (uint32_t)(z - era * 146097);
uint32_t yoe = (doe - doe/1460 + doe/36524 - doe/146096) / 365;
int32_t y = (int32_t)yoe + era * 400;
uint32_t doy = doe - (365*yoe + yoe/4 - yoe/100);
uint32_t mp = (5*doy + 2) / 153;
int32_t d = (int32_t)(doy - (153*mp + 2)/5 + 1);
int32_t m = (int32_t)(mp < 10 ? mp + 3 : mp - 9);
if (m <= 2) y++;
cw_printf(w, "%04d-%02d-%02d", y, m, d);
}
static void csv_write_time(csv_writer_t* w, int32_t ms) {
int32_t sign = ms < 0 ? -1 : 1;
uint32_t u = (ms == INT32_MIN) ? (uint32_t)INT32_MAX + 1u : (uint32_t)(sign == -1 ? -ms : ms);
uint32_t h = u / 3600000u;
uint32_t mi = (u % 3600000u) / 60000u;
uint32_t s = (u % 60000u) / 1000u;
uint32_t frac = u % 1000u;
if (sign == -1) cw_putc(w, '-');
if (frac) cw_printf(w, "%02u:%02u:%02u.%03u", h, mi, s, frac);
else cw_printf(w, "%02u:%02u:%02u", h, mi, s);
}
static void csv_write_timestamp(csv_writer_t* w, int64_t ns) {
const int64_t NS_PER_DAY = 86400000000000LL;
int64_t days = ns / NS_PER_DAY;
int64_t ns_in = ns % NS_PER_DAY;
if (ns_in < 0) { days--; ns_in += NS_PER_DAY; }
csv_write_date(w, (int32_t)days);
cw_putc(w, 'T');
uint64_t tns = (uint64_t)ns_in;
uint32_t h = (uint32_t)(tns / 3600000000000ULL);
uint32_t mi = (uint32_t)((tns % 3600000000000ULL) / 60000000000ULL);
uint32_t s = (uint32_t)((tns % 60000000000ULL) / 1000000000ULL);
uint32_t frac = (uint32_t)(tns % 1000000000ULL);
if (frac) cw_printf(w, "%02u:%02u:%02u.%09u", h, mi, s, frac);
else cw_printf(w, "%02u:%02u:%02u", h, mi, s);
}
static void csv_write_f64(csv_writer_t* w, double v) {
if (isnan(v)) { cw_puts(w, "nan"); return; }
if (isinf(v)) { cw_puts(w, v < 0 ? "-inf" : "inf"); return; }
cw_printf(w, "%.17g", v);
}
static void csv_write_guid(csv_writer_t* w, const uint8_t* g) {
cw_printf(w,
"%02x%02x%02x%02x-%02x%02x-%02x%02x-%02x%02x-%02x%02x%02x%02x%02x%02x",
g[0], g[1], g[2], g[3], g[4], g[5], g[6], g[7],
g[8], g[9], g[10], g[11], g[12], g[13], g[14], g[15]);
}
typedef struct csv_col_info_t {
ray_t* col;
ray_t* data_owner;
int64_t base_row;
const void* data;
int8_t type;
uint8_t attrs;
bool has_nulls;
} csv_col_info_t;
static void csv_col_info_init(csv_col_info_t* ci, ray_t* col) {
ci->col = col;
ci->data_owner = col;
ci->base_row = 0;
if (col && (col->attrs & RAY_ATTR_SLICE) && col->slice_parent) {
ci->data_owner = col->slice_parent;
ci->base_row = col->slice_offset;
}
ci->type = col ? col->type : 0;
ci->attrs = ci->data_owner ? ci->data_owner->attrs : 0;
ci->data = ci->data_owner ? ray_data(ci->data_owner) : NULL;
ci->has_nulls = ray_vec_may_have_nulls(col);
}
static void csv_write_cell(csv_writer_t* w, const csv_col_info_t* ci, int64_t r) {
if (!ci->col) return;
if (ci->has_nulls && ray_vec_is_null(ci->col, r)) return;
int64_t dr = ci->base_row + r;
int8_t t = ci->type;
const void* d = ci->data;
switch (t) {
case RAY_I64: case RAY_TIMESTAMP: break;
default: break;
}
switch (t) {
case RAY_I64:
cw_printf(w, "%" PRId64, ((const int64_t*)d)[dr]);
break;
case RAY_I32:
cw_printf(w, "%" PRId32, ((const int32_t*)d)[dr]);
break;
case RAY_I16:
cw_printf(w, "%d", (int)((const int16_t*)d)[dr]);
break;
case RAY_BOOL:
cw_puts(w, ((const uint8_t*)d)[dr] ? "true" : "false");
break;
case RAY_U8:
cw_printf(w, "%u", (unsigned)((const uint8_t*)d)[dr]);
break;
case RAY_F64:
csv_write_f64(w, ((const double*)d)[dr]);
break;
case RAY_F32:
csv_write_f64(w, (double)((const float*)d)[dr]);
break;
case RAY_DATE:
csv_write_date(w, ((const int32_t*)d)[dr]);
break;
case RAY_TIME:
csv_write_time(w, ((const int32_t*)d)[dr]);
break;
case RAY_TIMESTAMP:
csv_write_timestamp(w, ((const int64_t*)d)[dr]);
break;
case RAY_SYM: {
ray_t* s = ray_sym_vec_cell(ci->data_owner, dr);
if (s) csv_write_str(w, ray_str_ptr(s), ray_str_len(s));
break;
}
case RAY_STR: {
size_t slen = 0;
const char* sp = ray_str_vec_get(ci->col, r, &slen);
csv_write_str(w, sp ? sp : "", slen);
break;
}
case RAY_GUID:
csv_write_guid(w, (const uint8_t*)d + dr * 16);
break;
case RAY_LIST: {
ray_t** elems = (ray_t**)d;
ray_t* e = elems[dr];
if (!e || RAY_IS_ERR(e)) return;
ray_t* fmt = ray_fmt(e, false);
if (!fmt || RAY_IS_ERR(fmt)) return;
csv_write_str(w, ray_str_ptr(fmt), ray_str_len(fmt));
ray_release(fmt);
break;
}
default:
break;
}
}
ray_err_t ray_write_csv(ray_t* table, const char* path) {
if (!table || !path || path[0] == '\0') return RAY_ERR_TYPE;
int64_t ncols = ray_table_ncols(table);
int64_t nrows = ray_table_nrows(table);
if (ncols <= 0) return RAY_ERR_TYPE;
char tmp_path[1024];
if (snprintf(tmp_path, sizeof(tmp_path), "%s.tmp", path) >= (int)sizeof(tmp_path))
return RAY_ERR_IO;
FILE* fp = fopen(tmp_path, "wb");
if (!fp) return RAY_ERR_IO;
csv_writer_t w = { .fp = fp, .err = 0 };
ray_t* col_info_block = ray_alloc((size_t)ncols * sizeof(csv_col_info_t));
if (!col_info_block || RAY_IS_ERR(col_info_block)) {
fclose(fp);
remove(tmp_path);
return RAY_ERR_OOM;
}
csv_col_info_t* ci = (csv_col_info_t*)ray_data(col_info_block);
for (int64_t c = 0; c < ncols; c++)
csv_col_info_init(&ci[c], ray_table_get_col_idx(table, c));
for (int64_t c = 0; c < ncols; c++) {
if (c > 0) cw_putc(&w, ',');
int64_t name_id = ray_table_col_name(table, c);
ray_t* name_atom = ray_sym_str(name_id);
if (name_atom)
csv_write_str(&w, ray_str_ptr(name_atom), ray_str_len(name_atom));
}
cw_putc(&w, '\n');
for (int64_t r = 0; r < nrows && !w.err; r++) {
for (int64_t c = 0; c < ncols; c++) {
if (c > 0) cw_putc(&w, ',');
csv_write_cell(&w, &ci[c], r);
}
cw_putc(&w, '\n');
}
ray_free(col_info_block);
if (fflush(fp) != 0) w.err = 1;
int close_err = (fclose(fp) != 0);
if (close_err) w.err = 1;
if (w.err) {
remove(tmp_path);
return RAY_ERR_IO;
}
ray_fd_t fd = ray_file_open(tmp_path, RAY_OPEN_READ | RAY_OPEN_WRITE);
if (fd == RAY_FD_INVALID) { remove(tmp_path); return RAY_ERR_IO; }
ray_err_t sync_err = ray_file_sync(fd);
ray_file_close(fd);
if (sync_err != RAY_OK) { remove(tmp_path); return sync_err; }
ray_err_t rn_err = ray_file_rename(tmp_path, path);
if (rn_err != RAY_OK) { remove(tmp_path); return rn_err; }
return RAY_OK;
}