diff --git a/docs/THREADING.md b/docs/THREADING.md index c014a4d9..475389f2 100644 --- a/docs/THREADING.md +++ b/docs/THREADING.md @@ -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 | |---|---|---|---|---| @@ -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 @@ -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) @@ -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 @@ -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**. --- diff --git a/scripts/ci/check-threading-doc.sh b/scripts/ci/check-threading-doc.sh index d2b84c6f..c8d45d8d 100755 --- a/scripts/ci/check-threading-doc.sh +++ b/scripts/ci/check-threading-doc.sh @@ -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 diff --git a/tests/test_relation_mutation_set.c b/tests/test_relation_mutation_set.c index 50872ef6..1c7aa60d 100644 --- a/tests/test_relation_mutation_set.c +++ b/tests/test_relation_mutation_set.c @@ -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) { @@ -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}; diff --git a/wirelog/columnar/filter.c b/wirelog/columnar/filter.c index df94cc21..cd9676ad 100644 --- a/wirelog/columnar/filter.c +++ b/wirelog/columnar/filter.c @@ -602,21 +602,37 @@ wl_columnar_filter_status(wl_col_session_t *sess, col_rel_t *relation, return rc; } +/* Appends through @batch, which holds out's mutation lease for the whole + * fill: a lease per appended row made the filter operator's cost dominated + * by lease acquisition (main CI timeouts after #2037). The timestamp store + * below therefore also happens under the held writer gates. */ static int -wl_columnar_filter_append_selected(col_rel_t *out, const col_rel_t *src, - uint32_t src_row, const int64_t *row, wl_col_session_t *sess) +wl_columnar_filter_append_selected(col_rel_append_batch_t *batch, + col_rel_t *out, const col_rel_t *src, uint32_t src_row, + const int64_t *row, wl_col_session_t *sess) { - if (!out || !src || !row || src_row >= src->nrows + if (!batch || !out || !src || !row || src_row >= src->nrows || (src->timestamps && !out->timestamps)) return EINVAL; uint32_t dst_row = out->nrows; int rc = wl_columnar_filter_status(sess, out, - col_rel_append_row(out, row)); + col_rel_append_batch_row(batch, row)); if (rc == 0 && src->timestamps) out->timestamps[dst_row] = src->timestamps[src_row]; return rc; } +/* Every exit from a fill region ends its batch before out is destroyed or + * published: destroying a relation whose writer gates are still held fails, + * and the destroy result is ignored on these paths. */ +static void +wl_columnar_filter_discard_batch_out(col_rel_append_batch_t *batch, + col_rel_t *out) +{ + (void)col_rel_append_batch_end(batch, false); + (void)col_rel_destroy_checked(out); +} + static int wl_columnar_filter_dispose_input(eval_stack_t *stack, eval_entry_t *entry, int primary_rc) @@ -758,6 +774,15 @@ wl_columnar_filter_op(const wl_plan_op_t *op, eval_stack_t *stack, #ifdef __AVX2__ use_selection = true; #endif + col_rel_append_batch_t batch = { 0 }; + int batch_rc = wl_columnar_filter_status(sess, out, + col_rel_append_batch_begin(&batch, out)); + if (batch_rc != 0) { + free(row); + wl_columnar_filter_scratch_release(&scratch); + (void)col_rel_destroy_checked(out); + return wl_columnar_filter_dispose_input(stack, &e, batch_rc); + } if (use_selection) { uint32_t *sel = (uint32_t *)malloc( ((size_t)COL_FILTER_TILE + COL_FILTER_SEL_SLACK) @@ -765,7 +790,7 @@ wl_columnar_filter_op(const wl_plan_op_t *op, eval_stack_t *stack, if (!sel) { free(row); wl_columnar_filter_scratch_release(&scratch); - (void)col_rel_destroy_checked(out); + wl_columnar_filter_discard_batch_out(&batch, out); return wl_columnar_filter_dispose_input(stack, &e, ENOMEM); } @@ -784,13 +809,13 @@ wl_columnar_filter_op(const wl_plan_op_t *op, eval_stack_t *stack, uint32_t src_row = base + sel[i]; for (uint32_t c = 0; c < ncols; c++) row[c] = columns[c][src_row]; - int rc = wl_columnar_filter_append_selected(out, - e.rel, src_row, row, sess); + int rc = wl_columnar_filter_append_selected(&batch, + out, e.rel, src_row, row, sess); if (rc != 0) { free(sel); free(row); wl_columnar_filter_scratch_release(&scratch); - (void)col_rel_destroy_checked(out); + wl_columnar_filter_discard_batch_out(&batch, out); return wl_columnar_filter_dispose_input(stack, &e, rc); } @@ -807,18 +832,23 @@ wl_columnar_filter_op(const wl_plan_op_t *op, eval_stack_t *stack, continue; for (uint32_t c = 0; c < ncols; c++) row[c] = columns[c][r]; - int rc = wl_columnar_filter_append_selected(out, e.rel, r, - row, sess); + int rc = wl_columnar_filter_append_selected(&batch, out, + e.rel, r, row, sess); if (rc != 0) { free(row); wl_columnar_filter_scratch_release(&scratch); - (void)col_rel_destroy_checked(out); + wl_columnar_filter_discard_batch_out(&batch, out); return wl_columnar_filter_dispose_input(stack, &e, rc); } } } free(row); wl_columnar_filter_scratch_release(&scratch); + batch_rc = col_rel_append_batch_end(&batch, true); + if (batch_rc != 0) { + (void)col_rel_destroy_checked(out); + return wl_columnar_filter_dispose_input(stack, &e, batch_rc); + } int cleanup_rc = wl_columnar_filter_dispose_input(stack, &e, 0); if (cleanup_rc != 0) { (void)col_rel_destroy_checked(out); @@ -950,15 +980,21 @@ wl_columnar_filter_op(const wl_plan_op_t *op, eval_stack_t *stack, } } - /* Bulk-copy the passing rows into the output relation */ - for (uint32_t r = 0; r < nout; r++) { - int rc = col_rel_append_row(out, tmp + (size_t)r * ncols); - if (rc != 0) { - free(tmp); - wl_columnar_filter_scratch_release(&scratch); - col_rel_destroy(out); - return wl_columnar_filter_dispose_input(stack, &e, rc); - } + /* Bulk-copy the passing rows into the output relation, under one + * lease per storage epoch rather than one per row. */ + col_rel_append_batch_t batch = { 0 }; + int batch_rc = col_rel_append_batch_begin(&batch, out); + for (uint32_t r = 0; batch_rc == 0 && r < nout; r++) + batch_rc = col_rel_append_batch_row(&batch, + tmp + (size_t)r * ncols); + if (batch_rc == 0) + batch_rc = col_rel_append_batch_end(&batch, true); + if (batch_rc != 0) { + free(tmp); + wl_columnar_filter_scratch_release(&scratch); + (void)col_rel_append_batch_end(&batch, false); + col_rel_destroy(out); + return wl_columnar_filter_dispose_input(stack, &e, batch_rc); } free(tmp); wl_columnar_filter_scratch_release(&scratch); @@ -1008,6 +1044,16 @@ wl_columnar_filter_op(const wl_plan_op_t *op, eval_stack_t *stack, } int64_t *const row = row_rb.ptr; + col_rel_append_batch_t batch = { 0 }; + int batch_rc = wl_columnar_filter_status(sess, out, + col_rel_append_batch_begin(&batch, out)); + if (batch_rc != 0) { + col_row_buf_release(&row_rb); + wl_columnar_filter_scratch_release(&row_scratch); + wl_columnar_expr_compiled_free(ce); + (void)col_rel_destroy_checked(out); + return wl_columnar_filter_dispose_input(stack, &e, batch_rc); + } wl_columnar_expr_context_t expr_ctx = { .intern = sess ? sess->intern : NULL, .extensions = sess ? sess->base.extension_snapshot : NULL, @@ -1039,7 +1085,7 @@ wl_columnar_filter_op(const wl_plan_op_t *op, eval_stack_t *stack, col_row_buf_release(&row_rb); wl_columnar_filter_scratch_release(&row_scratch); wl_columnar_expr_compiled_free(ce); - (void)col_rel_destroy_checked(out); + wl_columnar_filter_discard_batch_out(&batch, out); /* A refused string allocation fails the step as a memory * error (Issue #1470). */ return wl_columnar_filter_dispose_input(stack, &e, @@ -1049,13 +1095,13 @@ wl_columnar_filter_op(const wl_plan_op_t *op, eval_stack_t *stack, pass = err == WL_COLUMNAR_EXPR_OK && val != 0; } if (pass) { - int rc = wl_columnar_filter_append_selected(out, e.rel, r, row, - sess); + int rc = wl_columnar_filter_append_selected(&batch, out, e.rel, + r, row, sess); if (rc != 0) { col_row_buf_release(&row_rb); wl_columnar_filter_scratch_release(&row_scratch); wl_columnar_expr_compiled_free(ce); - (void)col_rel_destroy_checked(out); + wl_columnar_filter_discard_batch_out(&batch, out); return wl_columnar_filter_dispose_input(stack, &e, rc); } } @@ -1063,6 +1109,11 @@ wl_columnar_filter_op(const wl_plan_op_t *op, eval_stack_t *stack, col_row_buf_release(&row_rb); wl_columnar_filter_scratch_release(&row_scratch); wl_columnar_expr_compiled_free(ce); + batch_rc = col_rel_append_batch_end(&batch, true); + if (batch_rc != 0) { + (void)col_rel_destroy_checked(out); + return wl_columnar_filter_dispose_input(stack, &e, batch_rc); + } int cleanup_rc = wl_columnar_filter_dispose_input(stack, &e, 0); if (cleanup_rc != 0) { @@ -1131,22 +1182,20 @@ fill_filtered_rel(const uint8_t *buf, uint32_t bsz, col_rel_t *rel, /* Fast path: simple colA CMP CONST or colA CMP colB predicate */ simple_filter_cmp_t cmp; + col_rel_append_batch_t batch = { 0 }; if (filter_is_simple_cmp(buf, bsz, &cmp)) { - for (uint32_t r = 0; r < rel->nrows; r++) { + int rc = wl_columnar_filter_status(sess, out, + col_rel_append_batch_begin(&batch, out)); + for (uint32_t r = 0; rc == 0 && r < rel->nrows; r++) { col_rel_row_copy_out(rel, r, row_buf); - if (col_filter_cmp_row(row_buf, rel->ncols, &cmp)) { - int rc = wl_columnar_filter_append_selected(out, rel, r, + if (col_filter_cmp_row(row_buf, rel->ncols, &cmp)) + rc = wl_columnar_filter_append_selected(&batch, out, rel, r, row_buf, sess); - if (rc != 0) { - col_row_buf_release(&rb); - wl_columnar_filter_scratch_release(&scratch); - return rc; - } - } } + int end_rc = col_rel_append_batch_end(&batch, rc == 0); col_row_buf_release(&rb); wl_columnar_filter_scratch_release(&scratch); - return 0; + return rc != 0 ? rc : end_rc; } /* Slow path: compile once, evaluate per row */ @@ -1157,6 +1206,14 @@ fill_filtered_rel(const uint8_t *buf, uint32_t bsz, col_rel_t *rel, col_row_buf_release(&rb); return compile_rc; } + int batch_rc = wl_columnar_filter_status(sess, out, + col_rel_append_batch_begin(&batch, out)); + if (batch_rc != 0) { + col_row_buf_release(&rb); + wl_columnar_filter_scratch_release(&scratch); + wl_columnar_expr_compiled_free(ce); + return batch_rc; + } for (uint32_t r = 0; r < rel->nrows; r++) { col_rel_row_copy_out(rel, r, row_buf); int pass; @@ -1174,9 +1231,10 @@ fill_filtered_rel(const uint8_t *buf, uint32_t bsz, col_rel_t *rel, pass = (err == 0) ? (val != 0 ? 1 : 0) : 0; /* fail-closed */ } if (pass) { - int rc = wl_columnar_filter_append_selected(out, rel, r, + int rc = wl_columnar_filter_append_selected(&batch, out, rel, r, row_buf, sess); if (rc != 0) { + (void)col_rel_append_batch_end(&batch, false); col_row_buf_release(&rb); wl_columnar_filter_scratch_release(&scratch); wl_columnar_expr_compiled_free(ce); @@ -1184,10 +1242,11 @@ fill_filtered_rel(const uint8_t *buf, uint32_t bsz, col_rel_t *rel, } } } + batch_rc = col_rel_append_batch_end(&batch, true); col_row_buf_release(&rb); wl_columnar_filter_scratch_release(&scratch); wl_columnar_expr_compiled_free(ce); - return 0; + return batch_rc; } /** diff --git a/wirelog/columnar/internal.h b/wirelog/columnar/internal.h index 1222d931..7133f29a 100644 --- a/wirelog/columnar/internal.h +++ b/wirelog/columnar/internal.h @@ -658,6 +658,16 @@ typedef struct wl_columnar_relation_terminal_sequence { bool consumed; } wl_columnar_relation_terminal_sequence_t; +/* Single-entry storage lives with the operation, including its role mapping. */ +typedef struct { + wl_columnar_relation_mutation_set_t set; + wl_columnar_relation_mutation_role_t role; + wl_columnar_relation_mutation_descriptor_t descriptor; + wl_columnar_relation_mutation_owner_t owners[2]; + wl_columnar_relation_mutation_lease_t lease; + wl_columnar_relation_mutation_initialization_t initialization; +} wl_columnar_relation_mutation_single_t; + int col_rel_mutation_set_acquire(wl_columnar_relation_mutation_set_t *set, const wl_columnar_relation_mutation_role_t *roles, size_t role_count, wl_columnar_relation_mutation_descriptor_t *descriptors, @@ -684,6 +694,7 @@ int wl_columnar_relation_terminal_sequence_finish( wl_columnar_relation_terminal_sequence_t *sequence); #ifdef WL_TEST_MUTATION_SET_HOOK void wl_columnar_relation_test_set_mutation_nonce(uint64_t next_nonce); +uint64_t wl_columnar_relation_test_mutation_nonce_peek(void); bool wl_columnar_relation_test_terminal_sequence_validate( const wl_columnar_relation_terminal_sequence_t *sequence); #endif @@ -2744,6 +2755,33 @@ int col_rel_cow_unshare_with_lease(col_rel_t *r, WL_MUST_CHECK int wl_columnar_relation_privatize_shared_view_with_lease(col_rel_t *r, wl_columnar_relation_mutation_lease_t *lease); +/* Append many rows to one relation under a held mutation lease instead of + * one per row (col_rel_append_row acquires and finishes a mutation set for + * every row). begin acquires the lease; row appends exactly as + * col_rel_append_row would, starting a fresh lease only before a row that + * replaces storage once this lease has already published a transition; end + * finishes it and is a no-op on an inactive batch. + * + * While a batch is active this thread holds the relation's descriptor and + * source writer gates, and for a shared view also its storage owner's + * source writer gate: readers of either get EBUSY, and destroying the + * relation fails, so end the batch before destroying, publishing or + * returning it. + * The struct holds address- and thread-bound tokens: keep it in the frame + * that runs the loop, pass it by pointer, never copy it. A failed row leaves + * the relation as that append's own contract says and the batch active; a + * failed lease renewal leaves it inactive. */ +typedef struct { + wl_columnar_relation_mutation_single_t single; + col_rel_t *relation; + bool active; +} col_rel_append_batch_t; + +int col_rel_append_batch_begin(col_rel_append_batch_t *batch, col_rel_t *r); +int col_rel_append_batch_row(col_rel_append_batch_t *batch, + const int64_t *row); +int col_rel_append_batch_end(col_rel_append_batch_t *batch, bool commit); + int col_rel_append_row_with_lease(col_rel_t *r, const int64_t *row, wl_columnar_relation_mutation_lease_t *lease); int col_rel_reserve_rows_with_lease(col_rel_t *r, uint32_t additional, diff --git a/wirelog/columnar/relation.c b/wirelog/columnar/relation.c index 24410692..3661abb1 100644 --- a/wirelog/columnar/relation.c +++ b/wirelog/columnar/relation.c @@ -34,6 +34,15 @@ wl_columnar_relation_test_set_mutation_nonce(uint64_t next_nonce) atomic_store_explicit(&wl_next_mutation_set_nonce, next_nonce, memory_order_relaxed); } + +/* The next nonce to be handed out: a delta across a single-threaded window + * counts the mutation sets acquired in it. */ +uint64_t +wl_columnar_relation_test_mutation_nonce_peek(void) +{ + return atomic_load_explicit(&wl_next_mutation_set_nonce, + memory_order_relaxed); +} #endif static int @@ -306,16 +315,6 @@ wl_columnar_relation_radix_bench_enabled(void) /* ---- COW helpers --------------------------------------------------------- */ -/* Single-entry storage lives with the operation, including its role mapping. */ -typedef struct { - wl_columnar_relation_mutation_set_t set; - wl_columnar_relation_mutation_role_t role; - wl_columnar_relation_mutation_descriptor_t descriptor; - wl_columnar_relation_mutation_owner_t owners[2]; - wl_columnar_relation_mutation_lease_t lease; - wl_columnar_relation_mutation_initialization_t initialization; -} wl_columnar_relation_mutation_single_t; - static int col_rel_mutation_single_acquire(col_rel_t *relation, wl_columnar_relation_mutation_single_t *single) @@ -4530,6 +4529,17 @@ col_rel_apply_compound_schema(col_rel_t *r, return rc; } +/* Whether appending one row to r replaces its storage: a shared view is + * copied on write, and full column or timestamp capacity grows. The append + * publishes a storage transition exactly when this holds, which is why the + * batch appender below uses it to start a fresh lease before such a row. */ +static bool +col_rel_append_replaces_storage(const col_rel_t *r) +{ + return r->col_shared || r->nrows >= r->capacity + || (r->timestamps && r->nrows >= r->timestamp_capacity); +} + static int col_rel_append_row_impl(col_rel_t *r, const int64_t *row, wl_columnar_source_access_writer_t *writer, bool writer_held, @@ -4572,8 +4582,7 @@ col_rel_append_row_impl(col_rel_t *r, const int64_t *row, : r->storage_generation; if (!col_rel_timestamp_shape_valid(r) || r->nrows > r->capacity) return EINVAL; - bool replaces_storage = r->col_shared || r->nrows >= r->capacity - || (r->timestamps && r->nrows >= r->timestamp_capacity); + bool replaces_storage = col_rel_append_replaces_storage(r); if (r->nrows == UINT32_MAX || r->view_generation >= WL_COLUMNAR_REL_GENERATION_INVALID - 1u || (replaces_storage && r->storage_generation @@ -4840,6 +4849,58 @@ col_rel_append_row_with_lease(col_rel_t *r, const int64_t *row, &lease->set->owners[lease->owner_slot].writer, true, lease); } +int +col_rel_append_batch_begin(col_rel_append_batch_t *batch, col_rel_t *r) +{ + if (!batch || !r || batch->active) + return EINVAL; + memset(batch, 0, sizeof(*batch)); + int rc = col_rel_mutation_single_acquire(r, &batch->single); + if (rc) { + memset(batch, 0, sizeof(*batch)); + return rc; + } + batch->relation = r; + batch->active = true; + return 0; +} + +int +col_rel_append_batch_row(col_rel_append_batch_t *batch, const int64_t *row) +{ + if (!batch || !batch->active || !row) + return EINVAL; + col_rel_t *r = batch->relation; + /* A lease publishes at most one storage transition. Once it has, a row + * that would replace storage again starts the next epoch's lease. */ + if (batch->single.lease.storage_transitioned + && col_rel_append_replaces_storage(r)) { + int rc = col_rel_mutation_set_finish(&batch->single.set, true); + if (rc == 0) { + memset(&batch->single, 0, sizeof(batch->single)); + rc = col_rel_mutation_single_acquire(r, &batch->single); + } + if (rc) { + batch->active = false; + return rc; + } + } + /* Each row's denial evidence starts clean, as for col_rel_append_row. */ + r->memory_budget_denial_pending = false; + return col_rel_append_row_impl(r, row, + &batch->single.owners[batch->single.lease.owner_slot].writer, + true, &batch->single.lease); +} + +int +col_rel_append_batch_end(col_rel_append_batch_t *batch, bool commit) +{ + if (!batch || !batch->active) + return 0; + batch->active = false; + return col_rel_mutation_set_finish(&batch->single.set, commit); +} + static int col_rel_capacity_for_rows(uint32_t current, uint32_t required, uint32_t *out_capacity)