diff --git a/tests/test_join_batch_resume.c b/tests/test_join_batch_resume.c index 82d14339..d6173aaa 100644 --- a/tests/test_join_batch_resume.c +++ b/tests/test_join_batch_resume.c @@ -16,6 +16,7 @@ #include "../wirelog/backend.h" #include "../wirelog/columnar/internal.h" #include "../wirelog/columnar/join_batch.h" +#include "../wirelog/util/log.h" #include "../wirelog/exec_plan_gen.h" #include "../wirelog/session.h" #include "../wirelog/wirelog-parser.h" @@ -26,6 +27,7 @@ #include #include #include +#include #ifdef _WIN32 static int @@ -43,6 +45,11 @@ wl_test_unsetenv_(const char *name) #define setenv wl_test_setenv_ #define unsetenv wl_test_unsetenv_ +#include +#define wl_test_getpid _getpid +#else +#include +#define wl_test_getpid getpid #endif static int tests_run; @@ -2052,6 +2059,294 @@ test_projected_output(void) fixture_fini(&f); } +/* #1909: the projected test above runs at seven rows per batch, below the + * 64-row floor, so projection never meets a scratch reallocation. Here the + * 90 projected matches cross the floor inside one batch. */ +static void +test_projected_output_across_grow(void) +{ + fixture_t f; + const int64_t keys[] = { 0, 1, 1 }; + uint32_t project[] = { 1, 3 }; + wl_columnar_continuation_t *cont = NULL; + col_rel_t *oracle = NULL; + col_rel_t *out2 = NULL; + uint32_t cap0, capN; + int rc; + + TEST("projected output across a scratch grow matches the oracle"); + if (!fixture_init(&f, 1u << 24, keys, 3, 2, 30)) { + FAIL("fixture"); + fixture_fini(&f); + return; + } + f.op.project_indices = project; + f.op.project_count = 2; + out2 = col_rel_new_auto("$join", 2); + if (!out2 || col_join_set_output_types(out2, f.left, f.right, &f.op) + != 0) { + FAIL("projected output"); + goto out; + } + rc = col_join_batch_producer_create(f.sess, &f.op, f.left, false, KEY0, + KEY0, 1, 512u * 16u, &cont); + if (rc != 0 || col_join_batch_rows_per_batch(cont) <= 90u) { + FAIL("producer create or rows_per_batch"); + goto out; + } + cap0 = col_join_batch_scratch_capacity(cont); + rc = col_join_batch_run_to_relation(cont, f.sess, out2); + capN = col_join_batch_scratch_capacity(cont); + oracle = run_oracle(f.sess, f.left, &f.op); + if (rc != 0 || !oracle || !same_rows(oracle, out2) || out2->nrows != 90u) + FAIL("projected result differs from the oracle across a grow"); + else if (cap0 != COL_REL_INIT_CAP || capN <= cap0) + FAIL("the projected scratch did not start at the floor and grow"); + else + PASS(); +out: + if (oracle) + col_rel_destroy(oracle); + if (out2) + col_rel_destroy(out2); + if (cont) + wl_columnar_continuation_destroy(cont); + fixture_fini(&f); +} + +/* #1909: a timestamped left gives the scratch a timestamp array, which the + * resize path copies by nrows just as it copies the columns. The producer + * writes zero timestamps, so the columns are what can be lost; the sink's + * timestamps must stay zero. */ +static void +test_timestamped_scratch_across_grow(void) +{ + fixture_t f; + const int64_t keys[] = { 0, 1, 2 }; + wl_columnar_continuation_t *cont = NULL; + col_rel_t *oracle = NULL; + uint32_t cap0, capN; + bool zero = true; + int rc; + + TEST("timestamped scratch across a grow matches the oracle"); + if (!fixture_init(&f, 1u << 24, keys, 3, 5, 47) + || col_rel_enable_timestamps(f.left) != 0 + || col_rel_enable_timestamps(f.out) != 0) { + FAIL("timestamped fixture"); + fixture_fini(&f); + return; + } + rc = col_join_batch_producer_create(f.sess, &f.op, f.left, false, KEY0, + KEY0, 1, 1u << 16, &cont); + if (rc != 0 || col_join_batch_rows_per_batch(cont) <= 3u * 47u) { + FAIL("producer create or rows_per_batch"); + goto out; + } + cap0 = col_join_batch_scratch_capacity(cont); + rc = col_join_batch_run_to_relation(cont, f.sess, f.out); + capN = col_join_batch_scratch_capacity(cont); + oracle = run_oracle(f.sess, f.left, &f.op); + for (uint32_t i = 0; f.out->timestamps && i < f.out->nrows; i++) { + const col_delta_timestamp_t *ts = &f.out->timestamps[i]; + if (ts->multiplicity != 0 || ts->iteration != 0 || ts->stratum != 0 + || ts->worker != 0) + zero = false; + } + if (rc != 0 || !oracle || !same_rows(oracle, f.out) + || f.out->nrows != 3u * 47u) + FAIL("timestamped result differs from the oracle across a grow"); + else if (!f.out->timestamps || !zero) + FAIL("the sink did not publish zero timestamps"); + else if (cap0 != COL_REL_INIT_CAP || capN <= cap0) + FAIL("the timestamped scratch did not start at the floor and grow"); + else + PASS(); +out: + if (oracle) + col_rel_destroy(oracle); + if (cont) + wl_columnar_continuation_destroy(cont); + fixture_fini(&f); +} + +/* #1909: the grow-refusal WARN is operator-facing, and no producer test + * could see it: the logger was not initialised until this binary's last + * test, and removing the warning, its once-per-producer gate, or rotating + * its reason strings survived the suite. Capture the JOIN section to a + * file and check the reason for a refused allocation and for an admission + * denial, and that two refused grows in one producer warn once. The + * caller's WL_LOG and WL_LOG_FILE are saved and restored. */ +static char warn_log_path[512]; +static char *warn_saved_log; +static char *warn_saved_file; + +static char * +warn_env_dup(const char *name) +{ + const char *value = getenv(name); + char *copy; + size_t len; + + if (!value) + return NULL; + len = strlen(value) + 1u; + copy = (char *)malloc(len); + if (copy) + memcpy(copy, value, len); + return copy; +} + +static void +warn_env_restore(const char *name, char **saved) +{ + if (*saved) + (void)setenv(name, *saved, 1); + else + (void)unsetenv(name); + free(*saved); + *saved = NULL; +} + +static bool +warn_log_begin(void) +{ + const char *dir = getenv("TMPDIR"); +#ifdef _WIN32 + if (!dir || !*dir) + dir = getenv("TEMP"); +#endif + if (!dir || !*dir) + dir = "."; + /* The process id keeps concurrent runs sharing a TMPDIR apart. */ + snprintf(warn_log_path, sizeof(warn_log_path), + "%s/wl_join_batch_warn_%ld_%lu.log", dir, (long)wl_test_getpid(), + (unsigned long)time(NULL)); + (void)remove(warn_log_path); + warn_saved_log = warn_env_dup("WL_LOG"); + warn_saved_file = warn_env_dup("WL_LOG_FILE"); + if (setenv("WL_LOG_FILE", warn_log_path, 1) != 0 + || setenv("WL_LOG", "JOIN:2", 1) != 0) + return false; + return wl_log_init() == 0; +} + +/* Stop capturing, read the file into buf and restore the caller's logging + * configuration. */ +static void +warn_log_end(char *buf, size_t size) +{ + FILE *fp; + size_t n = 0; + + wl_log_shutdown(); + fp = fopen(warn_log_path, "r"); + if (fp) { + n = fread(buf, 1, size - 1u, fp); + fclose(fp); + } + buf[n] = '\0'; + (void)remove(warn_log_path); + warn_env_restore("WL_LOG", &warn_saved_log); + warn_env_restore("WL_LOG_FILE", &warn_saved_file); + (void)wl_log_init(); +} + +static unsigned +count_substr(const char *haystack, const char *needle) +{ + unsigned count = 0; + for (const char *p = strstr(haystack, needle); p; + p = strstr(p + 1, needle)) + count++; + return count; +} + +static void +test_grow_refusal_warn_reason(void) +{ + for (int denied = 0; denied < 2; denied++) { + fixture_t f; + const int64_t keys[] = { 0, 1, 2 }; + wl_columnar_continuation_t *cont = NULL; + wl_columnar_continuation_sink_t sink; + col_join_batch_relation_sink_t sctx; + wl_columnar_continuation_status_t st = WL_COLUMNAR_CONTINUATION_OK; + wl_columnar_memory_governor_t *governor; + col_rel_t *oracle = NULL; + bool refused_ok = true; + char log[4096]; + + TEST(denied ? "a denied scratch grow warns: admission denied" + : "failed scratch grows warn once: allocation failed"); + /* The ceiling is a project-wide define, so it matches the library; + * 2 is WL_LOG_WARN, which #if cannot name. */ +#if WL_LOG_COMPILE_MAX_LEVEL < 2 + printf("SKIP (WARN compiled out) "); + PASS(); + continue; +#endif + if (!fixture_init(&f, 1u << 24, keys, 3, 5, 47) + || col_join_batch_producer_create(f.sess, &f.op, f.left, false, + KEY0, KEY0, 1, 512u * 32u, &cont) != 0 + || col_join_batch_relation_sink_init(&sctx, &sink, f.sess, + f.out) != 0 || !warn_log_begin()) { + FAIL("fixture, producer, sink or log capture"); + if (cont) + wl_columnar_continuation_destroy(cont); + fixture_fini(&f); + continue; + } + governor = wl_columnar_memory_governor_ref_get(f.sess->memory_governor); + if (denied) { + /* Room for a 64-row payload (2 KiB at four columns) but not + * for growing the scratch to 128 rows, which holds the old and + * the new buffers at once: the first batch is emitted short. */ + atomic_store_explicit(&governor->usable_bytes, + reserved_of(f.sess) + 64u * 32u * 2u, memory_order_release); + st = wl_columnar_continuation_publish(cont, &sink); + refused_ok = st == WL_COLUMNAR_CONTINUATION_OK + && f.out->nrows == COL_REL_INIT_CAP; + atomic_store_explicit(&governor->usable_bytes, 1u << 24, + memory_order_release); + } else { + /* Refuse the grow of two consecutive batches: both are emitted + * short, and the producer warns only for the first. */ + for (int b = 0; b < 2 && refused_ok; b++) { + wl_columnar_relation_test_fail_next_prepare_resize(); + st = wl_columnar_continuation_publish(cont, &sink); + refused_ok = st == WL_COLUMNAR_CONTINUATION_OK + && f.out->nrows == COL_REL_INIT_CAP * (uint32_t)(b + 1) + && col_join_batch_scratch_capacity(cont) + == COL_REL_INIT_CAP; + } + wl_columnar_relation_test_clear_prepare_resize(); + } + while (refused_ok && st == WL_COLUMNAR_CONTINUATION_OK) + st = wl_columnar_continuation_publish(cont, &sink); + warn_log_end(log, sizeof(log)); + oracle = run_oracle(f.sess, f.left, &f.op); + if (!refused_ok) + FAIL("a refused grow did not emit a short batch"); + else if (st != WL_COLUMNAR_CONTINUATION_DONE || !oracle + || !same_rows(oracle, f.out) || f.out->nrows != 3u * 47u) + FAIL("a refused grow lost or duplicated a row"); + else if (count_substr(log, "join batch scratch stuck") != 1u) + FAIL("the refusals did not warn exactly once"); + else if (count_substr(log, denied ? "admission denied" + : "allocation failed") != 1u + || count_substr(log, denied ? "allocation failed" + : "admission denied") != 0u) + FAIL("the warning names the wrong reason"); + else + PASS(); + if (oracle) + col_rel_destroy(oracle); + wl_columnar_continuation_destroy(cont); + fixture_fini(&f); + } +} + static int64_t double_bits(double d) { @@ -3157,6 +3452,9 @@ main(void) test_batch_payload_uses_right_governor_fallback(); test_batch_payload_denial_rolls_back_and_retries(); test_projected_output(); + test_projected_output_across_grow(); + test_timestamped_scratch_across_grow(); + test_grow_refusal_warn_reason(); test_float_key(); test_unsupported_budget_and_pooled_output(); test_row_cap_trips_after_a_committed_batch(); diff --git a/wirelog/columnar/internal.h b/wirelog/columnar/internal.h index 06dce90b..1222d931 100644 --- a/wirelog/columnar/internal.h +++ b/wirelog/columnar/internal.h @@ -2643,8 +2643,10 @@ col_join_output_width(const col_rel_t *left, const col_rel_t *right, int col_join_set_output_types(col_rel_t *out, const col_rel_t *left, const col_rel_t *right, const wl_plan_op_t *op); -/* Writes one joined row at @out_row with col_rel_set_raw; the caller - * publishes nrows and the view generation once per batch. */ +/* Writes one joined row at @out_row with col_rel_set_raw and publishes + * neither nrows nor a view generation; its callers arrange that. The + * batch producers publish nrows after every written row (#1909); the keyed + * fill workers leave a single publication to their coordinator. */ int col_join_write_pair_at(col_rel_t *out, uint64_t out_row, const col_rel_t *left, uint32_t lr, const col_rel_t *right, uint32_t rr, diff --git a/wirelog/columnar/join_batch.c b/wirelog/columnar/join_batch.c index 33108808..240ad506 100644 --- a/wirelog/columnar/join_batch.c +++ b/wirelog/columnar/join_batch.c @@ -173,13 +173,13 @@ producer_validate(void *context, * relation pre-allocated and clamping to rows_per_batch. Returns false when * the row cannot be written, which the caller turns into a short batch. * - * p->batch->nrows MUST be set to n before the reserve. producer_produce - * zeroes nrows while filling and only restores it at emit, and - * col_rel_prepare_resize copies exactly nrows rows -- so growing with nrows - * still 0 copies nothing, frees the old columns, and leaves every row written - * so far pointing at uninitialised malloc memory. That is a silently wrong - * join, and neither ASan nor UBSan reports it; only an oracle comparison - * across a grow does. #1481. + * col_rel_prepare_resize copies exactly nrows rows (and nrows timestamps), + * so p->batch->nrows must equal n here. A stale nrows copies too little, + * frees the old columns, and leaves rows written so far pointing at + * uninitialised malloc memory: a silently wrong join that neither ASan nor + * UBSan reports, only an oracle comparison across a grow (#1481). + * producer_produce publishes nrows after every written row (#1909), so no + * reallocation site has a sync point to forget. */ static bool grow_scratch(col_join_batch_producer_t *p, uint32_t n) @@ -196,7 +196,6 @@ grow_scratch(col_join_batch_producer_t *p, uint32_t n) * without growing whenever want <= capacity, so reporting its status * would be reporting the call rather than the outcome. */ want = cap < p->rows_per_batch / 2u ? cap * 2u : p->rows_per_batch; - p->batch->nrows = n; rc = col_rel_reserve_capacity_admitted(p->batch, want, &denied); if (rc == 0 && p->batch->capacity > n) return true; @@ -263,6 +262,9 @@ producer_produce(void *context, const wl_columnar_continuation_cursor_t *cursor, memset(&p->batch->timestamps[n], 0, sizeof(col_delta_timestamp_t)); n++; + /* Keep the live prefix structurally current, so every + * relocation preserves the row just written. */ + p->batch->nrows = n; if (n == p->rows_per_batch) { /* Batch full: park on the first unexamined candidate. * UINT32_MAX is stored only as "probe the next left