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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions src/ops/internal.h
Original file line number Diff line number Diff line change
Expand Up @@ -1604,6 +1604,9 @@ ray_t* exec_k_shortest(ray_graph_t* g, ray_op_t* op,
/* ── pivot_exec.c ── */
ray_t* exec_if(ray_graph_t* g, ray_op_t* op);

/* Is a descriptor view worth rebuilding over its own bytes (string.c)? */
bool ray_str_view_should_compact(uint64_t pooled_bytes, int64_t pool_len);

/* Shared-node memo around a sub-evaluation over a swapped g->table
* (exec.c): push sets the outer memo aside and arms one for the current
* table and sub-root; pop tears it down and restores the outer one. */
Expand Down
167 changes: 115 additions & 52 deletions src/ops/pivot.c
Original file line number Diff line number Diff line change
Expand Up @@ -510,9 +510,10 @@ static ray_t* if_scatter_str(ray_t* result, ray_t* value, int64_t* ids,
* STR vector with one row per id (or as many rows as the table, indexed by
* the id), or a broadcast scalar. The result takes the side's 16-byte
* descriptor at ids[j] instead of appending its bytes row by row: pooled
* strings keep pointing into their pool, two different pools are laid end
* to end (the second side's offsets shift), a pooled scalar's bytes go in
* once. Rows not in either id list stay the null descriptor the caller
* strings keep pointing into their pool when the rows keep most of it;
* otherwise (two different pools, a pooled scalar, or a small share of a
* big pool) the result gets a pool of exactly its own bytes. Rows not in
* either id list stay the null descriptor the caller
* zeroed. Returns NULL, the result untouched, for a side this cannot take
* (a SYM branch, a length that is neither) — the caller then falls back to
* the per-row scatter, which also reports the length error. */
Expand All @@ -523,6 +524,7 @@ typedef struct {
bool scalar;
bool full; /* vector of nrows: row ids[j] */
const ray_str_t* desc;
const char* bytes; /* its pool's bytes (NULL when inline-only) */
ray_t* pool;
const char* sp; /* scalar bytes */
size_t sl;
Expand All @@ -548,21 +550,49 @@ static bool if_str_side_init(ray_t* v, int64_t* ids, int64_t n, int64_t nrows,
if (v->len == n) s->full = false;
else if (v->len == nrows) s->full = true;
else return false;
const char* bytes = NULL;
str_resolve(v, &s->desc, &bytes);
str_resolve(v, &s->desc, &s->bytes);
s->pool = str_vec_pool_obj(v);
if (s->pool && RAY_IS_ERR(s->pool)) return false;
return true;
}

/* The side's descriptor for its j-th id. */
static inline ray_str_t if_str_side_desc(const if_str_side_t* s, int64_t j) {
return s->full ? s->desc[s->ids[j]] : s->desc[j];
}

/* Bytes the side's pooled descriptors point at. */
static uint64_t if_str_side_pooled_bytes(const if_str_side_t* s) {
if (!s->v || s->scalar) return 0;
uint64_t sum = 0;
for (int64_t j = 0; j < s->n; j++) {
ray_str_t d = if_str_side_desc(s, j);
if (!ray_str_is_inline(&d)) sum += d.len;
}
return sum;
}

/* Assign compact offsets to the side's pooled descriptors from *run. */
static void if_str_side_offsets(const if_str_side_t* s, uint32_t* newoff, uint64_t* run) {
for (int64_t j = 0; j < s->n; j++) {
ray_str_t d = if_str_side_desc(s, j);
if (ray_str_is_inline(&d)) continue;
newoff[j] = (uint32_t)*run;
*run += d.len;
}
}

typedef struct {
const int64_t* ids;
const ray_str_t* src;
bool scalar;
bool full;
ray_str_t sd; /* the scalar's descriptor */
uint32_t shift;
ray_str_t* dst;
/* compaction: pooled bytes move to dst_bytes at newoff[j] */
const uint32_t* newoff;
const char* src_bytes;
char* dst_bytes;
} if_str_scatter_ctx_t;

static void if_str_scatter_fn(void* vctx, uint32_t wid, int64_t start, int64_t end) {
Expand All @@ -574,17 +604,20 @@ static void if_str_scatter_fn(void* vctx, uint32_t wid, int64_t start, int64_t e
}
for (int64_t j = start; j < end; j++) {
ray_str_t d = c->full ? c->src[c->ids[j]] : c->src[j];
if (c->shift && !ray_str_is_inline(&d)) d.pool_off += c->shift;
if (c->newoff && !ray_str_is_inline(&d)) {
memcpy(c->dst_bytes + c->newoff[j], c->src_bytes + d.pool_off, d.len);
d.pool_off = c->newoff[j];
}
c->dst[c->ids[j]] = d;
}
}

static void if_str_scatter_side(const if_str_side_t* s, uint32_t shift, uint32_t scalar_off,
ray_str_t* dst) {
static void if_str_scatter_side(const if_str_side_t* s, const uint32_t* newoff, char* dst_bytes,
uint32_t scalar_off, ray_str_t* dst) {
if (!s->v) return;
if_str_scatter_ctx_t c = {
.ids = s->ids, .src = s->desc, .scalar = s->scalar, .full = s->full,
.shift = shift, .dst = dst,
.dst = dst, .newoff = newoff, .src_bytes = s->bytes, .dst_bytes = dst_bytes,
};
if (s->scalar) {
memset(&c.sd, 0, sizeof(c.sd));
Expand Down Expand Up @@ -615,32 +648,47 @@ static ray_t* if_scatter_str_desc(ray_t* result,
ray_t* pe = (e.v && !e.scalar) ? e.pool : NULL;
bool t_big = t.v && t.scalar && t.sl > RAY_STR_INLINE_MAX;
bool e_big = e.v && e.scalar && e.sl > RAY_STR_INLINE_MAX;
uint32_t e_shift = 0, ts_off = 0, es_off = 0;
if (!t_big && !e_big && (pt == pe || !pt || !pe)) {
ray_t* shared = pt ? pt : pe;
uint64_t t_ref = if_str_side_pooled_bytes(&t);
uint64_t e_ref = if_str_side_pooled_bytes(&e);
ray_str_t* dst = (ray_str_t*)ray_data(result);

/* One pool (or none) behind both sides, and the rows keep most of it:
* point into it as is. */
bool one_pool = (pt == pe) || !pt || !pe;
ray_t* shared = pt ? pt : pe;
if (!t_big && !e_big && one_pool &&
(!shared || !ray_str_view_should_compact(t_ref + e_ref, shared->len))) {
if (shared) { ray_retain(shared); result->str_pool = shared; }
} else {
int64_t tl = pt ? pt->len : 0;
int64_t el = (pe && pe != pt) ? pe->len : 0;
if (tl < 0 || el < 0) return NULL;
uint64_t total = (uint64_t)tl + (uint64_t)el
+ (t_big ? (uint64_t)t.sl : 0) + (e_big ? (uint64_t)e.sl : 0);
if (total > UINT32_MAX) return NULL;
ray_t* np = ray_alloc(total > 0 ? (size_t)total : 1);
if (!np || RAY_IS_ERR(np)) return NULL;
np->type = RAY_U8;
np->len = (int64_t)total;
char* dst = (char*)ray_data(np);
uint32_t off = 0;
if (tl) { memcpy(dst, ray_data(pt), (size_t)tl); off += (uint32_t)tl; }
if (pe && pe != pt && el) { memcpy(dst + off, ray_data(pe), (size_t)el); e_shift = off; off += (uint32_t)el; }
if (t_big) { memcpy(dst + off, t.sp, t.sl); ts_off = off; off += (uint32_t)t.sl; }
if (e_big) { memcpy(dst + off, e.sp, e.sl); es_off = off; off += (uint32_t)e.sl; }
result->str_pool = np;
if_str_scatter_side(&t, NULL, NULL, 0, dst);
if_str_scatter_side(&e, NULL, NULL, 0, dst);
return result;
}
ray_str_t* dst = (ray_str_t*)ray_data(result);
if_str_scatter_side(&t, 0, ts_off, dst);
if_str_scatter_side(&e, e_shift, es_off, dst);

/* Otherwise a pool of exactly the bytes the result points at: the two
* sides' pooled rows one after the other, then the pooled scalars. */
uint64_t total = t_ref + e_ref + (t_big ? (uint64_t)t.sl : 0) + (e_big ? (uint64_t)e.sl : 0);
if (total > UINT32_MAX) return NULL;
int64_t n_off = (t.v && !t.scalar ? t.n : 0) + (e.v && !e.scalar ? e.n : 0);
ray_t* off_hdr = NULL;
uint32_t* newoff = (uint32_t*)scratch_alloc(&off_hdr, (size_t)(n_off > 0 ? n_off : 1) * sizeof(uint32_t));
if (!newoff) return NULL;
uint32_t* t_off = (t.v && !t.scalar) ? newoff : NULL;
uint32_t* e_off = (e.v && !e.scalar) ? newoff + (t.v && !t.scalar ? t.n : 0) : NULL;
uint64_t run = 0;
if (t_off) if_str_side_offsets(&t, t_off, &run);
if (e_off) if_str_side_offsets(&e, e_off, &run);
uint32_t ts_off = 0, es_off = 0;
ray_t* np = ray_alloc(total > 0 ? (size_t)total : 1);
if (!np || RAY_IS_ERR(np)) { scratch_free(off_hdr); return NULL; }
np->type = RAY_U8;
np->len = (int64_t)total;
char* dst_bytes = (char*)ray_data(np);
if (t_big) { memcpy(dst_bytes + run, t.sp, t.sl); ts_off = (uint32_t)run; run += t.sl; }
if (e_big) { memcpy(dst_bytes + run, e.sp, e.sl); es_off = (uint32_t)run; run += e.sl; }
result->str_pool = np;
if_str_scatter_side(&t, t_off, dst_bytes, ts_off, dst);
if_str_scatter_side(&e, e_off, dst_bytes, es_off, dst);
scratch_free(off_hdr);
return result;
}

Expand Down Expand Up @@ -1051,20 +1099,24 @@ typedef struct {
const ray_str_t* t;
const ray_str_t* e;
ray_str_t* dst;
uint32_t e_shift;
/* two pools: the chosen row's pooled bytes move to dst_bytes at newoff[r] */
const uint32_t* newoff;
const char* t_bytes;
const char* e_bytes;
char* dst_bytes;
} if_str_desc_ctx_t;

static void if_str_desc_fn(void* vctx, uint32_t worker_id, int64_t start, int64_t end) {
(void)worker_id;
const if_str_desc_ctx_t* c = (const if_str_desc_ctx_t*)vctx;
for (int64_t r = start; r < end; r++) {
if (c->cond[r]) {
c->dst[r] = c->t[r];
} else {
ray_str_t d = c->e[r];
if (c->e_shift && !ray_str_is_inline(&d)) d.pool_off += c->e_shift;
c->dst[r] = d;
ray_str_t d = c->cond[r] ? c->t[r] : c->e[r];
if (c->newoff && !ray_str_is_inline(&d)) {
const char* src = c->cond[r] ? c->t_bytes : c->e_bytes;
memcpy(c->dst_bytes + c->newoff[r], src + d.pool_off, d.len);
d.pool_off = c->newoff[r];
}
c->dst[r] = d;
}
}

Expand Down Expand Up @@ -1137,8 +1189,8 @@ static ray_t* exec_if_eager(ray_graph_t* g, ray_op_t* op) {
/* Two STR vectors: the result is descriptors only. Each row takes
* its side's 16-byte descriptor; pooled strings keep pointing into
* their pool. One shared pool (or one side inline-only) is reused
* as is; two different pools are laid end to end in a new pool and
* the else side's offsets shift by the then pool's length. Nulls
* as is; two different pools give a pool of exactly the chosen
* rows' bytes. Nulls
* are empty descriptors and travel unchanged. No per-row append,
* no rehash; the fill runs on the worker pool. */
if (!then_scalar && !else_scalar &&
Expand All @@ -1153,7 +1205,8 @@ static ray_t* exec_if_eager(ray_graph_t* g, ray_op_t* op) {
str_resolve(then_v, &t_desc, &t_bytes);
str_resolve(else_v, &e_desc, &e_bytes);
bool ok = true;
uint32_t e_shift = 0;
ray_t* off_hdr = NULL;
uint32_t* newoff = NULL;
if (then_pool == else_pool || !then_pool || !else_pool) {
ray_t* out_pool = then_pool ? then_pool : else_pool;
if (out_pool && !RAY_IS_ERR(out_pool)) {
Expand All @@ -1163,33 +1216,43 @@ static ray_t* exec_if_eager(ray_graph_t* g, ray_op_t* op) {
} else if (RAY_IS_ERR(then_pool) || RAY_IS_ERR(else_pool)) {
ok = false;
} else {
int64_t tl = then_pool->len, el = else_pool->len;
if (tl < 0 || el < 0 || (uint64_t)tl + (uint64_t)el > UINT32_MAX) {
/* Two pools: a pool of exactly the chosen rows' bytes (a
* serial pass assigns the offsets, the fill copies). */
newoff = (uint32_t*)scratch_alloc(&off_hdr, (size_t)(len > 0 ? len : 1) * sizeof(uint32_t));
if (!newoff) {
ok = false;
} else {
ray_t* np = ray_alloc((size_t)(tl + el) > 0 ? (size_t)(tl + el) : 1);
uint64_t run = 0;
for (int64_t r = 0; r < len; r++) {
const ray_str_t* d = cond_p[r] ? &t_desc[r] : &e_desc[r];
if (ray_str_is_inline(d)) continue;
newoff[r] = (uint32_t)run;
run += d->len;
}
ray_t* np = (run <= UINT32_MAX) ? ray_alloc(run > 0 ? (size_t)run : 1) : NULL;
if (!np || RAY_IS_ERR(np)) {
ok = false;
} else {
np->type = RAY_U8;
np->len = tl + el;
if (tl) memcpy(ray_data(np), t_bytes, (size_t)tl);
if (el) memcpy((char*)ray_data(np) + tl, e_bytes, (size_t)el);
np->len = (int64_t)run;
result->str_pool = np;
e_shift = (uint32_t)tl;
}
}
if (!ok && off_hdr) { scratch_free(off_hdr); off_hdr = NULL; newoff = NULL; }
}
if (ok) {
if_str_desc_ctx_t dctx = {
.cond = cond_p, .t = t_desc, .e = e_desc,
.dst = (ray_str_t*)ray_data(result), .e_shift = e_shift,
.dst = (ray_str_t*)ray_data(result),
.newoff = newoff, .t_bytes = t_bytes, .e_bytes = e_bytes,
.dst_bytes = result->str_pool ? (char*)ray_data(result->str_pool) : NULL,
};
ray_pool_t* pool = ray_pool_get();
if (ray_pool_par_dispatch_ok(pool, len, RAY_PARALLEL_THRESHOLD))
ray_pool_dispatch(pool, if_str_desc_fn, &dctx, len);
else
if_str_desc_fn(&dctx, 0, 0, len);
if (off_hdr) scratch_free(off_hdr);
if (ray_vec_may_have_nulls(then_v) || ray_vec_may_have_nulls(else_v))
result->attrs |= RAY_ATTR_HAS_NULLS;
ray_release(cond_v); ray_release(then_v); ray_release(else_v);
Expand Down
25 changes: 25 additions & 0 deletions src/ops/string.c
Original file line number Diff line number Diff line change
Expand Up @@ -1053,12 +1053,14 @@ typedef struct {
substr_arg_t len;
_Atomic(uint32_t) any_null;
_Atomic(uint32_t) range_err;
_Atomic(uint64_t) pooled_bytes; /* bytes the pooled results point at */
} substr_view_ctx_t;

static void substr_view_fn(void* vctx, uint32_t worker_id, int64_t lo, int64_t hi) {
(void)worker_id;
substr_view_ctx_t* c = (substr_view_ctx_t*)vctx;
bool null_seen = false, range_seen = false;
uint64_t pooled = 0;
for (int64_t i = lo; i < hi; i++) {
ray_str_t* d = &c->dst[i];
memset(d, 0, sizeof(*d));
Expand All @@ -1084,13 +1086,25 @@ static void substr_view_fn(void* vctx, uint32_t worker_id, int64_t lo, int64_t h
if ((uint64_t)s->pool_off + (uint64_t)st > UINT32_MAX) { range_seen = true; d->len = 0; continue; }
memcpy(d->prefix, sp, 4);
d->pool_off = s->pool_off + (uint32_t)st;
pooled += (uint64_t)ln;
/* hash32 stays 0: a consumer that needs it computes it once
* (ray_str_t_hash32); hashing every substring here paid a pass
* over the bytes that most consumers never used. */
}
}
if (null_seen) atomic_store_explicit(&c->any_null, 1, memory_order_relaxed);
if (range_seen) atomic_store_explicit(&c->range_err, 1, memory_order_relaxed);
if (pooled) atomic_fetch_add_explicit(&c->pooled_bytes, pooled, memory_order_relaxed);
}

/* A view whose bytes are a small share (under an eighth) of the pool it
* points into is worth rebuilding over its own bytes: the copy costs the
* few bytes it keeps, the pool it would otherwise pin costs the rest. A
* larger share stays a view — the parent pool is usually alive anyway (a
* column, a sibling intermediate), and copying most of it would only add
* a second copy for the view's lifetime. */
bool ray_str_view_should_compact(uint64_t pooled_bytes, int64_t pool_len) {
return pool_len > 0 && pooled_bytes * 8 < (uint64_t)pool_len;
}

static ray_t* substr_str_view(ray_t* input, ray_t* start_v, ray_t* len_v) {
Expand All @@ -1114,6 +1128,7 @@ static ray_t* substr_str_view(ray_t* input, ray_t* start_v, ray_t* len_v) {
ctx.dst = (ray_str_t*)ray_data(result);
atomic_store_explicit(&ctx.any_null, 0, memory_order_relaxed);
atomic_store_explicit(&ctx.range_err, 0, memory_order_relaxed);
atomic_store_explicit(&ctx.pooled_bytes, 0, memory_order_relaxed);

ray_pool_t* pool = ray_pool_get();
if (ray_pool_par_dispatch_ok(pool, nrows, RAY_PARALLEL_THRESHOLD))
Expand All @@ -1126,6 +1141,16 @@ static ray_t* substr_str_view(ray_t* input, ray_t* start_v, ray_t* len_v) {
}
if (atomic_load_explicit(&ctx.any_null, memory_order_relaxed))
result->attrs |= RAY_ATTR_HAS_NULLS;
/* A result with no pooled descriptor (every substring fits inline)
* has nothing in the parent pool to keep alive. A view that does
* point into it stays a view: the column is alive anyway, and a
* sibling `if` over two such views can pick either side without
* copying (a compacted view would give it two different pools). */
if (result->str_pool &&
atomic_load_explicit(&ctx.pooled_bytes, memory_order_relaxed) == 0) {
ray_release(result->str_pool);
result->str_pool = NULL;
}
return result;
}

Expand Down
Loading
Loading