From f2cd6a3b03bca74be2700b895f9594d30e35f346 Mon Sep 17 00:00:00 2001 From: Justin Kim Date: Mon, 5 Oct 2026 16:19:27 +0900 Subject: [PATCH] feat(join): consume numeric expression MAPs in bounded join pipelines The bounded JOIN -> FILTER* -> MAP pipeline (#1475) accepted only projection MAPs. A MAP that computes anything, even `y + 1`, was reported as map-not-row-local, so the join was materialized in full before the MAP ran, or the step failed with ENOTSUP under WIRELOG_JOIN_BATCH_STRICT=1. The preflight now also accepts MAP expressions built only from numeric variables, numeric and boolean literals, integer and float arithmetic, and comparisons. Such an expression is a pure function of its row and reads no intern, extension or session state. Each batch is mapped by the ordinary col_op_map, so evaluation keeps that operator's semantics. String, digest, UUID, aggregate and extension expressions stay excluded, and FILTERs still admit no arithmetic. The sink used to report every FILTER or MAP failure as ENOMEM. It now records the operator's errno and returns it: ERANGE for a failed expression, as the one-shot MAP returns, and ENOSPC for a denied MAP admission. With a join output limit set, the order can differ: the pipeline counts mapped rows after each batch, so an early failing row yields ERANGE where the one-shot join, which counts join rows before the MAP runs, may stop first with EOVERFLOW. Either way the output relation is discarded and the left input is restored to the stack. Tests in test_join_pipeline.c: - The pipeline matches the one-shot oracle on a 256-row join, and the MAP never sees more than one 7-row batch. - A division by zero in a late batch returns ERANGE, as the oracle does. - A MAP failure injected on the fifth batch, the output limit, and a governor budget sweep that denies admission after committed batches each publish nothing and restore the left input. The sweep also checks that every governor reservation is released. - Stale input falls back to the one-shot path, or fails with ENOTSUP in strict mode. - An excluded digest MAP records a pipeline fallback, or fails with ENOTSUP in strict mode before it is evaluated. docs/MEMORY.md replaces the stale note that pipeline consumption "remains #1475" with the current eligibility and error contract. Fixes #1777 --- CHANGELOG.md | 11 + docs/MEMORY.md | 18 +- tests/meson.build | 1 + tests/test_join_pipeline.c | 591 ++++++++++++++++++++++++++++++- wirelog/columnar/join_pipeline.c | 74 +++- wirelog/columnar/join_pipeline.h | 10 +- 6 files changed, 682 insertions(+), 23 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index bc48f9de2..71cb4d970 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -19,6 +19,17 @@ All notable changes to wirelog are documented in this file. ### Changed +- **Bounded join pipelines consume expression MAPs** (#1777): with + `WIRELOG_JOIN_BATCH_BYTES` set, a `JOIN -> FILTER* -> MAP` pipeline + whose MAP computes integer or float arithmetic or comparisons is now + evaluated batch by batch instead of materializing the whole join first; + previously only projection MAPs qualified. String, digest, UUID, + aggregate and extension expressions still fall back to the one-shot + join (`ENOTSUP` under `WIRELOG_JOIN_BATCH_STRICT=1`). A FILTER or MAP + failure inside the pipeline now returns that operator's errno -- for + MAP, `ERANGE` for an expression error and `ENOSPC` for a denied + admission -- where every such failure previously became `ENOMEM`. + - **Bounded join scratch is sized lazily** (#1481): `col_join_batch_producer_create` reserved `WIRELOG_JOIN_BATCH_BYTES / row_bytes` rows of governed scratch at create, whatever the join's diff --git a/docs/MEMORY.md b/docs/MEMORY.md index 1379b8b8c..a2b6b7a63 100644 --- a/docs/MEMORY.md +++ b/docs/MEMORY.md @@ -1219,8 +1219,22 @@ join; the session records the fallback count and last reason and logs one `JOIN` warning per reason. `WIRELOG_JOIN_BATCH_STRICT=1` (ignored unless the bytes knob is set) turns a fallback into an `ENOTSUP` failure, which is the explicit bounded-mode result for unsupported shapes. Differential -keyed joins have the separate producer described below; pipeline consumption -of batches on the eval stack (JOIN -> FILTER* -> MAP) remains #1475. +keyed joins have the separate producer described below. + +A non-differential eligible join followed by `FILTER* -> MAP` is consumed +batch by batch on the eval stack (#1475), so the full join intermediate is +not materialized while the pipeline runs: each batch is filtered and mapped +into one output relation, which is published only after the last batch. The +FILTERs may use only numeric variables, numeric and boolean literals, and +comparisons. The MAP may project columns or evaluate numeric expressions -- +integer and float arithmetic and comparisons (#1777). A MAP with a string, +digest, UUID, aggregate or extension expression is not consumed; those have +no batch contract yet. A JOIN whose downstream operators are not +consumable, or whose input goes stale mid-stream, falls back to the one-shot +join with reason `pipeline-not-consumable` (`ENOTSUP` under +`WIRELOG_JOIN_BATCH_STRICT=1`). A FILTER or MAP failure inside the pipeline +returns that operator's errno -- for MAP, `ERANGE` for a failed expression +and `ENOSPC` for a denied admission -- and publishes no rows. Ownership and admission: diff --git a/tests/meson.build b/tests/meson.build index 06b855115..a92a82b30 100644 --- a/tests/meson.build +++ b/tests/meson.build @@ -2069,6 +2069,7 @@ test_join_pipeline_exe = executable( 'test_join_pipeline', files('test_join_pipeline.c'), include_directories: [wirelog_inc, wirelog_src_inc], + c_args: ['-DWL_SESSION_TEST_HOOKS=1'], dependencies: [nanoarrow_dep, threads_dep, xxhash_dep, mbedtls_dep, math_dep], link_with: [testlib_prod], ) diff --git a/tests/test_join_pipeline.c b/tests/test_join_pipeline.c index a4da1e6d7..1735368c0 100644 --- a/tests/test_join_pipeline.c +++ b/tests/test_join_pipeline.c @@ -1,9 +1,10 @@ -/* Conservative pipeline preflight for Issue #1475. */ +/* Conservative pipeline preflight for Issue #1475; expression MAPs #1777. */ #include "../wirelog/columnar/internal.h" #include "../wirelog/columnar/join_batch.h" #include "../wirelog/columnar/join_pipeline.h" #include "../wirelog/session.h" +#include #include #include #include @@ -108,8 +109,9 @@ test_conservative_exclusions(void) { static const char *const keys[] = { "col0" }; static const uint32_t project[] = { 0 }; - static uint8_t expr_data[] = { WL_PLAN_EXPR_CONST_INT, 0, 0, 0, 0, - 0, 0, 0, 0 }; + /* A string literal consults the session intern table, so this MAP + * stays excluded even though numeric expression MAPs are admitted. */ + static uint8_t expr_data[] = { WL_PLAN_EXPR_CONST_STR, 1, 0, 's' }; static uint8_t extension_filter[] = { WL_PLAN_EXPR_EXTENSION_CALL, 1, 0, 'f', 0, 0, 0, 0 }; @@ -144,7 +146,7 @@ test_conservative_exclusions(void) ops[1].map_expr_count = 1; CHECK(wl_columnar_join_pipeline_preflight(&plan, 0, &sess, NULL) == WL_COLUMNAR_JOIN_PIPELINE_EXCLUDED_MAP, - "expression MAP is excluded until transactional consumption exists"); + "intern-dependent string MAP is excluded"); ops[1].map_exprs = NULL; ops[1].map_expr_count = 0; ops[1].op = WL_PLAN_OP_FILTER; @@ -215,7 +217,7 @@ test_conservative_exclusions(void) ops[1].map_expr_count = 1; CHECK(wl_columnar_join_pipeline_preflight(&plan, 0, &sess, NULL) == WL_COLUMNAR_JOIN_PIPELINE_EXCLUDED_MAP, - "expression MAP is excluded until transactional consumption exists"); + "intern-dependent string MAP is excluded"); ops[1].map_exprs = NULL; ops[1].map_expr_count = 0; ops[1].op = WL_PLAN_OP_FILTER; @@ -398,12 +400,591 @@ test_pipeline_consumes_batches(void) pipeline_session_destroy(sess); } +/* ---- Issue #1777: row-local expression MAPs ---------------------------- */ + +#define PL_VAR(digit) WL_PLAN_EXPR_VAR, 4, 0, 'c', 'o', 'l', (digit) +#define PL_INT(byte) WL_PLAN_EXPR_CONST_INT, (byte), 0, 0, 0, 0, 0, 0, 0 + +static void +test_expression_map_preflight(void) +{ + static const char *const keys[] = { "col0" }; + static const uint32_t project[] = { 0, 1, 2 }; + static uint8_t sum[] = { PL_VAR('1'), PL_VAR('3'), WL_PLAN_EXPR_ARITH_ADD }; + static uint8_t bnot[] = { PL_VAR('0'), WL_PLAN_EXPR_ARITH_BNOT }; + static uint8_t float_div[] = { + WL_PLAN_EXPR_VAR_FLOAT, 4, 0, 'c', 'o', 'l', '1', + WL_PLAN_EXPR_CONST_FLOAT, 0, 0, 0, 0, 0, 0, 0xf0, 0x3f, + WL_PLAN_EXPR_ARITH_FLOAT_DIV, + }; + static uint8_t cmp[] = { PL_VAR('1'), PL_INT(2), WL_PLAN_EXPR_CMP_LT }; + static uint8_t var_string[] = { + WL_PLAN_EXPR_VAR_STRING, 4, 0, 'c', 'o', 'l', '1' + }; + static uint8_t extension[] = { + WL_PLAN_EXPR_EXTENSION_CALL, 1, 0, 'f', 0, 0, 0, 0 + }; + static uint8_t hash[] = { PL_VAR('1'), WL_PLAN_EXPR_ARITH_HASH }; + static uint8_t uuid4[] = { WL_PLAN_EXPR_ARITH_UUID4 }; + /* Leading operator: only its own depth guard rejects this, since the + * two trailing variables would otherwise end at depth one. */ + static uint8_t add_underflow[] = { + WL_PLAN_EXPR_ARITH_ADD, PL_VAR('1'), PL_VAR('3') + }; + /* The trailing literal restores the final depth to one, so only the + * unary operator's own operand check can reject this. */ + static uint8_t bnot_underflow[] = { WL_PLAN_EXPR_ARITH_BNOT, PL_INT(1) }; + static uint8_t residual[] = { PL_INT(1), PL_INT(2) }; + static uint8_t truncated[] = { WL_PLAN_EXPR_CONST_INT, 1, 0, 0 }; + static uint8_t arith_filter[] = { + PL_VAR('1'), PL_INT(1), WL_PLAN_EXPR_ARITH_ADD, PL_INT(2), + WL_PLAN_EXPR_CMP_GT, + }; + static const struct { + uint8_t *data; + uint32_t size; + wl_columnar_join_pipeline_eligibility_t expected; + const char *message; + } cases[] = { + { sum, sizeof(sum), WL_COLUMNAR_JOIN_PIPELINE_ELIGIBLE, + "integer arithmetic MAP is eligible" }, + { bnot, sizeof(bnot), WL_COLUMNAR_JOIN_PIPELINE_ELIGIBLE, + "unary integer MAP is eligible" }, + { float_div, sizeof(float_div), WL_COLUMNAR_JOIN_PIPELINE_ELIGIBLE, + "float arithmetic MAP is eligible" }, + { cmp, sizeof(cmp), WL_COLUMNAR_JOIN_PIPELINE_ELIGIBLE, + "comparison-valued MAP is eligible" }, + { var_string, sizeof(var_string), + WL_COLUMNAR_JOIN_PIPELINE_EXCLUDED_MAP, + "string variable MAP is excluded" }, + { extension, sizeof(extension), WL_COLUMNAR_JOIN_PIPELINE_EXCLUDED_MAP, + "extension MAP is excluded" }, + { hash, sizeof(hash), WL_COLUMNAR_JOIN_PIPELINE_EXCLUDED_MAP, + "digest MAP is excluded" }, + { uuid4, sizeof(uuid4), WL_COLUMNAR_JOIN_PIPELINE_EXCLUDED_MAP, + "non-deterministic MAP is excluded" }, + { add_underflow, sizeof(add_underflow), + WL_COLUMNAR_JOIN_PIPELINE_EXCLUDED_MAP, + "binary operator stack underflow is excluded" }, + { bnot_underflow, sizeof(bnot_underflow), + WL_COLUMNAR_JOIN_PIPELINE_EXCLUDED_MAP, + "unary operator stack underflow is excluded" }, + { residual, sizeof(residual), WL_COLUMNAR_JOIN_PIPELINE_EXCLUDED_MAP, + "residual MAP value is excluded" }, + { truncated, sizeof(truncated), WL_COLUMNAR_JOIN_PIPELINE_EXCLUDED_MAP, + "truncated MAP literal is excluded" }, + }; + wl_plan_expr_buffer_t exprs[2]; + wl_plan_op_t ops[2]; + wl_plan_relation_t plan = { .ops = ops, .op_count = 2 }; + wl_col_session_t sess; + uint32_t map_index; + + memset(&sess, 0, sizeof(sess)); + sess.join_batch_bytes = 4096; + memset(ops, 0, sizeof(ops)); + ops[0] = (wl_plan_op_t){ .op = WL_PLAN_OP_JOIN, .right_relation = "right", + .left_keys = keys, .right_keys = keys, + .key_count = 1 }; + ops[1] = (wl_plan_op_t){ .op = WL_PLAN_OP_MAP, .map_exprs = exprs, + .map_expr_count = 1, .project_count = 3 }; + for (size_t i = 0; i < sizeof(cases) / sizeof(cases[0]); i++) { + exprs[0] = (wl_plan_expr_buffer_t){ cases[i].data, cases[i].size }; + map_index = 7; + CHECK(wl_columnar_join_pipeline_preflight(&plan, 0, &sess, &map_index) + == cases[i].expected + && map_index == (cases[i].expected + == WL_COLUMNAR_JOIN_PIPELINE_ELIGIBLE ? 1u : UINT32_MAX), + cases[i].message); + } + + /* Columns without an expression fall back to projection, so a mixed + * MAP is judged by its non-empty expressions only. */ + exprs[0] = (wl_plan_expr_buffer_t){ sum, sizeof(sum) }; + exprs[1] = (wl_plan_expr_buffer_t){ NULL, 0 }; + ops[1].map_expr_count = 2; + ops[1].project_indices = project; + CHECK(wl_columnar_join_pipeline_preflight(&plan, 0, &sess, NULL) + == WL_COLUMNAR_JOIN_PIPELINE_ELIGIBLE, + "mixed expression/projection MAP is eligible"); + exprs[1] = (wl_plan_expr_buffer_t){ hash, sizeof(hash) }; + CHECK(wl_columnar_join_pipeline_preflight(&plan, 0, &sess, NULL) + == WL_COLUMNAR_JOIN_PIPELINE_EXCLUDED_MAP, + "one excluded expression excludes the whole MAP"); + ops[1].map_expr_count = 1; + ops[1].map_exprs = NULL; + CHECK(wl_columnar_join_pipeline_preflight(&plan, 0, &sess, NULL) + == WL_COLUMNAR_JOIN_PIPELINE_EXCLUDED_MAP, + "expression count without storage is excluded"); + ops[1].map_exprs = exprs; + ops[1].project_count = 0; + CHECK(wl_columnar_join_pipeline_preflight(&plan, 0, &sess, NULL) + == WL_COLUMNAR_JOIN_PIPELINE_EXCLUDED_MAP, + "zero-width expression MAP is excluded"); + + /* Admitting arithmetic in MAP must not widen the FILTER contract. */ + { + wl_plan_op_t filtered[3]; + wl_plan_relation_t filtered_plan = { .ops = filtered, .op_count = 3 }; + exprs[0] = (wl_plan_expr_buffer_t){ sum, sizeof(sum) }; + filtered[0] = ops[0]; + filtered[1] = (wl_plan_op_t){ .op = WL_PLAN_OP_FILTER, + .filter_expr = { arith_filter, + sizeof(arith_filter) } }; + filtered[2] = (wl_plan_op_t){ .op = WL_PLAN_OP_MAP, .map_exprs = exprs, + .map_expr_count = 1, + .project_count = 3 }; + CHECK(wl_columnar_join_pipeline_preflight(&filtered_plan, 0, &sess, + NULL) == WL_COLUMNAR_JOIN_PIPELINE_EXCLUDED_FILTER, + "arithmetic FILTER stays excluded"); + } +} + +#define PL_KEYS 4u +#define PL_FANOUT 64u +#define PL_JOIN_ROWS (PL_KEYS * PL_FANOUT) +/* join_batch_bytes = 7 * 32 bytes over four int64 join columns. */ +#define PL_BATCH_ROWS 7u +/* FILTER col3 >= 1002 drops key 0 and the first two rows of key 1. */ +#define PL_FILTERED_ROWS (PL_JOIN_ROWS - PL_FANOUT - 2u) + +static uint32_t pl_map_calls; +static uint32_t pl_map_max_rows; +static uint32_t pl_fail_map_call; +static col_rel_t *pl_touch_on_map; + +static void +pl_map_observer(eval_stack_t *stack, eval_entry_t *entry) +{ + (void)stack; + pl_map_calls++; + if (entry && entry->rel && entry->rel->nrows > pl_map_max_rows) + pl_map_max_rows = entry->rel->nrows; + if (pl_fail_map_call != 0 && pl_map_calls == pl_fail_map_call) + wl_columnar_ops_test_map_fail_output_alloc = true; + if (pl_touch_on_map) { + wl_columnar_relation_touch_view(pl_touch_on_map); + pl_touch_on_map = NULL; + } +} + +static void +pl_observe_reset(void) +{ + pl_map_calls = 0; + pl_map_max_rows = 0; + pl_fail_map_call = 0; + pl_touch_on_map = NULL; + wl_columnar_ops_test_map_fail_output_alloc = false; +} + +static wl_col_session_t * +expr_pipeline_session(void) +{ + static const char *const right_names[] = { "k", "r" }; + int64_t rows[PL_JOIN_ROWS * 2u]; + wl_col_session_t *sess; + col_rel_t *right; + + for (uint32_t k = 0; k < PL_KEYS; k++) { + for (uint32_t j = 0; j < PL_FANOUT; j++) { + rows[2u * (k * PL_FANOUT + j)] = k; + rows[2u * (k * PL_FANOUT + j) + 1u] = (int64_t)k * 1000 + j; + } + } + sess = pipeline_session(); + right = pipeline_relation("right", right_names, 2, rows, PL_JOIN_ROWS); + if (!sess || !right) { + col_rel_destroy(right); + pipeline_session_destroy(sess); + return NULL; + } + if (session_add_rel(sess, right) != 0) { + col_rel_destroy(right); + pipeline_session_destroy(sess); + return NULL; + } + if (!col_session_get_arrangement(&sess->base, "right", + (const uint32_t[]){ 0 }, 1)) { + pipeline_session_destroy(sess); + return NULL; + } + return sess; +} + +typedef struct { + int rc; + uint32_t stack_top; + uint32_t restored_left_rows; + uint32_t nrows; + int64_t (*rows)[3]; +} pl_run_t; + +static int +pl_row_cmp(const void *a, const void *b) +{ + const int64_t *x = (const int64_t *)a; + const int64_t *y = (const int64_t *)b; + for (uint32_t c = 0; c < 3; c++) { + if (x[c] != y[c]) + return x[c] < y[c] ? -1 : 1; + } + return 0; +} + +/* Evaluate @plan over a fresh four-row left input. batch_bytes == 0 runs + * the one-shot oracle path; otherwise the bounded pipeline is offered. */ +static pl_run_t +pl_run(wl_col_session_t *sess, const wl_plan_relation_t *plan, + uint32_t batch_bytes) +{ + static const char *const left_names[] = { "k", "v" }; + int64_t left_rows[PL_KEYS * 2u]; + pl_run_t run = { .rc = ENOMEM }; + eval_stack_t stack; + uint32_t saved = sess->join_batch_bytes; + col_rel_t *left; + + for (uint32_t k = 0; k < PL_KEYS; k++) { + left_rows[2u * k] = k; + left_rows[2u * k + 1u] = ((int64_t)k + 1) * 10; + } + left = pipeline_relation("left", left_names, 2, left_rows, PL_KEYS); + eval_stack_init(&stack); + if (!left || eval_stack_push(&stack, left, true) != 0) { + col_rel_destroy(left); + return run; + } + sess->join_batch_bytes = batch_bytes; + run.rc = col_eval_relation_plan(plan, &stack, sess); + sess->join_batch_bytes = saved; + run.stack_top = stack.top; + if (run.rc != 0 && stack.top == 1 && stack.items[0].rel) + run.restored_left_rows = stack.items[0].rel->nrows; + if (run.rc == 0 && stack.top == 1 && stack.items[0].rel + && stack.items[0].rel->ncols == 3) { + const col_rel_t *out = stack.items[0].rel; + run.nrows = out->nrows; + run.rows = calloc(out->nrows ? out->nrows : 1u, sizeof(*run.rows)); + if (run.rows) { + for (uint32_t r = 0; r < out->nrows; r++) { + for (uint32_t c = 0; c < 3; c++) + run.rows[r][c] = out->columns[c][r]; + } + qsort(run.rows, out->nrows, sizeof(*run.rows), pl_row_cmp); + } + } + (void)eval_stack_drain(&stack); + return run; +} + +static bool +pl_same_rows(const pl_run_t *a, const pl_run_t *b) +{ + return a->rc == 0 && b->rc == 0 && a->rows && b->rows + && a->nrows == b->nrows + && memcmp(a->rows, b->rows, (size_t)a->nrows * sizeof(*a->rows)) + == 0; +} + +/* JOIN(left.k = right.k) -> FILTER(col3 >= 1002) -> MAP(@exprs, 3 cols). */ +typedef struct { + wl_plan_op_t ops[3]; + wl_plan_relation_t plan; +} pl_plan_t; + +static void +pl_plan_init(pl_plan_t *p, wl_plan_expr_buffer_t *exprs, uint32_t nexprs) +{ + static const char *const keys[] = { "k" }; + static const uint32_t project[] = { 0, 1, 0 }; + static uint8_t filter[] = { + PL_VAR('3'), WL_PLAN_EXPR_CONST_INT, 0xea, 0x03, 0, 0, 0, 0, 0, 0, + WL_PLAN_EXPR_CMP_GTE, + }; + memset(p, 0, sizeof(*p)); + p->ops[0].op = WL_PLAN_OP_JOIN; + p->ops[0].right_relation = "right"; + p->ops[0].left_keys = keys; + p->ops[0].right_keys = keys; + p->ops[0].key_count = 1; + p->ops[1].op = WL_PLAN_OP_FILTER; + p->ops[1].filter_expr = (wl_plan_expr_buffer_t){ filter, sizeof(filter) }; + p->ops[2].op = WL_PLAN_OP_MAP; + p->ops[2].map_exprs = exprs; + p->ops[2].map_expr_count = nexprs; + p->ops[2].project_indices = project; + p->ops[2].project_count = 3; + p->plan.ops = p->ops; + p->plan.op_count = 3; +} + +/* col1 + col3, col3 * 3 - 1, col0 (projected) */ +static uint8_t pl_sum[] = { PL_VAR('1'), PL_VAR('3'), WL_PLAN_EXPR_ARITH_ADD }; +static uint8_t pl_scaled[] = { + PL_VAR('3'), PL_INT(3), WL_PLAN_EXPR_ARITH_MUL, PL_INT(1), + WL_PLAN_EXPR_ARITH_SUB, +}; + +static void +test_expression_pipeline_matches_oracle(void) +{ + wl_plan_expr_buffer_t exprs[] = { + { pl_sum, sizeof(pl_sum) }, { pl_scaled, sizeof(pl_scaled) }, + }; + wl_col_session_t *sess = expr_pipeline_session(); + pl_plan_t p; + pl_run_t got = { 0 }; + pl_run_t want = { 0 }; + uint32_t fallbacks; + + CHECK(sess != NULL, "expression pipeline fixture"); + if (!sess) + return; + pl_plan_init(&p, exprs, 2); + pl_observe_reset(); + wl_columnar_ops_test_before_map_dispose = pl_map_observer; + + fallbacks = sess->join_batch_fallback_count; + got = pl_run(sess, &p.plan, sess->join_batch_bytes); + CHECK(got.rc == 0 && got.nrows == PL_FILTERED_ROWS, + "expression pipeline produces every filtered row"); + CHECK(sess->join_batch_fallback_count == fallbacks, + "expression pipeline is consumed without fallback"); + CHECK(pl_map_calls >= (PL_JOIN_ROWS + PL_BATCH_ROWS - 1u) / PL_BATCH_ROWS + && pl_map_max_rows <= PL_BATCH_ROWS, + "expression MAP only ever sees one bounded join batch"); + + pl_observe_reset(); + want = pl_run(sess, &p.plan, 0); + CHECK(want.rc == 0 && pl_map_calls == 1 + && pl_map_max_rows == PL_FILTERED_ROWS, + "one-shot oracle materializes the filtered join for MAP"); + CHECK(pl_same_rows(&got, &want), + "expression pipeline matches the one-shot multiset"); + if (got.rows && got.nrows == PL_FILTERED_ROWS) { + /* Largest row: key 3, j 63 -> (40 + 3063, 3063 * 3 - 1, 3). */ + const int64_t *last = got.rows[got.nrows - 1u]; + CHECK(last[0] == 3103 && last[1] == 9188 && last[2] == 3, + "expression values are computed per joined row"); + } + + free(got.rows); + free(want.rows); + wl_columnar_ops_test_before_map_dispose = NULL; + pipeline_session_destroy(sess); +} + +static void +test_expression_pipeline_failures_roll_back(void) +{ + /* col1 / (col0 - 3): only key 3, the last left row, divides by zero. */ + static uint8_t late_div[] = { + PL_VAR('1'), PL_VAR('0'), PL_INT(3), WL_PLAN_EXPR_ARITH_SUB, + WL_PLAN_EXPR_ARITH_DIV, + }; + wl_plan_expr_buffer_t div_exprs[] = { { late_div, sizeof(late_div) } }; + wl_plan_expr_buffer_t exprs[] = { + { pl_sum, sizeof(pl_sum) }, { pl_scaled, sizeof(pl_scaled) }, + }; + wl_col_session_t *sess = expr_pipeline_session(); + pl_plan_t div_plan; + pl_plan_t p; + pl_run_t run; + + CHECK(sess != NULL, "expression failure fixture"); + if (!sess) + return; + pl_plan_init(&div_plan, div_exprs, 1); + pl_plan_init(&p, exprs, 2); + wl_columnar_ops_test_before_map_dispose = pl_map_observer; + + pl_observe_reset(); + run = pl_run(sess, &div_plan.plan, sess->join_batch_bytes); + CHECK(run.rc == ERANGE, + "expression error keeps the ordinary MAP ERANGE policy"); + CHECK(pl_map_calls > 1, + "expression error arrives after earlier batches were committed"); + CHECK(run.stack_top == 1 && run.restored_left_rows == PL_KEYS, + "expression error restores the left input and publishes nothing"); + free(run.rows); + pl_observe_reset(); + run = pl_run(sess, &div_plan.plan, 0); + CHECK(run.rc == ERANGE, "one-shot oracle reports the same ERANGE"); + free(run.rows); + + pl_observe_reset(); + pl_fail_map_call = 5; + run = pl_run(sess, &p.plan, sess->join_batch_bytes); + CHECK(run.rc == ENOMEM && pl_map_calls == 5, + "mid-stream MAP allocation failure stops at the fifth batch"); + CHECK(run.stack_top == 1 && run.restored_left_rows == PL_KEYS, + "mid-stream sink failure publishes nothing"); + free(run.rows); + + pl_observe_reset(); + sess->join_output_limit = 50; + run = pl_run(sess, &p.plan, sess->join_batch_bytes); + sess->join_output_limit = 0; + CHECK(run.rc == EOVERFLOW && run.stack_top == 1 + && run.restored_left_rows == PL_KEYS, + "expression pipeline enforces the join output limit"); + free(run.rows); + + { + pl_run_t want; + uint32_t fallbacks = sess->join_batch_fallback_count; + pl_observe_reset(); + want = pl_run(sess, &p.plan, 0); + pl_observe_reset(); + pl_touch_on_map = session_find_rel(sess, "right"); + run = pl_run(sess, &p.plan, sess->join_batch_bytes); + CHECK(run.rc == 0 && pl_same_rows(&run, &want), + "stale input abandons the pipeline for the ordinary path"); + CHECK(sess->join_batch_fallback_count == fallbacks + 1u + && sess->join_batch_last_reason + == COL_JOIN_BATCH_EXCLUDED_PIPELINE, + "stale pipeline records a fallback"); + free(run.rows); + free(want.rows); + + pl_observe_reset(); + sess->join_batch_strict = true; + pl_touch_on_map = session_find_rel(sess, "right"); + run = pl_run(sess, &p.plan, sess->join_batch_bytes); + sess->join_batch_strict = false; + CHECK(run.rc == ENOTSUP, "strict mode rejects a stale pipeline"); + free(run.rows); + } + + pl_observe_reset(); + wl_columnar_ops_test_before_map_dispose = NULL; + pipeline_session_destroy(sess); +} + +static void +test_expression_pipeline_budget_sweep(void) +{ + wl_plan_expr_buffer_t exprs[] = { + { pl_sum, sizeof(pl_sum) }, { pl_scaled, sizeof(pl_scaled) }, + }; + wl_col_session_t *sess = expr_pipeline_session(); + pl_plan_t p; + pl_run_t want; + bool saw_mid_stream_denial = false; + bool saw_success = false; + + CHECK(sess != NULL, "expression budget fixture"); + if (!sess) + return; + pl_plan_init(&p, exprs, 2); + pl_observe_reset(); + want = pl_run(sess, &p.plan, 0); + CHECK(want.rc == 0, "budget sweep oracle"); + wl_columnar_ops_test_before_map_dispose = pl_map_observer; + + for (uint64_t budget = 256; budget <= (1u << 20) && !saw_success; + budget += 256) { + wl_columnar_memory_resolution_t resolution = { 0 }; + wl_columnar_memory_governor_ref_t *ref; + wl_columnar_memory_governor_t *governor; + pl_run_t run; + + resolution.budget_bytes = budget; + resolution.usable_bytes = budget; + resolution.mode = WL_COLUMNAR_MEMORY_MODE_ENFORCING; + resolution.source = WL_COLUMNAR_MEMORY_SOURCE_ENV; + resolution.status = WL_COLUMNAR_MEMORY_OK; + ref = wl_columnar_memory_governor_ref_create(&resolution); + CHECK(ref != NULL, "budget sweep governor"); + if (!ref) + break; + governor = wl_columnar_memory_governor_ref_get(ref); + sess->memory_governor = ref; + pl_observe_reset(); + run = pl_run(sess, &p.plan, sess->join_batch_bytes); + sess->memory_governor = NULL; + if (run.rc == 0) { + saw_success = true; + CHECK(pl_same_rows(&run, &want), + "governed expression pipeline matches the oracle"); + } else { + CHECK(run.rc == ENOSPC, + "budget denial reports ENOSPC"); + CHECK(run.stack_top == 1 && run.restored_left_rows == PL_KEYS, + "budget denial restores the left input"); + if (pl_map_calls > 1) + saw_mid_stream_denial = true; + } + CHECK(wl_columnar_memory_reserved(governor) == 0, + "every pipeline reservation is released"); + free(run.rows); + wl_columnar_memory_governor_ref_release(ref); + } + CHECK(saw_mid_stream_denial, + "budget sweep denies a reservation after committed batches"); + CHECK(saw_success, "budget sweep reaches a sufficient budget"); + + free(want.rows); + pl_observe_reset(); + wl_columnar_ops_test_before_map_dispose = NULL; + pipeline_session_destroy(sess); +} + +static void +test_excluded_map_is_diagnosed(void) +{ + static uint8_t hash[] = { PL_VAR('1'), WL_PLAN_EXPR_ARITH_HASH }; + wl_plan_expr_buffer_t exprs[] = { { hash, sizeof(hash) } }; + wl_col_session_t *sess = expr_pipeline_session(); + pl_plan_t p; + pl_run_t run; + pl_run_t want; + uint32_t fallbacks; + + CHECK(sess != NULL, "excluded MAP fixture"); + if (!sess) + return; + pl_plan_init(&p, exprs, 1); + wl_columnar_ops_test_before_map_dispose = pl_map_observer; + + pl_observe_reset(); + want = pl_run(sess, &p.plan, 0); + fallbacks = sess->join_batch_fallback_count; + pl_observe_reset(); + run = pl_run(sess, &p.plan, sess->join_batch_bytes); + CHECK(run.rc == 0 && pl_same_rows(&run, &want) && pl_map_calls == 1, + "non-strict excluded MAP falls back to the one-shot path"); + CHECK(sess->join_batch_fallback_count == fallbacks + 1u + && sess->join_batch_last_reason == COL_JOIN_BATCH_EXCLUDED_PIPELINE, + "non-strict excluded MAP records a pipeline fallback"); + free(run.rows); + free(want.rows); + + pl_observe_reset(); + sess->join_batch_strict = true; + run = pl_run(sess, &p.plan, sess->join_batch_bytes); + sess->join_batch_strict = false; + CHECK(run.rc == ENOTSUP && pl_map_calls == 0, + "strict excluded MAP returns ENOTSUP before evaluating it"); + free(run.rows); + + pl_observe_reset(); + wl_columnar_ops_test_before_map_dispose = NULL; + pipeline_session_destroy(sess); +} + int main(void) { test_projection_pipeline(); test_conservative_exclusions(); test_pipeline_consumes_batches(); + test_expression_map_preflight(); + test_expression_pipeline_matches_oracle(); + test_expression_pipeline_failures_roll_back(); + test_expression_pipeline_budget_sweep(); + test_excluded_map_is_diagnosed(); if (failures) fprintf(stderr, "%d join pipeline preflight checks failed\n", failures); diff --git a/wirelog/columnar/join_pipeline.c b/wirelog/columnar/join_pipeline.c index 7274fe46f..127ffafe2 100644 --- a/wirelog/columnar/join_pipeline.c +++ b/wirelog/columnar/join_pipeline.c @@ -6,9 +6,13 @@ #include #include +/* Row-local numeric bytecode: variables, literals and comparisons, a pure + * function of the current row that reads no intern, extension or session + * state. MAP (Issue #1777) also admits integer and float arithmetic, whose + * checked failures surface as the ordinary MAP's ERANGE; FILTER does not. */ static bool -wl_columnar_join_pipeline_filter_is_row_local( - const wl_plan_expr_buffer_t *expr) +wl_columnar_join_pipeline_expr_is_row_local( + const wl_plan_expr_buffer_t *expr, bool allow_arithmetic) { const uint8_t *data; uint32_t pos = 0; @@ -63,10 +67,32 @@ wl_columnar_join_pipeline_filter_is_row_local( return false; depth--; break; + case WL_PLAN_EXPR_ARITH_ADD: + case WL_PLAN_EXPR_ARITH_SUB: + case WL_PLAN_EXPR_ARITH_MUL: + case WL_PLAN_EXPR_ARITH_DIV: + case WL_PLAN_EXPR_ARITH_MOD: + case WL_PLAN_EXPR_ARITH_BAND: + case WL_PLAN_EXPR_ARITH_BOR: + case WL_PLAN_EXPR_ARITH_BXOR: + case WL_PLAN_EXPR_ARITH_SHL: + case WL_PLAN_EXPR_ARITH_SHR: + case WL_PLAN_EXPR_ARITH_FLOAT_ADD: + case WL_PLAN_EXPR_ARITH_FLOAT_SUB: + case WL_PLAN_EXPR_ARITH_FLOAT_MUL: + case WL_PLAN_EXPR_ARITH_FLOAT_DIV: + if (!allow_arithmetic || depth < 2u) + return false; + depth--; + break; + case WL_PLAN_EXPR_ARITH_BNOT: + if (!allow_arithmetic || depth < 1u) + return false; + break; default: - /* Strings, arithmetic, aggregates, extensions and unknown tags - * stay on the materialized path until their batch contract is - * specified. */ + /* Strings, digests, UUIDs, aggregates, extensions and unknown + * tags stay on the materialized path until their batch contract + * is specified. */ return false; } } @@ -101,8 +127,8 @@ wl_columnar_join_pipeline_preflight(const wl_plan_relation_t *plan, while (i < plan->op_count && plan->ops[i].op == WL_PLAN_OP_FILTER) { if (plan->ops[i].materialized) return WL_COLUMNAR_JOIN_PIPELINE_EXCLUDED_FILTER; - if (!wl_columnar_join_pipeline_filter_is_row_local( - &plan->ops[i].filter_expr)) + if (!wl_columnar_join_pipeline_expr_is_row_local( + &plan->ops[i].filter_expr, false)) return WL_COLUMNAR_JOIN_PIPELINE_EXCLUDED_FILTER; i++; } @@ -110,10 +136,22 @@ wl_columnar_join_pipeline_preflight(const wl_plan_relation_t *plan, return WL_COLUMNAR_JOIN_PIPELINE_EXCLUDED_SHAPE; map = &plan->ops[i]; - if (map->map_expr_count != 0 || map->map_exprs != NULL - || map->project_count == 0 || !map->project_indices - || map->materialized) + if (map->project_count == 0 || map->materialized) return WL_COLUMNAR_JOIN_PIPELINE_EXCLUDED_MAP; + if (map->map_expr_count == 0) { + if (map->map_exprs != NULL || !map->project_indices) + return WL_COLUMNAR_JOIN_PIPELINE_EXCLUDED_MAP; + } else { + if (!map->map_exprs) + return WL_COLUMNAR_JOIN_PIPELINE_EXCLUDED_MAP; + /* An empty slot projects its column, as in col_op_map(). */ + for (uint32_t c = 0; c < map->map_expr_count; c++) { + const wl_plan_expr_buffer_t *expr = &map->map_exprs[c]; + if ((expr->data && expr->size > 0) + && !wl_columnar_join_pipeline_expr_is_row_local(expr, true)) + return WL_COLUMNAR_JOIN_PIPELINE_EXCLUDED_MAP; + } + } if (map_index) *map_index = i; @@ -144,6 +182,7 @@ typedef struct { uint32_t map_index; col_rel_t *out; uint32_t pending_begin; + int error; /* errno of a failed FILTER/MAP, reported unchanged */ bool begun; } wl_join_pipeline_sink_t; @@ -204,13 +243,16 @@ pipeline_sink_append(void *context, eval_stack_init(&batch_stack); rc = eval_stack_push(&batch_stack, (col_rel_t *)batch->payload, false); - if (rc != 0) + if (rc != 0) { + sink->error = rc; return WL_COLUMNAR_CONTINUATION_SINK_FAILURE; + } for (uint32_t i = sink->first_filter; i < sink->map_index; i++) { rc = wl_columnar_filter_op(&sink->plan->ops[i], &batch_stack, sink->sess); if (rc != 0) { (void)eval_stack_drain(&batch_stack); + sink->error = rc; return WL_COLUMNAR_CONTINUATION_SINK_FAILURE; } } @@ -218,6 +260,7 @@ pipeline_sink_append(void *context, sink->sess); if (rc != 0) { (void)eval_stack_drain(&batch_stack); + sink->error = rc; return WL_COLUMNAR_CONTINUATION_SINK_FAILURE; } rc = eval_stack_pop_relation(&batch_stack, &result); @@ -398,7 +441,14 @@ wl_columnar_join_pipeline_try_eval(const wl_plan_relation_t *plan, col_rel_destroy(out); if (eval_stack_repush_entry(stack, &left_entry) != 0) return EFAULT; - return status == WL_COLUMNAR_CONTINUATION_STALE ? EAGAIN : ENOMEM; + if (status == WL_COLUMNAR_CONTINUATION_STALE) + return EAGAIN; + /* A FILTER/MAP failure keeps the operator's own errno (for + * MAP, ERANGE for an expression error and ENOSPC for a denied + * admission) instead of collapsing into ENOMEM. An errno the + * caller treats as a fallback (EAGAIN, ENOENT, ENOTSUP) re-runs + * the restored input on the one-shot path. */ + return sink_ctx.error != 0 ? sink_ctx.error : ENOMEM; } if (col_join_output_limit_reached(sess, out)) { wl_columnar_continuation_cancel(continuation); diff --git a/wirelog/columnar/join_pipeline.h b/wirelog/columnar/join_pipeline.h index 78e5ea19f..d07c10e8b 100644 --- a/wirelog/columnar/join_pipeline.h +++ b/wirelog/columnar/join_pipeline.h @@ -14,9 +14,11 @@ typedef enum { WL_COLUMNAR_JOIN_PIPELINE_EXCLUDED_MAP, } wl_columnar_join_pipeline_eligibility_t; -/* Conservative preflight for the first continuation consumer. Projection - * MAPs are row-local without consulting intern or extension/session state; - * expression MAPs remain excluded until their transactional sink exists. */ +/* Conservative preflight for the continuation consumer. Projection MAPs + * and numeric expression MAPs (Issue #1777) are row-local without consulting + * intern or extension/session state; string, digest, UUID, aggregate and + * extension expressions remain excluded until their ownership and retry + * contracts are defined. */ wl_columnar_join_pipeline_eligibility_t wl_columnar_join_pipeline_preflight(const wl_plan_relation_t *plan, uint32_t join_index, const wl_col_session_t *sess, uint32_t *map_index); @@ -25,7 +27,7 @@ const char * wl_columnar_join_pipeline_eligibility_name( wl_columnar_join_pipeline_eligibility_t reason); -/* Consume an eligible JOIN -> FILTER* -> projection MAP pipeline directly +/* Consume an eligible JOIN -> FILTER* -> MAP pipeline directly * from a bounded join continuation. Returns 0 when the shape is not * eligible, 1 when the pipeline was evaluated and pushed, or an errno value * when an eligible pipeline could not be completed. The caller may then