#ifndef _POSIX_C_SOURCE
#define _POSIX_C_SOURCE 200809L
#endif
#include "parquet.h"
#include "core/platform.h"
#include "core/pool.h"
#include "core/crc32.h"
#include "mem/heap.h"
#include "table/sym.h"
#include "table/domain.h"
#include "vec/str.h"
#include "vec/vec.h"
#include "ops/idxop.h"
#include "ops/hash.h"
#include "ops/glob.h"
#include "ops/agg_registry.h"
#include "store/splay.h"
#include "store/stream.h"
#include "store/col.h"
#include "store/fileio.h"
#include "lang/eval.h"
#include "lang/internal.h"
#include <limits.h>
#include <math.h>
#include <stdio.h>
#include <string.h>
#include <sys/stat.h>
#ifndef RAY_OS_WINDOWS
#include <unistd.h>
#endif
#if defined(__SSE2__)
#include <emmintrin.h>
#endif
#define PQ_BATCH 65536
#define PQ_MAX_PART_ROWS 1048576
#define PQ_MAX_COLS 4096
#define PQ_MAX_PAGE (64u * 1024u * 1024u)
#define PQ_POOL_TARGET (64u * 1024u * 1024u)
#define PQ_POOL_MAX ((uint64_t)UINT32_MAX)
#define PQ_MAX_FOOTER (64u * 1024u * 1024u)
#define PQ_EPOCH_DAYS 10957
#define PQ_EPOCH_NS INT64_C(946684800000000000)
typedef struct { const uint8_t *p, *end; uint8_t type; } pq_span;
typedef struct { const uint8_t *p, *end; bool bad; } pq_cur;
static uint32_t pq_u32(const uint8_t* p) {
return (uint32_t)p[0] | (uint32_t)p[1]<<8 | (uint32_t)p[2]<<16 | (uint32_t)p[3]<<24;
}
static uint64_t pq_u64(const uint8_t* p) { return pq_u32(p) | (uint64_t)pq_u32(p+4)<<32; }
static uint8_t pq_byte(pq_cur* c) {
if (c->p == c->end) { c->bad = true; return 0; }
return *c->p++;
}
static bool pq_take(pq_cur* c, uint64_t n, const uint8_t** out) {
if (n > (uint64_t)(c->end - c->p)) { c->bad = true; return false; }
if (out) *out = c->p;
c->p += (size_t)n;
return true;
}
static uint64_t pq_var(pq_cur* c) {
uint64_t v = 0;
for (unsigned s = 0; s < 64; s += 7) {
uint8_t b = pq_byte(c);
if (c->bad || (s == 63 && b > 1)) { c->bad = true; return 0; }
v |= (uint64_t)(b & 127) << s;
if (!(b & 128)) return v;
}
c->bad = true; return 0;
}
static int64_t pq_int(pq_cur* c) {
uint64_t u = pq_var(c);
return (int64_t)(u >> 1) ^ -(int64_t)(u & 1);
}
static uint64_t pq_list(pq_cur* c, uint8_t* type) {
uint8_t h = pq_byte(c); *type = h & 15;
uint64_t n = h >> 4;
if (n == 15) n = pq_var(c);
if (!*type || *type > 12 || n > (uint64_t)(c->end-c->p)) c->bad = true;
return n;
}
static bool pq_skip(pq_cur* c, uint8_t t, unsigned depth, bool collection) {
if (depth > 32 || c->bad) { c->bad = true; return false; }
switch (t) {
case 1: case 2:
if (collection) { uint8_t b = pq_byte(c); if (b != 1 && b != 2) c->bad = true; }
break;
case 3: pq_byte(c); break;
case 4: case 5: case 6: pq_var(c); break;
case 7: pq_take(c, 8, NULL); break;
case 8: { uint64_t n = pq_var(c); pq_take(c, n, NULL); break; }
case 9: case 10: {
uint8_t et; uint64_t n = pq_list(c, &et);
for (uint64_t i = 0; i < n && !c->bad; i++) pq_skip(c, et, depth+1, true);
break;
}
case 11: {
uint64_t n = pq_var(c);
if (n > (uint64_t)(c->end-c->p)/2) { c->bad = true; break; }
uint8_t h = n ? pq_byte(c) : 0;
for (uint64_t i = 0; i < n && !c->bad; i++) {
pq_skip(c, h>>4, depth+1, true); pq_skip(c, h&15, depth+1, true);
}
break;
}
case 12:
while (!c->bad) {
uint8_t h = pq_byte(c);
if (!h) break;
if (!(h>>4)) pq_int(c);
pq_skip(c, h&15, depth+1, false);
}
break;
default: c->bad = true;
}
return !c->bad;
}
static bool pq_fields(pq_span s, pq_span* out, size_t n) {
memset(out, 0, n*sizeof(*out));
if (s.type != 12) return false;
pq_cur c = {s.p, s.end, false}; int64_t id = 0;
while (!c.bad) {
uint8_t h = pq_byte(&c);
if (!h) return !c.bad && c.p == c.end;
id = h>>4 ? id + (h>>4) : pq_int(&c);
if (id <= 0 || id > INT16_MAX) return false;
const uint8_t* p = c.p;
if (!pq_skip(&c, h&15, 0, false)) return false;
if ((uint64_t)id < n) {
if (out[id].type) return false;
out[id] = (pq_span){p, c.p, h&15};
}
}
return false;
}
static bool pq_num(pq_span s, int64_t* v) {
if (s.type < 4 || s.type > 6) return false;
pq_cur c = {s.p, s.end, false}; *v = pq_int(&c);
return !c.bad && c.p == c.end &&
(s.type != 4 || (*v >= INT16_MIN && *v <= INT16_MAX)) &&
(s.type != 5 || (*v >= INT32_MIN && *v <= INT32_MAX));
}
static int64_t pq_get(pq_span s, int64_t fallback) {
int64_t v; return pq_num(s, &v) ? v : fallback;
}
static bool pq_binary(pq_span s, const uint8_t** p, size_t* n) {
if (s.type != 8) return false;
pq_cur c = {s.p, s.end, false}; uint64_t len = pq_var(&c);
if (c.bad || !pq_take(&c, len, p) || c.p != c.end) return false;
*n = (size_t)len; return true;
}
static bool pq_element(pq_cur* c, uint8_t t, pq_span* s) {
const uint8_t* start = c->p;
if (!pq_skip(c, t, 0, true)) return false;
*s = (pq_span){start, c->p, t}; return true;
}
static ray_t* pq_error(const char* msg) { return ray_error("parquet", "%s", msg); }
bool ray_parquet_snappy(const uint8_t* src, size_t len, uint8_t* dst, size_t size) {
pq_cur c = {src, src+len, false}; uint64_t expected = pq_var(&c); size_t at = 0;
if (c.bad || expected != size) return false;
while (c.p < c.end && at < size) {
uint8_t tag = pq_byte(&c); size_t n, off;
if (!(tag&3)) {
n = tag>>2;
if (n < 60) n++;
else {
unsigned nb = (unsigned)n-59; const uint8_t* p;
if (!pq_take(&c, nb, &p)) return false;
uint32_t x = 0;
for (unsigned i = 0; i < nb; i++) x |= (uint32_t)p[i] << (8*i);
n = (size_t)x+1;
}
const uint8_t* p;
if (n > size-at || !pq_take(&c, n, &p)) return false;
memcpy(dst+at, p, n); at += n;
} else {
if ((tag&3) == 1) { n = 4+((tag>>2)&7); off = ((size_t)(tag&224)<<3) | pq_byte(&c); }
else {
const uint8_t* p; unsigned nb = (tag&3) == 2 ? 2 : 4;
if (!pq_take(&c, nb, &p)) return false;
off = nb == 2 ? (size_t)p[0] | (size_t)p[1]<<8 : pq_u32(p);
n = 1+(tag>>2);
}
if (c.bad || !off || off > at || n > size-at) return false;
size_t copied = 0;
while (copied < n) {
size_t step = n-copied < off ? n-copied : off;
memcpy(dst+at+copied,dst+at+copied-off,step);
copied += step;
if (copied >= off) off += off;
}
at += n;
}
}
return !c.bad && at == size && c.p == c.end;
}
typedef struct {
pq_cur c;
uint64_t left, bit;
const uint8_t* packed;
uint32_t value;
unsigned width;
} pq_rle;
static bool pq_rle_next(pq_rle* r, uint32_t* out) {
if (!r->left) {
uint64_t h = pq_var(&r->c);
if (r->c.bad || !h || r->width > 32) return false;
if (h&1) {
uint64_t groups = h>>1;
if (!groups || groups > INT32_MAX/8) return false;
r->left = groups*8; r->bit = 0;
if (!pq_take(&r->c, groups*r->width, &r->packed)) return false;
} else {
r->left = h>>1; r->packed = NULL; r->value = 0;
if (r->left > INT32_MAX) return false;
for (unsigned i = 0; i < (r->width+7)/8; i++) r->value |= (uint32_t)pq_byte(&r->c)<<(8*i);
if (r->c.bad || (r->width < 32 && r->value >= (UINT32_C(1)<<r->width))) return false;
}
}
if (r->packed) {
uint64_t v = 0; unsigned shift = (unsigned)(r->bit&7);
unsigned nb = (shift+r->width+7)/8;
for (unsigned i = 0; i < nb; i++) v |= (uint64_t)r->packed[r->bit/8+i]<<(8*i);
*out = (uint32_t)((v>>shift) & (r->width == 32 ? UINT32_MAX : (UINT32_C(1)<<r->width)-1));
r->bit += r->width;
} else *out = r->value;
r->left--; return true;
}
static bool pq_rle_done(const pq_rle* r) {
return !r->c.bad && r->c.p == r->c.end && (!r->left || (r->packed && r->left < 8));
}
typedef struct {
int64_t name;
int physical, optional, converted, type;
int64_t scale;
bool native_symbol, strict;
ray_sym_domain_t* import_domain;
} pq_schema;
typedef struct { const uint8_t* p; uint32_t n; } pq_string;
typedef struct {
const char* strings[8192];
size_t lengths[8192];
uint32_t hashes[8192];
int64_t ids[8192], positions[8192];
} pq_symbol_scratch;
typedef struct {
pq_span metadata;
pq_cur chunk, values;
pq_rle defs, ids;
uint8_t *page, *dict;
size_t page_cap, dict_cap;
pq_string* strings;
int64_t* symbol_ids;
pq_symbol_scratch* symbols;
int64_t dict_count, page_left, chunk_left, page_nulls, expected_nulls, pages_skipped;
uint64_t bool_bit;
int codec, encoding;
bool have_dict, optional;
bool bool_rle;
} pq_column;
static pq_symbol_scratch* pq_symbols(pq_column* c) {
if (!c->symbols) c->symbols = ray_alloc_raw(sizeof(*c->symbols));
return c->symbols;
}
typedef struct { int64_t lo, hi; } pq_interval;
struct ray_parquet {
uint8_t* map;
size_t size, data_end;
pq_schema* schema;
pq_span* groups;
int32_t *selected;
pq_column* cursors;
ray_sym_domain_t* scan_domain;
ray_t* text_pattern;
int64_t text_pos;
int64_t ncols, ngroups, nselected, noutput, rows, group, group_left, batch_rows;
int64_t group_batch;
int64_t filter_pos, filter_lo, filter_hi, skipped;
pq_interval* excluded;
int64_t nexcluded, exclude_pos, group_rows, bloom_skipped, parallel_batches;
bool failed, emitted, borrowed, filter_nulls;
};
static bool pq_reserve(uint8_t** buf, size_t* cap, size_t n) {
if (n > PQ_MAX_PAGE) return false;
if (!n) n = 1;
if (n <= *cap) return true;
size_t nc = *cap ? *cap : 4096;
while (nc < n) nc *= 2;
void* p = ray_realloc_raw(*buf, nc);
if (!p) return false;
*buf = p; *cap = nc; return true;
}
static bool pq_utf8(const uint8_t* p, size_t n) {
for (size_t i = 0; i < n;) {
uint32_t c = p[i++]; unsigned more; uint32_t min;
if (c < 0x80) continue;
if (c >= 0xc2 && c <= 0xdf) { more = 1; min = 0x80; c &= 31; }
else if (c >= 0xe0 && c <= 0xef) { more = 2; min = 0x800; c &= 15; }
else if (c >= 0xf0 && c <= 0xf4) { more = 3; min = 0x10000; c &= 7; }
else return false;
if (more > n-i) return false;
while (more--) { uint8_t b = p[i++]; if ((b&0xc0) != 0x80) return false; c = (c<<6)|(b&63); }
if (c < min || c > 0x10ffff || (c >= 0xd800 && c <= 0xdfff)) return false;
}
return true;
}
static bool pq_schema_type(pq_schema* s, pq_span* f) {
int64_t phys, rep;
if (!pq_num(f[1], &phys) || !pq_num(f[3], &rep) || rep < 0 || rep > 1 || f[5].type) return false;
s->physical = (int)phys; s->optional = (int)rep;
s->converted = (int)pq_get(f[6], -1); s->scale = 1;
switch (phys) {
case 0: s->type = rep ? RAY_I16 : RAY_BOOL; break;
case 1: s->type = RAY_I32; break;
case 2: s->type = RAY_I64; break;
case 4: s->type = RAY_F32; break;
case 5: s->type = RAY_F64; break;
case 6: s->type = RAY_STR; break;
default: return false;
}
int conv = s->converted;
if (f[10].type) {
pq_span logical[20];
if (!pq_fields(f[10], logical, 20)) return false;
int found = 0;
for (int i = 1; i < 20; i++) if (logical[i].type) { if (found) return false; found = i; }
if (found == 1) conv = 0;
else if (found == 6) conv = 6;
else if (found == 10) {
pq_span ints[3]; int64_t bits;
if (!pq_fields(logical[10], ints, 3) || ints[1].type != 3 || ints[1].end-ints[1].p != 1) return false;
bits = *ints[1].p;
if (bits != 8 && bits != 16 && bits != 32 && bits != 64) return false;
if (ints[2].type != 1 && ints[2].type != 2) return false;
conv = (ints[2].type == 1 ? 15 : 11) + (bits == 8 ? 0 : bits == 16 ? 1 : bits == 32 ? 2 : 3);
} else if (found == 8) {
pq_span ts[3], unit[4];
if (!pq_fields(logical[8], ts, 3) || (ts[1].type != 1 && ts[1].type != 2) ||
!pq_fields(ts[2], unit, 4) || phys != 2) return false;
int units = !!unit[1].type + !!unit[2].type + !!unit[3].type;
if (units != 1) return false;
s->scale = unit[1].type ? 1000000 : unit[2].type ? 1000 : 1;
conv = -2; s->type = RAY_TIMESTAMP;
} else return false;
}
if (conv == 0) { if (phys != 6) return false; }
else if (conv == 6) { if (phys != 1) return false; s->type = RAY_DATE; }
else if (conv == 9 || conv == 10) {
if (phys != 2) return false;
s->type = RAY_TIMESTAMP; s->scale = conv == 9 ? 1000000 : 1000;
} else if (conv >= 11 && conv <= 18) {
if (conv == 14) return false;
if (phys != (conv == 18 ? 2 : 1)) return false;
if (conv == 13) s->type = RAY_I64;
else if (conv == 11 || conv == 15) s->type = RAY_I16;
} else if (conv != -1 && conv != -2) return false;
s->converted = conv;
return true;
}
void ray_parquet_close(ray_parquet_t* r) {
if (!r) return;
if (r->cursors) for (int64_t i = 0; i < r->nselected; i++) {
ray_free_raw(r->cursors[i].page); ray_free_raw(r->cursors[i].dict);
ray_free_raw(r->cursors[i].strings); ray_free_raw(r->cursors[i].symbol_ids);
ray_free_raw(r->cursors[i].symbols);
}
ray_free_raw(r->cursors); ray_free_raw(r->selected); ray_free_raw(r->excluded);
if (!r->borrowed) {
if (r->scan_domain) ray_sym_domain_release(r->scan_domain);
if (r->text_pattern) ray_release(r->text_pattern);
ray_free_raw(r->schema); ray_free_raw(r->groups);
if (r->map) ray_vm_unmap_file(r->map, r->size);
}
ray_free_raw(r);
}
ray_t* ray_parquet_open(const char* path, ray_t* columns, int64_t batch_rows, ray_parquet_t** out) {
*out = NULL;
if (batch_rows < 1 || batch_rows > PQ_BATCH) return pq_error("batch size must be 1..65536");
if (columns && (columns->type != RAY_SYM || !columns->len)) return pq_error("columns must be a nonempty symbol vector");
ray_parquet_t* r = ray_calloc_raw(sizeof(*r));
if (!r) return ray_error("oom", NULL);
const char* err = "invalid Parquet footer";
r->batch_rows = batch_rows; r->group = -1; r->filter_pos = -1;
r->map = ray_vm_map_file(path, &r->size);
if (!r->map) { ray_parquet_close(r); return ray_error("io", "parquet: cannot map %s", path); }
if (r->size < 12 || memcmp(r->map, "PAR1", 4) || memcmp(r->map+r->size-4, "PAR1", 4)) goto fail;
uint32_t flen = pq_u32(r->map+r->size-8);
if (flen > r->size-12 || flen > PQ_MAX_FOOTER) goto fail;
r->data_end = r->size-8-flen;
pq_span f[10];
if (!pq_fields((pq_span){r->map+r->data_end,r->map+r->size-8,12}, f, 10) ||
pq_get(f[1], -1) < 1 || !pq_num(f[3], &r->rows) || r->rows < 0 || f[2].type != 9 || f[4].type != 9 || f[8].type || f[9].type) goto fail;
pq_cur sc = {f[2].p,f[2].end,false}; uint8_t st;
uint64_t ns = pq_list(&sc, &st);
if (sc.bad || st != 12 || ns < 2 || ns > PQ_MAX_COLS+1) goto fail;
r->ncols = (int64_t)ns-1;
r->schema = ray_calloc_raw((size_t)r->ncols*sizeof(*r->schema));
if (!r->schema) { err = "schema allocation failed"; goto fail; }
pq_span root, sf[11];
if (!pq_element(&sc, 12, &root) || !pq_fields(root,sf,11) || pq_get(sf[5],-1) != r->ncols || sf[1].type) goto fail;
err = "unsupported or malformed schema (only flat primitive columns are supported)";
for (int64_t i = 0; i < r->ncols; i++) {
pq_span se; const uint8_t* name; size_t len;
if (!pq_element(&sc,12,&se) || !pq_fields(se,sf,11) || !pq_binary(sf[4],&name,&len) ||
!len || len > 1024 || memchr(name,0,len) || !pq_schema_type(&r->schema[i],sf)) goto fail;
r->schema[i].name = ray_sym_intern((const char*)name,len);
if (r->schema[i].name < 0) goto fail;
for (int64_t j = 0; j < i; j++) if (r->schema[j].name == r->schema[i].name) goto fail;
}
if (sc.p != sc.end) goto fail;
pq_cur gc = {f[4].p,f[4].end,false}; uint8_t gt;
uint64_t ng = pq_list(&gc,>);
if (gc.bad || gt != 12 || ng > PQ_MAX_FOOTER/sizeof(pq_span)) goto fail;
r->ngroups = (int64_t)ng;
r->groups = ray_calloc_raw((size_t)(ng ? ng : 1)*sizeof(*r->groups));
if (!r->groups) goto fail;
int64_t total = 0;
err = "invalid row group metadata";
for (uint64_t i = 0; i < ng; i++) {
pq_span gf[8]; int64_t nr;
if (!pq_element(&gc,12,&r->groups[i]) || !pq_fields(r->groups[i],gf,8) ||
!pq_num(gf[3],&nr) || nr < 0 || nr > r->rows-total || gf[1].type != 9) goto fail;
pq_cur chunks = {gf[1].p,gf[1].end,false}; uint8_t chunk_type;
if (pq_list(&chunks,&chunk_type) != (uint64_t)r->ncols || chunks.bad || chunk_type != 12) goto fail;
total += nr;
}
if (gc.p != gc.end || total != r->rows) goto fail;
r->nselected = columns ? columns->len : r->ncols;
r->noutput = r->nselected;
if (r->nselected > r->ncols) goto fail;
r->selected = ray_alloc_raw((size_t)r->ncols*sizeof(*r->selected));
r->cursors = ray_calloc_raw((size_t)r->ncols*sizeof(*r->cursors));
if (!r->selected || !r->cursors) goto fail;
err = "unknown or duplicate projection column";
for (int64_t i = 0; i < r->nselected; i++) {
ray_t* cs = columns ? ray_sym_vec_cell(columns,i) : NULL;
int64_t name = columns ? (cs ? ray_sym_find(ray_str_ptr(cs),ray_str_len(cs)) : -1) : r->schema[i].name;
int64_t j = 0;
while (j < r->ncols && r->schema[j].name != name) j++;
if (j == r->ncols) goto fail;
for (int64_t k = 0; k < i; k++) if (r->selected[k] == j) goto fail;
r->selected[i] = (int32_t)j;
}
*out = r; return NULL;
fail:
ray_parquet_close(r); return pq_error(err);
}
ray_t* ray_parquet_range(ray_parquet_t* r, int64_t column, int64_t lo, int64_t hi) {
if (!r || r->group != -1 || r->emitted || r->filter_pos >= 0 || lo > hi)
return pq_error("range must be configured once before scanning, with lo <= hi");
int64_t c = 0;
while (c < r->ncols && r->schema[c].name != column) c++;
if (c == r->ncols) return pq_error("unknown range column");
pq_schema* s = &r->schema[c];
if ((s->type != RAY_I16 && s->type != RAY_I32 && s->type != RAY_I64) ||
(s->converted != -1 && (s->converted < 15 || s->converted > 18)))
return pq_error("range requires a signed integer column");
int64_t pos = 0;
while (pos < r->nselected && r->selected[pos] != c) pos++;
if (pos == r->nselected) r->selected[r->nselected++] = (int32_t)c;
r->filter_pos = pos; r->filter_lo = lo; r->filter_hi = hi;
return NULL;
}
int64_t ray_parquet_groups_skipped(const ray_parquet_t* r) { return r ? r->skipped : 0; }
static bool pq_disjoint(ray_parquet_t* r, pq_span metadata) {
pq_span mf[17], stats[10];
if (!pq_fields(metadata,mf,17) || !pq_fields(mf[12],stats,10)) return false;
const uint8_t *minp, *maxp; size_t minn, maxn;
pq_span low = stats[6].type ? stats[6] : stats[2];
pq_span high = stats[5].type ? stats[5] : stats[1];
if (!pq_binary(low,&minp,&minn) || !pq_binary(high,&maxp,&maxn)) return false;
int physical = r->schema[r->selected[r->filter_pos]].physical;
size_t width = physical == 1 ? 4 : 8;
if (minn != width || maxn != width) return false;
int64_t lo = width == 4 ? (int32_t)pq_u32(minp) : (int64_t)pq_u64(minp);
int64_t hi = width == 4 ? (int32_t)pq_u32(maxp) : (int64_t)pq_u64(maxp);
return lo <= hi && (hi < r->filter_lo || lo > r->filter_hi);
}
static bool pq_bounds_disjoint(ray_parquet_t* r, pq_span low, pq_span high) {
const uint8_t *a, *b; size_t an, bn;
size_t width = r->schema[r->selected[r->filter_pos]].physical == 1 ? 4 : 8;
if (!pq_binary(low,&a,&an) || !pq_binary(high,&b,&bn) || an != width || bn != width) return false;
int64_t lo = width == 4 ? (int32_t)pq_u32(a) : (int64_t)pq_u64(a);
int64_t hi = width == 4 ? (int32_t)pq_u32(b) : (int64_t)pq_u64(b);
return lo <= hi && (hi < r->filter_lo || lo > r->filter_hi);
}
static bool pq_index_span(ray_parquet_t* r, pq_span off, pq_span len, pq_span* out) {
int64_t o, n;
if (!pq_num(off,&o) || !pq_num(len,&n) || o < 4 || n <= 0 || n > PQ_MAX_FOOTER ||
(uint64_t)o > r->data_end || (uint64_t)n > r->data_end-(size_t)o) return false;
*out = (pq_span){r->map+o,r->map+o+n,12}; return true;
}
static uint64_t pq_rot(uint64_t x, unsigned n) { return (x<<n) | (x>>(64-n)); }
static uint64_t pq_integer_hash(uint64_t x, unsigned width) {
const uint64_t p1 = UINT64_C(11400714785074694791), p2 = UINT64_C(14029467366897019727);
const uint64_t p3 = UINT64_C(1609587929392839161), p4 = UINT64_C(9650029242287828579);
uint64_t h = UINT64_C(2870177450012600261)+width;
if (width == 8) { h ^= pq_rot(x*p2,31)*p1; h = pq_rot(h,27)*p1+p4; }
else { h ^= (uint64_t)(uint32_t)x*p1; h = pq_rot(h,23)*p2+p3; }
h ^= h>>33; h *= p2; h ^= h>>29; h *= p3; return h ^ (h>>32);
}
static bool pq_bloom_absent(ray_parquet_t* r, pq_span metadata) {
if (r->filter_lo != r->filter_hi) return false;
pq_span m[17], f[5]; int64_t offset;
if (!pq_fields(metadata,m,17) || !pq_num(m[14],&offset) || offset < 4 ||
(uint64_t)offset >= r->data_end) return false;
size_t avail = r->data_end-(size_t)offset;
if (avail > PQ_MAX_PAGE+1024) avail = PQ_MAX_PAGE+1024;
int64_t length;
if (m[15].type) {
if (!pq_num(m[15],&length) || length <= 0 || (uint64_t)length > avail) return false;
avail = (size_t)length;
}
pq_cur c = {r->map+offset,r->map+offset+avail,false};
const uint8_t* start = c.p;
if (!pq_skip(&c,12,0,false) || !pq_fields((pq_span){start,c.p,12},f,5)) return false;
int64_t n = pq_get(f[1],-1);
if (n < 32 || n > PQ_MAX_PAGE || n%32 || (uint64_t)n > (uint64_t)(c.end-c.p)) return false;
if (m[15].type && n != c.end-c.p) return false;
for (int i = 2; i <= 4; i++) {
static const uint8_t supported[] = {0x1c,0,0};
if (f[i].type != 12 || f[i].end-f[i].p != 3 || memcmp(f[i].p,supported,3)) return false;
}
unsigned width = r->schema[r->selected[r->filter_pos]].physical == 1 ? 4 : 8;
if (width == 4 && (r->filter_lo < INT32_MIN || r->filter_lo > INT32_MAX)) return true;
uint64_t h = pq_integer_hash((uint64_t)r->filter_lo,width);
const uint8_t* block = c.p+(((h>>32)*(uint64_t)(n/32))>>32)*32;
static const uint32_t salt[8] = {0x47b6137b,0x44974d91,0x8824ad5b,0xa2b7289d,0x705495c7,0x2df1424b,0x9efc4947,0x5c6bfb31};
for (int i = 0; i < 8; i++)
if (!(pq_u32(block+4*i) & (UINT32_C(1)<<(((uint32_t)h*salt[i])>>27)))) return true;
return false;
}
static void pq_page_intervals(ray_parquet_t* r, pq_span chunk) {
pq_span cf[10], os, cs, of[3], f[9];
if (!r->group_rows || !pq_fields(chunk,cf,10) ||
!pq_index_span(r,cf[4],cf[5],&os) || !pq_index_span(r,cf[6],cf[7],&cs) ||
!pq_fields(os,of,3) || !pq_fields(cs,f,9) || of[1].type != 9 ||
f[1].type != 9 || f[2].type != 9 || f[3].type != 9 || pq_get(f[4],-1) < 0 || pq_get(f[4],3) > 2) return;
pq_cur loc = {of[1].p,of[1].end,false}, nulls = {f[1].p,f[1].end,false};
pq_cur min = {f[2].p,f[2].end,false}, max = {f[3].p,f[3].end,false};
uint8_t lt, nt, at, bt;
uint64_t n = pq_list(&loc,<), nn = pq_list(&nulls,&nt), na = pq_list(&min,&at), nb = pq_list(&max,&bt);
if (loc.bad || nulls.bad || min.bad || max.bad || !n || n > PQ_MAX_FOOTER/sizeof(pq_interval) ||
n != nn || n != na || n != nb || lt != 12 || nt != 1 || at != 8 || bt != 8) return;
pq_interval* ranges = ray_alloc_raw((size_t)n*sizeof(*ranges));
if (!ranges) return;
int64_t expected_row = 0, nranges = 0;
const pq_column* col = &r->cursors[r->filter_pos];
const uint8_t* expected_page = NULL;
for (uint64_t i = 0; i < n; i++) {
pq_span ls, lf[4], low, high, hf[10], dh[9];
if (!pq_element(&loc,12,&ls) || !pq_fields(ls,lf,4) ||
!pq_element(&min,8,&low) || !pq_element(&max,8,&high)) goto invalid;
uint8_t isnull = pq_byte(&nulls);
int64_t off = pq_get(lf[1],-1), bytes = pq_get(lf[2],-1), row = pq_get(lf[3],-1);
if ((isnull != 1 && isnull != 2) || row != expected_row || off < 4 || bytes <= 0 ||
(uint64_t)off > r->data_end || (uint64_t)bytes > r->data_end-(size_t)off) goto invalid;
const uint8_t* p = r->map+off;
if (p < col->chunk.p || p+bytes > col->chunk.end || (expected_page && p != expected_page)) goto invalid;
if (!i) {
pq_span mf[17];
if (!pq_fields(col->metadata,mf,17) || off != pq_get(mf[9],-1)) goto invalid;
}
pq_cur page = {p,p+bytes,false};
if (!pq_skip(&page,12,0,false) || !pq_fields((pq_span){p,page.p,12},hf,10)) goto invalid;
int64_t kind = pq_get(hf[1],-1), count;
if ((kind != 0 && kind != 3) || pq_get(hf[3],-1) != page.end-page.p ||
!pq_fields(hf[kind == 0 ? 5 : 8],dh,9) || !pq_num(dh[1],&count) || count <= 0 ||
count > r->group_rows-row) goto invalid;
expected_row += count; expected_page = page.end;
if (isnull == 1 || pq_bounds_disjoint(r,low,high)) {
if (nranges && ranges[nranges-1].hi == row) ranges[nranges-1].hi += count;
else ranges[nranges++] = (pq_interval){row,row+count};
}
}
if (expected_row != r->group_rows || expected_page != col->chunk.end || loc.p != loc.end ||
nulls.bad || nulls.p != nulls.end || min.p != min.end || max.p != max.end) goto invalid;
r->excluded = ranges; r->nexcluded = nranges; return;
invalid:
ray_free_raw(ranges);
}
static bool pq_group_columns(ray_parquet_t* r, int64_t g, pq_span* cols, int64_t* rows, pq_span* filter_chunk) {
pq_span gf[8];
if (!pq_fields(r->groups[g],gf,8) || !pq_num(gf[3],rows) || gf[1].type != 9) return false;
pq_cur c = {gf[1].p,gf[1].end,false}; uint8_t type;
if (pq_list(&c,&type) != (uint64_t)r->ncols || c.bad || type != 12) return false;
for (int64_t i = 0; i < r->ncols; i++) {
pq_span ch, cf[10], mf[17];
if (!pq_element(&c,12,&ch) || !pq_fields(ch,cf,10) || cf[1].type || cf[8].type || cf[9].type ||
!pq_fields(cf[3],mf,17) || pq_get(mf[1],-1) != r->schema[i].physical ||
pq_get(mf[5],-1) != *rows || mf[3].type != 9) return false;
pq_cur path = {mf[3].p,mf[3].end,false}; uint8_t pt; pq_span ps; const uint8_t* p; size_t n;
if (pq_list(&path,&pt) != 1 || pt != 8 || !pq_element(&path,pt,&ps) ||
!pq_binary(ps,&p,&n) || path.p != path.end) return false;
const char* name = ray_str_ptr(ray_sym_str(r->schema[i].name));
if (!name || strlen(name) != n || memcmp(name,p,n)) return false;
cols[i] = cf[3];
if (filter_chunk && i == r->selected[r->filter_pos]) *filter_chunk = ch;
}
return c.p == c.end;
}
static int64_t pq_dict_max_len(ray_parquet_t* r, pq_span meta) {
pq_span mf[17];
if (!pq_fields(meta,mf,17)) return 0;
int64_t dict = pq_get(mf[11],-1), codec = pq_get(mf[4],-1);
if (dict < 4 || (uint64_t)dict >= r->data_end) return 0;
pq_cur c = {r->map+dict,r->map+r->data_end,false};
const uint8_t* start = c.p;
if (!pq_skip(&c,12,0,false)) return 0;
pq_span f[10];
if (!pq_fields((pq_span){start,c.p,12},f,10) || pq_get(f[1],-1) != 2) return 0;
int64_t raw = pq_get(f[2],-1), size = pq_get(f[3],-1);
if (raw < 0 || raw > PQ_MAX_PAGE || size < 0 || size > PQ_MAX_PAGE ||
(uint64_t)size > (uint64_t)(c.end-c.p)) return 0;
const uint8_t* p = c.p; uint8_t* buf = NULL;
if (codec == 1) {
buf = ray_alloc_raw((size_t)(raw ? raw : 1));
if (!buf || !ray_parquet_snappy(c.p,(size_t)size,buf,(size_t)raw)) { ray_free_raw(buf); return 0; }
p = buf;
} else if (raw != size) return 0;
int64_t max = 0;
for (const uint8_t* q = p, *end = p+raw; end-q >= 4; ) {
uint32_t len = pq_u32(q); q += 4;
if (len > (uint64_t)(end-q)) break;
if ((int64_t)len > max) max = len;
q += len;
}
ray_free_raw(buf);
return max;
}
static int64_t pq_group_batch(ray_parquet_t* r, int64_t g) {
int64_t batch = r->batch_rows, rows;
pq_span* cols = ray_alloc_raw((size_t)r->ncols*sizeof(*cols));
if (!cols) return batch;
if (pq_group_columns(r,g,cols,&rows,NULL)) {
for (int64_t i = 0; i < r->nselected; i++) {
const pq_schema* sc = &r->schema[r->selected[i]];
if (sc->physical != 6 || sc->type != RAY_STR || sc->import_domain) continue;
pq_span mf[17];
if (!pq_fields(cols[r->selected[i]],mf,17)) continue;
int64_t values = pq_get(mf[5],-1), bytes = pq_get(mf[6],-1);
if (values <= 0 || bytes <= 0) continue;
int64_t per_row = bytes/values + 1;
int64_t longest = pq_dict_max_len(r,cols[r->selected[i]]) + 1;
if (longest > per_row) per_row = longest;
int64_t cap = (int64_t)PQ_POOL_TARGET/per_row;
if (cap < 1) cap = 1;
if (cap < batch) batch = cap;
}
}
ray_free_raw(cols);
return batch;
}
static const char* pq_start_group(ray_parquet_t* r) {
pq_span* cols = ray_alloc_raw((size_t)r->ncols*sizeof(*cols));
if (!cols) return "row group allocation failed";
pq_span filter_chunk = {0};
ray_free_raw(r->excluded); r->excluded = NULL; r->nexcluded = r->exclude_pos = 0;
if (!pq_group_columns(r,r->group,cols,&r->group_left,r->filter_pos >= 0 ? &filter_chunk : NULL)) { ray_free_raw(cols); return "invalid column chunk metadata"; }
if (r->group_left == 0) { ray_free_raw(cols); return NULL; }
r->group_batch = pq_group_batch(r,r->group);
if (r->filter_pos >= 0 && !r->filter_nulls && pq_disjoint(r,cols[r->selected[r->filter_pos]])) {
r->group_left = 0; r->skipped++; ray_free_raw(cols); return NULL;
}
if (r->filter_pos >= 0 && !r->filter_nulls && pq_bloom_absent(r,cols[r->selected[r->filter_pos]])) {
r->group_left = 0; r->skipped++; r->bloom_skipped++; ray_free_raw(cols); return NULL;
}
r->group_rows = r->group_left;
const char* err = NULL;
for (int64_t i = 0; i < r->nselected; i++) {
pq_column* c = &r->cursors[i]; pq_span mf[17];
c->metadata = cols[r->selected[i]];
if (!pq_fields(c->metadata,mf,17)) { err = "invalid column metadata"; break; }
int64_t offset = pq_get(mf[9],-1), dict = pq_get(mf[11],-1), bytes = pq_get(mf[7],-1);
int64_t codec = pq_get(mf[4],-1);
if (codec != 0 && codec != 1) { err = "unsupported compression codec (supported: uncompressed, Snappy)"; break; }
if (dict >= 0) { if (dict > offset) { err = "dictionary follows data"; break; } offset = dict; }
if (offset < 4 || bytes < 0 || (uint64_t)offset > r->data_end || (uint64_t)bytes > r->data_end-(size_t)offset) {
err = "column chunk outside file data"; break;
}
c->chunk = (pq_cur){r->map+offset,r->map+offset+bytes,false};
c->codec = (int)codec; c->page_left = 0; c->chunk_left = r->group_left;
c->have_dict = false; c->dict_count = 0; c->optional = r->schema[r->selected[i]].optional;
ray_free_raw(c->strings); c->strings = NULL;
ray_free_raw(c->symbol_ids); c->symbol_ids = NULL;
}
if (!err && r->filter_pos >= 0 && !r->filter_nulls) pq_page_intervals(r,filter_chunk);
ray_free_raw(cols); return err;
}
static bool pq_unpack(pq_column* c, const uint8_t* p, size_t n, size_t decoded, bool compressed, const uint8_t** out) {
if (decoded > PQ_MAX_PAGE) return false;
if (!compressed || !c->codec) { if (n != decoded) return false; *out = p; return true; }
if (!pq_reserve(&c->page,&c->page_cap,decoded)) return false;
if (!ray_parquet_snappy(p,n,c->page,decoded)) return false;
*out = c->page; return true;
}
static bool pq_dictionary(pq_column* c, pq_schema* s, const uint8_t* data, size_t n, int64_t count) {
if (c->have_dict || count < 0 || (uint64_t)count > PQ_MAX_PAGE/sizeof(pq_string)) return false;
if (!pq_reserve(&c->dict,&c->dict_cap,n)) return false;
if (n) memcpy(c->dict,data,n);
c->dict_count = count; c->have_dict = true;
if (s->physical == 6) {
c->strings = ray_calloc_raw((size_t)(count ? count : 1)*sizeof(*c->strings));
if (!c->strings) return false;
pq_cur p = {c->dict,c->dict+n,false};
for (int64_t i = 0; i < count; i++) {
const uint8_t* len; const uint8_t* str;
if (!pq_take(&p,4,&len)) return false;
uint32_t size = pq_u32(len);
if (!pq_take(&p,size,&str)) return false;
if (s->strict && s->converted == 0 && !pq_utf8(str,size)) return false;
c->strings[i] = (pq_string){str,size};
}
if (p.p != p.end) return false;
if (s->import_domain) {
c->symbol_ids = ray_alloc_raw((size_t)(count ? count : 1)*sizeof(*c->symbol_ids));
if (!c->symbol_ids) return false;
pq_symbol_scratch* scratch = pq_symbols(c);
if (!scratch) return false;
const char** strings = scratch->strings; size_t* lengths = scratch->lengths;
uint32_t* hashes = scratch->hashes;
for (int64_t off = 0; off < count; off += 8192) {
int64_t n = count-off < 8192 ? count-off : 8192;
for (int64_t i = 0; i < n; i++) {
strings[i] = (const char*)c->strings[off+i].p; lengths[i] = c->strings[off+i].n;
hashes[i] = (uint32_t)ray_hash_bytes(strings[i],lengths[i]);
}
if (!ray_sym_domain_intern_batch(s->import_domain,n,strings,lengths,hashes,c->symbol_ids+off)) return false;
}
}
return true;
}
size_t width = s->physical == 1 || s->physical == 4 ? 4 : 8;
return s->physical != 0 && (uint64_t)count*width == n;
}
static const char* pq_page(pq_column* c, pq_schema* s, int64_t* skip) {
while (c->chunk.p < c->chunk.end) {
const uint8_t* start = c->chunk.p;
if (!pq_skip(&c->chunk,12,0,false)) return "invalid page header";
pq_span f[10];
if (!pq_fields((pq_span){start,c->chunk.p,12},f,10)) return "invalid page fields";
int64_t kind = pq_get(f[1],-1), raw = pq_get(f[2],-1), size = pq_get(f[3],-1);
if (raw < 0 || raw > PQ_MAX_PAGE || size < 0 || size > PQ_MAX_PAGE) return "invalid or oversized page (limit 64 MiB)";
const uint8_t* payload;
if (!pq_take(&c->chunk,(uint64_t)size,&payload)) return "truncated page";
if (skip && (kind == 0 || kind == 3)) {
pq_span h[9]; int64_t count;
if (!pq_fields(f[kind == 0 ? 5 : 8],h,9) || !pq_num(h[1],&count) ||
count <= 0 || count > c->chunk_left) return "invalid skipped page row count";
if (count <= *skip) {
*skip -= count; c->chunk_left -= count; c->pages_skipped++;
if (!c->chunk_left && c->chunk.p != c->chunk.end) return "excess column pages";
if (!*skip) return NULL;
continue;
}
}
if (f[4].type) {
int64_t expected;
if (!pq_num(f[4],&expected)) return "invalid page checksum";
if (ray_crc32(0,payload,(size_t)size) != (uint32_t)expected) return "page checksum mismatch";
}
if (kind == 1) continue;
const uint8_t* data = NULL; pq_span h[9]; int64_t count, encoding;
c->bool_bit = 0; c->bool_rle = false; c->page_nulls = 0; c->expected_nulls = -1;
if (kind == 2) {
if (!pq_fields(f[7],h,9) || !pq_num(h[1],&count) ||
(pq_get(h[2],-1) != 0 && pq_get(h[2],-1) != 2) ||
!pq_unpack(c,payload,(size_t)size,(size_t)raw,true,&data) ||
!pq_dictionary(c,s,data,(size_t)raw,count)) return "invalid dictionary page";
continue;
}
if (kind == 0) {
if (!pq_fields(f[5],h,9) || !pq_num(h[1],&count) || !pq_num(h[2],&encoding) ||
(c->optional && pq_get(h[3],-1) != 3) ||
!pq_unpack(c,payload,(size_t)size,(size_t)raw,true,&data)) return "invalid v1 data page";
c->values = (pq_cur){data,data+raw,false};
if (c->optional) {
const uint8_t* p; const uint8_t* levels;
if (!pq_take(&c->values,4,&p) || !pq_take(&c->values,pq_u32(p),&levels)) return "truncated definition levels";
c->defs = (pq_rle){.c={levels,c->values.p,false},.width=1};
}
} else if (kind == 3) {
int64_t dl, rl, nulls, nr;
if (!pq_fields(f[8],h,9) || !pq_num(h[1],&count) || !pq_num(h[2],&nulls) ||
!pq_num(h[3],&nr) || nr != count || nulls < 0 || nulls > count ||
(!c->optional && nulls) || !pq_num(h[4],&encoding) || !pq_num(h[5],&dl) ||
!pq_num(h[6],&rl) || dl < 0 || rl != 0 || dl > size || dl > raw ||
(!c->optional && dl) || (h[7].type && h[7].type != 1 && h[7].type != 2)) return "invalid v2 data page";
c->defs = (pq_rle){.c={payload,payload+dl,false},.width=1};
if (!pq_unpack(c,payload+dl,(size_t)(size-dl),(size_t)(raw-dl),h[7].type != 2,&data)) return "invalid v2 compressed values";
c->values = (pq_cur){data,data+raw-dl,false};
c->expected_nulls = nulls;
} else return "unsupported page type";
if (count <= 0 || count > c->chunk_left) return "invalid page row count";
if (encoding == 3 && s->physical == 0) {
const uint8_t* lp; const uint8_t* body;
if (!pq_take(&c->values,4,&lp) || !pq_take(&c->values,pq_u32(lp),&body) ||
c->values.p != c->values.end) return "invalid RLE boolean data";
c->ids = (pq_rle){.c={body,c->values.end,false},.width=1};
c->bool_rle = true; c->encoding = 0; c->page_left = count;
return NULL;
}
if (encoding != 0 && encoding != 2 && encoding != 8) return "unsupported value encoding (supported: plain, dictionary, RLE booleans)";
c->encoding = (int)encoding; c->page_left = count;
if (encoding) {
if (!c->have_dict || s->physical == 0) return "missing or invalid dictionary";
unsigned width = pq_byte(&c->values);
if (c->values.bad || width > 32) return "invalid dictionary bit width";
c->ids = (pq_rle){.c=c->values,.width=width};
}
return NULL;
}
return "missing data page";
}
#define PQ_PUT_FAIL "value conversion or allocation failed"
static const char* pq_put_value(ray_t** vp, int64_t row, pq_schema* s, const uint8_t* p, uint32_t len) {
ray_t* v = *vp;
switch (s->physical) {
case 0:
if (v->type == RAY_BOOL) ((uint8_t*)ray_data(v))[row] = *p;
else ((int16_t*)ray_data(v))[row] = *p;
break;
case 1: case 2: {
int64_t x = s->physical == 1 ? (int32_t)pq_u32(p) : (int64_t)pq_u64(p);
if (s->physical == 1 && s->converted >= 11 && s->converted <= 13) x = pq_u32(p);
if ((s->converted == 11 && x > UINT8_MAX) ||
(s->converted == 12 && x > UINT16_MAX) ||
(s->converted == 15 && (x < INT8_MIN || x > INT8_MAX)) ||
(s->converted == 16 && (x < INT16_MIN || x > INT16_MAX))) return PQ_PUT_FAIL;
if (s->type == RAY_DATE && __builtin_sub_overflow(x,PQ_EPOCH_DAYS,&x)) return PQ_PUT_FAIL;
int64_t raw = x;
if (s->type == RAY_TIMESTAMP &&
(__builtin_sub_overflow(x,PQ_EPOCH_NS/s->scale,&x) || __builtin_mul_overflow(x,s->scale,&x))) {
if (s->strict) return "timestamp outside the native nanosecond range (strict mode)";
x = raw >= PQ_EPOCH_NS/s->scale ? INT64_MAX : INT64_MIN+1;
}
if (s->type == RAY_I16) {
if (x < INT16_MIN || x > INT16_MAX) return PQ_PUT_FAIL;
((int16_t*)ray_data(v))[row] = (int16_t)x;
} else if (s->type == RAY_I64 || s->type == RAY_TIMESTAMP) ((int64_t*)ray_data(v))[row] = x;
else {
if (x < INT32_MIN || x > INT32_MAX) return PQ_PUT_FAIL;
((int32_t*)ray_data(v))[row] = (int32_t)x;
}
break;
}
case 4: { uint32_t u = pq_u32(p); memcpy((float*)ray_data(v)+row,&u,4); break; }
case 5: { uint64_t u = pq_u64(p); memcpy((double*)ray_data(v)+row,&u,8); break; }
case 6: {
if (s->strict && s->converted == 0 && !pq_utf8(p,len)) return PQ_PUT_FAIL;
ray_t* next = ray_str_vec_set(v,row,(const char*)p,len);
if (!next || RAY_IS_ERR(next)) { if (next) ray_release(next); return PQ_PUT_FAIL; }
*vp = v = next; break;
}
default: return PQ_PUT_FAIL;
}
uint8_t attrs = v->attrs;
v->attrs |= RAY_ATTR_HAS_NULLS;
if (!ray_vec_is_null(v,row)) v->attrs = attrs;
else if (s->strict) return PQ_PUT_FAIL;
return NULL;
}
static void pq_plain_int(ray_t* v, int64_t row, const uint8_t* p, int64_t n) {
int64_t i = 0;
size_t width = v->type == RAY_I32 ? 4 : 8;
uint8_t* dst = (uint8_t*)ray_data(v)+(size_t)row*width;
#if defined(__SSE2__)
__m128i sentinel = width == 4 ? _mm_set1_epi32(INT32_MIN) : _mm_set_epi32(INT32_MIN,0,INT32_MIN,0);
for (; i+16/(int64_t)width <= n; i += 16/(int64_t)width) {
__m128i x = _mm_loadu_si128((const __m128i*)(p+(size_t)i*width));
_mm_storeu_si128((__m128i*)(dst+(size_t)i*width),x);
int mask = _mm_movemask_epi8(_mm_cmpeq_epi32(x,sentinel));
if ((width == 4 && mask) || (width == 8 && ((mask&255) == 255 || (mask>>8) == 255)))
v->attrs |= RAY_ATTR_HAS_NULLS;
}
#endif
for (; i < n; i++) {
if (width == 4) {
int32_t x = (int32_t)pq_u32(p+(size_t)i*4); memcpy(dst+(size_t)i*4,&x,4);
if (x == INT32_MIN) v->attrs |= RAY_ATTR_HAS_NULLS;
} else {
int64_t x = (int64_t)pq_u64(p+(size_t)i*8); memcpy(dst+(size_t)i*8,&x,8);
if (x == INT64_MIN) v->attrs |= RAY_ATTR_HAS_NULLS;
}
}
}
static const char* pq_decode_symbols(pq_column* c, pq_schema* s, ray_t** out, int64_t rows) {
pq_symbol_scratch* scratch = pq_symbols(c);
if (!scratch) return "symbol scratch allocation failed";
ray_t* v = ray_sym_vec_new(RAY_SYM_W32,rows);
if (!v || RAY_IS_ERR(v)) { if (v) ray_release(v); return "symbol vector allocation failed"; }
v->sym_domain = s->import_domain; ray_sym_domain_retain(v->sym_domain);
v->len = rows; *out = v; uint32_t* dst = ray_data(v);
const char** strings = scratch->strings; size_t* lengths = scratch->lengths;
uint32_t* hashes = scratch->hashes;
int64_t* ids = scratch->ids; int64_t* positions = scratch->positions;
for (int64_t at = 0; at < rows;) {
if (ray_interrupted()) return "scan interrupted";
if (!c->page_left) { const char* err = pq_page(c,s,NULL); if (err) return err; }
int64_t n = rows-at < c->page_left ? rows-at : c->page_left, count = 0;
if (n > 8192) n = 8192;
for (int64_t i = 0; i < n; i++) {
if (!c->optional && c->encoding && c->ids.left && !c->ids.packed) {
uint32_t id = c->ids.value;
if (id >= c->dict_count || !c->symbol_ids) return "invalid symbol dictionary index";
int64_t pos = c->symbol_ids[id];
if (pos < 0 || (uint64_t)pos >= UINT32_MAX) return "symbol domain exceeds W32";
if (!pos && s->strict) return "nonnull empty text collides with native null";
int64_t count = n-i < (int64_t)c->ids.left ? n-i : (int64_t)c->ids.left;
for (int64_t j = 0; j < count; j++) dst[at+i+j] = (uint32_t)pos;
if (!pos) v->attrs |= RAY_ATTR_HAS_NULLS;
c->ids.left -= (uint64_t)count; i += count-1; continue;
}
uint32_t present = 1;
if (c->optional && (!pq_rle_next(&c->defs,&present) || present > 1)) return "invalid definition levels";
if (!present) { dst[at+i] = 0; v->attrs |= RAY_ATTR_HAS_NULLS; c->page_nulls++; continue; }
if (c->encoding) {
uint32_t id;
if (!pq_rle_next(&c->ids,&id) || id >= c->dict_count || !c->symbol_ids) return "invalid symbol dictionary index";
int64_t pos = c->symbol_ids[id];
if (pos < 0 || (uint64_t)pos >= UINT32_MAX) return "symbol domain exceeds W32";
if (!pos && s->strict) return "nonnull empty text collides with native null";
dst[at+i] = (uint32_t)pos;
if (!pos) v->attrs |= RAY_ATTR_HAS_NULLS;
} else {
const uint8_t *lp, *p;
if (!pq_take(&c->values,4,&lp)) return "truncated string length";
uint32_t len = pq_u32(lp);
if (!pq_take(&c->values,len,&p)) return "truncated string data";
if (s->strict && (!len || (s->converted == 0 && !pq_utf8(p,len))))
return "nonnull text is empty or invalid UTF-8";
strings[count] = (const char*)p; lengths[count] = len;
hashes[count] = (uint32_t)ray_hash_bytes(p,len); positions[count++] = at+i;
}
}
if (count) {
if (!ray_sym_domain_intern_batch(s->import_domain,count,strings,lengths,hashes,ids)) return "symbol domain allocation failed";
for (int64_t i = 0; i < count; i++) {
if (ids[i] < 0 || (uint64_t)ids[i] >= UINT32_MAX) return "symbol domain exceeds W32";
dst[positions[i]] = (uint32_t)ids[i]; if (!ids[i]) v->attrs |= RAY_ATTR_HAS_NULLS;
}
}
at += n; c->page_left -= n; c->chunk_left -= n;
if (!c->page_left) {
if (c->optional && !pq_rle_done(&c->defs)) return "excess definition levels";
if (c->expected_nulls >= 0 && c->page_nulls != c->expected_nulls) return "v2 null count mismatch";
if (c->encoding) { if (!pq_rle_done(&c->ids)) return "excess dictionary indices"; }
else if (c->values.p != c->values.end) return "excess string values";
if (!c->chunk_left && c->chunk.p != c->chunk.end) return "excess column pages";
}
}
return NULL;
}
static const char* pq_decode(pq_column* c, pq_schema* s, ray_t** out, int64_t rows) {
if (s->import_domain) return pq_decode_symbols(c,s,out,rows);
ray_t* v = ray_vec_new((int8_t)s->type,rows);
if (!v || RAY_IS_ERR(v)) { if (v) ray_release(v); return "vector allocation failed"; }
v->len = rows; *out = v;
if (v->type == RAY_STR) memset(ray_data(v),0,(size_t)rows*sizeof(ray_str_t));
for (int64_t i = 0; i < rows; i++) {
if (!(i&1023) && ray_interrupted()) return "scan interrupted";
if (!c->page_left) { const char* err = pq_page(c,s,NULL); if (err) return err; }
if (!c->optional && c->encoding && c->ids.left && !c->ids.packed) {
uint32_t id = c->ids.value;
if (id >= c->dict_count) return "dictionary index out of range";
int64_t n = rows-i < c->page_left ? rows-i : c->page_left;
if ((uint64_t)n > c->ids.left) n = (int64_t)c->ids.left;
if (n > 4096) n = 4096;
const uint8_t* p = s->physical == 6 ? c->strings[id].p :
c->dict+(size_t)id*((s->physical == 1 || s->physical == 4) ? 4 : 8);
uint32_t len = s->physical == 6 ? c->strings[id].n : 0;
if (len > PQ_MAX_PAGE || (s->physical == 6 && v->str_pool && (uint64_t)v->str_pool->len+len > PQ_POOL_MAX)) return PQ_PUT_FAIL;
const char* perr = pq_put_value(&v,i,s,p,len);
if (perr) return perr;
*out = v;
if (ray_elem_size(v->type) == 2) {
int16_t* dst = ray_data(v); int16_t x = dst[i];
for (int64_t j = 1; j < n; j++) dst[i+j] = x;
} else if (ray_elem_size(v->type) == 4) {
uint32_t* dst = ray_data(v); uint32_t x; memcpy(&x,dst+i,4);
for (int64_t j = 1; j < n; j++) memcpy(dst+i+j,&x,4);
} else if (ray_elem_size(v->type) == 8) {
uint64_t* dst = ray_data(v); uint64_t x; memcpy(&x,dst+i,8);
for (int64_t j = 1; j < n; j++) memcpy(dst+i+j,&x,8);
} else {
size_t width = ray_elem_size(v->type); uint8_t* dst = ray_data(v);
for (int64_t j = 1; j < n; j++) memcpy(dst+(size_t)(i+j)*width,dst+(size_t)i*width,width);
}
c->ids.left -= (uint64_t)n; c->page_left -= n-1; c->chunk_left -= n-1; i += n-1;
} else if (!c->optional && !c->encoding &&
((s->physical == 1 && s->type == RAY_I32 && (s->converted == -1 || s->converted == 17)) ||
(s->physical == 2 && s->type == RAY_I64))) {
int64_t n = rows-i < c->page_left ? rows-i : c->page_left;
if (n > 4096) n = 4096;
const uint8_t* p;
if (!pq_take(&c->values,(uint64_t)n*(s->physical == 1 ? 4 : 8),&p)) return "truncated numeric data";
pq_plain_int(v,i,p,n);
if (s->strict && (v->attrs & RAY_ATTR_HAS_NULLS)) return "nonnull integer collides with native null";
c->page_left -= n-1; c->chunk_left -= n-1; i += n-1;
} else {
uint32_t present = 1;
if (c->optional && (!pq_rle_next(&c->defs,&present) || present > 1)) return "invalid definition levels";
if (!present) {
if (ray_vec_set_null_checked(v,i,true) != RAY_OK) return "unrepresentable null";
c->page_nulls++;
} else {
const uint8_t* p; uint32_t len = 0; uint8_t b = 0;
if (c->encoding) {
uint32_t id;
if (!pq_rle_next(&c->ids,&id) || id >= c->dict_count) return "dictionary index out of range";
if (s->physical == 6) { p = c->strings[id].p; len = c->strings[id].n; }
else p = c->dict+(size_t)id*((s->physical == 1 || s->physical == 4) ? 4 : 8);
} else if (c->bool_rle) {
uint32_t bit;
if (!pq_rle_next(&c->ids,&bit)) return "truncated RLE boolean data";
b = (uint8_t)bit; p = &b;
} else if (s->physical == 0) {
if (c->bool_bit/8 >= (uint64_t)(c->values.end-c->values.p)) return "truncated boolean data";
b = (c->values.p[c->bool_bit/8]>>(c->bool_bit&7))&1; c->bool_bit++; p = &b;
} else if (s->physical == 6) {
const uint8_t* lp;
if (!pq_take(&c->values,4,&lp)) return "truncated string length";
len = pq_u32(lp);
if (!pq_take(&c->values,len,&p)) return "truncated string data";
} else if (!pq_take(&c->values,(s->physical == 1 || s->physical == 4) ? 4 : 8,&p)) return "truncated numeric data";
if (s->physical == 6 && v->str_pool && (uint64_t)v->str_pool->len+len > PQ_POOL_MAX) return "batch string pool exceeds 4 GiB";
if (len > PQ_MAX_PAGE) return PQ_PUT_FAIL;
const char* perr = pq_put_value(&v,i,s,p,len);
if (perr) return perr;
*out = v;
}
}
c->page_left--; c->chunk_left--;
if (!c->page_left) {
if (c->optional && !pq_rle_done(&c->defs)) return "excess definition levels";
if (c->expected_nulls >= 0 && c->page_nulls != c->expected_nulls) return "v2 null count mismatch";
if (c->encoding) { if (!pq_rle_done(&c->ids)) return "excess dictionary indices"; }
else if (c->bool_rle) { if (!pq_rle_done(&c->ids)) return "excess RLE boolean data"; }
else if (s->physical == 0) { if ((c->bool_bit+7)/8 != (uint64_t)(c->values.end-c->values.p)) return "excess boolean data"; }
else if (c->values.p != c->values.end) return "excess plain data";
if (!c->chunk_left && c->chunk.p != c->chunk.end) return "excess column pages";
}
}
return NULL;
}
static const char* pq_discard(pq_column* c, pq_schema* s, int64_t rows) {
while (rows) {
if (ray_interrupted()) return "scan interrupted";
if (!c->page_left) {
const char* err = pq_page(c,s,&rows);
if (err) return err;
if (!rows) break;
}
int64_t n = rows < c->page_left ? rows : c->page_left;
if (n > PQ_BATCH) n = PQ_BATCH;
ray_t* discard = NULL;
const char* err = pq_decode(c,s,&discard,n);
if (discard) ray_release(discard);
if (err) return err;
rows -= n;
}
return NULL;
}
typedef struct {
ray_parquet_t* reader;
ray_t** columns;
const char** errors;
int64_t rows;
} pq_work;
static const char* pq_decode_column(ray_parquet_t* r, int64_t i, ray_t** out, int64_t rows) {
pq_schema schema = r->schema[r->selected[i]];
if (r->filter_pos >= 0 || r->text_pattern) schema.import_domain = NULL;
return pq_decode(&r->cursors[i],&schema,out,rows);
}
static const char* pq_intern_selected(pq_column* c, ray_t** vector, ray_sym_domain_t* domain) {
pq_symbol_scratch* scratch = pq_symbols(c);
if (!scratch) return "symbol scratch allocation failed";
ray_t* src = *vector;
ray_t* dst = ray_sym_vec_new(RAY_SYM_W32,src->len);
if (!dst || RAY_IS_ERR(dst)) { if (dst) ray_release(dst); return "symbol allocation failed"; }
dst->sym_domain = domain; ray_sym_domain_retain(domain); dst->len = src->len;
const char** strings = scratch->strings; size_t* lengths = scratch->lengths;
uint32_t* hashes = scratch->hashes; int64_t* ids = scratch->ids;
for (int64_t at = 0; at < src->len; at += 8192) {
int64_t n = src->len-at < 8192 ? src->len-at : 8192;
for (int64_t i = 0; i < n; i++) {
strings[i] = ray_str_vec_get(src,at+i,&lengths[i]);
hashes[i] = (uint32_t)ray_hash_bytes(strings[i],lengths[i]);
}
if (!ray_sym_domain_intern_batch(domain,n,strings,lengths,hashes,ids)) {
ray_release(dst); return "symbol domain allocation failed";
}
for (int64_t i = 0; i < n; i++) {
if (ids[i] < 0 || (uint64_t)ids[i] >= UINT32_MAX) { ray_release(dst); return "symbol domain exceeds W32"; }
((uint32_t*)ray_data(dst))[at+i] = (uint32_t)ids[i];
if (!ids[i]) dst->attrs |= RAY_ATTR_HAS_NULLS;
}
}
ray_release(src); *vector = dst; return NULL;
}
static void pq_decode_task(void* ptr, uint32_t worker, int64_t start, int64_t end) {
(void)worker;
pq_work* w = ptr; ray_parquet_t* r = w->reader;
for (int64_t i = start; i < end; i++) if (!w->columns[i])
w->errors[i] = pq_decode_column(r,i,&w->columns[i],w->rows);
}
int64_t ray_parquet_pages_skipped(const ray_parquet_t* r) {
int64_t n = 0;
if (r) for (int64_t i = 0; i < r->nselected; i++) n += r->cursors[i].pages_skipped;
return n;
}
int64_t ray_parquet_bloom_skipped(const ray_parquet_t* r) { return r ? r->bloom_skipped : 0; }
int64_t ray_parquet_parallel_batches(const ray_parquet_t* r) { return r ? r->parallel_batches : 0; }
ray_t* ray_parquet_next(ray_parquet_t* r) {
if (!r || r->failed) return pq_error("reader is closed or failed");
if (ray_interrupted()) { r->failed = true; return ray_error("cancel","parquet scan interrupted"); }
const char* err = NULL;
int64_t rows;
for (;;) {
while (!r->group_left && r->group+1 < r->ngroups) {
r->group++;
err = pq_start_group(r);
if (err) { r->failed = true; return pq_error(err); }
}
if (!r->group_left && r->emitted) return NULL;
int64_t cap = r->group_batch > 0 ? r->group_batch : r->batch_rows;
rows = r->group_left < cap ? r->group_left : cap;
if (r->group_left && r->exclude_pos < r->nexcluded) {
pq_interval range = r->excluded[r->exclude_pos];
int64_t at = r->group_rows-r->group_left;
if (at == range.lo) {
for (int64_t i = 0; i < r->nselected; i++) {
pq_schema schema = r->schema[r->selected[i]]; schema.import_domain = NULL;
err = pq_discard(&r->cursors[i],&schema,range.hi-range.lo);
if (err) { r->failed = true; return pq_error(err); }
}
r->group_left -= range.hi-range.lo; r->exclude_pos++;
continue;
}
if (rows > range.lo-at) rows = range.lo-at;
}
break;
}
ray_t* tbl = ray_table_new(r->noutput);
if (!tbl || RAY_IS_ERR(tbl)) return tbl ? tbl : ray_error("oom",NULL);
ray_t** cols = ray_calloc_raw((size_t)r->nselected*sizeof(*cols));
uint8_t* keep = NULL;
const char** errors = NULL;
int64_t bad_col = -1;
if (!cols) { ray_release(tbl); r->failed = true; return ray_error("oom",NULL); }
ray_pool_t* pool = rows >= 4096 && r->nselected > 1 ? ray_pool_get() : NULL;
if (ray_pool_par_dispatch_ok(pool,rows,4096)) {
errors = ray_calloc_raw((size_t)r->nselected*sizeof(*errors));
if (!errors) { err = "task allocation failed"; goto fail; }
pq_work work = {r,cols,errors,rows};
ray_pool_dispatch_n(pool,pq_decode_task,&work,(uint32_t)r->nselected);
r->parallel_batches++;
for (int64_t i = 0; i < r->nselected; i++) {
if (errors[i] || !cols[i]) { err = errors[i] ? errors[i] : "scan interrupted"; bad_col = i; goto fail; }
}
} else for (int64_t i = 0; i < r->nselected; i++) {
err = pq_decode_column(r,i,&cols[i],rows);
if (err) { bad_col = i; goto fail; }
}
int64_t kept = rows;
if (r->filter_pos >= 0 && rows) {
keep = ray_alloc_raw((size_t)rows);
if (!keep) { err = "selection allocation failed"; goto fail; }
ray_t* f = cols[r->filter_pos]; kept = 0;
for (int64_t i = 0; i < rows; i++) {
int64_t value = f->type == RAY_I16 ? ((int16_t*)ray_data(f))[i] :
f->type == RAY_I32 ? ((int32_t*)ray_data(f))[i] : ((int64_t*)ray_data(f))[i];
keep[i] = ray_vec_is_null(f,i) ? r->filter_nulls : value >= r->filter_lo && value <= r->filter_hi;
kept += keep[i];
}
}
if (r->text_pattern && rows) {
if (!keep) keep = ray_alloc_raw((size_t)rows);
if (!keep) { err = "selection allocation failed"; goto fail; }
const char* pat = ray_str_ptr(r->text_pattern); size_t pn = ray_str_len(r->text_pattern);
ray_glob_compiled_t compiled = ray_glob_compile(pat,pn);
kept = 0;
for (int64_t i = 0; i < rows; i++) {
size_t len; const char* str = ray_str_vec_get(cols[r->text_pos],i,&len);
keep[i] = compiled.shape == RAY_GLOB_SHAPE_NONE ? ray_glob_match(str,len,pat,pn) :
ray_glob_match_compiled(&compiled,str,len);
kept += keep[i];
}
}
for (int64_t c = 0; c < r->noutput; c++) {
ray_t* v = cols[c];
if (kept != rows) {
size_t esz = ray_sym_elem_size(v->type,v->attrs); uint8_t* data = ray_data(v); int64_t dst = 0;
for (int64_t i = 0; i < rows; i++) if (keep[i]) {
if (dst != i) memcpy(data+(size_t)dst*esz,data+(size_t)i*esz,esz);
dst++;
}
v->len = kept;
}
if (v->type == RAY_STR && r->schema[r->selected[c]].import_domain) {
err = pq_intern_selected(&r->cursors[c],&cols[c],r->schema[r->selected[c]].import_domain);
if (err) goto fail;
v = cols[c];
}
tbl = ray_table_add_col(tbl,r->schema[r->selected[c]].name,v);
if (!tbl || RAY_IS_ERR(tbl)) { err = "table allocation failed"; goto fail; }
}
r->group_left -= rows; r->emitted = true;
for (int64_t c = 0; c < r->nselected; c++) ray_release(cols[c]);
ray_free_raw(cols); ray_free_raw(keep); ray_free_raw(errors);
return tbl;
fail:
for (int64_t c = 0; c < r->nselected; c++) if (cols[c]) ray_release(cols[c]);
ray_free_raw(cols); ray_free_raw(keep); ray_free_raw(errors);
if (tbl) ray_release(tbl);
r->failed = true;
if (ray_interrupted()) return ray_error("cancel","parquet scan interrupted");
if (bad_col >= 0) {
ray_t* name = ray_sym_str(r->schema[r->selected[bad_col]].name);
return ray_error("parquet","%s (column %s)",err,name ? ray_str_ptr(name) : "?");
}
return pq_error(err);
}
static ray_t* pq_option(ray_t* options, const char* name) {
ray_t* key = ray_sym(ray_sym_intern(name,strlen(name)));
if (!key || RAY_IS_ERR(key)) return key;
ray_t* result = ray_dict_get(options,key); ray_release(key); return result;
}
static ray_t* pq_import_types(ray_parquet_t* r, ray_t* types);
static ray_t* pq_strict_option(ray_parquet_t* r, ray_t* options) {
ray_t* strict = pq_option(options,"strict");
if (!strict) return NULL;
bool valid = strict->type == -RAY_BOOL;
if (valid) for (int64_t c = 0; c < r->ncols; c++) r->schema[c].strict = strict->b8 != 0;
ray_release(strict);
return valid ? NULL : pq_error("strict must be a boolean");
}
static ray_t* pq_scan_types(ray_parquet_t* r, ray_t* types) {
if (types && types->type != RAY_SYM) return pq_error("types must be a symbol vector");
ray_t* err = pq_import_types(r,types);
if (err) return err;
for (int64_t c = 0; c < r->ncols; c++) if (r->schema[c].native_symbol) {
if (!r->scan_domain) r->scan_domain = ray_sym_domain_new();
if (!r->scan_domain) return ray_error("oom",NULL);
r->schema[c].import_domain = r->scan_domain;
}
return NULL;
}
static ray_t* pq_open_options(const char* path, ray_t* options, ray_parquet_t** out) {
if (!options || options->type != RAY_DICT) return ray_parquet_open(path,options,PQ_BATCH,out);
*out = NULL;
ray_t* keys = ray_dict_keys(options);
if (!keys || keys->type != RAY_SYM) return pq_error("scan option keys must be symbols");
for (int64_t i = 0; i < keys->len; i++) {
ray_t* key = ray_sym_vec_cell(keys,i); const char* k = key ? ray_str_ptr(key) : NULL;
if (!k || (strcmp(k,"columns") && strcmp(k,"range") && strcmp(k,"types") && strcmp(k,"strict"))) return pq_error("unknown scan option");
}
ray_t* columns = pq_option(options,"columns");
ray_t* range = pq_option(options,"range");
ray_t* types = pq_option(options,"types");
ray_t* err = NULL;
if ((columns && RAY_IS_ERR(columns)) || (range && RAY_IS_ERR(range))) { err = pq_error("invalid scan options"); goto done; }
err = ray_parquet_open(path,columns,PQ_BATCH,out);
if (!err) err = pq_scan_types(*out,types);
if (!err) err = pq_strict_option(*out,options);
if (err || !range) goto done;
if (range->type != RAY_DICT) { err = pq_error("range must be {column: name min: lo max: hi}"); goto done; }
ray_t* name = pq_option(range,"column"); ray_t* lo = pq_option(range,"min"); ray_t* hi = pq_option(range,"max");
if (!name || !lo || !hi || name->type != -RAY_SYM || lo->type != -RAY_I64 || hi->type != -RAY_I64)
err = pq_error("range expects a symbol column and two I64 bounds");
else err = ray_parquet_range(*out,name->i64,lo->i64,hi->i64);
if (name) ray_release(name);
if (lo) ray_release(lo);
if (hi) ray_release(hi);
done:
if (columns) ray_release(columns);
if (range) ray_release(range);
if (types) ray_release(types);
if (err && *out) { ray_parquet_close(*out); *out = NULL; }
return err;
}
static ray_parquet_t* pq_group_reader(const ray_parquet_t* parent, int64_t group) {
ray_parquet_t* r = ray_calloc_raw(sizeof(*r));
if (!r) return NULL;
r->borrowed = true; r->map = parent->map; r->size = parent->size; r->data_end = parent->data_end;
r->schema = parent->schema; r->groups = parent->groups; r->ncols = parent->ncols;
r->ngroups = group+1; r->group = group-1; r->batch_rows = parent->batch_rows;
r->nselected = parent->nselected; r->noutput = parent->noutput;
r->filter_pos = parent->filter_pos; r->filter_lo = parent->filter_lo; r->filter_hi = parent->filter_hi;
r->filter_nulls = parent->filter_nulls;
r->text_pattern = parent->text_pattern; r->text_pos = parent->text_pos;
r->selected = ray_alloc_raw((size_t)r->nselected*sizeof(*r->selected));
r->cursors = ray_calloc_raw((size_t)r->nselected*sizeof(*r->cursors));
if (!r->selected || !r->cursors) { ray_parquet_close(r); return NULL; }
memcpy(r->selected,parent->selected,(size_t)r->nselected*sizeof(*r->selected));
return r;
}
static ray_pool_t* pq_group_pool(ray_parquet_t* r) {
ray_pool_t* pool = r->ngroups > 1 ? ray_pool_get() : NULL;
return ray_pool_par_dispatch_ok(pool,r->ngroups,2) ? pool : NULL;
}
static ray_t* pq_append_table(ray_t** vectors, ray_t* batch, int64_t ncols) {
int64_t nr = ray_table_nrows(batch);
for (int64_t c = 0; c < ncols; c++) {
ray_t* src = ray_table_get_col_idx(batch,c);
if (src->type == RAY_STR) {
for (int64_t i = 0; i < nr; i++) {
size_t n; const char* p = ray_str_vec_get(src,i,&n);
ray_t* v = ray_str_vec_append(vectors[c],p,n);
if (!v || RAY_IS_ERR(v)) return v ? v : ray_error("oom",NULL);
vectors[c] = v;
}
} else {
ray_t* v = ray_vec_append_raw(vectors[c],ray_data(src),nr);
if (!v || RAY_IS_ERR(v)) return v ? v : ray_error("oom",NULL);
vectors[c] = v;
}
vectors[c]->attrs |= src->attrs & RAY_ATTR_HAS_NULLS;
}
return NULL;
}
static ray_t* pq_materialize(ray_parquet_t* r);
typedef struct { ray_parquet_t* parent; int64_t first; ray_t** tables; } pq_read_work;
static void pq_read_group(void* ptr, uint32_t worker, int64_t start, int64_t end) {
(void)worker; pq_read_work* w = ptr;
for (int64_t i = start; i < end; i++) {
ray_parquet_t* r = pq_group_reader(w->parent,w->first+i);
w->tables[i] = r ? pq_materialize(r) : ray_error("oom",NULL);
}
}
static ray_t* pq_materialize(ray_parquet_t* r) {
ray_t* err = NULL;
ray_t** vectors = ray_calloc_raw((size_t)r->noutput*sizeof(*vectors));
if (!vectors) { ray_parquet_close(r); return ray_error("oom",NULL); }
for (int64_t c = 0; c < r->noutput; c++) {
pq_schema* schema = &r->schema[r->selected[c]];
vectors[c] = schema->import_domain ? ray_sym_vec_new(RAY_SYM_W32,0) : ray_vec_new((int8_t)schema->type,0);
if (!vectors[c] || RAY_IS_ERR(vectors[c])) {
err = vectors[c] ? vectors[c] : ray_error("oom",NULL); vectors[c] = NULL; goto done;
}
if (schema->import_domain) {
vectors[c]->sym_domain = schema->import_domain;
ray_sym_domain_retain(schema->import_domain);
}
}
ray_pool_t* pool = pq_group_pool(r);
if (pool) {
uint32_t width = ray_pool_total_workers(pool);
if (width > RAY_POOL_INIT_TASKS) width = RAY_POOL_INIT_TASKS;
ray_t** tables = ray_calloc_raw((size_t)width*sizeof(*tables));
if (!tables) { err = ray_error("oom",NULL); goto done; }
for (int64_t first = 0; first < r->ngroups && !err; first += width) {
uint32_t n = r->ngroups-first < width ? (uint32_t)(r->ngroups-first) : width;
memset(tables,0,(size_t)n*sizeof(*tables));
pq_read_work work = {r,first,tables};
ray_pool_dispatch_n(pool,pq_read_group,&work,n);
for (uint32_t i = 0; i < n; i++) {
ray_t* batch = tables[i];
if (!err) {
if (!batch) err = ray_error("cancel","parquet scan interrupted");
else if (RAY_IS_ERR(batch)) { err = batch; batch = NULL; }
else err = pq_append_table(vectors,batch,r->noutput);
}
if (batch) ray_release(batch);
}
}
ray_free_raw(tables);
if (err) goto done;
} else for (;;) {
ray_t* batch = ray_parquet_next(r);
if (!batch) break;
if (RAY_IS_ERR(batch)) { err = batch; goto done; }
err = pq_append_table(vectors,batch,r->noutput);
ray_release(batch);
if (err) goto done;
}
err = ray_table_new(r->noutput);
for (int64_t c = 0; c < r->noutput && err && !RAY_IS_ERR(err); c++)
err = ray_table_add_col(err,r->schema[r->selected[c]].name,vectors[c]);
done:
for (int64_t c = 0; c < r->noutput; c++) if (vectors[c]) ray_release(vectors[c]);
ray_free_raw(vectors); ray_parquet_close(r);
return err ? err : ray_error("oom",NULL);
}
ray_t* ray_parquet_read(const char* path, ray_t* columns) {
ray_parquet_t* r = NULL;
ray_t* err = pq_open_options(path,columns,&r);
return err ? err : pq_materialize(r);
}
ray_t* ray_parquet_metadata(const char* path) {
ray_parquet_t* r = NULL; ray_t* err = ray_parquet_open(path,NULL,PQ_BATCH,&r);
if (err) return err;
const char* names[] = {"row_group","column","physical_type","ray_type","codec","rows","compressed_bytes","uncompressed_bytes","statistics","page_index","bloom_filter"};
enum { N = 11 };
ray_t* v[N] = {0}; pq_span* cols = NULL;
int64_t count = r->ngroups*r->ncols;
if ((uint64_t)count > PQ_MAX_FOOTER/(N*sizeof(int64_t))) {
ray_parquet_close(r); return pq_error("metadata result exceeds 64 MiB");
}
for (int i = 0; i < N; i++) {
v[i] = i == 1 || i == 3 ? ray_sym_vec_new(RAY_SYM_W64,count) : ray_vec_new(RAY_I64,count);
if (!v[i] || RAY_IS_ERR(v[i])) { err = v[i]; v[i] = NULL; goto done; }
v[i]->len = count;
}
cols = ray_alloc_raw((size_t)r->ncols*sizeof(*cols));
if (!cols) goto done;
for (int64_t g = 0; g < r->ngroups; g++) {
int64_t nr;
if (!pq_group_columns(r,g,cols,&nr,NULL)) { err = pq_error("invalid row group metadata"); goto done; }
pq_span gf[8];
if (!pq_fields(r->groups[g],gf,8)) { err = pq_error("invalid row group"); goto done; }
pq_cur cc = {gf[1].p,gf[1].end,false}; uint8_t ct; pq_list(&cc,&ct);
for (int64_t c = 0; c < r->ncols; c++) {
int64_t row = g*r->ncols+c; pq_span mf[17], ch, cf[10];
if (!pq_fields(cols[c],mf,17) || !pq_element(&cc,12,&ch) || !pq_fields(ch,cf,10)) {
err = pq_error("invalid column metadata"); goto done;
}
const char* tn = ray_type_name((int8_t)r->schema[c].type);
int64_t cells[N] = {g,r->schema[c].name,r->schema[c].physical,ray_sym_intern(tn,strlen(tn)),
pq_get(mf[4],-1),nr,pq_get(mf[7],-1),pq_get(mf[6],-1),!!mf[12].type,
!!cf[4].type && !!cf[6].type,!!mf[14].type};
for (int i = 0; i < N; i++) ((int64_t*)ray_data(v[i]))[row] = cells[i];
}
}
err = ray_table_new(N);
for (int i = 0; i < N && err && !RAY_IS_ERR(err); i++)
err = ray_table_add_col(err,ray_sym_intern(names[i],strlen(names[i])),v[i]);
done:
for (int i = 0; i < N; i++) if (v[i]) ray_release(v[i]);
ray_free_raw(cols); ray_parquet_close(r);
return err ? err : ray_error("oom",NULL);
}
static bool pq_safe_name(const char* s) {
if (!s || !*s || *s == '.') return false;
for (; *s; s++) if (*s == '/' || *s == '\\' || *s == ':') return false;
return true;
}
static ray_t* pq_import_types(ray_parquet_t* r, ray_t* types) {
if (!types) return NULL;
if (types->type == RAY_DICT) {
ray_t* keys = ray_dict_keys(types);
if (!keys || keys->type != RAY_SYM) return pq_error("import option keys must be symbols");
for (int64_t i = 0; i < keys->len; i++) {
ray_t* key = ray_sym_vec_cell(keys,i); const char* k = key ? ray_str_ptr(key) : NULL;
if (!k || (strcmp(k,"types") && strcmp(k,"strict") && strcmp(k,"rows"))) return pq_error("unknown import option");
}
ray_t* err = pq_strict_option(r,types);
ray_t* rows = pq_option(types,"rows");
if (!err && rows) {
if (rows->type != -RAY_I64 || rows->i64 < 1 || rows->i64 > PQ_MAX_PART_ROWS)
err = pq_error("partition rows must be 1..1048576");
else r->batch_rows = rows->i64;
}
if (rows) ray_release(rows);
ray_t* schema = pq_option(types,"types");
if (!err && schema && schema->type != RAY_SYM) err = pq_error("types must be a symbol vector");
if (!err) err = pq_import_types(r,schema);
if (schema) ray_release(schema);
return err;
}
if (types->type != RAY_SYM || types->len != r->ncols) return pq_error("import types must match the Parquet columns");
for (int64_t c = 0; c < r->ncols; c++) {
ray_t* atom = ray_sym_vec_cell(types,c); const char* name = atom ? ray_str_ptr(atom) : NULL;
pq_schema* s = &r->schema[c];
if (!name) return pq_error("invalid import type");
if (!strcmp(name,"SYM") && s->type == RAY_STR) { s->native_symbol = true; continue; }
if (!strcmp(name,ray_type_name((int8_t)s->type))) continue;
bool integer = (s->type == RAY_I16 || s->type == RAY_I32 || s->type == RAY_I64);
if (!integer) return pq_error("unsupported import type conversion");
if (!strcmp(name,"I16")) s->type = RAY_I16;
else if (!strcmp(name,"I32")) s->type = RAY_I32;
else if (!strcmp(name,"I64")) s->type = RAY_I64;
else if (!strcmp(name,"UNIX_DATE")) s->type = RAY_DATE;
else {
int64_t scale = !strcmp(name,"UNIX_SECONDS") ? 1000000000 :
!strcmp(name,"UNIX_MILLIS") ? 1000000 : !strcmp(name,"UNIX_MICROS") ? 1000 :
!strcmp(name,"UNIX_NANOS") ? 1 : 0;
if (!scale) return pq_error("unsupported import type; raw integer timestamps require explicit UNIX_* units");
s->type = RAY_TIMESTAMP; s->scale = scale;
}
}
return NULL;
}
typedef struct {
ray_parquet_t* parent;
const char *root, *table;
int64_t *offsets;
ray_t** errors;
bool durable;
} pq_native_work;
static ray_err_t pq_save_partition(ray_parquet_t* r, ray_t* batch, const char* leaf,
const char* sym, bool durable) {
ray_sym_domain_t* domain = NULL;
for (int64_t c = 0; c < r->ncols; c++) if (r->schema[c].import_domain) domain = r->schema[c].import_domain;
if (!domain) return durable ? ray_splay_save(batch,leaf,sym) : ray_splay_save_staged_bulk(batch,leaf,sym);
ray_err_t e = ray_mkdir_p(leaf); if (e != RAY_OK) return e;
ray_t* schema = ray_vec_new(RAY_STR,r->ncols);
if (!schema || RAY_IS_ERR(schema)) { if (schema) ray_release(schema); return RAY_ERR_OOM; }
char path[1600];
for (int64_t c = 0; c < r->ncols; c++) {
ray_t* name = ray_sym_str(r->schema[c].name);
int n = snprintf(path,sizeof(path),"%s/%s",leaf,ray_str_ptr(name));
if (n < 0 || (size_t)n >= sizeof(path)) { e = RAY_ERR_RANGE; break; }
ray_t* col = ray_table_get_col_idx(batch,c);
e = col->type == RAY_SYM ? ray_col_save_sym_encoded(col,path,domain,durable) :
durable ? ray_col_save(col,path) : ray_col_save_bulk(col,path);
if (e != RAY_OK) break;
ray_t* next = ray_str_vec_append(schema,ray_str_ptr(name),ray_str_len(name));
if (!next || RAY_IS_ERR(next)) { if (next) ray_release(next); e = RAY_ERR_OOM; break; }
schema = next;
}
if (e == RAY_OK) {
int n = snprintf(path,sizeof(path),"%s/.d",leaf);
e = n < 0 || (size_t)n >= sizeof(path) ? RAY_ERR_RANGE :
durable ? ray_col_save(schema,path) : ray_col_save_bulk(schema,path);
}
ray_release(schema); return e;
}
static ray_t* pq_write_reader(ray_parquet_t* r, const char* root, const char* table,
int64_t part, bool durable) {
char sym[1200]; snprintf(sym,sizeof(sym),"%s/.sym",root);
for (;;) {
ray_t* batch = ray_parquet_next(r);
if (!batch || RAY_IS_ERR(batch)) return batch;
if (durable) for (int64_t c = 0; c < r->ncols; c++) {
ray_t* col = ray_table_get_col_idx(batch,c);
if (col->type == RAY_STR || col->type == RAY_SYM || col->type == RAY_F32) continue;
ray_retain(col);
ray_t* e = col->len >= PQ_BATCH ? ray_index_attach_chunk_zone(&col,16) : ray_index_attach_zone(&col);
if (e && RAY_IS_ERR(e)) { ray_release(col); ray_release(batch); return e; }
ray_table_set_col_idx(batch,c,col); ray_release(col);
}
char leaf[1400];
int n = snprintf(leaf,sizeof(leaf),"%s/%lld/%s",root,(long long)part++,table);
ray_err_t e = n < 0 || (size_t)n >= sizeof(leaf) ? RAY_ERR_RANGE :
pq_save_partition(r,batch,leaf,sym,false);
ray_release(batch);
if (e != RAY_OK) return ray_error(ray_err_code_str(e),"parquet: native partition write failed");
}
}
static void pq_write_group(void* ptr, uint32_t worker, int64_t start, int64_t end) {
(void)worker; pq_native_work* w = ptr;
for (int64_t g = start; g < end; g++) {
ray_parquet_t* r = pq_group_reader(w->parent,g);
if (!r) { w->errors[g] = ray_error("oom",NULL); continue; }
w->errors[g] = pq_write_reader(r,w->root,w->table,w->offsets[g],w->durable);
ray_parquet_close(r);
}
}
static ray_t* pq_write_groups(ray_parquet_t* r, const char* root, const char* table,
bool durable, int64_t* parts) {
*parts = 0;
int64_t ng = r->ngroups;
int64_t* offsets = ray_alloc_raw((size_t)(ng+1)*sizeof(*offsets));
ray_t** errors = ray_calloc_raw((size_t)(ng+1)*sizeof(*errors));
ray_t* err = NULL;
if (!offsets || !errors) { err = ray_error("oom",NULL); goto done; }
for (int64_t g = 0; g < ng; g++) {
pq_span fields[8];
if (!pq_fields(r->groups[g],fields,8)) { err = pq_error("invalid row group"); goto done; }
int64_t rows = pq_get(fields[3],-1);
if (rows < 0) { err = pq_error("invalid row group row count"); goto done; }
offsets[g] = *parts;
int64_t n = rows ? 1+(rows-1)/pq_group_batch(r,g) : 1;
if (*parts > INT64_MAX-n) { err = pq_error("too many partitions"); goto done; }
*parts += n;
}
if (!ng) { *parts = 1; err = pq_write_reader(r,root,table,0,durable); goto done; }
pq_native_work work = {r,root,table,offsets,errors,durable};
ray_pool_t* pool = pq_group_pool(r);
if (pool) {
ray_pool_dispatch_n(pool,pq_write_group,&work,(uint32_t)ng);
} else pq_write_group(&work,0,0,ng);
for (int64_t g = 0; g < ng; g++) if (errors[g]) {
if (!err) err = errors[g]; else ray_release(errors[g]);
}
if (!err && ray_interrupted()) err = ray_error("cancel","parquet conversion interrupted");
done:
ray_free_raw(offsets); ray_free_raw(errors); return err;
}
static bool pq_dest(const char* in, char* out, size_t size) {
size_t len = in ? strlen(in) : 0;
while (len > 1 && (in[len-1] == '/' || in[len-1] == '\\')) len--;
if (!len || len >= size) return false;
memcpy(out,in,len); out[len] = 0;
return true;
}
static ray_t* pq_stage(ray_parquet_t* r, const char* root, char* staging, size_t size, bool parted) {
struct stat st;
int n = snprintf(staging,size,"%s.parquet-partial",root);
if (n < 0 || (size_t)n >= size || stat(root,&st) == 0)
return pq_error("destination exists or path is too long");
for (int64_t c = 0; c < r->ncols; c++) {
const char* name = ray_str_ptr(ray_sym_str(r->schema[c].name));
if (!pq_safe_name(name) || (parted && !strcmp(name,"part")))
return pq_error("column name is unsafe or conflicts with native partition key 'part'");
}
#ifdef RAY_OS_WINDOWS
bool created = CreateDirectoryA(staging,NULL) != 0;
#else
bool created = mkdir(staging,0755) == 0;
#endif
return created ? NULL : ray_error("parquet","cannot create staging directory %s (the parent must exist; one left by an earlier failed import is not overwritten, remove it to retry)",staging);
}
static ray_t* pq_publish(const char* staging, const char* root, int64_t rows) {
struct stat st;
if (ray_interrupted()) return ray_error("cancel","parquet conversion interrupted");
char child[1100];
int n = snprintf(child,sizeof(child),"%s/.d",staging);
if (n < 0 || (size_t)n >= sizeof(child) || ray_file_sync_dir(child) != RAY_OK)
return pq_error("cannot sync staging directory; native destination not published");
if (stat(root,&st) == 0 || ray_file_rename_new(staging,root) != RAY_OK)
return pq_error("cannot publish native destination; staging directory retained");
if (ray_file_sync_dir(root) != RAY_OK) return pq_error("native destination published but directory sync failed");
return ray_i64(rows);
}
static ray_t* pq_sync_partitions(ray_parquet_t* r, const char* root, const char* table, int64_t parts) {
char path[1600], leaf[1400];
for (int64_t part = 0; part < parts; part++) {
if (ray_interrupted()) return ray_error("cancel","parquet conversion interrupted");
int n = snprintf(leaf,sizeof(leaf),"%s/%lld/%s",root,(long long)part,table);
if (n < 0 || (size_t)n >= sizeof(leaf)) return pq_error("partition path is too long");
for (int64_t c = 0; c <= r->ncols; c++) {
const char* name = c == r->ncols ? ".d" : ray_str_ptr(ray_sym_str(r->schema[c].name));
n = snprintf(path,sizeof(path),"%s/%s",leaf,name);
if (n < 0 || (size_t)n >= sizeof(path)) return pq_error("column path is too long");
ray_fd_t fd = ray_file_open(path,RAY_OPEN_READ | RAY_OPEN_WRITE);
if (fd == RAY_FD_INVALID) return pq_error("cannot open staged column for sync");
ray_err_t e = ray_file_sync(fd); ray_file_close(fd);
if (e != RAY_OK) return ray_error(ray_err_code_str(e),"parquet: staged column sync failed");
}
if (ray_file_sync_dir(path) != RAY_OK || ray_file_sync_dir(leaf) != RAY_OK)
return pq_error("cannot sync partition directories");
}
return NULL;
}
ray_t* ray_parquet_parted_typed(const char* path, const char* root, const char* table, ray_t* types) {
#ifdef RAY_FUZZING
return ray_error("restricted","Parquet imports disabled under fuzzing");
#endif
if (!pq_safe_name(table)) return pq_error("invalid destination table name");
char dest[1024];
if (!pq_dest(root,dest,sizeof(dest))) return pq_error("invalid or too long destination path");
root = dest;
ray_parquet_t* r = NULL; ray_t* err = ray_parquet_open(path,NULL,PQ_BATCH,&r);
if (err) return err;
char staging[1100], sym[1200]; int64_t parts;
ray_sym_domain_t* domain = NULL;
err = pq_import_types(r,types);
if (!err) err = pq_stage(r,root,staging,sizeof(staging),true);
if (err) goto done;
for (int64_t c = 0; c < r->ncols; c++) if (r->schema[c].native_symbol) {
if (!domain) {
snprintf(sym,sizeof(sym),"%s/.sym",staging);
domain = ray_sym_domain_open_or_create(sym);
if (!domain || ray_sym_domain_intern(domain,"",0) != 0) {
err = ray_error("oom",NULL); goto done;
}
}
r->schema[c].import_domain = domain;
}
err = pq_write_groups(r,staging,table,true,&parts);
if (!err && domain) {
ray_err_t e = ray_sym_domain_flush(domain,true);
if (e != RAY_OK) err = ray_error(ray_err_code_str(e),"parquet: cannot flush symbol file");
}
if (!err) err = pq_sync_partitions(r,staging,table,parts);
if (!err) err = pq_publish(staging,root,r->rows);
done:
if (domain) ray_sym_domain_release(domain);
ray_parquet_close(r); return err;
}
ray_t* ray_parquet_parted(const char* path, const char* root, const char* table) {
return ray_parquet_parted_typed(path,root,table,NULL);
}
static bool pq_remove_dir(const char* path) {
#ifdef RAY_OS_WINDOWS
return RemoveDirectoryA(path) != 0;
#else
return rmdir(path) == 0;
#endif
}
typedef struct {
ray_parquet_t* parent;
ray_col_stream_t* writers;
int64_t* offsets;
_Atomic uint32_t* nulls;
_Atomic uint32_t* locks;
ray_t** errors;
} pq_direct_work;
static void pq_write_direct_group(void* ptr, uint32_t worker, int64_t start, int64_t end) {
(void)worker; pq_direct_work* w = ptr;
for (int64_t task = start; task < end; task++) {
int64_t g = task/w->parent->ncols, c = task%w->parent->ncols;
ray_parquet_t* r = pq_group_reader(w->parent,g);
if (!r) { w->errors[task] = ray_error("oom",NULL); continue; }
r->selected[0] = (int32_t)c; r->nselected = r->noutput = 1;
int64_t row = w->offsets[g];
ray_col_stream_t local = {0};
local.type = w->writers[c].type; local.dom = w->writers[c].dom;
size_t size = local.type == RAY_SYM ? 4 : ray_elem_size(local.type);
local.fp = fopen(w->writers[c].tmp_path,"r+b");
if (!local.fp) { w->errors[task] = pq_error("cannot open native column"); ray_parquet_close(r); continue; }
int64_t offset = 32+row*(int64_t)size;
#ifdef RAY_OS_WINDOWS
bool seek = _fseeki64(local.fp,offset,SEEK_SET) == 0;
#else
bool seek = fseeko(local.fp,(off_t)offset,SEEK_SET) == 0;
#endif
if (!seek) w->errors[task] = pq_error("cannot seek native output range");
if (!w->errors[task]) {
ray_err_t ie = ray_col_stream_index_begin(&local,row);
if (ie != RAY_OK) w->errors[task] = ray_error(ray_err_code_str(ie),"parquet: zone accumulator");
}
while (!w->errors[task]) {
ray_t* batch = ray_parquet_next(r);
if (!batch) break;
if (RAY_IS_ERR(batch)) { w->errors[task] = batch; break; }
ray_err_t err = ray_col_stream_append(&local,ray_table_get_col_idx(batch,0));
row += ray_table_nrows(batch); ray_release(batch);
if (err != RAY_OK) { w->errors[task] = ray_error(ray_err_code_str(err),"parquet: direct column write failed"); break; }
}
if (local.zone && !w->errors[task]) {
while (atomic_exchange_explicit(&w->locks[c],1,memory_order_acquire)) RAY_CPU_RELAX();
ray_err_t me = ray_col_stream_index_merge(&w->writers[c],&local);
atomic_store_explicit(&w->locks[c],0,memory_order_release);
if (me != RAY_OK && !w->errors[task]) w->errors[task] = ray_error(ray_err_code_str(me),"parquet: zone merge");
}
if (local.zone) { ray_zone_acc_free(local.zone); ray_free_raw(local.zone); local.zone = NULL; }
if (fclose(local.fp) && !w->errors[task]) w->errors[task] = pq_error("native column close failed");
if (local.had_nulls) atomic_store_explicit(&w->nulls[c],1,memory_order_relaxed);
if (!w->errors[task] && row != w->offsets[g+1]) w->errors[task] = pq_error("row group ended before its assigned output range");
ray_parquet_close(r);
}
}
static ray_t* pq_write_direct(ray_parquet_t* r, ray_col_stream_t* writers) {
int64_t tasks = r->ngroups*r->ncols;
if (tasks > UINT32_MAX) return pq_error("too many column chunk tasks");
int64_t* offsets = ray_calloc_raw((size_t)(r->ngroups+1)*sizeof(*offsets));
_Atomic uint32_t* nulls = ray_calloc_raw((size_t)r->ncols*sizeof(*nulls));
_Atomic uint32_t* locks = ray_calloc_raw((size_t)r->ncols*sizeof(*locks));
ray_t** errors = ray_calloc_raw((size_t)(tasks+1)*sizeof(*errors));
ray_t* err = NULL;
if (!offsets || !nulls || !locks || !errors) { err = ray_error("oom",NULL); goto done; }
for (int64_t c = 0; c < r->ncols; c++) {
size_t size = writers[c].type == RAY_SYM ? 4 : ray_elem_size(writers[c].type);
if (r->rows > (INT64_MAX-32)/(int64_t)size) { err = pq_error("native column exceeds file offset range"); goto done; }
if (fflush(writers[c].fp)) { err = pq_error("cannot flush column header"); goto done; }
atomic_init(&nulls[c],0); atomic_init(&locks[c],0);
if (r->schema[c].native_symbol) r->schema[c].import_domain = writers[c].dom;
}
for (int64_t g = 0; g < r->ngroups; g++) {
pq_span fields[8];
if (!pq_fields(r->groups[g],fields,8)) { err = pq_error("invalid row group"); goto done; }
offsets[g+1] = offsets[g]+pq_get(fields[3],0);
}
pq_direct_work work = {r,writers,offsets,nulls,locks,errors};
ray_pool_t* pool = ray_pool_get();
if (ray_pool_par_dispatch_ok(pool,tasks,2)) ray_pool_dispatch_n(pool,pq_write_direct_group,&work,(uint32_t)tasks);
else pq_write_direct_group(&work,0,0,tasks);
for (int64_t g = 0; g < tasks; g++) if (errors[g]) {
if (!err) err = errors[g]; else ray_release(errors[g]);
}
if (!err && ray_interrupted()) err = ray_error("cancel","parquet conversion interrupted");
for (int64_t c = 0; c < r->ncols; c++) {
writers[c].rows = r->rows;
writers[c].had_nulls = atomic_load_explicit(&nulls[c],memory_order_relaxed) != 0;
}
done:
for (int64_t c = 0; c < r->ncols; c++) r->schema[c].import_domain = NULL;
ray_free_raw(offsets); ray_free_raw(nulls); ray_free_raw(locks); ray_free_raw(errors); return err;
}
ray_t* ray_parquet_splayed_typed(const char* path, const char* dir, ray_t* types) {
#ifdef RAY_FUZZING
return ray_error("restricted","Parquet imports disabled under fuzzing");
#endif
char dest[1024];
if (!pq_dest(dir,dest,sizeof(dest))) return pq_error("invalid or too long destination path");
dir = dest;
ray_parquet_t* r = NULL; ray_t* err = ray_parquet_open(path,NULL,PQ_BATCH,&r);
if (err) return err;
char staging[1100], spool[1200], leaf[1500], file[1700];
ray_col_stream_t* writers = NULL; int64_t opened = 0, parts = 0;
ray_err_t e = RAY_OK; ray_sym_domain_t* domain = NULL;
err = pq_import_types(r,types);
if (err) goto done;
err = pq_stage(r,dir,staging,sizeof(staging),false);
if (err) goto done;
writers = ray_calloc_raw((size_t)r->ncols*sizeof(*writers));
if (!writers) { err = ray_error("oom",NULL); goto done; }
for (int64_t c = 0; c < r->ncols; c++) if (r->schema[c].native_symbol && !domain) {
snprintf(file,sizeof(file),"%s/.sym",staging);
domain = ray_sym_domain_open_or_create(file);
if (!domain) { err = ray_error("oom",NULL); goto done; }
if (ray_sym_domain_intern(domain,"",0) != 0) { err = pq_error("cannot initialize native symbol domain"); goto done; }
}
for (int64_t c = 0; c < r->ncols; c++) {
opened++;
e = ray_col_stream_open(&writers[c],staging,r->schema[c].name,r->schema[c].native_symbol ? RAY_SYM : (int8_t)r->schema[c].type,domain);
if (e != RAY_OK) goto io_fail;
}
for (int64_t c = 0; c < r->ncols; c++) {
e = ray_col_stream_index_begin(&writers[c],0);
if (e != RAY_OK) goto io_fail;
}
bool direct = true;
for (int64_t c = 0; c < r->ncols; c++) if (writers[c].type == RAY_STR) direct = false;
if (direct) {
err = pq_write_direct(r,writers);
if (err) goto done;
goto finish_columns;
}
snprintf(spool,sizeof(spool),"%s/.groups",staging);
e = ray_mkdir(spool); if (e != RAY_OK) goto io_fail;
err = pq_write_groups(r,spool,"batch",false,&parts);
if (err) goto done;
for (int64_t part = 0; part < parts; part++) {
if (ray_interrupted()) { e = RAY_ERR_CANCEL; goto io_fail; }
snprintf(leaf,sizeof(leaf),"%s/%lld/batch",spool,(long long)part);
ray_t* batch = ray_read_splayed(leaf,NULL);
if (!batch || RAY_IS_ERR(batch)) { err = batch ? batch : ray_error("oom",NULL); goto done; }
for (int64_t c = 0; c < r->ncols; c++) {
e = ray_col_stream_append(&writers[c],ray_table_get_col_idx(batch,c));
if (e != RAY_OK) break;
}
ray_release(batch);
if (e != RAY_OK) goto io_fail;
for (int64_t c = 0; c <= r->ncols; c++) {
const char* name = c == r->ncols ? ".d" : ray_str_ptr(ray_sym_str(r->schema[c].name));
snprintf(file,sizeof(file),"%s/%s",leaf,name);
if (remove(file)) { e = RAY_ERR_IO; goto io_fail; }
}
if (!pq_remove_dir(leaf)) { e = RAY_ERR_IO; goto io_fail; }
snprintf(leaf,sizeof(leaf),"%s/%lld",spool,(long long)part);
if (!pq_remove_dir(leaf)) { e = RAY_ERR_IO; goto io_fail; }
}
if (!pq_remove_dir(spool)) { e = RAY_ERR_IO; goto io_fail; }
finish_columns:
if (domain) { e = ray_sym_domain_flush(domain,true); if (e != RAY_OK) goto io_fail; }
e = ray_col_stream_close_all(writers,r->ncols,false,NULL);
if (e != RAY_OK) goto io_fail;
if (domain) {
ray_sym_domain_release(domain); domain = NULL;
for (int64_t c = 0; c < opened; c++) writers[c].dom = NULL;
}
ray_col_stream_hash_all(writers,r->ncols,NULL);
ray_t* schema = ray_vec_new(RAY_STR,r->ncols);
if (!schema || RAY_IS_ERR(schema)) { err = schema ? schema : ray_error("oom",NULL); goto done; }
for (int64_t c = 0; c < r->ncols; c++) {
ray_t* name = ray_sym_str(r->schema[c].name);
ray_t* next = ray_str_vec_append(schema,ray_str_ptr(name),ray_str_len(name));
if (!next || RAY_IS_ERR(next)) { ray_release(schema); err = next ? next : ray_error("oom",NULL); goto done; }
schema = next;
}
snprintf(file,sizeof(file),"%s/.d",staging);
e = ray_col_save_bulk(schema,file); ray_release(schema);
if (e != RAY_OK) goto io_fail;
for (int64_t c = 0; c <= r->ncols; c++) {
snprintf(file,sizeof(file),"%s/.d",staging);
const char* path = c == r->ncols ? file : writers[c].path;
ray_fd_t fd = ray_file_open(path,RAY_OPEN_READ | RAY_OPEN_WRITE);
if (fd == RAY_FD_INVALID) { e = RAY_ERR_IO; goto io_fail; }
e = ray_file_sync(fd); ray_file_close(fd);
if (e != RAY_OK) goto io_fail;
}
e = ray_file_sync_dir(file);
if (e != RAY_OK) goto io_fail;
err = pq_publish(staging,dir,r->rows);
goto done;
io_fail:
err = ray_error(ray_err_code_str(e),"parquet: splayed write failed; staging directory retained");
done:
for (int64_t c = 0; c < opened; c++) ray_col_stream_abort(&writers[c]);
ray_free_raw(writers); if (domain) ray_sym_domain_release(domain); ray_parquet_close(r); return err;
}
ray_t* ray_parquet_splayed(const char* path, const char* dir) {
return ray_parquet_splayed_typed(path,dir,NULL);
}
static const char* pq_path_arg(ray_t* arg) {
if (!arg || arg->type != -RAY_STR) return NULL;
const char* p = ray_str_ptr(arg); size_t n = ray_str_len(arg);
return p && n && !memchr(p,0,n) ? p : NULL;
}
ray_t* ray_parquet_metadata_fn(ray_t* path) {
const char* p = pq_path_arg(path);
return p ? ray_parquet_metadata(p) : ray_error("type","parquet.meta: path must be a nonempty string");
}
ray_t* ray_parquet_read_fn(ray_t** args, int64_t n) {
if (n != 1 && n != 2) return ray_error("arity","parquet.read: [columns] path");
const char* p = pq_path_arg(args[n-1]);
return p ? ray_parquet_read(p,n == 2 ? args[0] : NULL) : ray_error("type","parquet.read: path must be a nonempty string");
}
ray_t* ray_parquet_splayed_fn(ray_t** args, int64_t n) {
if (n != 2 && n != 3) return ray_error("arity","parquet.splayed: [types] path directory");
const char* path = pq_path_arg(args[n-2]); const char* dir = pq_path_arg(args[n-1]);
return path && dir ? ray_parquet_splayed_typed(path,dir,n == 3 ? args[0] : NULL) : ray_error("type","parquet.splayed: expected two paths");
}
ray_t* ray_parquet_parted_fn(ray_t** args, int64_t n) {
if (n != 3 && n != 4) return ray_error("arity","parquet.parted: [types] path root table-name");
const char* p = pq_path_arg(args[n-3]); const char* root = pq_path_arg(args[n-2]);
ray_t* name = args[n-1]->type == -RAY_SYM ? ray_sym_str(args[n-1]->i64) : args[n-1];
const char* table = pq_path_arg(name);
return p && root && table ? ray_parquet_parted_typed(p,root,table,n == 4 ? args[0] : NULL) : ray_error("type","parquet.parted: expected paths and table name");
}
ray_t* ray_parquet_each_fn(ray_t** args, int64_t n) {
if (n != 2 && n != 3) return ray_error("arity","parquet.each: [columns] path callback");
const char* p = pq_path_arg(args[n-2]); ray_t* fn = args[n-1];
if (!p || (fn->type != RAY_LAMBDA && fn->type != RAY_UNARY && fn->type != RAY_VARY))
return ray_error("type","parquet.each: expected path and callable");
ray_parquet_t* r = NULL; ray_t* err = pq_open_options(p,n == 3 ? args[0] : NULL,&r);
if (err) return err;
int64_t rows = 0;
for (;;) {
ray_t* batch = ray_parquet_next(r);
if (!batch) break;
if (RAY_IS_ERR(batch)) { err = batch; break; }
rows += ray_table_nrows(batch);
ray_t* result = call_fn1(fn,batch); ray_release(batch);
if (!result || RAY_IS_ERR(result)) { err = result ? result : ray_error("oom",NULL); break; }
ray_release(result);
}
ray_parquet_close(r); return err ? err : ray_i64(rows);
}
ray_t* ray_parquet_scan_fn(ray_t** args, int64_t n) {
if (n < 1 || n > 2) return ray_error("arity","parquet.scan: [options] path");
ray_t* path = args[n-1];
if (!pq_path_arg(path)) return ray_error("type","parquet.scan: path must be a nonempty string");
if (n == 2) {
ray_t* opts = args[0];
ray_t* keys = opts->type == RAY_DICT ? ray_dict_keys(opts) : NULL;
if (!keys || keys->type != RAY_SYM) return pq_error("lazy scan options must be {types: [...]}");
for (int64_t i = 0; i < keys->len; i++) {
ray_t* k = ray_sym_vec_cell(keys,i);
if (!k || (strcmp(ray_str_ptr(k),"types") && strcmp(ray_str_ptr(k),"strict"))) return pq_error("unknown lazy scan option");
}
}
int64_t names[] = {ray_sym_intern(".parquet.source",15),ray_sym_intern(".parquet.options",16)};
ray_t* keys = ray_vec_from_raw(RAY_SYM,names,n);
ray_t* vals = ray_list_new(n);
if (!keys || RAY_IS_ERR(keys) || !vals || RAY_IS_ERR(vals)) {
if (keys) ray_release(keys);
if (vals) ray_release(vals);
return ray_error("oom",NULL);
}
ray_retain(path); ((ray_t**)ray_data(vals))[0] = path; vals->len = 1;
if (n == 2) { ray_retain(args[0]); ((ray_t**)ray_data(vals))[1] = args[0]; vals->len = 2; }
return ray_dict_new(keys,vals);
}
static const char* pq_symbol_name(ray_t* x) {
return x && x->type == -RAY_SYM ? ray_str_ptr(ray_sym_str(x->i64)) : NULL;
}
static const char* pq_builtin_name(ray_t* x) {
const char* name = pq_symbol_name(x);
ray_t* value = name ? ray_env_get(x->i64) : NULL;
return value && value->type == RAY_LAMBDA ? NULL : name;
}
static bool pq_reserved(const char* k) {
return !strcmp(k,"from") || !strcmp(k,"where") || !strcmp(k,"by") ||
!strcmp(k,"asc") || !strcmp(k,"desc") || !strcmp(k,"take") || !strcmp(k,"nearest");
}
static int64_t pq_column_id(ray_parquet_t* r, ray_t* x) {
if (x && x->type == -RAY_SYM)
for (int64_t c = 0; c < r->ncols; c++) if (r->schema[c].name == x->i64) return c;
return -1;
}
static bool pq_references(ray_parquet_t* r, ray_t* e, uint8_t* keep, unsigned depth) {
if (!e || depth > 64) return false;
if (e->type == -RAY_SYM) {
int64_t c = pq_column_id(r,e); if (c >= 0) keep[c] = 1;
return true;
}
if (e->type == RAY_SYM) {
for (int64_t i = 0; i < e->len; i++) {
ray_t* cell = ray_sym_vec_cell(e,i);
const char* name = cell ? ray_str_ptr(cell) : NULL;
if (!name) return false;
for (int64_t c = 0; c < r->ncols; c++)
if (!strcmp(name,ray_str_ptr(ray_sym_str(r->schema[c].name)))) keep[c] = 1;
}
return true;
}
if (e->type == RAY_DICT) {
ray_t* vals = ray_dict_vals(e);
if (!vals || vals->type != RAY_LIST) return false;
for (int64_t i = 0; i < vals->len; i++)
if (!pq_references(r,ray_list_get(vals,i),keep,depth+1)) return false;
return true;
}
if (e->type == RAY_LIST) {
if (!e->len) return true;
const char* fn = pq_builtin_name(ray_list_get(e,0));
static const char* const transparent[] = {
"+","-","*","/","mod","%","div","==","!=","<","<=",">",">=","and","or","not",
"sum","count","avg","min","max","first","last","distinct","all","any",
"abs","sqrt","floor","ceil","round","log","exp","like","within","in",
"strlen","lower","upper","year","month","day","hour","minute","second","xbar",
"if","nil?","substr","str-find"
};
if (fn && !strcmp(fn,"let")) {
if (e->len != 4) return false;
const char* binding = pq_symbol_name(ray_list_get(e,1));
if (!binding || !strcmp(binding,"let")) return false;
for (size_t i = 0; i < sizeof(transparent)/sizeof(*transparent); i++)
if (!strcmp(binding,transparent[i])) return false;
return pq_references(r,ray_list_get(e,2),keep,depth+1) &&
pq_references(r,ray_list_get(e,3),keep,depth+1);
}
bool known = false;
for (size_t i = 0; fn && i < sizeof(transparent)/sizeof(*transparent); i++)
if (!strcmp(fn,transparent[i])) { known = true; break; }
if (!known) return false;
if (pq_column_id(r,ray_list_get(e,0)) >= 0) return false;
for (int64_t i = 1; i < e->len; i++)
if (!pq_references(r,ray_list_get(e,i),keep,depth+1)) return false;
return true;
}
return e->type < RAY_TABLE;
}
static bool pq_query_range(ray_parquet_t* r, ray_t* e, int64_t* column, int64_t* lo, int64_t* hi, unsigned depth) {
if (!e || e->type != RAY_LIST || e->len < 3 || depth > 64) return false;
const char* fn = pq_builtin_name(ray_list_get(e,0));
if (!fn || pq_column_id(r,ray_list_get(e,0)) >= 0) return false;
if (!strcmp(fn,"and")) {
for (int64_t i = 1; i < e->len; i++)
if (!pq_query_range(r,ray_list_get(e,i),column,lo,hi,depth+1)) return false;
return true;
}
if (e->len != 3) return false;
ray_t* a = ray_list_get(e,1); ray_t* b = ray_list_get(e,2);
int64_t c = pq_column_id(r,a); bool reverse = false;
if (c < 0) { c = pq_column_id(r,b); b = a; reverse = true; }
bool within = !strcmp(fn,"within") && !reverse && b->type == RAY_I64 && b->len == 2;
if (c < 0 || (!within && (b->type != -RAY_I64 || b->i64 == INT64_MIN))) return false;
pq_schema* s = &r->schema[c];
if ((s->type != RAY_I16 && s->type != RAY_I32 && s->type != RAY_I64) ||
(s->converted != -1 && (s->converted < 15 || s->converted > 18))) return false;
int64_t l = INT64_MIN, h = INT64_MAX, x = within ? 0 : b->i64;
if (within) {
l = ((int64_t*)ray_data(b))[0]; h = ((int64_t*)ray_data(b))[1];
if (l == INT64_MIN || h == INT64_MIN || l > h) return false;
} else if (!strcmp(fn,"==")) l = h = x;
else if (!strcmp(fn,reverse ? ">=" : "<=")) h = x;
else if (!strcmp(fn,reverse ? ">" : "<")) h = x-1;
else if (!strcmp(fn,reverse ? "<=" : ">=")) l = x;
else if (!strcmp(fn,reverse ? "<" : ">")) { if (x == INT64_MAX) return false; l = x+1; }
else return false;
if (*column < 0) *column = c;
if (*column == c) { if (l > *lo) *lo = l; if (h < *hi) *hi = h; }
return true;
}
typedef struct {
int64_t name, pos;
const agg_vtable_t* kernel;
void* state;
bool count;
} pq_aggregate;
static bool pq_exact_range(ray_parquet_t* r, ray_t* where) {
if (!where) return true;
if (r->filter_pos < 0 || where->type != RAY_LIST) return false;
const char* fn = pq_symbol_name(ray_list_get(where,0));
if (fn && !strcmp(fn,"and")) {
for (int64_t i = 1; i < where->len; i++) if (!pq_exact_range(r,ray_list_get(where,i))) return false;
return true;
}
int64_t c = r->selected[r->filter_pos];
return where->len == 3 && (pq_column_id(r,ray_list_get(where,1)) == c || pq_column_id(r,ray_list_get(where,2)) == c);
}
typedef struct { int64_t pos, lo, hi; } pq_predicate;
static bool pq_predicates(ray_parquet_t* r, ray_t* where, pq_predicate* predicates, int64_t* n) {
if (!where) return true;
if (where->type != RAY_LIST) return false;
const char* fn = pq_symbol_name(ray_list_get(where,0));
if (fn && !strcmp(fn,"and")) {
for (int64_t i = 1; i < where->len; i++)
if (!pq_predicates(r,ray_list_get(where,i),predicates,n)) return false;
return true;
}
int64_t column = -1, lo = INT64_MIN, hi = INT64_MAX, pos = 0;
if (*n == PQ_MAX_COLS || !pq_query_range(r,where,&column,&lo,&hi,0) || column < 0) return false;
while (pos < r->noutput && r->selected[pos] != column) pos++;
if (pos == r->noutput) return false;
predicates[(*n)++] = (pq_predicate){pos,lo,hi}; return true;
}
static void pq_filter_integers(ray_t* batch, const pq_predicate* predicates, int64_t n) {
int64_t rows = ray_table_nrows(batch), dst = 0, nc = ray_table_ncols(batch);
for (int64_t i = 0; i < rows; i++) {
bool keep = true;
for (int64_t p = 0; p < n && keep; p++) {
ray_t* v = ray_table_get_col_idx(batch,predicates[p].pos);
int64_t x = v->type == RAY_I16 ? ((int16_t*)ray_data(v))[i] :
v->type == RAY_I32 ? ((int32_t*)ray_data(v))[i] : ((int64_t*)ray_data(v))[i];
keep = ray_vec_is_null(v,i) ? predicates[p].lo == INT64_MIN :
x >= predicates[p].lo && x <= predicates[p].hi;
}
if (!keep) continue;
if (dst != i) for (int64_t c = 0; c < nc; c++) {
ray_t* v = ray_table_get_col_idx(batch,c); size_t size = ray_sym_elem_size(v->type,v->attrs);
memcpy((uint8_t*)ray_data(v)+(size_t)dst*size,(uint8_t*)ray_data(v)+(size_t)i*size,size);
}
dst++;
}
for (int64_t c = 0; c < nc; c++) ray_table_get_col_idx(batch,c)->len = dst;
}
static ray_t* pq_accumulate(ray_parquet_t* r, pq_aggregate* aggs, int64_t n,
const uint32_t* gids, const pq_predicate* predicates, int64_t np) {
for (;;) {
ray_t* batch = ray_parquet_next(r);
if (!batch || RAY_IS_ERR(batch)) return batch;
if (np) pq_filter_integers(batch,predicates,np);
int64_t rows = ray_table_nrows(batch);
for (int64_t a = 0; a < n; a++) {
ray_t* v = ray_table_get_col_idx(batch,aggs[a].pos);
ray_valid_t valid = {ray_data(v),v->type,ray_vec_may_have_nulls(v)};
aggs[a].kernel->update_batch(aggs[a].state,aggs[a].kernel->state_size,gids,ray_data(v),&valid,rows,NULL);
}
ray_release(batch);
}
}
typedef struct {
ray_parquet_t* parent;
int64_t first, n, np;
pq_aggregate *templates, **states;
const uint32_t* gids;
const pq_predicate* predicates;
ray_t** errors;
} pq_aggregate_work;
static void pq_aggregate_group(void* ptr, uint32_t worker, int64_t start, int64_t end) {
(void)worker; pq_aggregate_work* w = ptr;
for (int64_t i = start; i < end; i++) {
ray_parquet_t* r = pq_group_reader(w->parent,w->first+i);
pq_aggregate* aggs = ray_calloc_raw((size_t)w->n*sizeof(*aggs));
w->states[i] = aggs;
if (!r || !aggs) { w->errors[i] = ray_error("oom",NULL); ray_parquet_close(r); continue; }
for (int64_t a = 0; a < w->n; a++) {
aggs[a] = w->templates[a]; aggs[a].state = ray_alloc_raw(aggs[a].kernel->state_size);
if (!aggs[a].state) { w->errors[i] = ray_error("oom",NULL); break; }
aggs[a].kernel->init(aggs[a].state);
}
if (!w->errors[i]) w->errors[i] = pq_accumulate(r,aggs,w->n,w->gids,w->predicates,w->np);
ray_parquet_close(r);
}
}
static ray_t* pq_aggregate_groups(ray_parquet_t* r, pq_aggregate* aggs, int64_t n,
const uint32_t* gids, const pq_predicate* predicates, int64_t np) {
ray_pool_t* pool = pq_group_pool(r);
if (!pool) return pq_accumulate(r,aggs,n,gids,predicates,np);
uint32_t width = ray_pool_total_workers(pool);
if (width > RAY_POOL_INIT_TASKS) width = RAY_POOL_INIT_TASKS;
pq_aggregate** states = ray_calloc_raw((size_t)width*sizeof(*states));
ray_t** errors = ray_calloc_raw((size_t)width*sizeof(*errors));
ray_t* err = NULL;
if (!states || !errors) { err = ray_error("oom",NULL); goto done; }
for (int64_t first = 0; first < r->ngroups && !err; first += width) {
uint32_t count = r->ngroups-first < width ? (uint32_t)(r->ngroups-first) : width;
memset(states,0,(size_t)count*sizeof(*states)); memset(errors,0,(size_t)count*sizeof(*errors));
pq_aggregate_work work = {r,first,n,np,aggs,states,gids,predicates,errors};
ray_pool_dispatch_n(pool,pq_aggregate_group,&work,count);
for (uint32_t i = 0; i < count; i++) {
if (errors[i]) { if (!err) err = errors[i]; else ray_release(errors[i]); }
if (!states[i]) { if (!err) err = ray_error("cancel","parquet scan interrupted"); continue; }
for (int64_t a = 0; a < n; a++) {
if (!err) aggs[a].kernel->merge(aggs[a].state,states[i][a].state,NULL);
ray_free_raw(states[i][a].state);
}
ray_free_raw(states[i]);
}
}
done:
ray_free_raw(states); ray_free_raw(errors); return err;
}
static ray_t* pq_stream_aggregates(ray_parquet_t* r, ray_t* query, ray_t* where) {
if (where && r->filter_pos < 0 && !r->text_pattern) return NULL;
bool residual = !r->text_pattern && !pq_exact_range(r,where);
pq_predicate predicates[PQ_MAX_COLS]; int64_t np = 0;
if (residual && !pq_predicates(r,where,predicates,&np)) return NULL;
ray_t* keys = ray_dict_keys(query); ray_t* vals = ray_dict_vals(query);
if (!keys || keys->type != RAY_SYM || !vals || vals->type != RAY_LIST || keys->len > PQ_MAX_COLS) return NULL;
pq_aggregate* aggs = ray_calloc_raw((size_t)keys->len*sizeof(*aggs));
if (!aggs) return ray_error("oom",NULL);
int64_t n = 0; bool all_count = true;
ray_t* result = NULL; uint32_t* gids = NULL;
for (int64_t i = 0; i < keys->len; i++) {
ray_t* key = ray_sym_vec_cell(keys,i); const char* k = key ? ray_str_ptr(key) : NULL;
if (!k) goto done;
if (!strcmp(k,"from") || !strcmp(k,"where")) continue;
if (pq_reserved(k)) goto done;
ray_t* e = ray_list_get(vals,i);
if (!e || e->type != RAY_LIST || e->len != 2) goto done;
const char* fn = pq_builtin_name(ray_list_get(e,0));
int64_t col = pq_column_id(r,ray_list_get(e,1));
if (!fn || col < 0 || pq_column_id(r,ray_list_get(e,0)) >= 0) goto done;
uint16_t op = !strcmp(fn,"count") ? OP_COUNT : !strcmp(fn,"sum") ? OP_SUM :
!strcmp(fn,"avg") ? OP_AVG : !strcmp(fn,"min") ? OP_MIN : !strcmp(fn,"max") ? OP_MAX : 0;
int8_t type = (int8_t)r->schema[col].type;
if (r->schema[col].native_symbol) {
if (op != OP_COUNT) goto done;
type = RAY_I32;
}
if (!op || ((op == OP_SUM || op == OP_AVG) && type != RAY_I16 && type != RAY_I32 && type != RAY_I64)) goto done;
const agg_vtable_t* kernel = agg_resolve(op,type);
if (!kernel || kernel->kind != ACC_STREAMING || !kernel->finalize_value || !kernel->update_batch || !kernel->merge) goto done;
int64_t pos = 0;
while (pos < r->noutput && r->selected[pos] != col) pos++;
if (pos == r->noutput) goto done;
aggs[n] = (pq_aggregate){ray_sym_intern(k,strlen(k)),pos,kernel,NULL,op == OP_COUNT};
if (op != OP_COUNT) all_count = false;
n++;
}
if (!n) goto done;
for (int64_t a = 0; a < n; a++) {
aggs[a].state = ray_alloc_raw(aggs[a].kernel->state_size);
if (!aggs[a].state) { result = ray_error("oom",NULL); goto done; }
aggs[a].kernel->init(aggs[a].state);
}
if (!(all_count && !where)) {
gids = ray_calloc_raw(PQ_BATCH*sizeof(*gids));
if (!gids) { result = ray_error("oom",NULL); goto done; }
result = pq_aggregate_groups(r,aggs,n,gids,predicates,np);
if (result) goto done;
}
result = ray_table_new(n);
if (!result) result = ray_error("oom",NULL);
if (RAY_IS_ERR(result)) goto done;
for (int64_t a = 0; a < n; a++) {
ray_t* v = ray_vec_new(aggs[a].kernel->out_type,1);
if (!v || RAY_IS_ERR(v)) { ray_release(result); result = v ? v : ray_error("oom",NULL); goto done; }
v->len = 1;
if (all_count && !where) ((int64_t*)ray_data(v))[0] = r->rows;
else if (aggs[a].kernel->finalize_value(aggs[a].state,ray_data(v))) v->attrs |= RAY_ATTR_HAS_NULLS;
result = ray_table_add_col(result,aggs[a].name,v); ray_release(v);
if (!result) result = ray_error("oom",NULL);
if (RAY_IS_ERR(result)) break;
}
done:
for (int64_t a = 0; a < n; a++) ray_free_raw(aggs[a].state);
ray_free_raw(aggs); ray_free_raw(gids);
return result;
}
ray_t* ray_parquet_select_source(ray_t* source, ray_t* query, bool* complete) {
*complete = false;
if (!source || source->type != RAY_DICT || ray_dict_len(source) < 1 || ray_dict_len(source) > 2) return NULL;
ray_t* path = pq_option(source,".parquet.source");
if (!path) return NULL;
if (ray_eval_get_restricted()) { ray_release(path); return ray_error("access","restricted"); }
const char* p = pq_path_arg(path);
if (!p) { ray_release(path); return pq_error("invalid lazy source path"); }
ray_parquet_t* r = NULL;
ray_t* err = ray_parquet_open(p,NULL,PQ_BATCH,&r); ray_release(path);
if (err) return err;
ray_t* options = pq_option(source,".parquet.options");
if (options) {
ray_t* types = options->type == RAY_DICT ? pq_option(options,"types") : NULL;
err = options->type == RAY_DICT ? pq_scan_types(r,types) : pq_error("invalid lazy scan options");
if (!err) err = pq_strict_option(r,options);
if (types) ray_release(types);
ray_release(options);
if (err) { ray_parquet_close(r); return err; }
}
uint8_t* keep = ray_calloc_raw((size_t)r->ncols);
if (!keep) { ray_parquet_close(r); return ray_error("oom",NULL); }
ray_t* keys = ray_dict_keys(query); ray_t* vals = ray_dict_vals(query);
bool project = keys && keys->type == RAY_SYM && vals && vals->type == RAY_LIST;
bool explicit = false;
for (int64_t i = 0; project && i < keys->len; i++) {
ray_t* key = ray_sym_vec_cell(keys,i); const char* k = key ? ray_str_ptr(key) : NULL;
if (!k) { project = false; break; }
if (!strcmp(k,"from")) continue;
if (!pq_reserved(k)) explicit = true;
project = pq_references(r,ray_list_get(vals,i),keep,0);
}
if (project && explicit) {
int64_t n = 0;
for (int64_t c = 0; c < r->ncols; c++) if (keep[c]) r->selected[n++] = (int32_t)c;
if (!n) { r->selected[0] = 0; n = 1; }
r->nselected = r->noutput = n;
}
ray_free_raw(keep);
ray_t* where = pq_option(query,"where");
int64_t col = -1, lo = INT64_MIN, hi = INT64_MAX;
if (pq_query_range(r,where,&col,&lo,&hi,0) && col >= 0 && lo <= hi) {
err = ray_parquet_range(r,r->schema[col].name,lo,hi);
if (!err) r->filter_nulls = lo == INT64_MIN;
}
if (!err && where && where->type == RAY_LIST && where->len == 3) {
const char* fn = pq_builtin_name(ray_list_get(where,0));
int64_t c = pq_column_id(r,ray_list_get(where,1));
ray_t* pat = ray_list_get(where,2);
if (fn && !strcmp(fn,"like") && pq_column_id(r,ray_list_get(where,0)) < 0 &&
c >= 0 && r->schema[c].type == RAY_STR && pat && pat->type == -RAY_STR) {
for (int64_t i = 0; i < r->nselected; i++) if (r->selected[i] == c) {
r->text_pos = i; r->text_pattern = pat; ray_retain(pat); break;
}
}
}
ray_t* streamed = err ? NULL : pq_stream_aggregates(r,query,where);
if (where) ray_release(where);
if (err) { ray_parquet_close(r); return err; }
if (streamed) { *complete = true; ray_parquet_close(r); return streamed; }
return pq_materialize(r);
}