From 13f655645080213b49f51c2dc689c2d15596d2ee Mon Sep 17 00:00:00 2001 From: Justin Kim Date: Mon, 5 Oct 2026 22:29:06 +0900 Subject: [PATCH] fix(session): re-derive from scratch on steps without a delta callback A session driven by wirelog_easy_step() with no delta callback installed re-derived every rule on every step, on top of the rows the previous evaluation had already derived. Every step appended another copy of each derived row. A removal never retracted the rows it had derived. A step with nothing pending re-ran every rule, so a value-position @call ran once per derived row on each idle step. With three `ready` rows, the derived relation held 3, 6, 9 and 12 rows after one to four steps. A snapshot taken afterwards emitted that state as the stable model, which is the "step then snapshot duplicates rows" pattern wirelog-easy.h warned about. Inserts and removals made without a callback already force a full evaluation, so a plain step is meant to re-derive everything. It now does so correctly: - The early return for a step with nothing pending no longer requires a delta callback. Every write entry point that changes input sets pending_input_change, and a new session starts with it set. - A plain step evaluates every stratum. Once the session has evaluated, it first clears the derived relations with col_session_clear_idb_rows and resets every stratum frontier, under the same has_evaluated condition the snapshot path uses before a full re-evaluation. Without the stratum reset, the transitive closure in the new test lost (1,4) after an edge was added. A resumed step keeps what its first attempt already merged. - The full mask is forced even when the pending input recorded only some strata, as an insert staged with a callback that is then removed does. Clearing alone would leave the other strata empty. - Before clearing, the step sets pending_full_input_eval, which only its commit resets. If the step then fails, the next snapshot evaluates in full instead of re-deriving only the strata the pending input reaches over the emptied relations. - A plain step no longer seeds retraction deltas ($r$) left by a removal staged with a callback. Seeded, the first iteration read the removed rows in place of the relation and re-derived paths through them. Clearing derived rows also cleared input. A rule head can also hold input: an inline fact such as `reach(1).` beside `reach(y) :- reach(x), edge(x, y).`, or a host row inserted into it. It lost that input whenever a full re-evaluation reset it. On main a second snapshot after an insert already returned `reach` empty, and four-worker recursive steps did too. A rule head now keeps its input rows in a private `$in$` shadow: - The first insert into a rule head registers its shadow; inline facts are loaded through that same insert. Creating shadows at session creation would raise the creation memory floor that test_session pins, and creating them when the first evaluation starts would leave a reservation behind after a refused step. - A direct-mapped cache of relation identities answers "not a rule head" for plain input relations. Without it, a million single-row inserts took 0.71 s against main's 0.51 s; with it the two match within noise, also when two relations alternate. - Inserts write the shadow first and truncate it if the relation's own append fails. - Removals plan the shadow change before the relation is touched and apply it once the removal has committed. The shadow records input, so a removal applies to it even when the relation no longer holds the row; a callback-mode step can already have dropped it. Such a removal leaves the model as it is, so it schedules no re-evaluation. - wl_columnar_session_restore_seed appends the shadow back after both resets: col_session_clear_idb_rows and the multi-worker recursive reset in eval.c. Tests: - New test_easy_plain_step: - idle steps leave one copy of each row; - inserts and removals in two strata hold exactly the model; - a recursive closure is extended and then cut; - a step after removing the callback re-derives everything, after a staged insert and after a staged removal; - inline facts survive plain steps on one and four workers, and snapshots alone; - host rows in a rule head without inline facts survive the same three modes; - host rows and an inline fact removed with a callback installed, or after it was removed, stay removed; - a host row inserted twice is gone after two removals. - test_extension_replay: - a no-callback case pins five invocations to load, zero per idle step, five for an insert into an unread relation and four after a retraction. Snapshots in that case leave the invocation counter unchanged, so they read the stable model rather than repairing it. - a plain step whose addon fails after the clear is followed by a snapshot that returns the full model. - removing input the model already lacks, with a callback installed, invokes no addon in another stratum. - Mutation checks: each of fourteen mutants fails a named check: - restore the callback requirement; - drop the clear, the stratum reset, the forced mask, or pending_full_input_eval; - re-enable the retraction seed; - skip the shadow on insert, on its creation, on plain removal, on the incremental removal hit or miss, on restore, or in the eval.c reset; - schedule a re-evaluation for an input-only removal. Docs: - wirelog-easy.h documents both step modes and drops the warning against a snapshot after a step. - docs/SEMANTICS.md specifies the no-callback invocation count. Its publication-cutoff state map now lists col_session_step_impl as a reader of has_evaluated, and session_seed_shadow_truncate as a reader and writer of base_nrows. A selftest anchor that the first edit made ambiguous now names the one row it edits. - Example 12's comment and README no longer claim that step followed by snapshot double-counts; checked on that program, it returns the same five rows. - CHANGELOG. Not changed: with a delta callback, a rule head's own input is mishandled. The first step publishes an inline fact as a retraction and leaves the relation empty, and a host row inserted into a rule head is also published as a retraction. Main does the same. Fixes #1994 --- CHANGELOG.md | 22 + docs/SEMANTICS.md | 18 +- examples/12-snapshot-vs-delta/README.md | 7 +- examples/12-snapshot-vs-delta/snapshot_demo.c | 8 +- scripts/ci/test-check-state-map-anchors.py | 6 +- tests/meson.build | 12 + tests/test_easy_plain_step.c | 620 ++++++++++++++++++ tests/test_extension_replay.c | 284 +++++++- wirelog/columnar/eval.c | 3 + wirelog/columnar/internal.h | 11 + wirelog/columnar/session.c | 327 ++++++++- wirelog/wirelog-easy.h | 20 +- 12 files changed, 1292 insertions(+), 46 deletions(-) create mode 100644 tests/test_easy_plain_step.c diff --git a/CHANGELOG.md b/CHANGELOG.md index 71cb4d970..44bc31043 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -51,6 +51,28 @@ All notable changes to wirelog are documented in this file. ### Fixed +- **Steps without a delta callback no longer pile up derived rows** + (#1994): `wirelog_easy_step()` on a session with no delta callback + re-derived every rule on every step on top of what the previous step + had derived. Each step appended another copy of every derived row, a + removal never retracted the rows it had derived, and a step with + nothing pending re-ran every rule, invoking each value-position + `@call` once per derived row. A step with an insert or removal + pending now discards the derived rows and re-derives every rule, as a + full snapshot does, and a step with nothing pending does nothing. A + snapshot taken after such a step therefore reads exactly the model the + step derived, and `wirelog-easy.h` no longer warns against that + combination. + +- **Full re-evaluation keeps a rule head's own input rows** (#1994): a + relation that is a rule head and also holds input -- an inline fact + such as `reach(1).`, or a host row -- lost that input whenever a full + re-evaluation discarded the derived rows. A second snapshot after an + insert returned `reach` empty, and so did the multi-worker recursive + path. Such a relation now keeps its input rows in a private + `$in$` copy, maintained on every insert and removal and restored + after each discard. + - **Every configure left the worktree dirty** (#1814): `.gitignore` ignored `subprojects/xxHash-0.8.3/` while `subprojects/xxhash.wrap` pins `directory = xxHash-0.8.4`, so the rule named a directory meson diff --git a/docs/SEMANTICS.md b/docs/SEMANTICS.md index d0997f0a5..885d0811c 100644 --- a/docs/SEMANTICS.md +++ b/docs/SEMANTICS.md @@ -87,11 +87,15 @@ engine decides that at stratum granularity rather than per rule -- reaches a particular rule is a property of how the program stratifies, and mutating a relation the rule does not read is not on its own a guarantee. -How many invocations a session performs under other configurations differs and -is not specified here; issue #1994 tracks the one difference that is known. -Choosing a configuration for its invocation count would rest on something a -release may change, so a host should not do it -- and a delta callback is -installed because the host wants deltas, not as a tuning knob. +Without a delta callback, `wirelog_easy_step` is not incremental (#1994). A +step with an insert or removal pending since the last step or snapshot that +evaluated discards the derived rows and re-derives every rule, so the addon +callback runs once for every row its rule then derives, whether or not the +mutation reaches that rule. A step with nothing pending derives nothing and +invokes nothing. Other configurations are not specified here. Choosing a +configuration for its invocation count would rest on something a release may +change, so a host should not do it -- and a delta callback is installed because +the host wants deltas, not as a tuning knob. The delta stream carries the set difference, so a re-derivation that reproduces a row publishes nothing. Inserting a sixth row into a five-row @@ -840,7 +844,7 @@ an earlier revision summarised it here and got it wrong within one round. | `last_inserted_relation` | `col_session_snapshot_impl`, `col_session_step_impl`, `col_worker_session_create`, `session_note_inserted_input` | `col_eval_stratum_tdd_recursive`, `col_session_snapshot_impl`, `col_session_step_impl`, `session_note_inserted_input` | PRESERVE | | `pending_input_change` | `col_session_create_internal`, `col_session_remove`, `col_session_remove_incremental`, `col_session_snapshot_impl`, `col_session_step_impl`, `session_note_inserted_input` | `col_session_snapshot_impl`, `col_session_step_impl` | PRESERVE | | `pending_full_input_eval` | `col_session_insert`, `col_session_snapshot_impl`, `col_session_step_impl`, `session_note_inserted_input` | `col_session_snapshot_impl`, `col_session_step_impl` | PRESERVE | -| `has_evaluated` | `col_session_snapshot_impl`, `col_session_step_impl` | `col_eval_stratum_tdd_recursive`, `col_session_snapshot_impl` | PRESERVE | +| `has_evaluated` | `col_session_snapshot_impl`, `col_session_step_impl` | `col_eval_stratum_tdd_recursive`, `col_session_snapshot_impl`, `col_session_step_impl` | PRESERVE | | `snapshot_stable_valid` | `col_session_remove`, `col_session_remove_incremental`, `col_session_snapshot_impl`, `col_session_step_impl`, `session_note_inserted_input` | `col_session_snapshot_impl` | PRESERVE | | `delta_seeded` | `col_session_snapshot_impl`, `tdd_worker_subpass_fn` | `col_op_variable`, `col_session_snapshot_impl`, `has_empty_forced_delta`, `tdd_worker_subpass_fn`, `wl_columnar_eval_nonrec_relation_parallel`, `wl_columnar_join_select_right`, `wl_columnar_eval_tdd_plan_prepare_inputs` | PRESERVE | | `last_removed_relation` | `col_session_remove_incremental`, `col_session_step_impl`, `col_worker_session_create` | `col_session_remove_incremental`, `col_session_snapshot_impl`, `col_session_step_impl` | PRESERVE | @@ -849,7 +853,7 @@ an earlier revision summarised it here and got it wrong within one round. | `plain_step_completion_phase` | `col_eval_stratum_tdd_nonrecursive`, `col_session_step_impl`, `col_worker_session_create`, `wl_columnar_eval_resume_nonrecursive_completion` | `col_session_step_impl`, `wl_columnar_eval_resume_nonrecursive_completion` | PRESERVE | | `plain_step_completion_step_context` | `col_session_snapshot_impl`, `col_session_step_impl`, `col_worker_session_create` | `col_eval_stratum_tdd_nonrecursive`, `col_session_step_impl` | PRESERVE | | `plain_step_completion_active` | `col_eval_stratum_tdd_nonrecursive`, `col_session_step_impl`, `col_worker_session_create`, `wl_columnar_eval_resume_nonrecursive_completion` | `col_session_snapshot_impl`, `col_session_step_impl` | **COMMIT** on the step path -- a re-entrancy latch, not progress. While it stays set beside `plain_step_completion_pending`, the entry guard of both `col_session_step_impl` and `col_session_snapshot_impl` returns `EBUSY` above every line that would clear it. The snapshot path never writes it, so the snapshot cutoff has nothing to perform here. | -| `col_rel_t::base_nrows` | `col_rel_compact_impl`, `col_rel_deep_copy`, `col_rel_deep_copy_governed_impl`, `col_rel_install_shared_view_unprotected`, `col_rel_reset_rows_locked`, `col_session_remove`, `col_session_remove_incremental`, `col_session_snapshot_impl`, `col_session_step_impl`, `col_stratum_step_retraction_nonrecursive`, `tdd_empty_relation_candidate`, `tdd_seed_bdx_coordinator_idb`, `wl_columnar_eval_serial_canonicalize_aggregate_locked`, `wl_retraction_restore` | `col_rel_compact_impl`, `col_rel_deep_copy`, `col_rel_deep_copy_governed_impl`, `col_rel_install_shared_view_unprotected`, `col_rel_mutable_image_validate`, `col_session_remove`, `col_session_remove_incremental`, `col_session_snapshot_impl`, `col_stratum_step_retraction_nonrecursive`, `wl_retraction_stage_prepare` | PRESERVE -- the snapshot delta pre-seed skips a relation whose `base_nrows` is zero or whose `nrows <= base_nrows`, so committing the bookkeeping's `base_nrows = nrows` changes what the next attempt pre-seeds, including whether it pre-seeds at all. That is the whole verified consequence; see the note below before adding another | +| `col_rel_t::base_nrows` | `col_rel_compact_impl`, `col_rel_deep_copy`, `col_rel_deep_copy_governed_impl`, `col_rel_install_shared_view_unprotected`, `col_rel_reset_rows_locked`, `col_session_remove`, `col_session_remove_incremental`, `col_session_snapshot_impl`, `col_session_step_impl`, `col_stratum_step_retraction_nonrecursive`, `session_seed_shadow_truncate`, `tdd_empty_relation_candidate`, `tdd_seed_bdx_coordinator_idb`, `wl_columnar_eval_serial_canonicalize_aggregate_locked`, `wl_retraction_restore` | `col_rel_compact_impl`, `col_rel_deep_copy`, `col_rel_deep_copy_governed_impl`, `col_rel_install_shared_view_unprotected`, `col_rel_mutable_image_validate`, `col_session_remove`, `col_session_remove_incremental`, `col_session_snapshot_impl`, `col_stratum_step_retraction_nonrecursive`, `session_seed_shadow_truncate`, `wl_retraction_stage_prepare` | PRESERVE -- the snapshot delta pre-seed skips a relation whose `base_nrows` is zero or whose `nrows <= base_nrows`, so committing the bookkeeping's `base_nrows = nrows` changes what the next attempt pre-seeds, including whether it pre-seeds at all. That is the whole verified consequence; see the note below before adding another | Two non-field actions sit in the same region and need their own verdicts. `col_session_reclaim_quiescent` is **COMMIT**: the step cutoff's unwind calls diff --git a/examples/12-snapshot-vs-delta/README.md b/examples/12-snapshot-vs-delta/README.md index 381aeb625..3e08190cc 100644 --- a/examples/12-snapshot-vs-delta/README.md +++ b/examples/12-snapshot-vs-delta/README.md @@ -62,9 +62,10 @@ simpler when you just need the current answer without tracking history. - **Semantic equivalence** -- the incremental delta path produces the same derived facts as full re-evaluation via snapshot. -- **Two independent sessions** -- because `wirelog_easy_snapshot` is an - evaluating call, the driver uses separate sessions to avoid - double-counting. See `wirelog-easy.h` for the full contract. +- **Two independent sessions** -- the snapshot side runs in its own + session, so it is a full evaluation that shares no state with the + incremental side it is compared against. See `wirelog-easy.h` for how + a snapshot behaves after a step on the same session. - **Deterministic comparison** -- both result sets are sorted before comparison so the test is stable regardless of backend evaluation order. diff --git a/examples/12-snapshot-vs-delta/snapshot_demo.c b/examples/12-snapshot-vs-delta/snapshot_demo.c index 2f13a4943..46a6c2cb6 100644 --- a/examples/12-snapshot-vs-delta/snapshot_demo.c +++ b/examples/12-snapshot-vs-delta/snapshot_demo.c @@ -23,11 +23,9 @@ * wirelog_easy_step (delta callback, streaming) and wirelog_easy_snapshot * (one-shot batch) and verifies that they produce identical results. * - * IMPORTANT: wirelog_easy_snapshot() is an evaluating call -- calling - * wirelog_easy_step() followed by wirelog_easy_snapshot() on the same insert - * batch would double-count derived tuples. This driver therefore - * uses two independent sessions: one for the delta path and one for - * the snapshot path. + * The driver uses two independent sessions, one for the delta path and + * one for the snapshot path, so the snapshot side is a full evaluation + * that shares no state with the incremental side it is checked against. * * Build: meson compile -C build snapshot_demo * Run: ./build/examples/12-snapshot-vs-delta/snapshot_demo diff --git a/scripts/ci/test-check-state-map-anchors.py b/scripts/ci/test-check-state-map-anchors.py index 79c621292..54d45e600 100755 --- a/scripts/ci/test-check-state-map-anchors.py +++ b/scripts/ci/test-check-state-map-anchors.py @@ -399,9 +399,11 @@ def row_edit(name: str, old: str, new: str, want: int = 1) -> None: # The attribution check, read column: a function that exists, is # uniquely defined, and never touches the field beside it. row_edit("a real function that never reads the field fails", - "| `col_eval_stratum_tdd_recursive`, `col_session_snapshot_impl`,", + "| `col_eval_stratum_tdd_recursive`, `col_session_snapshot_impl`," + " `col_session_step_impl`, `session_note_inserted_input` |", "| `arr_build_full`, `col_eval_stratum_tdd_recursive`," - " `col_session_snapshot_impl`,") + " `col_session_snapshot_impl`, `col_session_step_impl`," + " `session_note_inserted_input` |") # The same for the write column, so `writes()` is not free to be # `return True` at this level either. diff --git a/tests/meson.build b/tests/meson.build index db7267604..c784ef6ae 100644 --- a/tests/meson.build +++ b/tests/meson.build @@ -5754,6 +5754,18 @@ test_extension_replay_exe = executable( test('extension_replay', test_extension_replay_exe) +# Issue #1994: a session stepped without a delta callback must hold exactly +# the current model -- no copy per step, no residue after a retraction. +test_easy_plain_step_exe = executable( + 'test_easy_plain_step', + files('test_easy_plain_step.c'), + include_directories: [wirelog_inc, wirelog_src_inc], + dependencies: [nanoarrow_dep, threads_dep, xxhash_dep, mbedtls_dep, math_dep], + link_with: [testlib_prod], +) + +test('easy_plain_step', test_easy_plain_step_exe) + # ============================================================================ # I/O Context Accessor Tests (#452) # ============================================================================ diff --git a/tests/test_easy_plain_step.c b/tests/test_easy_plain_step.c new file mode 100644 index 000000000..3470bd920 --- /dev/null +++ b/tests/test_easy_plain_step.c @@ -0,0 +1,620 @@ +/* + * test_easy_plain_step.c - Issue #1994 regression test. + * + * Copyright (C) CleverPlant + * Licensed under LGPL-3.0 + * For commercial licenses, contact: inquiry@cleverplant.com + * + * A session driven by wirelog_easy_step() with no delta callback installed + * re-derives every rule in full on a step with pending input. It used to + * do so on top of the rows the previous step derived, so each step appended + * another copy of every derived row, a retraction never removed the row it + * had derived, and a step with nothing pending re-derived everything again. + * + * Every check reads the model through wirelog_easy_snapshot() after a + * completed step. With nothing pending that snapshot emits the stable model + * without evaluating, so it observes what the steps left behind rather than + * repairing it; test_plain_step_rederives_every_rule in + * test_extension_replay.c proves that by counting addon invocations around + * such a snapshot. + * + * The SEEDED cases cover a relation that is a rule head and also holds + * input -- an inline fact or a host row. Discarding derived rows must keep + * that input, on plain steps, on snapshots alone, and in the multi-worker + * recursive path. + */ + +#include "wirelog/wirelog-easy.h" + +#include +#include +#include +#include + +static int tests_run = 0; +static int tests_passed = 0; +static int tests_failed = 0; + +#define TEST(name) \ + do { \ + tests_run++; \ + printf(" [%d] %s", tests_run, name); \ + } while (0) +#define PASS() \ + do { \ + tests_passed++; \ + printf(" ... PASS\n"); \ + } while (0) +#define FAIL(msg) \ + do { \ + tests_failed++; \ + printf(" ... FAIL: %s\n", (msg)); \ + } while (0) + +#define MAX_ROWS 64 + +typedef struct { + int64_t col[MAX_ROWS][2]; + uint32_t count; + int overflowed; +} rows_t; + +static void +collect(const char *relation, const int64_t *row, uint32_t ncols, + void *user_data) +{ + rows_t *rows = (rows_t *)user_data; + (void)relation; + if (rows->count >= MAX_ROWS) { + rows->overflowed = 1; + return; + } + rows->col[rows->count][0] = ncols > 0 ? row[0] : 0; + rows->col[rows->count][1] = ncols > 1 ? row[1] : 0; + rows->count++; +} + +static char message[256]; + +/* Snapshot @relation and require exactly @nwant rows, each of @want once. + * An overflow or a repeated row fails rather than reading as a match. */ +static const char * +expect_rows(wirelog_easy_session_t *s, const char *relation, + const int64_t (*want)[2], uint32_t nwant, const char *when) +{ + rows_t rows; + + memset(&rows, 0, sizeof(rows)); + if (wirelog_easy_snapshot(s, relation, collect, &rows) != WIRELOG_OK) { + snprintf(message, sizeof(message), "%s: snapshot failed", when); + return message; + } + if (rows.overflowed || rows.count != nwant) { + snprintf(message, sizeof(message), "%s: %s has %u rows, want %u", + when, relation, rows.overflowed ? MAX_ROWS + 1u : rows.count, + nwant); + return message; + } + for (uint32_t w = 0; w < nwant; w++) { + uint32_t seen = 0; + for (uint32_t r = 0; r < rows.count; r++) + if (rows.col[r][0] == want[w][0] && rows.col[r][1] == want[w][1]) + seen++; + if (seen != 1) { + snprintf(message, sizeof(message), + "%s: %s(%lld,%lld) seen %u times, want once", when, relation, + (long long)want[w][0], (long long)want[w][1], seen); + return message; + } + } + return NULL; +} + +static const char * +step(wirelog_easy_session_t *s, const char *when) +{ + if (wirelog_easy_step(s) != WIRELOG_OK) { + snprintf(message, sizeof(message), "%s: step failed", when); + return message; + } + return NULL; +} + +static const char * +insert1(wirelog_easy_session_t *s, const char *relation, int64_t a) +{ + return wirelog_easy_insert(s, relation, &a, 1) == WIRELOG_OK + ? NULL : "insert failed"; +} + +static const char * +insert2(wirelog_easy_session_t *s, const char *relation, int64_t a, int64_t b) +{ + int64_t row[2] = { a, b }; + return wirelog_easy_insert(s, relation, row, 2) == WIRELOG_OK + ? NULL : "insert failed"; +} + +static const char * +remove1(wirelog_easy_session_t *s, const char *relation, int64_t a) +{ + return wirelog_easy_remove(s, relation, &a, 1) == WIRELOG_OK + ? NULL : "remove failed"; +} + +static const char * +remove2(wirelog_easy_session_t *s, const char *relation, int64_t a, int64_t b) +{ + int64_t row[2] = { a, b }; + return wirelog_easy_remove(s, relation, row, 2) == WIRELOG_OK + ? NULL : "remove failed"; +} + +#define TRY(expr) \ + do { \ + const char *why_ = (expr); \ + if (why_) { \ + FAIL(why_); \ + goto done; \ + } \ + } while (0) + +/* Two strata: `twice` reads only `ready`, `echo` only `other`. */ +static const char *const TWO_STRATA = + ".decl ready(id: int64)\n" + ".decl other(k: int64)\n" + ".decl twice(id: int64, v: int64)\n" + ".decl echo(k: int64)\n" + "twice(id, id * 2) :- ready(id).\n" + "echo(k) :- other(k).\n"; + +static void +test_idle_steps_keep_one_copy(void) +{ + static const int64_t want[][2] = { { 1, 2 }, { 2, 4 }, { 3, 6 } }; + wirelog_easy_session_t *s = NULL; + + TEST("idle plain steps leave one copy of every derived row"); + if (wirelog_easy_open(TWO_STRATA, &s) != WIRELOG_OK) { + FAIL("open"); + return; + } + for (int64_t i = 1; i <= 3; i++) + TRY(insert1(s, "ready", i)); + TRY(step(s, "load")); + for (int i = 0; i < 4; i++) + TRY(step(s, "idle")); + TRY(expect_rows(s, "twice", want, 3, "after four idle steps")); + PASS(); +done: + wirelog_easy_close(s); +} + +static void +test_mutating_steps_replace_the_model(void) +{ + static const int64_t loaded[][2] = { { 1, 2 }, { 2, 4 }, { 3, 6 } }; + static const int64_t after_remove[][2] = { { 2, 4 }, { 3, 6 } }; + static const int64_t after_insert[][2] = { + { 2, 4 }, { 3, 6 }, { 9, 18 } + }; + static const int64_t one_echo[][2] = { { 7, 0 } }; + wirelog_easy_session_t *s = NULL; + + TEST("plain steps after inserts and removals hold exactly the model"); + if (wirelog_easy_open(TWO_STRATA, &s) != WIRELOG_OK) { + FAIL("open"); + return; + } + for (int64_t i = 1; i <= 3; i++) + TRY(insert1(s, "ready", i)); + TRY(step(s, "load")); + TRY(insert1(s, "other", 7)); + TRY(step(s, "insert other")); + TRY(expect_rows(s, "twice", loaded, 3, "unrelated insert")); + TRY(expect_rows(s, "echo", one_echo, 1, "unrelated insert")); + TRY(remove1(s, "other", 7)); + TRY(step(s, "remove other")); + TRY(expect_rows(s, "echo", NULL, 0, "retracted other")); + TRY(expect_rows(s, "twice", loaded, 3, "retracted other")); + TRY(remove1(s, "ready", 1)); + TRY(step(s, "remove ready")); + TRY(expect_rows(s, "twice", after_remove, 2, "retracted ready")); + TRY(insert1(s, "ready", 9)); + TRY(step(s, "insert ready")); + TRY(expect_rows(s, "twice", after_insert, 3, "inserted ready")); + PASS(); +done: + wirelog_easy_close(s); +} + +static void +test_recursive_rederivation(void) +{ + static const char *const TC = + ".decl edge(a: int64, b: int64)\n" + ".decl path(a: int64, b: int64)\n" + "path(a, b) :- edge(a, b).\n" + "path(a, c) :- path(a, b), edge(b, c).\n"; + static const int64_t chain[][2] = { + { 1, 2 }, { 2, 3 }, { 1, 3 } + }; + static const int64_t extended[][2] = { + { 1, 2 }, { 2, 3 }, { 3, 4 }, { 1, 3 }, { 2, 4 }, { 1, 4 } + }; + static const int64_t cut[][2] = { { 1, 2 }, { 3, 4 } }; + wirelog_easy_session_t *s = NULL; + + TEST("recursive plain steps re-derive the closure without residue"); + if (wirelog_easy_open(TC, &s) != WIRELOG_OK) { + FAIL("open"); + return; + } + TRY(insert2(s, "edge", 1, 2)); + TRY(insert2(s, "edge", 2, 3)); + TRY(step(s, "load")); + TRY(step(s, "idle")); + TRY(expect_rows(s, "path", chain, 3, "chain")); + TRY(insert2(s, "edge", 3, 4)); + TRY(step(s, "extend")); + TRY(step(s, "idle")); + TRY(expect_rows(s, "path", extended, 6, "extended")); + TRY(remove2(s, "edge", 2, 3)); + TRY(step(s, "cut")); + TRY(expect_rows(s, "path", cut, 2, "cut")); + PASS(); +done: + wirelog_easy_close(s); +} + +static void +ignore_delta(const char *relation, const int64_t *row, uint32_t ncols, + int32_t diff, void *user_data) +{ + (void)relation; + (void)row; + (void)ncols; + (void)diff; + (void)user_data; +} + +/* An insert made while a delta callback was installed records only the + * strata it reaches. If the callback is removed before the next step, that + * step is a plain one and discards every derived relation, so it must also + * re-derive every stratum, not just the recorded ones. */ +static void +test_callback_removed_before_step(void) +{ + static const int64_t loaded[][2] = { { 1, 2 }, { 2, 4 }, { 3, 6 } }; + static const int64_t one_echo[][2] = { { 7, 0 } }; + wirelog_easy_session_t *s = NULL; + + TEST("a step after the delta callback is removed re-derives everything"); + if (wirelog_easy_open(TWO_STRATA, &s) != WIRELOG_OK) { + FAIL("open"); + return; + } + if (wirelog_easy_set_delta_cb(s, ignore_delta, NULL) != WIRELOG_OK) { + FAIL("set_delta_cb"); + goto done; + } + for (int64_t i = 1; i <= 3; i++) + TRY(insert1(s, "ready", i)); + TRY(step(s, "load with callback")); + TRY(insert1(s, "other", 7)); + if (wirelog_easy_set_delta_cb(s, NULL, NULL) != WIRELOG_OK) { + FAIL("clear delta_cb"); + goto done; + } + TRY(step(s, "first plain step")); + TRY(expect_rows(s, "twice", loaded, 3, "after callback removal")); + TRY(expect_rows(s, "echo", one_echo, 1, "after callback removal")); + PASS(); +done: + wirelog_easy_close(s); +} + +/* The removal half of the same transition: a removal staged with a callback + * leaves its retraction relation behind, which a plain step must not read + * in place of the relation itself. */ +static void +test_callback_removal_before_plain_step(void) +{ + static const char *const TC = + ".decl edge(a: int64, b: int64)\n" + ".decl path(a: int64, b: int64)\n" + "path(a, b) :- edge(a, b).\n" + "path(a, c) :- path(a, b), edge(b, c).\n"; + static const int64_t cut[][2] = { { 1, 2 }, { 3, 4 } }; + wirelog_easy_session_t *s = NULL; + + TEST("a removal staged with a callback is honoured by a plain step"); + if (wirelog_easy_open(TC, &s) != WIRELOG_OK) { + FAIL("open"); + return; + } + if (wirelog_easy_set_delta_cb(s, ignore_delta, NULL) != WIRELOG_OK) { + FAIL("set_delta_cb"); + goto done; + } + TRY(insert2(s, "edge", 1, 2)); + TRY(insert2(s, "edge", 2, 3)); + TRY(insert2(s, "edge", 3, 4)); + TRY(step(s, "load with callback")); + TRY(remove2(s, "edge", 2, 3)); + if (wirelog_easy_set_delta_cb(s, NULL, NULL) != WIRELOG_OK) { + FAIL("clear delta_cb"); + goto done; + } + TRY(step(s, "first plain step")); + TRY(expect_rows(s, "path", cut, 2, "after staged removal")); + PASS(); +done: + wirelog_easy_close(s); +} + +/* `reach` is a rule head that also carries an inline fact. A full + * re-evaluation discards what `reach` derived; it must keep the fact. */ +static const char *const SEEDED = + ".decl edge(a: int64, b: int64)\n" + ".decl reach(x: int64)\n" + "reach(1).\n" + "reach(y) :- reach(x), edge(x, y).\n"; + +typedef enum { SEED_PLAIN_STEP, SEED_SNAPSHOT_ONLY } seed_mode_t; + +static const char * +seed_sequence(wirelog_easy_session_t *s, seed_mode_t mode) +{ + static const int64_t r12[][2] = { { 1, 0 }, { 2, 0 } }; + static const int64_t r123[][2] = { { 1, 0 }, { 2, 0 }, { 3, 0 } }; + static const int64_t r1[][2] = { { 1, 0 } }; + static const int64_t host[][2] = { + { 1, 0 }, { 5, 0 }, { 6, 0 } + }; + const char *why; + +#define SEED_STEP(when) \ + do { \ + if (mode != SEED_SNAPSHOT_ONLY \ + && (why = step(s, (when))) != NULL) \ + return why; \ + } while (0) +#define SEED_EXPECT(want, n, when) \ + do { \ + if ((why = expect_rows(s, "reach", (want), (n), (when))) \ + != NULL) \ + return why; \ + } while (0) + + if ((why = insert2(s, "edge", 1, 2)) != NULL) + return why; + SEED_STEP("first"); + SEED_EXPECT(r12, 2, "first"); + if ((why = insert2(s, "edge", 2, 3)) != NULL) + return why; + SEED_STEP("extend"); + SEED_EXPECT(r123, 3, "extend"); + if ((why = remove2(s, "edge", 1, 2)) != NULL) + return why; + SEED_STEP("cut"); + SEED_EXPECT(r1, 1, "cut keeps the fact"); + SEED_STEP("idle"); + SEED_EXPECT(r1, 1, "idle keeps the fact"); + /* A host row in the same relation is input too, and so is its removal. */ + if ((why = insert1(s, "reach", 5)) != NULL + || (why = insert2(s, "edge", 5, 6)) != NULL) + return why; + SEED_STEP("host row"); + SEED_EXPECT(host, 3, "host row derives"); + if ((why = remove1(s, "reach", 5)) != NULL) + return why; + SEED_STEP("host removal"); + SEED_EXPECT(r1, 1, "host removal retracts"); +#undef SEED_STEP +#undef SEED_EXPECT + return NULL; +} + +/* `reach` has no inline fact, so its only input is what the host inserts + * into it once evaluation has registered it. */ +static const char *const HOST_SEEDED = + ".decl src(x: int64)\n" + ".decl edge(a: int64, b: int64)\n" + ".decl reach(x: int64)\n" + "reach(x) :- src(x).\n" + "reach(y) :- reach(x), edge(x, y).\n"; + +static const char * +host_seed_sequence(wirelog_easy_session_t *s, seed_mode_t mode) +{ + static const int64_t r12[][2] = { { 1, 0 }, { 2, 0 } }; + static const int64_t r1278[][2] = { + { 1, 0 }, { 2, 0 }, { 7, 0 }, { 8, 0 } + }; + const char *why; + + if ((why = insert1(s, "src", 1)) != NULL + || (why = insert2(s, "edge", 1, 2)) != NULL) + return why; + if (mode != SEED_SNAPSHOT_ONLY && (why = step(s, "first")) != NULL) + return why; + if ((why = expect_rows(s, "reach", r12, 2, "first")) != NULL) + return why; + if ((why = insert1(s, "reach", 7)) != NULL + || (why = insert2(s, "edge", 7, 8)) != NULL) + return why; + if (mode != SEED_SNAPSHOT_ONLY && (why = step(s, "host row")) != NULL) + return why; + if ((why = expect_rows(s, "reach", r1278, 4, "host row derives")) != NULL) + return why; + if ((why = remove1(s, "reach", 7)) != NULL) + return why; + if (mode != SEED_SNAPSHOT_ONLY + && (why = step(s, "host removal")) != NULL) + return why; + return expect_rows(s, "reach", r12, 2, "host removal retracts"); +} + +static void +run_seed_case(const char *name, seed_mode_t mode, uint32_t workers, + bool host_rows) +{ + wirelog_easy_open_opts_t opts = WIRELOG_EASY_OPEN_OPTS_INIT; + wirelog_easy_session_t *s = NULL; + const char *why; + + TEST(name); + opts.num_workers = workers; + if (wirelog_easy_open_opts(host_rows ? HOST_SEEDED : SEEDED, &opts, &s) + != WIRELOG_OK) { + FAIL("open"); + return; + } + why = host_rows ? host_seed_sequence(s, mode) : seed_sequence(s, mode); + if (why) + FAIL(why); + else + PASS(); + wirelog_easy_close(s); +} + +/* A host row removed while a delta callback is installed goes through the + * incremental removal path; the next plain step must not bring it back. + * With @plain_removal the callback is cleared first, so the plain removal + * path takes the row out. Either way the relation need not hold the row: + * a delta-callback step can already have dropped it. */ +static void +run_callback_host_row_case(const char *name, bool plain_removal) +{ + static const int64_t r12[][2] = { { 1, 0 }, { 2, 0 } }; + wirelog_easy_session_t *s = NULL; + + TEST(name); + if (wirelog_easy_open(HOST_SEEDED, &s) != WIRELOG_OK) { + FAIL("open"); + return; + } + if (wirelog_easy_set_delta_cb(s, ignore_delta, NULL) != WIRELOG_OK) { + FAIL("set_delta_cb"); + goto done; + } + TRY(insert1(s, "src", 1)); + TRY(insert2(s, "edge", 1, 2)); + TRY(step(s, "load with callback")); + TRY(insert1(s, "reach", 7)); + TRY(insert2(s, "edge", 7, 8)); + TRY(step(s, "host row with callback")); + if (!plain_removal) { + TRY(remove1(s, "reach", 7)); + TRY(step(s, "host removal with callback")); + } + if (wirelog_easy_set_delta_cb(s, NULL, NULL) != WIRELOG_OK) { + FAIL("clear delta_cb"); + goto done; + } + if (plain_removal) + TRY(remove1(s, "reach", 7)); + TRY(insert1(s, "src", 1)); + TRY(step(s, "plain step")); + TRY(expect_rows(s, "reach", r12, 2, "after the plain step")); + PASS(); +done: + wirelog_easy_close(s); +} + +/* The removal finds the fact in the relation, before any step, while a + * callback is installed: the incremental path's own commit. */ +static void +test_callback_removal_of_an_inline_fact(void) +{ + wirelog_easy_session_t *s = NULL; + + TEST("an inline fact removed with a callback stays removed"); + if (wirelog_easy_open(SEEDED, &s) != WIRELOG_OK) { + FAIL("open"); + return; + } + if (wirelog_easy_set_delta_cb(s, ignore_delta, NULL) != WIRELOG_OK) { + FAIL("set_delta_cb"); + goto done; + } + TRY(remove1(s, "reach", 1)); + if (wirelog_easy_set_delta_cb(s, NULL, NULL) != WIRELOG_OK) { + FAIL("clear delta_cb"); + goto done; + } + TRY(insert2(s, "edge", 1, 2)); + TRY(step(s, "first plain step")); + TRY(step(s, "second plain step")); + TRY(insert2(s, "edge", 2, 3)); + TRY(step(s, "re-deriving step")); + TRY(expect_rows(s, "reach", NULL, 0, "after the removed fact")); + PASS(); +done: + wirelog_easy_close(s); +} + +/* Each insert is input in its own right: the relation holds both copies of + * a row inserted twice, and each removal must take one copy out of the + * input as well, or the next re-derivation restores it. */ +static void +test_twice_inserted_host_row(void) +{ + static const int64_t r12[][2] = { { 1, 0 }, { 2, 0 } }; + wirelog_easy_session_t *s = NULL; + + TEST("a host row inserted twice is gone after two removals"); + if (wirelog_easy_open(HOST_SEEDED, &s) != WIRELOG_OK) { + FAIL("open"); + return; + } + TRY(insert1(s, "src", 1)); + TRY(insert2(s, "edge", 1, 2)); + TRY(step(s, "load")); + TRY(insert1(s, "reach", 7)); + TRY(insert1(s, "reach", 7)); + TRY(step(s, "two host rows")); + TRY(remove1(s, "reach", 7)); + TRY(step(s, "first removal")); + TRY(remove1(s, "reach", 7)); + TRY(step(s, "second removal")); + TRY(insert1(s, "src", 1)); + TRY(step(s, "re-deriving step")); + TRY(expect_rows(s, "reach", r12, 2, "after both removals")); + PASS(); +done: + wirelog_easy_close(s); +} + +int +main(void) +{ + printf("test_easy_plain_step (Issue #1994)\n"); + test_idle_steps_keep_one_copy(); + test_mutating_steps_replace_the_model(); + test_recursive_rederivation(); + test_callback_removed_before_step(); + test_callback_removal_before_plain_step(); + run_seed_case("inline facts and host rows survive plain steps", + SEED_PLAIN_STEP, 1, false); + run_seed_case("inline facts and host rows survive snapshots alone", + SEED_SNAPSHOT_ONLY, 1, false); + run_seed_case("inline facts survive plain steps on four workers", + SEED_PLAIN_STEP, 4, false); + run_seed_case("host rows in an unseeded rule head survive plain steps", + SEED_PLAIN_STEP, 1, true); + run_seed_case("host rows in an unseeded rule head survive snapshots", + SEED_SNAPSHOT_ONLY, 1, true); + run_seed_case("host rows in an unseeded rule head on four workers", + SEED_PLAIN_STEP, 4, true); + run_callback_host_row_case( + "a host row removed with a callback stays removed", false); + run_callback_host_row_case( + "a host row removed after the callback stays removed", true); + test_callback_removal_of_an_inline_fact(); + test_twice_inserted_host_row(); + printf("%d run, %d passed, %d failed\n", tests_run, tests_passed, + tests_failed); + return tests_failed ? 1 : 0; +} diff --git a/tests/test_extension_replay.c b/tests/test_extension_replay.c index 6b3577fe4..1447ee063 100644 --- a/tests/test_extension_replay.c +++ b/tests/test_extension_replay.c @@ -21,12 +21,12 @@ * docs/SEMANTICS.md. If you change the engine on purpose, update these * numbers and re-read that section to check it still holds. * - * Every case installs a delta callback and drives the session with - * wirelog_easy_step(). Without a delta callback the addon runs once per - * derived row on every step, including steps with nothing pending, but that - * count is unspecified and a stepped session has no non-perturbing way to - * observe anything through this API, so it is deliberately not pinned here: - * pinning it would freeze behavior nobody has characterised. Issue #1994. + * Every case drives the session with wirelog_easy_step(), and all but two + * of them, named below, do so with a delta callback installed throughout. + * test_plain_step_rederives_every_rule drives the same rule with no + * callback (Issue #1994): there a step with pending input re-derives every + * rule in full and an idle step does nothing. The second-to-last case + * removes its callback before a step that fails. */ #include "wirelog/wirelog-easy.h" @@ -77,10 +77,9 @@ typedef struct { int32_t diff; } event_t; -/* Host-side z-set mirror of a derived relation, rebuilt from the deltas. - * wirelog_easy_snapshot() is itself an evaluating call and must not follow - * wirelog_easy_step() on the same session, so the deltas are the only - * non-perturbing way to observe what the engine published. */ +/* Host-side z-set mirror of a derived relation, rebuilt from the deltas, + * so the delta-callback cases check what the engine published rather than + * what it holds. */ typedef struct { int64_t col[MAX_ROWS][2]; int32_t mult[MAX_ROWS]; @@ -181,6 +180,18 @@ stable_action(const wirelog_extension_value_t *args, uint32_t nargs, return 0; } +/* Stable addon that fails while flaky_fail is set. */ +static int flaky_fail; + +static int +flaky_action(const wirelog_extension_value_t *args, uint32_t nargs, + wirelog_extension_value_t *result, void *user_data) +{ + if (flaky_fail) + return -1; + return stable_action(args, nargs, result, user_data); +} + /* Unstable addon: a fresh value on every invocation, which is what a host * handing out tickets does. Declares no PURE or DETERMINISTIC bit, and * nothing in the engine asks for one. */ @@ -223,6 +234,18 @@ static const char *const JOIN_PROGRAM = "started(id, @call(\"replay.action\", id)) :- ready(id)," " enabled(id).\n"; +/* REPLAY_PROGRAM's call site beside a recursive rule head the host also + * writes into. `reach` shares no relation with `started`. */ +static const char *const REACH_PROGRAM = + ".decl ready(id: int64)\n" + ".decl started(id: int64, ticket: int64)\n" + ".decl src(x: int64)\n" + ".decl edge(a: int64, b: int64)\n" + ".decl reach(x: int64)\n" + "started(id, @call(\"replay.action\", id)) :- ready(id).\n" + "reach(x) :- src(x).\n" + "reach(y) :- reach(x), edge(x, y).\n"; + typedef struct { wirelog_extension_registry_t *registry; wirelog_extension_snapshot_t *snapshot; @@ -232,7 +255,7 @@ typedef struct { static char message[256]; static const char * -fixture_open(fixture_t *fx, const char *program, uint32_t policy, +fixture_open_plain(fixture_t *fx, const char *program, uint32_t policy, wirelog_extension_scalar_fn fn) { uint32_t types[] = { WIRELOG_EXTENSION_VALUE_INT64 }; @@ -273,6 +296,16 @@ fixture_open(fixture_t *fx, const char *program, uint32_t policy, snprintf(message, sizeof(message), "open_opts returned %d", (int)rc); return message; } + return NULL; +} + +static const char * +fixture_open(fixture_t *fx, const char *program, uint32_t policy, + wirelog_extension_scalar_fn fn) +{ + const char *why = fixture_open_plain(fx, program, policy, fn); + if (why) + return why; if (wirelog_easy_set_delta_cb(fx->session, on_delta, NULL) != WIRELOG_OK) return "set_delta_cb"; return NULL; @@ -382,6 +415,32 @@ expect_one_event(const char *relation, int64_t a, int64_t b, int32_t diff, return NULL; } +/* Counts the `started` rows a snapshot emits, each of which must carry the + * stable addon's value for its id. A repeated row counts twice. */ +static void +on_started_row(const char *relation, const int64_t *row, uint32_t ncols, + void *user_data) +{ + unsigned *rows = (unsigned *)user_data; + (void)relation; + if (ncols == 2 && row[1] == row[0] * 10) + (*rows)++; + else + *rows += 1000; /* a malformed row can never read as a match */ +} + +static const char * +expect_snapshot(fixture_t *fx, unsigned want, const char *what) +{ + unsigned rows = 0; + if (wirelog_easy_snapshot(fx->session, "started", on_started_row, &rows) + != WIRELOG_OK) { + snprintf(message, sizeof(message), "%s: snapshot failed", what); + return message; + } + return expect_count(rows, want, what); +} + /* ======================================================================== */ /* Cases */ /* ======================================================================== */ @@ -801,6 +860,206 @@ test_selective_body_narrows_the_count(void) PASS(); } +/* Issue #1994: the same rule with no delta callback. A step with pending + * input re-derives every rule in full -- a host insert or removal without a + * callback is not incremental -- so the count follows the derivation of + * every rule, including one the mutation does not reach. A step with + * nothing pending does no work at all. The snapshots read the stable model + * a completed step leaves behind: with nothing pending they do not evaluate, + * which the unchanged invocation counter around them proves. */ +static void +test_plain_step_rederives_every_rule(void) +{ + static const int64_t loaded[] = { 1, 2, 3, 4, 5 }; + static const int64_t survivors[] = { 2, 3, 4, 5 }; + fixture_t fx; + const char *why; + unsigned n = 0; + unsigned before; + int64_t value; + + TEST("without a delta callback only steps with pending input invoke"); + calls = 0; + why = fixture_open_plain(&fx, REPLAY_PROGRAM, + WIRELOG_EXTENSION_CALLBACK_THREAD_SAFE, stable_action); + for (unsigned i = 0; !why && i < 5; i++) { + value = (int64_t)(i + 1); + if (wirelog_easy_insert(fx.session, "ready", &value, 1) + != WIRELOG_OK) + why = "insert ready"; + } + if (!why) + why = step_once(&fx, &n); + if (!why) + why = expect_count(n, 5, "load"); + if (!why) + why = expect_args(loaded, 5, "load args"); + for (unsigned i = 0; !why && i < 3; i++) { + why = step_once(&fx, &n); + if (!why) + why = expect_count(n, 0, "idle step"); + } + before = calls; + if (!why) + why = expect_snapshot(&fx, 5, "after idle steps"); + if (!why) + why = expect_count(calls - before, 0, "stable snapshot"); + + value = 7; + if (!why && wirelog_easy_insert(fx.session, "other", &value, 1) + != WIRELOG_OK) + why = "insert other"; + if (!why) + why = step_once(&fx, &n); + if (!why) + why = expect_count(n, 5, "insert into an unread relation"); + + value = 1; + if (!why && wirelog_easy_remove(fx.session, "ready", &value, 1) + != WIRELOG_OK) + why = "remove ready"; + if (!why) + why = step_once(&fx, &n); + if (!why) + why = expect_count(n, 4, "retraction with four survivors"); + if (!why) + why = expect_args(survivors, 4, "retraction args"); + before = calls; + if (!why) + why = expect_snapshot(&fx, 4, "after retraction"); + if (!why) + why = expect_count(calls - before, 0, "stable snapshot"); + if (!why) + why = expect_count(calls, 14, "total invocations"); + { + const char *teardown = fixture_close(&fx); + if (!why) + why = teardown; + } + if (why) + FAIL(why); + else + PASS(); +} + +/* Issue #1994: a plain step discards the derived relations before it + * re-derives them. If it then fails, nothing it discarded may stay lost: + * the next evaluation must still be a full one. The pending insert below + * is staged with a callback installed, so on its own it would let a + * snapshot evaluate only the stratum it reaches. */ +static void +test_failed_plain_step_leaves_a_full_evaluation_pending(void) +{ + fixture_t fx; + const char *why; + unsigned n = 0; + int64_t value; + + TEST("a plain step that fails after discarding still re-derives later"); + calls = 0; + flaky_fail = 0; + why = fixture_open(&fx, REPLAY_PROGRAM, + WIRELOG_EXTENSION_CALLBACK_THREAD_SAFE, flaky_action); + for (unsigned i = 0; !why && i < 3; i++) { + value = (int64_t)(i + 1); + if (wirelog_easy_insert(fx.session, "ready", &value, 1) + != WIRELOG_OK) + why = "insert ready"; + } + if (!why) + why = step_once(&fx, &n); + if (!why) + why = expect_count(n, 3, "load"); + value = 7; + if (!why && wirelog_easy_insert(fx.session, "other", &value, 1) + != WIRELOG_OK) + why = "insert other"; + if (!why && wirelog_easy_set_delta_cb(fx.session, NULL, NULL) + != WIRELOG_OK) + why = "clear delta_cb"; + if (!why) { + flaky_fail = 1; + if (wirelog_easy_step(fx.session) == WIRELOG_OK) + why = "a step whose addon fails reported success"; + flaky_fail = 0; + } + if (!why) + why = expect_snapshot(&fx, 3, "snapshot after the failed step"); + if (!why) + why = step_once(&fx, &n); + if (!why) + why = expect_snapshot(&fx, 3, "idle step after the snapshot"); + { + const char *teardown = fixture_close(&fx); + if (!why) + why = teardown; + } + if (why) + FAIL(why); + else + PASS(); +} + +/* Issue #1994: removing a host row from a rule head also removes it from the + * head's input record. When the relation no longer holds the row -- a + * delta-callback step can drop it -- the model does not change, so the next + * step must leave a rule in another stratum alone. */ +static void +test_input_only_removal_leaves_other_rules_alone(void) +{ + fixture_t fx; + const char *why; + unsigned n = 0; + int64_t value; + int64_t edge[2] = { 1, 2 }; + + TEST("removing input the model already lacks re-runs no other rule"); + calls = 0; + why = fixture_open(&fx, REACH_PROGRAM, + WIRELOG_EXTENSION_CALLBACK_THREAD_SAFE, stable_action); + for (unsigned i = 0; !why && i < 3; i++) { + value = (int64_t)(i + 1); + if (wirelog_easy_insert(fx.session, "ready", &value, 1) + != WIRELOG_OK) + why = "insert ready"; + } + value = 1; + if (!why && (wirelog_easy_insert(fx.session, "src", &value, 1) + != WIRELOG_OK + || wirelog_easy_insert(fx.session, "edge", edge, 2) != WIRELOG_OK)) + why = "insert reach inputs"; + if (!why) + why = step_once(&fx, &n); + if (!why) + why = expect_count(n, 3, "load"); + value = 7; + if (!why && wirelog_easy_insert(fx.session, "reach", &value, 1) + != WIRELOG_OK) + why = "insert reach(7)"; + if (!why) + why = step_once(&fx, &n); + if (!why) + why = expect_count(n, 0, "insert into the other stratum"); + if (!why && wirelog_easy_remove(fx.session, "reach", &value, 1) + != WIRELOG_OK) + why = "remove reach(7)"; + if (!why) + why = step_once(&fx, &n); + if (!why) + why = expect_count(n, 0, "removal the model already reflects"); + if (!why) + why = expect_count(nevents, 0, "events for that removal"); + { + const char *teardown = fixture_close(&fx); + if (!why) + why = teardown; + } + if (why) + FAIL(why); + else + PASS(); +} + int main(void) { @@ -811,6 +1070,9 @@ main(void) test_selective_body_narrows_the_count(); test_purity_bits_do_not_change_the_count(); test_unstable_value_republishes_unchanged_rows(); + test_plain_step_rederives_every_rule(); + test_failed_plain_step_leaves_a_full_evaluation_pending(); + test_input_only_removal_leaves_other_rules_alone(); printf("\n"); printf("Passed: %d/%d\n", tests_passed, tests_run); diff --git a/wirelog/columnar/eval.c b/wirelog/columnar/eval.c index cf4a490f1..68ce87b74 100644 --- a/wirelog/columnar/eval.c +++ b/wirelog/columnar/eval.c @@ -6840,6 +6840,9 @@ col_eval_stratum_tdd_recursive(const wl_plan_stratum_t *sp, col_rel_t *r = session_find_rel(coord, sp->relations[ri].name); if (r && r->nrows > 0) { rc = tdd_reset_coord_relation(r); + /* Issue #1994: keep the relation's own input rows. */ + if (rc == 0) + rc = wl_columnar_session_restore_seed(coord, r); if (rc != 0) { coord->tdd_total_ns += now_ns() - tdd_total_t0; return rc; diff --git a/wirelog/columnar/internal.h b/wirelog/columnar/internal.h index a240f8e78..9bf1a5bab 100644 --- a/wirelog/columnar/internal.h +++ b/wirelog/columnar/internal.h @@ -2352,12 +2352,23 @@ typedef struct wl_col_session_t { * reused. The report flag prevents retry from double-counting metrics. */ bool teardown_started; bool teardown_reported; + /* Issue #1994: identities of relations a host insert found to be no + * rule head, direct-mapped by identity, so inserts into input relations + * skip the seed-shadow lookup. Rule-head membership is fixed by the + * plan, and relation identities are never reused, so an entry cannot go + * stale; a collision only costs a lookup. */ + uint64_t seed_none_identity[16]; } wl_col_session_t; bool wl_columnar_session_budget_denied(const void *session); void wl_columnar_session_budget_denial_clear(wl_col_session_t *sess); +/* Issue #1994: append the input rows a rule-head relation holds in its + * `$in$` seed shadow back into @r after @r was reset for a full + * re-evaluation. Returns 0 when @r has no shadow. */ +int +wl_columnar_session_restore_seed(wl_col_session_t *sess, col_rel_t *r); typedef struct wl_columnar_session_hash_registry_image wl_columnar_session_hash_registry_image_t; diff --git a/wirelog/columnar/session.c b/wirelog/columnar/session.c index 206cf147c..caf7e7e73 100644 --- a/wirelog/columnar/session.c +++ b/wirelog/columnar/session.c @@ -219,6 +219,212 @@ session_note_inserted_input(wl_col_session_t *sess, const col_rel_t *relation, sess->snapshot_stable_valid = false; } +/* + * Seed shadows (Issue #1994). + * + * A rule head can also receive input -- inline facts such as `reach(1).`, or + * a host insert -- and then holds its input rows and its derived rows in one + * relation. A full re-evaluation discards the derived rows by resetting the + * whole relation, which would discard the input too. Every rule head + * therefore has a private copy of its input rows in `$in$`, registered + * on the first insert into it (see session_seed_shadow_for_insert), kept in + * step with every insert and removal, and appended back after each reset. + * No plan operator names a `$in$` relation, so evaluation never reads it. + */ +#define WL_SEED_SHADOW_PREFIX "$in$" + +static int +session_seed_shadow_name(const char *relation, char *out, size_t cap) +{ + int len = snprintf(out, cap, WL_SEED_SHADOW_PREFIX "%s", relation); + return len < 0 || (size_t)len >= cap ? EOVERFLOW : 0; +} + +static col_rel_t * +session_seed_shadow(wl_col_session_t *sess, const char *relation) +{ + char name[256]; + if (!sess || !relation + || session_seed_shadow_name(relation, name, sizeof(name)) != 0) + return NULL; + return session_find_rel(sess, name); +} + +static bool +session_plan_has_rule_head(const wl_plan_t *plan, const char *relation) +{ + for (uint32_t si = 0; plan && si < plan->stratum_count; si++) + for (uint32_t ri = 0; ri < plan->strata[si].relation_count; ri++) + if (plan->strata[si].relations[ri].name + && strcmp(plan->strata[si].relations[ri].name, relation) + == 0) + return true; + return false; +} + +static int +session_seed_shadow_create(wl_col_session_t *sess, const char *relation) +{ + char name[256]; + col_rel_t *seed = NULL; + int rc = session_seed_shadow_name(relation, name, sizeof(name)); + if (rc == 0) + rc = wl_columnar_relation_alloc_governed(&seed, name, + sess->memory_governor); + if (rc != 0) + return rc; + rc = session_add_rel(sess, seed); + if (rc != 0) { + rc = rc == ENOMEM && seed->memory_budget_denial_pending ? ENOSPC : rc; + col_rel_destroy(seed); + } + return rc; +} + +/* The shadow a host insert into @r must also write, created on the first + * insert into a rule head -- inline facts are loaded through this same + * insert path. *seed_out is NULL for a relation that is not a rule head. + * On the common path -- inserts into plain input relations -- + * seed_none_identity answers without a lookup. */ +static int +session_seed_shadow_for_insert(wl_col_session_t *sess, const col_rel_t *r, + col_rel_t **seed_out) +{ + uint64_t *none = &sess->seed_none_identity[r->relation_identity + % (sizeof(sess->seed_none_identity) + / sizeof(sess->seed_none_identity[0]))]; + + *seed_out = NULL; + if (r->relation_identity != 0 && r->relation_identity == *none) + return 0; + *seed_out = session_seed_shadow(sess, r->name); + if (*seed_out) + return 0; + if (!session_plan_has_rule_head(sess->plan, r->name)) { + *none = r->relation_identity; + return 0; + } + int rc = session_seed_shadow_create(sess, r->name); + if (rc == 0) + *seed_out = session_seed_shadow(sess, r->name); + return rc; +} + +/* Drop rows [nrows, seed->nrows) appended by a batch whose relation write + * then failed. The shadow is private to the session, so nothing reads it + * between the two writes. */ +static void +session_seed_shadow_truncate(col_rel_t *seed, uint32_t nrows) +{ + if (!seed || seed->nrows <= nrows) + return; + seed->nrows = nrows; + if (seed->base_nrows > nrows) + seed->base_nrows = nrows; + seed->sorted_nrows = 0; + seed->run_count = 0; + memset(seed->run_ends, 0, sizeof(seed->run_ends)); + wl_columnar_eval_dedup_set_clear(seed); + wl_columnar_relation_touch_view(seed); +} + +/* Mark, without changing the shadow, one shadow row for each requested row + * it holds. *mask_out is NULL when there is no shadow or nothing matched; + * otherwise the caller passes it to session_seed_shadow_apply_remove once the + * relation's own removal has committed, or frees it. */ +static int +session_seed_shadow_plan_remove(col_rel_t *seed, const int64_t *data, + uint32_t num_rows, uint32_t num_cols, uint8_t **mask_out) +{ + int64_t row_stack[COL_STACK_MAX]; + int64_t *row_buf = row_stack; + uint8_t *mask = NULL; + bool any = false; + + *mask_out = NULL; + if (!seed || seed->nrows == 0 || seed->ncols != num_cols) + return 0; + mask = (uint8_t *)calloc(seed->nrows, sizeof(*mask)); + if (!mask) + return ENOMEM; + if (num_cols > COL_STACK_MAX) { + row_buf = (int64_t *)malloc((size_t)num_cols * sizeof(*row_buf)); + if (!row_buf) { + free(mask); + return ENOMEM; + } + } + for (uint32_t di = 0; di < num_rows; di++) { + const int64_t *del = data + (size_t)di * num_cols; + for (uint32_t ri = 0; ri < seed->nrows; ri++) { + if (mask[ri]) + continue; + col_rel_row_copy_out(seed, ri, row_buf); + if (memcmp(row_buf, del, (size_t)num_cols * sizeof(*del)) == 0) { + mask[ri] = 1; + any = true; + break; + } + } + } + if (row_buf != row_stack) + free(row_buf); + if (!any) { + free(mask); + return 0; + } + *mask_out = mask; + return 0; +} + +/* Apply a planned removal. No allocation or fallible operation. */ +static void +session_seed_shadow_apply_remove(col_rel_t *seed, uint8_t *mask) +{ + uint32_t out_r = 0; + + if (!seed || !mask) + return; + for (uint32_t ri = 0; ri < seed->nrows; ri++) { + if (mask[ri]) + continue; + if (out_r != ri) + col_rel_row_move_raw(seed, out_r, ri); + out_r++; + } + session_seed_shadow_truncate(seed, out_r); + free(mask); +} + +int +wl_columnar_session_restore_seed(wl_col_session_t *sess, col_rel_t *r) +{ + col_rel_t *seed = r ? session_seed_shadow(sess, r->name) : NULL; + int64_t *rows = NULL; + uint64_t cells = 0; + uint64_t bytes = 0; + bool denied = false; + int rc; + + if (!seed || seed->nrows == 0) + return 0; + if (r->ncols != 0 && r->ncols != seed->ncols) + return EINVAL; + if (!wl_columnar_memory_size_mul(seed->nrows, seed->ncols, &cells) + || !wl_columnar_memory_size_mul(cells, sizeof(int64_t), &bytes) + || bytes > SIZE_MAX) + return EOVERFLOW; + rows = (int64_t *)malloc((size_t)bytes); + if (!rows) + return ENOMEM; + for (uint32_t i = 0; i < seed->nrows; i++) + col_rel_row_copy_out(seed, i, rows + (size_t)i * seed->ncols); + rc = col_rel_append_rows_atomic(r, rows, seed->nrows, seed->ncols, + &denied); + free(rows); + return rc == ENOMEM && denied ? ENOSPC : rc; +} + /* Shared views borrow a root relation's columns. Teardown must release the * aliases before their root; otherwise the root's owner metadata would point * at freed storage. A root owned by another session is left for that owner @@ -3865,11 +4071,23 @@ col_session_insert(wl_session_t *session, const char *relation, return EINVAL; /* column count mismatch */ } + /* Issue #1994: the seed shadow takes the batch first, so a failure on + * either side leaves the two in step. */ + col_rel_t *seed = NULL; + int rc = session_seed_shadow_for_insert(COL_SESSION(session), r, &seed); + if (rc != 0) + return rc; + uint32_t seed_rows = seed ? seed->nrows : 0; bool denied = false; - int rc = col_rel_append_rows_atomic(r, data, num_rows, num_cols, - &denied); + rc = seed ? col_rel_append_rows_atomic(seed, data, num_rows, num_cols, + &denied) : 0; if (rc != 0) return rc == ENOMEM && denied ? ENOSPC : rc; + rc = col_rel_append_rows_atomic(r, data, num_rows, num_cols, &denied); + if (rc != 0) { + session_seed_shadow_truncate(seed, seed_rows); + return rc == ENOMEM && denied ? ENOSPC : rc; + } session_note_inserted_input(sess, r, false); /* The non-incremental API must force a full epoch evaluation even when @@ -4045,11 +4263,23 @@ col_session_insert_incremental(wl_session_t *session, const char *relation, /* Append the complete batch atomically; frontier[] is intentionally NOT * modified by the incremental path. */ + /* Issue #1994: the seed shadow takes the batch first, so a failure on + * either side leaves the two in step. */ + col_rel_t *seed = NULL; + int rc = session_seed_shadow_for_insert(COL_SESSION(session), r, &seed); + if (rc != 0) + return rc; + uint32_t seed_rows = seed ? seed->nrows : 0; bool denied = false; - int rc = col_rel_append_rows_atomic(r, data, num_rows, num_cols, - &denied); + rc = seed ? col_rel_append_rows_atomic(seed, data, num_rows, num_cols, + &denied) : 0; if (rc != 0) return rc == ENOMEM && denied ? ENOSPC : rc; + rc = col_rel_append_rows_atomic(r, data, num_rows, num_cols, &denied); + if (rc != 0) { + session_seed_shadow_truncate(seed, seed_rows); + return rc == ENOMEM && denied ? ENOSPC : rc; + } wl_col_session_t *sess = COL_SESSION(session); session_note_inserted_input(sess, r, true); @@ -4104,6 +4334,8 @@ col_session_remove(wl_session_t *session, const char *relation, return writer_rc; uint8_t *matched = NULL; + uint8_t *seed_mask = NULL; + col_rel_t *seed = NULL; int64_t row_stack[COL_STACK_MAX]; int64_t *row_buf = row_stack; uint32_t remove_count = 0; @@ -4160,6 +4392,14 @@ col_session_remove(wl_session_t *session, const char *relation, goto remove_release; } + /* Issue #1994: the seed shadow records input, so a request removes from + * it even when the relation itself no longer holds the row. */ + seed = session_seed_shadow(sess, r->name); + writer_rc = session_seed_shadow_plan_remove(seed, data, num_rows, + num_cols, &seed_mask); + if (writer_rc != 0) + goto remove_release; + /* No allocation or fallible operation follows the first row write. */ if (remove_count != 0) { uint32_t out_r = 0; @@ -4189,6 +4429,10 @@ col_session_remove(wl_session_t *session, const char *relation, wl_columnar_eval_dedup_set_clear(r); wl_columnar_relation_touch_view(r); } + /* A row only the shadow still holds is already absent from the model, + * so taking it out of the input changes nothing to re-evaluate. */ + session_seed_shadow_apply_remove(seed, seed_mask); + seed_mask = NULL; if (remove_count != 0) { session_invalidate_relation_caches(sess, r->name); sess->pending_input_change = true; @@ -4198,6 +4442,7 @@ col_session_remove(wl_session_t *session, const char *relation, remove_release: if (row_buf != row_stack) free(row_buf); + free(seed_mask); free(matched); if (writer.owner) { int release_rc = wl_columnar_source_access_writer_release(&writer); @@ -4263,6 +4508,8 @@ col_session_remove_incremental(wl_session_t *session, const char *relation, uint32_t source_row; } *match_plan = NULL; uint8_t *matched = NULL; + uint8_t *seed_mask = NULL; + col_rel_t *seed = NULL; col_rel_t *rdelta = NULL; col_rel_t *previous_delta = NULL; wl_columnar_source_access_reader_t previous_reader = { 0 }; @@ -4326,9 +4573,21 @@ col_session_remove_incremental(wl_session_t *session, const char *relation, } } - /* A miss is a true no-op: it must not replace an earlier staged retraction - * or disturb any input/session publication state. */ + /* Issue #1994: the seed shadow records input, so a request removes from + * it even when the relation itself no longer holds the row. */ + seed = session_seed_shadow(sess, r->name); + writer_rc = session_seed_shadow_plan_remove(seed, data, num_rows, + num_cols, &seed_mask); + if (writer_rc != 0) + goto incremental_release; + + /* A miss is a no-op for the model: it must not replace an earlier + * staged retraction or disturb any session publication state. It still + * removes the input the seed shadow holds -- a row the model already + * lacks -- so a later full re-evaluation does not bring it back. */ if (match_count == 0) { + session_seed_shadow_apply_remove(seed, seed_mask); + seed_mask = NULL; writer_rc = 0; goto incremental_release; } @@ -4463,6 +4722,8 @@ col_session_remove_incremental(wl_session_t *session, const char *relation, wl_columnar_relation_touch_view(r); session_rel_registration_commit(sess, rdelta, ®istration); rdelta = NULL; + session_seed_shadow_apply_remove(seed, seed_mask); + seed_mask = NULL; session_invalidate_relation_caches(sess, r->name); sess->last_removed_relation = r->name; sess->outer_epoch++; @@ -4483,6 +4744,7 @@ col_session_remove_incremental(wl_session_t *session, const char *relation, if (row_buf != row_stack) free(row_buf); free(match_plan); + free(seed_mask); free(matched); if (writer.owner) { int release_rc = wl_columnar_source_access_writer_release(&writer); @@ -4591,6 +4853,9 @@ col_session_publication_cutoff(wl_col_session_t *sess) ? wl_evaluation_control_charge(sess->base.evaluation_control, 0) : 0; } +static int +col_session_clear_idb_rows(const wl_plan_t *plan, wl_col_session_t *sess); + static int col_session_step_impl(wl_session_t *session) { @@ -4642,7 +4907,11 @@ col_session_step_impl(wl_session_t *session) } sess->extension_expr_status = 0; - if (!completing_plain_step && !sess->delta_observer && sess->delta_cb + /* A step with nothing pending has nothing to derive, with or without a + * delta callback (Issue #1994). Every write that invalidates the stable + * model also sets pending_input_change, and a session starts with it + * set, so the first step still evaluates. */ + if (!completing_plain_step && !sess->delta_observer && !sess->pending_input_change && sess->last_inserted_relation == NULL && sess->last_removed_relation == NULL){ @@ -4673,9 +4942,34 @@ col_session_step_impl(wl_session_t *session) session, sess->last_removed_relation); } } + /* Issue #1994: a step with no delta callback is not incremental. It + * re-derives every stratum from scratch, so once the session has evaluated + * it first discards what that evaluation derived and resets every stratum + * frontier, under the same has_evaluated condition a snapshot uses before + * a full re-evaluation; the rule frontiers are reset below. Without the + * clear every step appended another copy of each derived row and a + * retraction never reached the rows it had derived; without the stratum + * reset a stale frontier skips iterations a recursive stratum needs after + * the clear. A resumed step (completing_plain_step) keeps what its first + * attempt already merged. */ if (sess->plain_step_completion_step_context - && !completing_plain_step) + && !completing_plain_step) { + affected_mask = UINT64_MAX; + if (sess->has_evaluated) { + /* Until this step commits, the derived relations are empty or + * partial, so whatever evaluates next must clear and evaluate in + * full; only the commit below resets this. */ + sess->pending_full_input_eval = true; + int clear_rc = col_session_clear_idb_rows(plan, sess); + if (clear_rc != 0) + return clear_rc; + for (uint32_t si = 0; si < plan->stratum_count && si < MAX_STRATA; + si++) + sess->frontier_ops->reset_stratum_frontier(sess, si, + sess->outer_epoch); + } sess->plain_step_completion_mask = affected_mask; + } if (!completing_plain_step && sess->delta_observer) { affected_mask = wl_columnar_eval_delta_observer_mask(sess); @@ -4693,8 +4987,11 @@ col_session_step_impl(wl_session_t *session) * set (removal via col_session_remove_incremental), check if $r$ * exists for retractions, and set retraction_seeded. This block does * not narrow affected_mask; the removal's contribution to it is unioned - * in above (Issue #1031). */ - if (sess->last_removed_relation != NULL) { + * in above (Issue #1031). A plain step re-derives from scratch instead + * (Issue #1994): seeding there would make the first iteration read the + * removed rows in place of the relation. */ + if (sess->last_removed_relation != NULL + && !sess->plain_step_completion_step_context) { char rname[256]; if (retraction_rel_name(sess->last_removed_relation, rname, sizeof(rname)) @@ -5031,6 +5328,7 @@ col_session_clear_idb_rows(const wl_plan_t *plan, wl_col_session_t *sess) wl_columnar_source_access_writer_t *writers = NULL; uint32_t target_count = 0; uint32_t owner_count = 0; + uint32_t reset_count = 0; uint32_t capacity = sess ? (sess->nrels ? sess->nrels : 1) : 0; int rc = 0; @@ -5129,6 +5427,7 @@ col_session_clear_idb_rows(const wl_plan_t *plan, wl_col_session_t *sess) rc = col_rel_reset_rows_locked(targets[i], writer); if (rc != 0) goto cleanup; + reset_count = i + 1u; col_session_invalidate_arrangements(&sess->base, targets[i]->name); } @@ -5141,6 +5440,14 @@ col_session_clear_idb_rows(const wl_plan_t *plan, wl_col_session_t *sess) for (uint32_t i = owner_count; i > 0; i--) if (writers[i - 1].owner) (void)wl_columnar_source_access_writer_release(&writers[i - 1]); + /* Issue #1994: a reset relation that also holds input gets its input + * back. This appends under the relation's own writer, so it runs only + * once every writer above is released. */ + for (uint32_t i = 0; i < reset_count; i++) { + int seed_rc = wl_columnar_session_restore_seed(sess, targets[i]); + if (rc == 0 && seed_rc != 0) + rc = seed_rc; + } free((void *)targets); free((void *)owners); free(writers); diff --git a/wirelog/wirelog-easy.h b/wirelog/wirelog-easy.h index 6f3e0e352..b1ffa92da 100644 --- a/wirelog/wirelog-easy.h +++ b/wirelog/wirelog-easy.h @@ -290,6 +290,13 @@ wirelog_easy_remove_sym(wirelog_easy_session_t *s, const char *relation, ...); * * Advance the session by one step, lazily building the plan if needed. * + * With a delta callback installed (wirelog_easy_set_delta_cb()) a step is + * incremental and reports what changed to the callback. Without one a step + * is a full re-evaluation: when anything was inserted or removed since the + * last step or snapshot that evaluated, it discards the derived rows and + * re-derives every rule, keeping the input rows of a relation that is also + * a rule head. Either way a step with nothing pending does no evaluation. + * * Returns: WIRELOG_OK on success, WIRELOG_ERR_EXEC on failure. */ WIRELOG_API wirelog_error_t @@ -351,14 +358,11 @@ wirelog_easy_banner(const char *label); * Take a snapshot of the current state and forward only tuples whose * relation name matches @relation to @cb. * - * IMPORTANT: wirelog_easy_snapshot() is an *evaluating* call — the underlying - * columnar backend re-evaluates every stratum and emits the resulting IDB - * rows. Do NOT call wirelog_easy_step() followed by wirelog_easy_snapshot() on the - * same insert batch: step() already derives and appends the IDB rows, and - * a subsequent snapshot() will re-derive and append again, producing - * duplicated tuples. Choose one mode per batch: - * - Incremental / delta mode: wirelog_easy_set_delta_cb() + wirelog_easy_step() - * - Query mode: wirelog_easy_snapshot() (no prior step) + * wirelog_easy_snapshot() evaluates when input is pending: if anything was + * inserted or removed since the last evaluation, it evaluates before it + * emits. After a completed wirelog_easy_step() with nothing inserted or + * removed since, it emits the model that step left without evaluating + * again, so a snapshot may follow a step on the same session. * * Returns: WIRELOG_OK on success, WIRELOG_ERR_EXEC on failure. */