diff --git a/CHANGELOG.md b/CHANGELOG.md index bc48f9de..71cb4d97 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 1379b8b8..a2b6b7a6 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 06b85511..a92a82b3 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 a4da1e6d..1735368c 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 7274fe46..127ffafe 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 78e5ea19..d07c10e8 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