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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
26 changes: 14 additions & 12 deletions docs/MEMORY.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down Expand Up @@ -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

Expand Down
158 changes: 139 additions & 19 deletions tests/test_join_limit.c
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -27,7 +27,9 @@
#include "../wirelog/wirelog.h"
#include "plan_fixture.h"

#include <errno.h>
#include <stdio.h>
#include <stdint.h>
#include <stdlib.h>
#include <string.h>

Expand Down Expand Up @@ -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");
Expand All @@ -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);
}
Expand Down Expand Up @@ -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"
Expand Down Expand Up @@ -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);
Expand All @@ -315,19 +371,83 @@ 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 */
/* ======================================================================== */

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);
Expand Down
15 changes: 3 additions & 12 deletions wirelog/columnar/internal.h
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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;
Expand Down
62 changes: 21 additions & 41 deletions wirelog/columnar/session.c
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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. */
{
Expand Down
Loading