diff --git a/docs/MEMORY.md b/docs/MEMORY.md index 2719c72f..1379b8b8 100644 --- a/docs/MEMORY.md +++ b/docs/MEMORY.md @@ -364,9 +364,9 @@ memory use, how to read the numbers, what they cost, and the baselines that bounded-memory work (Issue #1367 and its sub-issues) starts from. It covers Issue #1380. -Nothing here enforces a limit. The ledger is an accounting instrument; the -budget it carries (`WIRELOG_MEMORY_BUDGET`) only drives the join operator's -backpressure poll described in ยง6. +The ledger is an accounting instrument, while the memory governor uses the +resolved budget to admit covered allocation classes before they grow. Coverage +is incremental; the sections below describe which paths participate. ## 1. Units and conventions @@ -543,15 +543,17 @@ and host Job Object limits; when no finite source exists, the resolver is advisory/unbounded. Explicit `0`, malformed values, overflow, and values below 256 MiB are invalid. This resolver does not use physical RAM as an enforcing fallback. The legacy ledger worker-share hint remains separate from governor -admission; fixed eval-arena, delta-pool, and compound-arena backing storage -are now admitted. Joins and other allocation classes remain follow-up work. - -The only consumer is the join operator: when RELATION reaches 80% of its -share (`wl_mem_ledger_should_backpressure(RELATION, 80)`), a worker session -stops generating rows for the current join and reports the condition -upstream. The join row cap (`WIRELOG_JOIN_OUTPUT_LIMIT`, Issue #221) is a -separate mechanism. Replacing both with admission control is #1367 -(foundation in #1368). +admission. Fixed eval-arena, delta-pool, compound-arena backing storage, and +the covered join allocation paths use admission; LFTJ, auxiliary metadata, +and eval-entry segment allocations remain follow-up work. + +The ledger's RELATION worker-share hint can still trigger join backpressure +at 80% (`wl_mem_ledger_should_backpressure(RELATION, 80)`). The governor also +admits covered allocations and reports refusals upstream. The legacy +`WIRELOG_JOIN_OUTPUT_LIMIT` row cap is disabled when unset, empty, or zero; a +positive decimal keeps the optional per-join cap, clamped to `UINT32_MAX`. +Malformed values fail session creation. This row cap is not the memory budget +and does not replace allocation admission. ## 7. Baselines diff --git a/tests/test_join_limit.c b/tests/test_join_limit.c index 08fcaf5d..27bd8c05 100644 --- a/tests/test_join_limit.c +++ b/tests/test_join_limit.c @@ -5,13 +5,13 @@ * Licensed under LGPL-3.0 * For commercial licenses, contact: inquiry@cleverplant.com * - * Validates that join_output_limit is correctly computed at session init: - * - Default: auto-detected from physical RAM / num_workers - * - Env override: WIRELOG_JOIN_OUTPUT_LIMIT sets exact value - * - Disable: WIRELOG_JOIN_OUTPUT_LIMIT=0 disables the limit - * - Scaling: 8-worker limit is ~1/8 of 1-worker limit + * Validates optional legacy join row-cap configuration: + * - Unset, empty and zero disable the cap + * - Env override: WIRELOG_JOIN_OUTPUT_LIMIT sets a positive cap + * - Malformed values fail session creation + * - Disabled caps still preserve join output high-water accounting * - * Issue #221: Dynamic join output limit based on available memory + * Issue #1369: remove the implicit physical-memory-derived row cap */ #define _GNU_SOURCE @@ -27,7 +27,9 @@ #include "../wirelog/wirelog.h" #include "plan_fixture.h" +#include #include +#include #include #include @@ -97,13 +99,13 @@ build_plan(const char *src) } /* ======================================================================== */ -/* Test 1: Default limit is auto-detected (> 0) */ +/* Test 1: Unset limit leaves the optional row cap disabled */ /* ======================================================================== */ static int test_default_limit(void) { - TEST("Default join_output_limit > 0 (auto-detected from RAM)"); + TEST("Unset WIRELOG_JOIN_OUTPUT_LIMIT disables the row cap"); /* Remove env override if present */ unsetenv("WIRELOG_JOIN_OUTPUT_LIMIT"); @@ -125,12 +127,12 @@ test_default_limit(void) } wl_col_session_t *sess = (wl_col_session_t *)session; - bool ok = sess->join_output_limit > 0; + bool ok = sess->join_output_limit == 0; if (!ok) { char msg[128]; snprintf(msg, sizeof(msg), - "join_output_limit == 0 (expected > 0), got %llu", + "expected join_output_limit=0, got %llu", (unsigned long long)sess->join_output_limit); FAIL(msg); } @@ -239,16 +241,72 @@ test_env_disable(void) return ok ? 0 : 1; } +static int +test_empty_env_is_disabled(void) +{ + TEST("empty WIRELOG_JOIN_OUTPUT_LIMIT disables the row cap"); + setenv("WIRELOG_JOIN_OUTPUT_LIMIT", "", 1); + wl_plan_t *plan = build_plan(".decl a(x: int32)\n" + ".decl r(x: int32)\n" + "r(x) :- a(x).\n"); + if (!plan) { + unsetenv("WIRELOG_JOIN_OUTPUT_LIMIT"); + FAIL("could not generate plan"); + return 1; + } + wl_session_t *session = NULL; + int rc = wl_session_create(wl_backend_columnar(), plan, 1, &session); + bool ok = rc == 0 && session != NULL + && ((wl_col_session_t *)session)->join_output_limit == 0; + if (session) + wl_session_destroy(session); + wl_plan_free(plan); + unsetenv("WIRELOG_JOIN_OUTPUT_LIMIT"); + if (ok) + PASS(); + else + FAIL("empty value did not disable the row cap"); + return ok ? 0 : 1; +} + +static int +test_env_clamps_to_row_width(void) +{ + TEST("positive row cap is clamped to UINT32_MAX"); + setenv("WIRELOG_JOIN_OUTPUT_LIMIT", "18446744073709551615", 1); + wl_plan_t *plan = build_plan(".decl a(x: int32)\n" + ".decl r(x: int32)\n" + "r(x) :- a(x).\n"); + if (!plan) { + unsetenv("WIRELOG_JOIN_OUTPUT_LIMIT"); + FAIL("could not generate plan"); + return 1; + } + wl_session_t *session = NULL; + int rc = wl_session_create(wl_backend_columnar(), plan, 1, &session); + bool ok = rc == 0 && session != NULL + && ((wl_col_session_t *)session)->join_output_limit == UINT32_MAX; + if (session) + wl_session_destroy(session); + wl_plan_free(plan); + unsetenv("WIRELOG_JOIN_OUTPUT_LIMIT"); + if (ok) + PASS(); + else + FAIL("positive cap was not clamped to UINT32_MAX"); + return ok ? 0 : 1; +} + /* ======================================================================== */ -/* Test 4: Limit is constant regardless of num_workers (Issue #404) */ +/* Test 4: Disabled cap is independent of configured worker count */ /* ======================================================================== */ static int test_limit_constant_across_workers(void) { - TEST("join_output_limit is constant across W=1, W=4, W=8 (Issue #404)"); + TEST("disabled row cap stays zero at W=1, W=4, W=8"); - /* Remove env override to exercise auto-detection path */ + /* Remove env override to exercise unset configuration. */ unsetenv("WIRELOG_JOIN_OUTPUT_LIMIT"); wl_plan_t *plan = build_plan(".decl a(x: int32)\n" @@ -290,15 +348,13 @@ test_limit_constant_across_workers(void) uint64_t limit4 = ((wl_col_session_t *)sess4)->join_output_limit; uint64_t limit8 = ((wl_col_session_t *)sess8)->join_output_limit; - /* Global per-join cap: limit must be identical for W=1, W=4, W=8. - * Regression check for Issue #404: commit 6929689 divided by num_workers, - * causing silent data loss in multi-worker mode. */ - bool ok = (limit1 > 0 && limit1 == limit4 && limit1 == limit8); + /* These simple plans exercise coordinator session configuration only. */ + bool ok = (limit1 == 0 && limit4 == 0 && limit8 == 0); if (!ok) { char msg[256]; snprintf(msg, sizeof(msg), - "limit1=%llu limit4=%llu limit8=%llu: expected all equal", + "limit1=%llu limit4=%llu limit8=%llu: expected all zero", (unsigned long long)limit1, (unsigned long long)limit4, (unsigned long long)limit8); @@ -315,6 +371,66 @@ test_limit_constant_across_workers(void) return ok ? 0 : 1; } +static int +test_invalid_env_values(void) +{ + static const char *const invalid[] = { + "-1", "+1", " 1", "1 ", "1x", "18446744073709551616" + }; + TEST("malformed row-cap values fail before session allocation"); + wl_plan_t *plan = build_plan(".decl a(x: int32)\n" + ".decl r(x: int32)\n" + "r(x) :- a(x).\n"); + if (!plan) { + unsetenv("WIRELOG_JOIN_OUTPUT_LIMIT"); + FAIL("could not generate plan"); + return 1; + } + for (size_t i = 0; i < sizeof(invalid) / sizeof(invalid[0]); i++) { + setenv("WIRELOG_JOIN_OUTPUT_LIMIT", invalid[i], 1); + wl_session_t *session = (wl_session_t *)(uintptr_t)1; + int rc = wl_session_create(wl_backend_columnar(), plan, 1, &session); + if (rc != EINVAL || session != NULL) { + if (session && session != (wl_session_t *)(uintptr_t)1) + wl_session_destroy(session); + wl_plan_free(plan); + unsetenv("WIRELOG_JOIN_OUTPUT_LIMIT"); + FAIL("invalid value did not return EINVAL with NULL output"); + return 1; + } + } + setenv("WIRELOG_JOIN_OUTPUT_LIMIT", "0", 1); + wl_session_t *session = NULL; + int rc = wl_session_create(wl_backend_columnar(), plan, 1, &session); + bool ok = rc == 0 && session != NULL + && ((wl_col_session_t *)session)->join_output_limit == 0; + if (session) + wl_session_destroy(session); + wl_plan_free(plan); + unsetenv("WIRELOG_JOIN_OUTPUT_LIMIT"); + if (ok) + PASS(); + else + FAIL("valid creation failed after malformed configuration"); + return ok ? 0 : 1; +} + +static int +test_disabled_limit_records_high_water(void) +{ + TEST("disabled row cap still records the output high-water mark"); + wl_col_session_t session = {0}; + col_rel_t output = {0}; + output.nrows = 17; + bool stopped = col_join_output_limit_reached(&session, &output); + bool ok = !stopped && session.join_output_peak == 17; + if (ok) + PASS(); + else + FAIL("zero cap stopped output or lost its high-water mark"); + return ok ? 0 : 1; +} + /* ======================================================================== */ /* main */ /* ======================================================================== */ @@ -322,12 +438,16 @@ test_limit_constant_across_workers(void) int main(void) { - printf("=== test_join_limit (Issue #221) ===\n"); + printf("=== test_join_limit (Issue #1369) ===\n"); test_default_limit(); test_env_override(); test_env_disable(); + test_empty_env_is_disabled(); + test_env_clamps_to_row_width(); test_limit_constant_across_workers(); + test_invalid_env_values(); + test_disabled_limit_records_high_water(); printf("\nPassed: %d/%d\n", tests_passed, tests_run); printf("Failed: %d/%d\n", tests_failed, tests_run); diff --git a/wirelog/columnar/internal.h b/wirelog/columnar/internal.h index 1222d931..49898bd9 100644 --- a/wirelog/columnar/internal.h +++ b/wirelog/columnar/internal.h @@ -195,14 +195,6 @@ now_ns(void) #define COL_ARR_CACHE_MAX 128u #define COL_ARR_CACHE_LIMIT_BYTES (256ULL * 1024ULL * 1024ULL) /* 256 MB default */ -/* Default maximum output rows per single join operation (issue #218, #221). - * Prevents unbounded memory growth from cardinality explosion in - * cross-product-heavy joins (e.g., DOOP VarPointsTo). When exceeded, - * the join returns EOVERFLOW. Set to 0 to disable the limit. - * At runtime, the session uses a dynamically computed limit based on - * available physical memory (see wl_col_session_t.join_output_limit). */ -#define COL_JOIN_OUTPUT_LIMIT_DEFAULT (50u * 1000u * 1000u) /* 50M rows */ - typedef enum wl_columnar_internal_tdd_fallback_reason { WL_COLUMNAR_INTERNAL_TDD_FALLBACK_NONE = 0, WL_COLUMNAR_INTERNAL_TDD_FALLBACK_NON_RECURSIVE, @@ -2113,10 +2105,9 @@ typedef struct wl_col_session_t { uint32_t callback_active_workers; bool callback_parallel_execution; const void *callback_session_key; - /* Dynamic join output limit (Issue #221). - * Maximum output rows per single join operation. Computed at session init - * based on available physical memory, num_workers, and estimated row width. - * 0 = disabled (no limit). Overridable via WIRELOG_JOIN_OUTPUT_LIMIT env var. */ + /* Optional legacy maximum output rows per single join, configured through + * WIRELOG_JOIN_OUTPUT_LIMIT. Zero disables the row cap; memory admission + * is enforced separately. */ uint64_t join_output_limit; /* Bounded keyed-join sub-batches (Issue #1446). join_batch_bytes is * the per-producer batch payload budget from WIRELOG_JOIN_BATCH_BYTES; diff --git a/wirelog/columnar/session.c b/wirelog/columnar/session.c index a0e517a6..206cf147 100644 --- a/wirelog/columnar/session.c +++ b/wirelog/columnar/session.c @@ -2446,11 +2446,29 @@ col_session_create_internal(const wl_plan_t *plan, uint32_t num_workers, { wl_columnar_memory_governor_ref_t *memory_governor; const wl_columnar_memory_governor_t *governor; + uint64_t join_output_limit = 0; bool intern_attached_here = false; int relation_create_rc = ENOMEM; if (!plan || !out) return EINVAL; + *out = NULL; + { + const char *env = getenv("WIRELOG_JOIN_OUTPUT_LIMIT"); + if (env && env[0] != '\0') { + uint64_t value = 0; + for (const unsigned char *p = (const unsigned char *)env; + *p != '\0'; p++) { + if (*p < '0' || *p > '9' + || value > (UINT64_MAX - (*p - '0')) / 10) + return EINVAL; + value = value * 10 + (*p - '0'); + } + join_output_limit = value; + } + if (join_output_limit > UINT32_MAX) + join_output_limit = UINT32_MAX; + } if (options && options->memory_governor) { /* #1473: an injected governor replaces environment and host * resolution. The session holds its own reference; the caller's @@ -2649,47 +2667,9 @@ col_session_create_internal(const wl_plan_t *plan, uint32_t num_workers, sess->callback_parallel_execution = false; sess->callback_session_key = sess; - /* Dynamic join output limit (Issue #221) */ - { - const char *join_limit_env = getenv("WIRELOG_JOIN_OUTPUT_LIMIT"); - bool env_valid = false; - if (join_limit_env && join_limit_env[0] != '\0') { - char *endp = NULL; - errno = 0; - uint64_t val = strtoull(join_limit_env, &endp, 10); - if (endp != join_limit_env && *endp == '\0' && errno != ERANGE) { - sess->join_output_limit = val; - env_valid = true; - } - } - if (!env_valid) { - uint64_t phys = col_detect_physical_memory(); - if (phys > 0) { - /* Global per-join cap: 25% of RAM / (8 bytes * 3 avg cols). - * Each K-fusion worker processes a 1/K data partition, so - * per-partition join output does NOT scale with K. Dividing - * by num_workers here was a regression (commit 6929689) that - * caused silent data loss in multi-worker mode (Issue #404). - * Dynamic mem_ledger backpressure handles runtime coordination. - * - * Avg-col assumption was 5 (Issue #221, /40), but DOOP-class - * points-to analyses dominate the recursive recursive-join - * cost on 2-3 column intermediates (SubtypeOf, MethodLookup, - * VarPointsTo, CallGraphEdge); /40 was conservatively low and - * caused DOOP to hit the cap at ~105M intermediate rows on - * 16 GB hosts (Issue #791). /24 is a 67 % headroom bump that - * preserves the 25 % RAM safety margin and matches the - * narrow-join workloads that dominate the v0.43 portfolio. */ - sess->join_output_limit = (phys / 4) / 24ULL; - } else { - sess->join_output_limit - = (uint64_t)COL_JOIN_OUTPUT_LIMIT_DEFAULT; - } - } - /* Clamp to UINT32_MAX since nrows is uint32_t */ - if (sess->join_output_limit > UINT32_MAX) - sess->join_output_limit = UINT32_MAX; - } + /* Zero leaves the legacy row cap disabled; memory admission governs + * allocations independently. A configured positive cap remains optional. */ + sess->join_output_limit = join_output_limit; /* Bounded keyed-join sub-batches (Issue #1446): strict decimals like * the siblings above; anything else leaves the mode off. */ {