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
298 changes: 298 additions & 0 deletions tests/test_join_batch_resume.c
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand All @@ -26,6 +27,7 @@
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <time.h>

#ifdef _WIN32
static int
Expand All @@ -43,6 +45,11 @@ wl_test_unsetenv_(const char *name)

#define setenv wl_test_setenv_
#define unsetenv wl_test_unsetenv_
#include <process.h>
#define wl_test_getpid _getpid
#else
#include <unistd.h>
#define wl_test_getpid getpid
#endif

static int tests_run;
Expand Down Expand Up @@ -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)
{
Expand Down Expand Up @@ -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();
Expand Down
6 changes: 4 additions & 2 deletions wirelog/columnar/internal.h
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
18 changes: 10 additions & 8 deletions wirelog/columnar/join_batch.c
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand All @@ -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;
Expand Down Expand Up @@ -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
Expand Down
Loading