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
101 changes: 101 additions & 0 deletions tests/test_tdd_inline_workers.c
Original file line number Diff line number Diff line change
Expand Up @@ -57,6 +57,33 @@ wl_columnar_session_get_tdd_decision_stats(wl_session_t *, uint32_t *,
static bool owner_partial_submit_armed;
static unsigned owner_partial_submit_count;

/* #1977: when armed, every TDD sub-pass runs through a trampoline that
* checks the worker restores delta_seeded on its successful fall-through
* exit. Only the seeding sub-passes (force_diff, outbound_only,
* eff_iter > 0) are checked; each worker writes its own slot, and the
* coordinator reads the slots only after the workqueue barrier. */
#define SEEDED_PROBE_SLOTS 64u
static bool seeded_probe_armed;
static void (*seeded_probe_fn)(void *);
static uint32_t seeded_probe_checked[SEEDED_PROBE_SLOTS];
static uint32_t seeded_probe_leaked[SEEDED_PROBE_SLOTS];

static void
seeded_probe_trampoline(void *arg)
{
col_eval_tdd_worker_ctx_t *ctx = (col_eval_tdd_worker_ctx_t *)arg;
wl_col_session_t *worker = ctx->worker_sess;
bool before = worker->delta_seeded;
seeded_probe_fn(arg);
if (!ctx->force_diff || !ctx->outbound_only || ctx->eff_iter == 0
|| ctx->rc != 0 || ctx->all_empty_delta
|| worker->worker_id >= SEEDED_PROBE_SLOTS)
return;
seeded_probe_checked[worker->worker_id]++;
if (worker->delta_seeded != before)
seeded_probe_leaked[worker->worker_id]++;
}

/* Enable the production TDD submission path while keeping this test's
* dispatch hook deterministic. */
int
Expand All @@ -66,6 +93,10 @@ wl_columnar_eval_test_submit(wl_work_queue_t *wq, void (*fn)(void *),
if (owner_partial_submit_armed
&& ++owner_partial_submit_count == 2)
return ENOMEM;
if (seeded_probe_armed) {
seeded_probe_fn = fn;
return wl_workqueue_submit(wq, seeded_probe_trampoline, ctx);
}
return wl_workqueue_submit(wq, fn, ctx);
}

Expand Down Expand Up @@ -405,6 +436,74 @@ run_scalar_dispatch_probe(void)
return 1;
}

/* #1977: a worker whose seeding sub-pass completes normally must restore
* delta_seeded on its clone, as its early exits after the save already
* do. The program joins reach on both columns, so no exchange key aligns
* and the stratum takes global-read mode, which seeds. */
static int
run_delta_seeded_restore_probe(void)
{
const char *source = ".decl edge(x: int32, y: int32)\n"
".decl reach(x: int32, y: int32)\n"
"reach(x,y) :- edge(x,y).\n"
"reach(x,z) :- reach(x,y), edge(y,z).\n"
"reach(x,z) :- edge(x,y), reach(y,z).\n";
wirelog_error_t error;
wirelog_program_t *program = wirelog_parse_string(source, &error);
wl_plan_t *plan = NULL;
wl_session_t *session = NULL;
if (!program)
return 1;
wl_fusion_apply(program, NULL);
wl_jpp_apply(program, NULL);
wl_sip_apply(program, NULL);
int rc = wl_plan_from_program(program, &plan);
if (rc != 0 || !plan)
goto fail;
plan_fixture_hold(program);
rc = wl_session_create(wl_backend_columnar(), plan, 8, &session);
if (rc != 0)
goto fail;
int64_t rows[200];
for (uint32_t i = 0; i < 100; i++) {
rows[i * 2] = (int64_t)i;
rows[i * 2 + 1] = (int64_t)i + 1;
}
memset(seeded_probe_checked, 0, sizeof(seeded_probe_checked));
memset(seeded_probe_leaked, 0, sizeof(seeded_probe_leaked));
seeded_probe_armed = true;
rc = wl_session_insert(session, "edge", rows, 100, 2);
uint64_t count = 0;
if (rc == 0)
rc = wl_session_snapshot(session, count_rows, &count);
seeded_probe_armed = false;
uint32_t checked = 0, leaked = 0;
for (uint32_t w = 0; w < SEEDED_PROBE_SLOTS; w++) {
checked += seeded_probe_checked[w];
leaked += seeded_probe_leaked[w];
}
if (rc == 0 && (count != 5050 || checked == 0 || leaked != 0)) {
fprintf(stderr,
"delta_seeded probe rc=%d count=%" PRIu64 " checked=%u "
"leaked=%u (checked=0: no global-read seeding sub-pass ran; "
"is WIRELOG_TDD_GLOBAL_READ=0 set?)\n", rc, count, checked,
leaked);
rc = EIO;
}
wl_session_destroy(session);
wl_plan_free(plan);
return rc == 0 ? 0 : 1;
fail:
seeded_probe_armed = false;
if (session)
wl_session_destroy(session);
if (plan)
wl_plan_free(plan);
else if (program)
wirelog_program_free(program);
return 1;
}

static int
run_lifecycle(uint32_t workers, bool require_dispatch)
{
Expand Down Expand Up @@ -823,6 +922,8 @@ main(void)
rc = run_lifecycle(8, true);
if (rc == 0)
rc = run_scalar_dispatch_probe();
if (rc == 0)
rc = run_delta_seeded_restore_probe();
if (rc == 0)
rc = run_teardown_refusal_gate();
if (rc == 0)
Expand Down
8 changes: 4 additions & 4 deletions wirelog/columnar/eval.c
Original file line number Diff line number Diff line change
Expand Up @@ -2594,10 +2594,10 @@ tdd_worker_subpass_fn(void *arg)
}

ctx->any_new = any_new;
sess->tdd_subpass_active = saved_tdd_subpass;
sess->tdd_outbound_only_active = saved_outbound_only;
sess->diff_operators_active = saved_diff;
ctx->runtime_ns = now_ns() - worker_t0;
/* Leave through the same restore as the early exits after the flags
* were saved, so the fall-through cannot skip one (#1977:
* delta_seeded). */
TDD_WORKER_RETURN();
#undef TDD_WORKER_RETURN
}

Expand Down
Loading