From ed0ae202ab36b6c47d5106f086f468287030b193 Mon Sep 17 00:00:00 2001 From: Justin Kim Date: Sat, 3 Oct 2026 23:07:51 +0900 Subject: [PATCH] fix(columnar): restore delta_seeded on the TDD worker sub-pass success path tdd_worker_subpass_fn saves four session flags and restores them through TDD_WORKER_RETURN() at its early exits after the save, but its fall-through exit restored only tdd_subpass_active, tdd_outbound_only_active and diff_operators_active. A worker that took the seeding branch (force_diff, outbound_only and eff_iter > 0) and completed normally left delta_seeded set on its clone, which survives into the next sub-pass. The setter's condition does not change between sub-passes today, so the stale value was immediately set again, but letting force_diff or outbound_only vary per sub-pass would make the stale flag select a delta where the full relation was meant, which can under-derive. End the function with TDD_WORKER_RETURN() so the fall-through takes the same restore as those early exits. The test wraps every TDD sub-pass submitted through the test submission hook and fails if a seeding sub-pass that returned successfully changed delta_seeded on its worker, or if no seeding sub-pass ran. Its program, a transitive closure that joins reach on both columns, puts the stratum in global-read mode, the mode that sets both force_diff and outbound_only. Before the fix it observed 800 seeding sub-passes and 8 leaks, one per worker. Fixes #1977 --- tests/test_tdd_inline_workers.c | 101 ++++++++++++++++++++++++++++++++ wirelog/columnar/eval.c | 8 +-- 2 files changed, 105 insertions(+), 4 deletions(-) diff --git a/tests/test_tdd_inline_workers.c b/tests/test_tdd_inline_workers.c index e5a2473b9..5d2bdd415 100644 --- a/tests/test_tdd_inline_workers.c +++ b/tests/test_tdd_inline_workers.c @@ -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 @@ -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); } @@ -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) { @@ -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) diff --git a/wirelog/columnar/eval.c b/wirelog/columnar/eval.c index 3c96ab607..5b4e0b5e7 100644 --- a/wirelog/columnar/eval.c +++ b/wirelog/columnar/eval.c @@ -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 }