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
9 changes: 5 additions & 4 deletions docs/THREADING.md
Original file line number Diff line number Diff line change
Expand Up @@ -277,7 +277,7 @@ The 64-byte padding between `tail` and `head`
cache-line ping-pong between producer and consumer collapses
throughput by 2-10x.

### 5.3 Init and relation identity/nonce allocators (8 rows)
### 5.3 Init and relation identity/nonce allocators (9 rows)

| Anchor (`file:function[#N]`) | Field | Op | Order | Justification |
|---|---|---|---|---|
Expand All @@ -289,6 +289,7 @@ throughput by 2-10x.
| `relation.c:col_rel_mutation_set_nonce_allocate` | `wl_next_mutation_set_nonce` | `atomic_load_explicit` | `relaxed` | Read the candidate nonce before the nonwrapping CAS reservation loop; the counter only supplies unique admission provenance and publishes no other state |
| `relation.c:col_rel_mutation_set_nonce_allocate#2` | `wl_next_mutation_set_nonce` | `atomic_compare_exchange_weak_explicit` | `relaxed`/`relaxed` | Reserve a unique nonzero mutation-set nonce and retry with the observed value after a lost race; the RMW provides uniqueness without publishing payload state |
| `relation.c:wl_columnar_relation_test_set_mutation_nonce` | `wl_next_mutation_set_nonce` | `atomic_store_explicit` | `relaxed` | Test-only seam selects allocator exhaustion or retry states before test admissions; tests do not race this reset with nonce allocation |
| `relation.c:wl_columnar_relation_test_mutation_nonce_peek` | `wl_next_mutation_set_nonce` | `atomic_load_explicit` | `relaxed` | Test-only seam reads the next nonce so a single-threaded window can count the mutation sets acquired in it; only the delta matters, and no other state is read through it |

Only the two `io_adapter.c` rows use non-explicit atomic APIs, which default
to `memory_order_seq_cst`. The relation identity and mutation-set nonce
Expand Down Expand Up @@ -466,7 +467,7 @@ measured by `bench/bench_intern.c`; baselines are in `docs/INTERN_PERF.md`

### 5.12 Existing inventory total

21 + 4 + 8 + 19 + 2 + 2 + 3 + 37 + 5 + 7 + 3 = **111 atomic call sites**
21 + 4 + 9 + 19 + 2 + 2 + 3 + 37 + 5 + 7 + 3 = **112 atomic call sites**
before the source-access contract below.

### 5.13 `wirelog/columnar/source_access.h` — relation source gate (21 rows)
Expand Down Expand Up @@ -561,7 +562,7 @@ concurrent alias removals cannot underflow the count.
| `session.c:session_pool_rel_promote#2` | `src->retained_reservation.owner_bits` | `atomic_load_explicit` | acquire | Promote a committed reservation only when the pool slot still owns it |
| `session.c:session_pool_rel_promote#3` | `src->storage_alias_borrows` | `atomic_store_explicit` | relaxed | Leave the closed pool tombstone with no child aliases |

111 + 21 + 30 = **162 atomic call sites**.
112 + 21 + 30 = **163 atomic call sites**.

The `#N` suffix counts all atomic sites in a symbol, regardless of operation;
the first site remains unsuffixed. `scripts/ci/check-threading-doc.sh` uses
Expand Down Expand Up @@ -686,7 +687,7 @@ the committed token after publication and before a growth transaction.
| `eval_dedup.c:wl_columnar_eval_dedup_test_fail_next_growth_alloc` | test-only fault flag | `atomic_store_explicit` | release | Arm one allocation refusal before a test invokes dedup growth; excluded from the production library |
| `eval_dedup.c:wl_columnar_eval_dedup_set_grow` | test-only fault flag | `atomic_exchange_explicit` | acquire-release | Consume the one-shot fault safely when test workers grow dedup tables; excluded from the production library |

The complete source audit now contains **219 atomic call sites**.
The complete source audit now contains **220 atomic call sites**.

---

Expand Down
2 changes: 1 addition & 1 deletion scripts/ci/check-threading-doc.sh
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,7 @@ fi
rows="$tmp_dir/rows"
sed -nE 's/^\| `([^`]+:[A-Za-z_][A-Za-z0-9_]*(#[0-9]+)?)` \| [^|]* \| `([^`]*)` \|.*/\1\t\3/p' "$doc" >"$rows"
row_count=$(wc -l <"$rows")
expected_rows="${WIRELOG_THREADING_EXPECTED_ROWS:-219}"
expected_rows="${WIRELOG_THREADING_EXPECTED_ROWS:-220}"
[ "$row_count" -eq "$expected_rows" ] || {
echo "check-threading-doc: FAIL: expected $expected_rows audit rows, found $row_count" >&2
exit 1
Expand Down
215 changes: 215 additions & 0 deletions tests/test_relation_mutation_set.c
Original file line number Diff line number Diff line change
Expand Up @@ -562,6 +562,218 @@ invalid_tests(void)
assert(col_rel_storage_alias_borrow_count(&root) == 1);
}

