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: 9 additions & 0 deletions crates/omnigraph-core/src/branch_control.rs
Original file line number Diff line number Diff line change
Expand Up @@ -776,6 +776,14 @@ decide_seam! {
pub static BRANCH_CREATE_POST_NATIVE = ("branch_create.post_native", BranchCreate, [Fail]);
}

decide_seam! {
/// After the namespace inventory (refs listed, path collision checked,
/// a ref-less tree reclaimed), before the native create. No CAS covers
/// this window: the caller's exclusive schema permit is what keeps a
/// sibling create from passing its own inventory meanwhile.
pub static BRANCH_CREATE_POST_INVENTORY_PRE_NATIVE = ("branch_create.post_inventory_pre_native", BranchCreate, [Fail]);
}

/// Archived legacy ancestors still own their native path, even without refs.
/// Flat generated names need no archive probe; slash names cost one per ancestor.
async fn refuse_archived_path_ancestor(dataset: &Dataset, branch: &str) -> Result<()> {
Expand Down Expand Up @@ -848,6 +856,7 @@ pub async fn create_branch_recoverably(
{
return Err(authority_appeared_after_absence(source, branch).await?);
}
fail(&BRANCH_CREATE_POST_INVENTORY_PRE_NATIVE)?;

for attempt in 0..2 {
let native_error = match crate::lance_clone::create_branch(source, branch, source_version)
Expand Down
191 changes: 182 additions & 9 deletions crates/omnigraph-dst/src/concurrent.rs

Large diffs are not rendered by default.

4 changes: 3 additions & 1 deletion crates/omnigraph-dst/src/lance_faults.rs
Original file line number Diff line number Diff line change
Expand Up @@ -256,7 +256,8 @@ fn seam_scheduler() -> Option<(Arc<crate::concurrent::SeamScheduler>, usize)> {

/// Thread-name attribution: the actors' OS threads carry their
/// identities — `dst-writer-N` (scheduler id N), `dst-branch-actor` (id =
/// writers), `dst-maintenance` (id = writers+1). Lance-realm calls executed
/// writers), `dst-maintenance` (id = writers+1), `dst-schema-actor` (id =
/// writers+2). Lance-realm calls executed
/// INLINE on an actor's thread inherit its name and take turns; calls from
/// Lance's own pool threads (lance-cpu, lance-io) carry other names and run
/// UNGATED — the measured coverage gap (`note_unattributed`), never a
Expand All @@ -270,6 +271,7 @@ fn actor_from_thread(writers: usize) -> Option<usize> {
match name {
"dst-branch-actor" => Some(writers),
"dst-maintenance" => Some(writers + 1),
"dst-schema-actor" => Some(writers + 2),
_ => None,
}
}
Expand Down
153 changes: 152 additions & 1 deletion crates/omnigraph-dst/tests/scenarios.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4581,6 +4581,7 @@ fn dst_concurrent_two_writers_first_contact() {
writers: 2,
ops_per_writer: 12,
maintenance_ops: 0,
schema_ops: 0,
kill_writer: None,
branch_cycles: 0,
readers: 0,
Expand Down Expand Up @@ -4630,6 +4631,7 @@ fn dst_concurrent_contention_hunt() {
writers: 4,
ops_per_writer: 20,
maintenance_ops: 0,
schema_ops: 0,
kill_writer: None,
branch_cycles: 0,
readers: 0,
Expand Down Expand Up @@ -4669,6 +4671,7 @@ fn dst_maintenance_actor_first_contact() {
writers: 2,
ops_per_writer: 12,
maintenance_ops: 8,
schema_ops: 0,
kill_writer: None,
branch_cycles: 0,
readers: 0,
Expand Down Expand Up @@ -4698,6 +4701,128 @@ fn dst_maintenance_actor_first_contact() {
}
}

/// Schema apply racing live writers — the deterministic coverage the
/// shared/exclusive schema gate ships with (RFC 2026-09-18-shared-schema-gate;
/// before it, writers-racing-apply had NO DST coverage). A dedicated schema
/// actor performs monotone additive applies (each takes the gate's EXCLUSIVE
/// side, draining every writer's shared permit) while two data writers race
/// under shared permits. Oracles: every writer claim commits (no wedge under
/// schema contention), every apply commits (writers cannot starve the
/// exclusive side), each apply lands exactly one empty-person-diff era
/// commit, and — in every seed — the writers genuinely interleaved
/// (`alternations` ≥ 1) and at least one apply landed between two writer
/// commits (`era_commits_between_data` ≥ 1); otherwise the green is
/// vacuous. Plain mode only:
/// an apply's table rewrite runs on the single `lance-cpu` pool thread,
/// which the seam arbiter deliberately cannot see, so under the seam
/// scheduler its stall budget trips on a loaded machine (measured: 0, 4,
/// 12, 18 or 27 escapes across runs of one seed). The strict-replay claim
/// for this arm (`sched_escapes == 0`) is therefore the hunt's, run
/// explicitly on an idle machine; the permits' turn/epoch protocol itself
/// is pinned by `dst_seam_scheduler_bite_and_replay`.
#[test]
#[serial]
fn dst_schema_apply_racing_writers_first_contact() {
use omnigraph_dst::concurrent::{ConcurrentScenario, run_concurrent_universe};
for seed in dst_seeds(&[24_301, 24_302, 24_303]) {
let root = format!("shared-memory://dst-s24-schema-{seed}");
let sc = ConcurrentScenario {
seed,
writers: 2,
ops_per_writer: 12,
maintenance_ops: 0,
schema_ops: 3,
kill_writer: None,
branch_cycles: 0,
readers: 1,
writer_fault_pct: 0,
seam_schedule: false,
park_deleter_hold: false,
};
let report = run_concurrent_universe(&root, &sc);
assert_eq!(
report.committed, 24,
"every data write must commit despite schema contention"
);
assert_eq!(
report.schema_committed, 3,
"every schema apply must commit; writers cannot starve the exclusive side"
);
assert_eq!(
report.maintenance_commits, 3,
"each apply lands exactly one empty-person-diff era commit"
);
assert!(
report.alternations >= 1,
"seed {seed}: the writers never interleaved — a vacuous green for the \
concurrency claim"
);
assert!(
report.era_commits_between_data >= 1,
"seed {seed}: no schema apply landed between two writer commits — the \
applies never contended with the writers"
);
println!(
"dst s24 schema [seed={seed}]: committed={} occ_retries={} \
schema(committed={} retries={}) alternations={} applies_between_writes={}",
report.committed,
report.occ_retries,
report.schema_committed,
report.schema_retries,
report.alternations,
report.era_commits_between_data
);
}
}

/// The schema-arm hunt instrument: wider seeds, scheduler on, faults on —
/// run explicitly when hunting interleavings around the shared/exclusive
/// boundary (`OMNIGRAPH_DST_SEEDS` widens the search). No readers: a reader
/// opens read-only handles, which take the schema gate's exclusive side with
/// no arbiter hook, so their gate transitions would fall outside the turns
/// that `sched_escapes == 0` certifies.
#[test]
#[serial]
#[ignore = "hunt: schema-apply-vs-writers interleaving search — run explicitly"]
fn dst_schema_apply_racing_writers_hunt() {
use omnigraph_dst::concurrent::{ConcurrentScenario, run_concurrent_universe};
// 24_304 is the first-contact scenario's own shape under the scheduler
// (see the pin for why strict replay is a hunt claim for this arm).
for seed in dst_seeds(&[
24_304, 24_310, 24_311, 24_312, 24_313, 24_314, 24_315, 24_316, 24_317,
]) {
let root = format!("shared-memory://dst-s24-schema-hunt-{seed}");
let sc = ConcurrentScenario {
seed,
writers: 3,
ops_per_writer: 10,
maintenance_ops: 0,
schema_ops: 4,
kill_writer: None,
branch_cycles: 0,
readers: 0,
writer_fault_pct: 10,
seam_schedule: true,
park_deleter_hold: false,
};
let report = run_concurrent_universe(&root, &sc);
assert_eq!(report.committed, 30);
assert_eq!(report.schema_committed, 4);
assert_eq!(report.sched_escapes, 0, "strict replay must hold");
println!(
"dst s24 schema-hunt [seed={seed}]: committed={} schema(committed={} \
retries={}) faults={} alternations={} sched(turns={} escapes={})",
report.committed,
report.schema_committed,
report.schema_retries,
report.writer_faults_injected,
report.alternations,
report.sched_turns,
report.sched_escapes
);
}
}

/// ARM 2 — crash one writer mid-op while the other keeps racing:
/// writer 0's adapter-realm storage dies at its k-th write-class call
/// (post-mortem refusal, no revive — the one-participant process-death
Expand All @@ -4719,6 +4844,7 @@ fn dst_crash_one_writer_first_contact() {
writers: 2,
ops_per_writer: 12,
maintenance_ops: 0,
schema_ops: 0,
kill_writer: Some((0, kill_at)),
branch_cycles: 0,
readers: 0,
Expand Down Expand Up @@ -4772,6 +4898,7 @@ fn dst_branch_actor_first_contact() {
writers: 2,
ops_per_writer: 12,
maintenance_ops: 0,
schema_ops: 0,
kill_writer: None,
branch_cycles: 4,
readers: 0,
Expand Down Expand Up @@ -4821,6 +4948,7 @@ fn dst_concurrent_fleet() {
writers: 3,
ops_per_writer: 10,
maintenance_ops: 0,
schema_ops: 0,
kill_writer: None,
branch_cycles: 0,
// Readers in EVERY fleet arm — live differential reads during
Expand All @@ -4830,8 +4958,26 @@ fn dst_concurrent_fleet() {
seam_schedule: seam,
park_deleter_hold: false,
};
let arms: [(&str, ConcurrentScenario); 6] = [
let arms: [(&str, ConcurrentScenario); 8] = [
("race", base.clone()),
// Schema arms: plain mode only (the first-contact pin says why).
(
"schema",
ConcurrentScenario {
schema_ops: 3,
seam_schedule: false,
..base.clone()
},
),
(
"schema+maint",
ConcurrentScenario {
schema_ops: 3,
maintenance_ops: 4,
seam_schedule: false,
..base.clone()
},
),
(
"maint",
ConcurrentScenario {
Expand All @@ -4842,6 +4988,7 @@ fn dst_concurrent_fleet() {
(
"crash",
ConcurrentScenario {
schema_ops: 0,
kill_writer: Some((0, 7 + (seed as usize % 17))),
..base.clone()
},
Expand Down Expand Up @@ -4915,6 +5062,7 @@ fn dst_seam_scheduler_bite_and_replay() {
writers: 2,
ops_per_writer: 8,
maintenance_ops: 0,
schema_ops: 0,
kill_writer: None,
branch_cycles: 0,
readers: 0,
Expand Down Expand Up @@ -5016,6 +5164,7 @@ fn dst_optimize_races_branch_delete() {
writers: 3,
ops_per_writer: 10,
maintenance_ops: 4,
schema_ops: 0,
kill_writer: None,
branch_cycles: 3,
readers: 0,
Expand Down Expand Up @@ -5084,6 +5233,7 @@ fn dst_optimize_races_branch_delete_seed_search() {
writers: 2,
ops_per_writer: 6,
maintenance_ops: 0,
schema_ops: 0,
kill_writer: None,
branch_cycles: 3,
readers: 0,
Expand Down Expand Up @@ -5180,6 +5330,7 @@ fn dst_optimize_races_branch_delete_directed_hold() {
writers: 2,
ops_per_writer: 6,
maintenance_ops: 0,
schema_ops: 0,
kill_writer: None,
branch_cycles: 3,
readers: 0,
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,36 @@
# issue: none
--- runner
timeout_ms: 10000
environments:
- target: omnigraph-engine-dst
storage: in-memory-object-store
seeds: [0, 42]

--- schema
node Person {
name: String @key
}
--- seed
{"type":"Person","data":{"name":"alice"}}
--- mutate
branch create feature
--- expect ok
--- concurrent
w1: query add_bob() { insert Person { name: "bob" } }
w2: query add_carol() { insert Person { name: "carol" } }
w3 on feature: query add_dave() { insert Person { name: "dave" } }
r1: query all() { match { $p: Person } return { $p.name } }
order: w2 park put nodes/, w3 park put _versions/, w1, w2 put nodes/, r1 start, r1, w2, w3 put _versions/, w3
--- expect
w1: ok
w2: ok
w3: ok
r1: ok
--- query
query all() { match { $p: Person } return { $p.name } }
--- expect unordered
{"p.name":"alice"}
{"p.name":"bob"}
{"p.name":"carol"}
--- expect shape
p.name: String
5 changes: 3 additions & 2 deletions crates/omnigraph/benches/scenarios/concurrent_writes.rs
Original file line number Diff line number Diff line change
Expand Up @@ -14,8 +14,9 @@
//!
//! Workload: insert-only `Chunk` rows with disjoint keys per worker
//! (`cw-w{worker}-{seq}`), the shape that exercises the write path's real
//! serialization — the process-global write queue and the exclusive schema
//! gate every writer crosses in `commit_all` — without manufacturing key
//! serialization — the process-global write queue and the schema gate
//! every writer crosses in `commit_all` (shared since RFC
//! 2026-09-18-shared-schema-gate) — without manufacturing key
//! conflicts. `Omnigraph::mutate` replays a typed read-set/authority
//! conflict (`ReadSetChanged`: the graph head moved under a concurrent
//! writer) itself, up to `MAX_PRE_EFFECT_REPREPARES` times for an
Expand Down
Loading
Loading