diff --git a/docs/THREADING.md b/docs/THREADING.md index 7f46cbd6..962ca13e 100644 --- a/docs/THREADING.md +++ b/docs/THREADING.md @@ -700,6 +700,35 @@ excluded from the production library. The complete source audit now contains **222 atomic call sites**. +### 5.20 `wirelog/columnar/ops.c` — LFTJ output growth admission test hook (3 rows) + +The test-only growth hook saves the configured governor limit, sets +`usable_bytes` to the already-reserved total to force a growth denial, then +restores the limit after the append attempt. These operations are excluded +from the production library. + +| Anchor (file:function[#N]) | Field | Op | Order | Justification | +|---|---|---|---|---| +| `ops.c:lftj_test_before_output_growth` | `memory_governor->usable_bytes` | `atomic_load_explicit` | relaxed | Save the configured limit before the test hook forces output-growth admission to fail; test-only and excluded from the production library | +| `ops.c:lftj_test_before_output_growth#2` | `memory_governor->usable_bytes` | `atomic_store_explicit` | relaxed | Temporarily set the limit to the current reserved total so the output-growth reservation is denied; test-only and excluded from the production library | +| `ops.c:lftj_test_after_output_append` | `memory_governor->usable_bytes` | `atomic_store_explicit` | relaxed | Restore the saved governor limit after the append attempt; test-only and excluded from the production library | + +The complete source audit now contains **225 atomic call sites**. + +### 5.21 `wirelog/columnar/ops.c` — LFTJ output construction admission test hook (3 rows) + +The test-only constructor hook records and temporarily lowers the governor +limit to exercise denied initial output allocation, then restores the limit. +These operations are excluded from the production library. + +| Anchor (file:function[#N]) | Field | Op | Order | Justification | +|---|---|---|---|---| +| `ops.c:lftj_test_before_output_construction` | `memory_governor->usable_bytes` | `atomic_load_explicit` | relaxed | Save the configured limit before the test hook forces initial output-construction admission to fail; test-only and excluded from the production library | +| `ops.c:lftj_test_before_output_construction#2` | `memory_governor->usable_bytes` | `atomic_store_explicit` | relaxed | Temporarily lower the limit so initial output construction is denied; test-only and excluded from the production library | +| `ops.c:lftj_test_after_output_construction` | `memory_governor->usable_bytes` | `atomic_store_explicit` | relaxed | Restore the saved limit after the constructor attempt; test-only and excluded from the production library | + +The complete source audit now contains **228 atomic call sites**. + --- ## 6. Lock-free SPSC delta queue diff --git a/scripts/ci/check-threading-doc.sh b/scripts/ci/check-threading-doc.sh index 035fe6bf..32b71390 100755 --- a/scripts/ci/check-threading-doc.sh +++ b/scripts/ci/check-threading-doc.sh @@ -19,7 +19,7 @@ fi rows="$tmp_dir/rows" sed -nE 's/^\| `([^`]+:[A-Za-z_][A-Za-z0-9_]*(#[0-9]+)?)` \| [^|]* \| `([^`]*)` \|.*/\1\t\3/p' "$doc" >"$rows" row_count=$(wc -l <"$rows") -expected_rows="${WIRELOG_THREADING_EXPECTED_ROWS:-222}" +expected_rows="${WIRELOG_THREADING_EXPECTED_ROWS:-228}" [ "$row_count" -eq "$expected_rows" ] || { echo "check-threading-doc: FAIL: expected $expected_rows audit rows, found $row_count" >&2 exit 1 diff --git a/tests/test_lftj_integration.c b/tests/test_lftj_integration.c index cab9e65f..e50df6b2 100644 --- a/tests/test_lftj_integration.c +++ b/tests/test_lftj_integration.c @@ -650,6 +650,294 @@ test_lftj_inner_denial_propagates_and_retries(void) PASS(); } +static void +test_lftj_output_growth_denial_propagates_and_retries(void) +{ + TEST("LFTJ output growth denial and allocator failure restore state"); + char source[8192]; + size_t used = 0; + const char *prefix = ".decl r1(x: int32, a: int32)\n" + ".decl r2(x: int32, b: int32)\n" + ".decl r3(x: int32, c: int32)\n" + ".decl out(x: int32, a: int32, b: int32, c: int32)\n"; + int written = snprintf(source, sizeof(source), "%s", prefix); + if (written < 0 || (size_t)written >= sizeof(source)) + FAIL("query prefix did not fit the test buffer"); + used = (size_t)written; + for (uint32_t i = 1; i <= 100; i++) { + written = snprintf(source + used, sizeof(source) - used, + "r1(%u, %u). r2(%u, %u). r3(%u, %u).\n", + i, i + 100u, i, i + 200u, i, i + 300u); + if (written < 0 || (size_t)written >= sizeof(source) - used) + FAIL("generated input facts did not fit the test buffer"); + used += (size_t)written; + } + written = snprintf(source + used, sizeof(source) - used, + "out(x, a, b, c) :- r1(x, a), r2(x, b), r3(x, c).\n"); + if (written < 0 || (size_t)written >= sizeof(source) - used) + FAIL("query rule did not fit the test buffer"); + + wirelog_program_t *prog = NULL; + wl_plan_t *plan = NULL; + wl_session_t *session = NULL; + wl_columnar_memory_governor_ref_t *ref + = test_governor(UINT64_C(1) << 30); + wl_session_options_t options; + eval_stack_t stack; + bool stack_initialized = false; + wl_columnar_lftj_output_test_hook_state_t hook_state = {0}; + wl_col_session_t *columnar = NULL; + const wl_plan_op_t *lftj_op = NULL; + const char *failure = NULL; + uint64_t baseline_reserved = 0; + uint64_t original_usable = UINT64_C(1) << 30; + int rc; + + wl_columnar_lftj_test_clear_hooks(); + wl_columnar_lftj_test_clear_output_hooks(); + if (!ref) { + failure = "governor creation failed"; + goto cleanup; + } + if (make_plan_no_opt(source, &plan, &prog) != 0) { + failure = "plan creation failed"; + goto cleanup; + } + wl_session_options_init(&options); + options.memory_governor = ref; + if (wl_session_create_with_options(wl_backend_columnar(), plan, 1, + &options, &session) != 0 || !session) { + failure = "session creation failed"; + goto cleanup; + } + if (wl_session_load_facts(session, prog) != 0) { + failure = "loading EDB facts failed"; + goto cleanup; + } + for (uint32_t s = 0; s < plan->stratum_count && !lftj_op; s++) { + const wl_plan_stratum_t *stratum = &plan->strata[s]; + for (uint32_t r = 0; r < stratum->relation_count && !lftj_op; r++) { + const wl_plan_relation_t *relation = &stratum->relations[r]; + for (uint32_t o = 0; o < relation->op_count; o++) { + if (relation->ops[o].op == WL_PLAN_OP_LFTJ) { + lftj_op = &relation->ops[o]; + break; + } + } + } + } + if (!lftj_op) { + failure = "plan has no LFTJ operation"; + goto cleanup; + } + columnar = COL_SESSION(session); + original_usable = atomic_load_explicit( + &wl_columnar_memory_governor_ref_get(ref)->usable_bytes, + memory_order_relaxed); + eval_stack_init(&stack); + stack_initialized = true; + + rc = col_op_lftj(lftj_op, &stack, columnar); + if (rc != 0 || stack.top != 1 || !stack.items[0].rel + || stack.items[0].rel->nrows != 100) { + failure = "warmup did not publish all 100 expected rows"; + goto cleanup; + } + eval_stack_drain(&stack); + baseline_reserved = wl_columnar_memory_reserved( + wl_columnar_memory_governor_ref_get(ref)); + if (baseline_reserved == 0 || columnar->sarr_active_pins != 0) { + failure = "warm arrangements did not reach a stable baseline"; + goto cleanup; + } + + columnar->memory_budget_denied = false; + wl_columnar_lftj_test_deny_next_output_construction(); + rc = col_op_lftj(lftj_op, &stack, columnar); + wl_columnar_lftj_test_get_output_hook_state(&hook_state); + if (hook_state.deny_output_construction_consumed) { + atomic_store_explicit( + &wl_columnar_memory_governor_ref_get(ref)->usable_bytes, + hook_state.previous_constructor_usable_bytes, + memory_order_relaxed); + } + if (!hook_state.deny_output_construction_consumed + || hook_state.reserved_before_construction <= baseline_reserved + || rc != ENOSPC || !columnar->memory_budget_denied || stack.top != 0) { + failure = "output constructor denial did not set session state"; + goto cleanup; + } + if (wl_columnar_memory_reserved( + wl_columnar_memory_governor_ref_get(ref)) != baseline_reserved) { + failure = "output constructor denial leaked reservations"; + goto cleanup; + } + if (columnar->sarr_active_pins != 0) { + failure = "output constructor denial leaked an arrangement probe"; + goto cleanup; + } + for (uint32_t i = 0; i < columnar->sarr_count; i++) { + if (columnar->sarr_entries[i].sarr.pin_count != 0) { + failure = "output constructor denial left an arrangement pinned"; + goto cleanup; + } + } + for (uint32_t i = 0; i < 3; i++) { + static const char *const names[] = { "r1", "r2", "r3" }; + col_rel_t *rel = session_find_rel(columnar, names[i]); + if (!rel || wl_columnar_source_access_gate_busy(&rel->source_access)) { + failure = "output constructor denial leaked a source reader"; + goto cleanup; + } + } + + wl_columnar_lftj_test_clear_output_hooks(); + columnar->memory_budget_denied = false; + rc = col_op_lftj(lftj_op, &stack, columnar); + if (rc != 0 || columnar->memory_budget_denied || stack.top != 1 + || !stack.items[0].rel || stack.items[0].rel->nrows != 100) { + failure = "retry after output constructor denial failed"; + goto cleanup; + } + eval_stack_drain(&stack); + if (wl_columnar_memory_reserved( + wl_columnar_memory_governor_ref_get(ref)) != baseline_reserved) { + failure = "output constructor retry leaked reservations"; + goto cleanup; + } + + columnar->memory_budget_denied = false; + wl_columnar_lftj_test_deny_next_output_growth(); + rc = col_op_lftj(lftj_op, &stack, columnar); + wl_columnar_lftj_test_get_output_hook_state(&hook_state); + if (hook_state.deny_output_growth_consumed) { + atomic_store_explicit( + &wl_columnar_memory_governor_ref_get(ref)->usable_bytes, + hook_state.previous_usable_bytes, memory_order_relaxed); + } + if (!hook_state.deny_output_growth_consumed + || hook_state.output_rows_at_growth == 0 + || hook_state.output_rows_at_growth + < hook_state.output_capacity_at_growth + || hook_state.reserved_before_growth < baseline_reserved + || rc != ENOSPC || !columnar->memory_budget_denied || stack.top != 0) { + failure = "output growth denial was not propagated as ENOSPC"; + goto cleanup; + } + if (wl_columnar_memory_reserved( + wl_columnar_memory_governor_ref_get(ref)) != baseline_reserved) { + failure = "output growth denial leaked reservation bytes"; + goto cleanup; + } + if (columnar->sarr_active_pins != 0) { + failure = "output growth denial leaked an arrangement probe"; + goto cleanup; + } + for (uint32_t i = 0; i < columnar->sarr_count; i++) { + if (columnar->sarr_entries[i].sarr.pin_count != 0) { + failure = "output growth denial left an arrangement pinned"; + goto cleanup; + } + } + for (uint32_t i = 0; i < 3; i++) { + static const char *const names[] = { "r1", "r2", "r3" }; + col_rel_t *rel = session_find_rel(columnar, names[i]); + if (!rel || wl_columnar_source_access_gate_busy(&rel->source_access)) { + failure = "output growth denial leaked a source reader"; + goto cleanup; + } + } + + wl_columnar_lftj_test_clear_output_hooks(); + columnar->memory_budget_denied = false; + rc = col_op_lftj(lftj_op, &stack, columnar); + if (rc != 0 || columnar->memory_budget_denied || stack.top != 1 + || !stack.items[0].rel || stack.items[0].rel->nrows != 100) { + failure = "retry after output growth denial failed"; + goto cleanup; + } + eval_stack_drain(&stack); + if (wl_columnar_memory_reserved( + wl_columnar_memory_governor_ref_get(ref)) != baseline_reserved) { + failure = "successful retry leaked output reservation bytes"; + goto cleanup; + } + + columnar->memory_budget_denied = false; + wl_columnar_lftj_test_fail_next_output_growth(); + rc = col_op_lftj(lftj_op, &stack, columnar); + wl_columnar_lftj_test_get_output_hook_state(&hook_state); + if (!hook_state.fail_output_growth_consumed + || hook_state.output_rows_at_growth == 0 + || hook_state.output_rows_at_growth + < hook_state.output_capacity_at_growth + || rc != ENOMEM + || hook_state.failure_pending_flag || columnar->memory_budget_denied + || stack.top != 0) { + failure = "ordinary output allocator failure was misclassified"; + goto cleanup; + } + if (wl_columnar_memory_reserved( + wl_columnar_memory_governor_ref_get(ref)) != baseline_reserved) { + failure = "ordinary output failure leaked reservations"; + goto cleanup; + } + if (columnar->sarr_active_pins != 0) { + failure = "ordinary output failure leaked an arrangement probe"; + goto cleanup; + } + for (uint32_t i = 0; i < columnar->sarr_count; i++) { + if (columnar->sarr_entries[i].sarr.pin_count != 0) { + failure = "ordinary output failure left an arrangement pinned"; + goto cleanup; + } + } + for (uint32_t i = 0; i < 3; i++) { + static const char *const names[] = { "r1", "r2", "r3" }; + col_rel_t *rel = session_find_rel(columnar, names[i]); + if (!rel || wl_columnar_source_access_gate_busy(&rel->source_access)) { + failure = "ordinary output failure leaked a source reader"; + goto cleanup; + } + } + wl_columnar_lftj_test_clear_output_hooks(); + rc = col_op_lftj(lftj_op, &stack, columnar); + if (rc != 0 || stack.top != 1 || !stack.items[0].rel + || stack.items[0].rel->nrows != 100) { + failure = "retry after ordinary output failure failed"; + goto cleanup; + } + eval_stack_drain(&stack); + if (wl_columnar_memory_reserved( + wl_columnar_memory_governor_ref_get(ref)) != baseline_reserved) { + failure = "ordinary output retry leaked reservation bytes"; + goto cleanup; + } + +cleanup: + wl_columnar_lftj_test_clear_hooks(); + wl_columnar_lftj_test_clear_output_hooks(); + if (columnar && ref) { + atomic_store_explicit( + &wl_columnar_memory_governor_ref_get(ref)->usable_bytes, + original_usable, memory_order_relaxed); + columnar->memory_budget_denied = false; + } + if (stack_initialized) + eval_stack_drain(&stack); + if (session) + wl_session_destroy(session); + if (plan) + wl_plan_free(plan); + if (prog) + wirelog_program_free(prog); + if (ref) + wl_columnar_memory_governor_ref_release(ref); + if (failure) + FAIL(failure); + PASS(); +} + /* Test 5: IDB relation in chain prevents LFTJ rewrite. */ static void test_idb_not_rewritten(void) @@ -860,6 +1148,7 @@ main(void) test_idb_not_rewritten(); test_lftj_stack_overflow_releases_output(); test_lftj_inner_denial_propagates_and_retries(); + test_lftj_output_growth_denial_propagates_and_retries(); test_tdd_lftj_metadata_is_conservative(); printf("\nResults: %d/%d passed", pass_count, test_count); diff --git a/wirelog/columnar/lftj.h b/wirelog/columnar/lftj.h index 55f764e5..f4b5c6dc 100644 --- a/wirelog/columnar/lftj.h +++ b/wirelog/columnar/lftj.h @@ -126,6 +126,24 @@ typedef struct { wl_columnar_memory_admission_status_t inner_status; } wl_columnar_lftj_test_hook_state_t; +typedef struct { + bool deny_output_construction_pending; + bool deny_output_construction_consumed; + bool deny_output_construction_restore_pending; + bool deny_output_growth_pending; + bool deny_output_growth_consumed; + bool deny_output_growth_restore_pending; + bool fail_output_growth_pending; + bool fail_output_growth_consumed; + bool failure_pending_flag; + uint32_t output_rows_at_growth; + uint32_t output_capacity_at_growth; + uint64_t previous_usable_bytes; + uint64_t reserved_before_growth; + uint64_t previous_constructor_usable_bytes; + uint64_t reserved_before_construction; +} wl_columnar_lftj_output_test_hook_state_t; + void wl_columnar_lftj_test_fail_next_iters_alloc(void); void @@ -135,6 +153,18 @@ wl_columnar_lftj_test_get_hook_state( wl_columnar_lftj_test_hook_state_t *out); void wl_columnar_lftj_test_clear_hooks(void); + +void +wl_columnar_lftj_test_deny_next_output_growth(void); +void +wl_columnar_lftj_test_deny_next_output_construction(void); +void +wl_columnar_lftj_test_fail_next_output_growth(void); +void +wl_columnar_lftj_test_get_output_hook_state( + wl_columnar_lftj_output_test_hook_state_t *out); +void +wl_columnar_lftj_test_clear_output_hooks(void); #endif #endif /* WL_COLUMNAR_LFTJ_H */ diff --git a/wirelog/columnar/ops.c b/wirelog/columnar/ops.c index fa1440de..fae819fc 100644 --- a/wirelog/columnar/ops.c +++ b/wirelog/columnar/ops.c @@ -1037,6 +1037,142 @@ col_op_reduce_weighted(const col_rel_t *src, col_rel_t *dst) /* LFTJ Operator (Issue #195) */ /* ======================================================================== */ +#ifdef WL_TEST_LFTJ_ADMISSION_HOOKS +#if defined(_MSC_VER) +#define WL_COLUMNAR_LFTJ_OUTPUT_TEST_TLS __declspec(thread) +#else +#define WL_COLUMNAR_LFTJ_OUTPUT_TEST_TLS _Thread_local +#endif +static WL_COLUMNAR_LFTJ_OUTPUT_TEST_TLS +wl_columnar_lftj_output_test_hook_state_t lftj_output_test_hook_state; + +void +wl_columnar_lftj_test_deny_next_output_growth(void) +{ + lftj_output_test_hook_state.deny_output_growth_pending = true; + lftj_output_test_hook_state.deny_output_growth_consumed = false; +} + +void +wl_columnar_lftj_test_deny_next_output_construction(void) +{ + lftj_output_test_hook_state.deny_output_construction_pending = true; + lftj_output_test_hook_state.deny_output_construction_consumed = false; +} + +void +wl_columnar_lftj_test_fail_next_output_growth(void) +{ + lftj_output_test_hook_state.fail_output_growth_pending = true; + lftj_output_test_hook_state.fail_output_growth_consumed = false; +} + +void +wl_columnar_lftj_test_get_output_hook_state( + wl_columnar_lftj_output_test_hook_state_t *out) +{ + if (out) + *out = lftj_output_test_hook_state; +} + +void +wl_columnar_lftj_test_clear_output_hooks(void) +{ + memset(&lftj_output_test_hook_state, 0, + sizeof(lftj_output_test_hook_state)); +} + +static void +lftj_test_before_output_construction( + wl_columnar_memory_governor_ref_t *governor) +{ + wl_columnar_lftj_output_test_hook_state_t *state + = &lftj_output_test_hook_state; + if (!state->deny_output_construction_pending || !governor) + return; + wl_columnar_memory_governor_t *memory_governor + = wl_columnar_memory_governor_ref_get(governor); + if (!memory_governor) + return; + state->previous_constructor_usable_bytes = atomic_load_explicit( + &memory_governor->usable_bytes, memory_order_relaxed); + state->reserved_before_construction + = wl_columnar_memory_reserved(memory_governor); + atomic_store_explicit(&memory_governor->usable_bytes, + state->reserved_before_construction, memory_order_relaxed); + state->deny_output_construction_pending = false; + state->deny_output_construction_consumed = true; + state->deny_output_construction_restore_pending = true; +} + +static void +lftj_test_after_output_construction( + wl_columnar_memory_governor_ref_t *governor) +{ + wl_columnar_lftj_output_test_hook_state_t *state + = &lftj_output_test_hook_state; + if (!state->deny_output_construction_restore_pending || !governor) + return; + wl_columnar_memory_governor_t *memory_governor + = wl_columnar_memory_governor_ref_get(governor); + if (memory_governor) + atomic_store_explicit(&memory_governor->usable_bytes, + state->previous_constructor_usable_bytes, memory_order_relaxed); + state->deny_output_construction_restore_pending = false; +} + +static bool +lftj_test_before_output_growth(col_rel_t *out, + wl_columnar_memory_governor_ref_t *governor) +{ + wl_columnar_lftj_output_test_hook_state_t *state + = &lftj_output_test_hook_state; + if (!out || out->nrows == 0 || out->nrows < out->capacity) + return false; + if (!state->fail_output_growth_pending + && !state->deny_output_growth_pending) + return false; + state->output_rows_at_growth = out->nrows; + state->output_capacity_at_growth = out->capacity; + if (state->fail_output_growth_pending) { + state->fail_output_growth_pending = false; + state->fail_output_growth_consumed = true; + return true; + } + if (state->deny_output_growth_pending && governor) { + wl_columnar_memory_governor_t *memory_governor + = wl_columnar_memory_governor_ref_get(governor); + if (memory_governor) { + state->previous_usable_bytes = atomic_load_explicit( + &memory_governor->usable_bytes, memory_order_relaxed); + state->reserved_before_growth + = wl_columnar_memory_reserved(memory_governor); + atomic_store_explicit(&memory_governor->usable_bytes, + state->reserved_before_growth, memory_order_relaxed); + state->deny_output_growth_pending = false; + state->deny_output_growth_consumed = true; + state->deny_output_growth_restore_pending = true; + } + } + return false; +} + +static void +lftj_test_after_output_append(wl_columnar_memory_governor_ref_t *governor) +{ + wl_columnar_lftj_output_test_hook_state_t *state + = &lftj_output_test_hook_state; + if (!state->deny_output_growth_restore_pending || !governor) + return; + wl_columnar_memory_governor_t *memory_governor + = wl_columnar_memory_governor_ref_get(governor); + if (memory_governor) + atomic_store_explicit(&memory_governor->usable_bytes, + state->previous_usable_bytes, memory_order_relaxed); + state->deny_output_growth_restore_pending = false; +} +#endif + /* * lftj_binary_ctx_t: callback context for col_op_lftj. * @@ -1059,6 +1195,7 @@ typedef struct { int64_t *tmp; /* scratch row buffer */ col_rel_t *out; /* destination relation */ int rc; /* first error code encountered; 0 = ok */ + wl_columnar_memory_governor_ref_t *governor; } lftj_binary_ctx_t; /* @@ -1093,7 +1230,20 @@ lftj_binary_cb(const int64_t *row, uint32_t lftj_ncols, void *user) ctx->tmp[bo + c] = val; } } +#ifdef WL_TEST_LFTJ_ADMISSION_HOOKS + if (lftj_test_before_output_growth(ctx->out, ctx->governor)) { + lftj_output_test_hook_state.failure_pending_flag + = ctx->out->memory_budget_denial_pending; + ctx->rc = ENOMEM; + return; + } +#endif int rc = col_rel_append_row(ctx->out, ctx->tmp); +#ifdef WL_TEST_LFTJ_ADMISSION_HOOKS + lftj_test_after_output_append(ctx->governor); +#endif + if (rc == ENOMEM && ctx->out->memory_budget_denial_pending) + rc = ENOSPC; if (rc) ctx->rc = rc; } @@ -1261,16 +1411,29 @@ col_op_lftj(const wl_plan_op_t *op, eval_stack_t *stack, wl_col_session_t *sess) } { - col_rel_t *out = col_rel_pool_new_auto(sess->delta_pool, - sess->eval_arena, "$lftj", - total_binary_ncols); + col_rel_t *out = NULL; + if (sess->memory_governor) { +#ifdef WL_TEST_LFTJ_ADMISSION_HOOKS + lftj_test_before_output_construction(sess->memory_governor); +#endif + rc = wl_columnar_relation_new_auto_governed("$lftj", + total_binary_ncols, 0, false, sess->memory_governor, &out); +#ifdef WL_TEST_LFTJ_ADMISSION_HOOKS + lftj_test_after_output_construction(sess->memory_governor); +#endif + } else { + out = col_rel_pool_new_auto(sess->delta_pool, sess->eval_arena, + "$lftj", total_binary_ncols); + rc = out ? 0 : ENOMEM; + } int64_t *tmp = (int64_t *)malloc( (total_binary_ncols ? total_binary_ncols : 1u) * sizeof(int64_t)); if (!out || !tmp) { free(tmp); if (out) col_rel_destroy(out); - rc = ENOMEM; + if (rc == 0) + rc = ENOMEM; goto cleanup_arrays; } @@ -1282,15 +1445,14 @@ col_op_lftj(const wl_plan_op_t *op, eval_stack_t *stack, wl_col_session_t *sess) total_binary_ncols, tmp, out, - 0 }; + 0, + sess->memory_governor }; rc = wl_columnar_lftj_join_typed_governed(inputs, key_type, k, lftj_binary_cb, &ctx, sess ? sess->memory_governor : NULL); if (rc == 0) rc = ctx.rc; - if (rc == ENOSPC && sess) - sess->memory_budget_denied = true; free(tmp); if (rc != 0) { @@ -1360,6 +1522,8 @@ col_op_lftj(const wl_plan_op_t *op, eval_stack_t *stack, wl_col_session_t *sess) free(binary_offsets); WL_LFTJ_RELEASE_SCRATCH(); #undef WL_LFTJ_RELEASE_SCRATCH + if (rc == ENOSPC && sess) + sess->memory_budget_denied = true; return rc; }