/* A held-lease batch appender (main CI regression after #2037): one
* mutation set per storage transition instead of one per row, with the
* same rows, capacities and generations as per-row col_rel_append_row. */
static col_rel_t *
batch_relation(const char *name)
{
static const wirelog_column_type_t types[2] = {
WIRELOG_TYPE_INT64, WIRELOG_TYPE_FLOAT
};
col_rel_t *r = col_rel_new_auto(name, 2);
assert(r);
assert(col_rel_set_column_types(r, types, 2) == 0);
assert(col_rel_enable_timestamps(r) == 0);
return r;
}

static void
batch_row_values(uint32_t i, int64_t row[2])
{
double value = (i % 7u == 0u) ? -0.0 : (double)i;
row[0] = (int64_t)i;
memcpy(&row[1], &value, sizeof(value));
}

static void
append_batch_tests(void)
{
enum { ROWS = 1000 };
col_rel_t *ref = batch_relation("batch_ref");
col_rel_t *out = batch_relation("batch_out");
int64_t row[2];
uint64_t ref_storage_before = ref->storage_generation;
uint64_t ref_view_before = ref->view_generation;
for (uint32_t i = 0; i < ROWS; i++) {
batch_row_values(i, row);
assert(col_rel_append_row(ref, row) == 0);
}

wl_columnar_relation_test_set_mutation_nonce(1000);
uint64_t storage_before = out->storage_generation;
uint64_t view_before = out->view_generation;
col_rel_append_batch_t batch;
memset(&batch, 0, sizeof(batch));
assert(col_rel_append_batch_row(&batch, row) == EINVAL);
assert(col_rel_append_batch_end(&batch, true) == 0);
uint64_t nonce_before = wl_columnar_relation_test_mutation_nonce_peek();
assert(col_rel_append_batch_begin(&batch, out) == 0);
assert(col_rel_append_batch_begin(&batch, out) == EINVAL);
for (uint32_t i = 0; i < ROWS; i++) {
batch_row_values(i, row);
assert(col_rel_append_batch_row(&batch, row) == 0);
if (i == 10) {
/* The held gates exclude readers until the batch ends. */
wl_columnar_source_access_reader_t reader = {0};
assert(col_rel_source_reader_acquire(out, &reader) == EBUSY);
}
}
assert(col_rel_append_batch_end(&batch, true) == 0);
uint64_t acquisitions
= wl_columnar_relation_test_mutation_nonce_peek() - nonce_before;
assert(col_rel_append_batch_end(&batch, true) == 0);
assert(col_rel_append_batch_row(&batch, row) == EINVAL);

/* Same storage shape and publications as the per-row reference. */
assert(out->nrows == ref->nrows && out->capacity == ref->capacity
&& out->timestamp_capacity == ref->timestamp_capacity);
for (uint32_t c = 0; c < 2; c++)
assert(memcmp(out->columns[c], ref->columns[c],
ROWS * sizeof(int64_t)) == 0);
uint64_t transitions = out->storage_generation - storage_before;
assert(transitions == ref->storage_generation - ref_storage_before);
assert(transitions > 1);
assert(out->view_generation - view_before
== ref->view_generation - ref_view_before);
/* One set per storage epoch: at least one, never more than the
* transitions plus the first, and far below one per row. */
assert(acquisitions >= 1 && acquisitions <= transitions + 1);

wl_columnar_source_access_reader_t reader = {0};
assert(col_rel_source_reader_acquire(out, &reader) == 0);
assert(col_rel_source_reader_release(&reader) == 0);
assert(col_rel_destroy_checked(out) == 0);
assert(col_rel_destroy_checked(ref) == 0);
}

/* The filter operator appends its selected rows under the batch appender:
* a right-side filter of 4096 timestamped rows acquires a handful of
* mutation sets (one per storage epoch plus construction), not one per
* selected row, and produces exactly the per-row result. */
static uint8_t batch_filter_simple[] = {
WL_PLAN_EXPR_VAR, 4, 0, 'c', 'o', 'l', '1',
WL_PLAN_EXPR_CONST_INT, 150, 0, 0, 0, 0, 0, 0, 0,
WL_PLAN_EXPR_CMP_GT
};
static uint8_t batch_filter_compiled[] = {
WL_PLAN_EXPR_VAR, 4, 0, 'c', 'o', 'l', '1',
WL_PLAN_EXPR_CONST_INT, 1, 0, 0, 0, 0, 0, 0, 0,
WL_PLAN_EXPR_ARITH_ADD,
WL_PLAN_EXPR_CONST_INT, 151, 0, 0, 0, 0, 0, 0, 0,
WL_PLAN_EXPR_CMP_GT
};

static void
filter_batch_tests(void)
{
enum { ROWS = 4096, KEPT = ROWS - 151 };
col_rel_t *src = col_rel_new_auto("filter_src", 2);
delta_pool_t *pool = delta_pool_create(32, sizeof(col_rel_t), 4096);
assert(src && pool && col_rel_enable_timestamps(src) == 0);
for (uint32_t i = 0; i < ROWS; i++) {
int64_t row[2] = { (int64_t)(i * 3u), (int64_t)i };
assert(col_rel_append_row(src, row) == 0);
src->timestamps[i] = (col_delta_timestamp_t){
.iteration = i, .stratum = 7u, .worker = i % 5u,
.multiplicity = (int64_t)(i % 3u) + 1
};
}
wl_plan_expr_buffer_t exprs[2] = {
{ batch_filter_simple, sizeof(batch_filter_simple) },
{ batch_filter_compiled, sizeof(batch_filter_compiled) }
};
for (int e = 0; e < 2; e++) {
wl_columnar_relation_test_set_mutation_nonce(1000);
uint64_t before = wl_columnar_relation_test_mutation_nonce_peek();
col_rel_t *out = wl_columnar_filter_apply_right_filter(&exprs[e],
src, pool, NULL);
uint64_t acquisitions
= wl_columnar_relation_test_mutation_nonce_peek() - before;
assert(out && out->nrows == KEPT && out->timestamps);
/* Per-row acquisition would be at least KEPT; growing 64 -> 4096
* takes six storage epochs. */
assert(acquisitions >= 1 && acquisitions <= 16);
for (uint32_t i = 0; i < KEPT; i++) {
uint32_t s = i + 151u;
assert(out->columns[0][i] == (int64_t)(s * 3u)
&& out->columns[1][i] == (int64_t)s
&& out->timestamps[i].iteration == s
&& out->timestamps[i].stratum == 7u
&& out->timestamps[i].worker == s % 5u
&& out->timestamps[i].multiplicity == (int64_t)(s % 3u) + 1);
}
/* The batch has ended: the result is readable and destroyable. */
wl_columnar_source_access_reader_t reader = {0};
assert(col_rel_source_reader_acquire(out, &reader) == 0);
assert(col_rel_source_reader_release(&reader) == 0);
assert(col_rel_destroy_checked(out) == 0);
}
delta_pool_destroy(pool);
assert(col_rel_destroy_checked(src) == 0);
}

/* The filter operator itself, on its three fills: timestamped simple
* predicate, untimestamped bulk copy, and the compiled slow path. These are
* the paths the evaluator runs; each must take a handful of mutation sets,
* not one per selected row. */
static void
filter_op_batch_tests(void)
{
enum { ROWS = 4096, KEPT = ROWS - 151 };
struct {
bool timestamped;
uint8_t *expr;
uint32_t size;
} cases[] = {
{ true, batch_filter_simple, sizeof(batch_filter_simple) },
{ false, batch_filter_simple, sizeof(batch_filter_simple) },
{ true, batch_filter_compiled, sizeof(batch_filter_compiled) },
};
for (size_t k = 0; k < sizeof(cases) / sizeof(cases[0]); k++) {
col_rel_t *src = col_rel_new_auto("filter_op_src", 2);
assert(src);
if (cases[k].timestamped)
assert(col_rel_enable_timestamps(src) == 0);
for (uint32_t i = 0; i < ROWS; i++) {
int64_t row[2] = { (int64_t)(i * 3u), (int64_t)i };
assert(col_rel_append_row(src, row) == 0);
if (cases[k].timestamped)
src->timestamps[i] = (col_delta_timestamp_t){
.iteration = i, .stratum = 9u, .worker = 1u,
.multiplicity = 1
};
}
wl_plan_op_t op;
memset(&op, 0, sizeof(op));
op.filter_expr.data = cases[k].expr;
op.filter_expr.size = cases[k].size;
eval_stack_t stack;
eval_stack_init(&stack);
assert(eval_stack_push(&stack, src, true) == 0);
wl_columnar_relation_test_set_mutation_nonce(1000);
uint64_t before = wl_columnar_relation_test_mutation_nonce_peek();
assert(wl_columnar_filter_op(&op, &stack, &(wl_col_session_t){ 0 })
== 0);
uint64_t acquisitions
= wl_columnar_relation_test_mutation_nonce_peek() - before;
eval_entry_t result = eval_stack_pop(&stack);
col_rel_t *out = result.rel;
assert(out && out->nrows == KEPT
&& (out->timestamps != NULL) == cases[k].timestamped);
assert(acquisitions >= 1 && acquisitions <= 16);
for (uint32_t i = 0; i < KEPT; i++) {
uint32_t s = i + 151u;
assert(out->columns[0][i] == (int64_t)(s * 3u)
&& out->columns[1][i] == (int64_t)s);
if (cases[k].timestamped)
assert(out->timestamps[i].iteration == s
&& out->timestamps[i].stratum == 9u);
}
assert(col_rel_destroy_checked(out) == 0);
}
}

int
main(void)
{
Expand All @@ -572,6 +784,9 @@ main(void)
rollback_tests();
invalid_tests();
metadata_final_access_test();
append_batch_tests();
filter_batch_tests();
filter_op_batch_tests();
{
col_rel_t root;
fixture_t f = {0};
Expand Down
Loading
Loading