Skip to content

Commit 934e7ea

Browse files
committed
engine: schema apply stages detached rewrites and publishes its contract without a sidecar (RFC 0067 step 4)
Schema apply is the fourth writer moved to detached table commits. An existing-table rewrite is a detached Overwrite of the promoted HEAD, published as a pin one past the published version and promoted after the manifest commit. An added type stays a linear version-one create at its identity path; that path is a deterministic function of the accepted identity allocator (the RFC's "fresh incarnation path" was wrong), so an attempt that died after creating the dataset left it exactly where the retry creates it, and the retry reclaims the unregistered leftover under the schema sentinel instead of refusing to claim unowned physical state. The schema contract is staged with the graph commit that publishes it recorded in `__schema_state.json.staging`, and the writer installs the live contract from memory after the commit. No SchemaApply sidecar is written: a failure before the manifest commit returns the plain error and leaves detached versions, a created dataset and a staged contract that the next read-write open discards; a failure after it reports `RecoveryRequired` naming the published commit, and the next read-write open, or the same handle's next write, installs the contract because that commit is in main's lineage. A read-only open refuses a published but uninstalled contract and serves an unpublished staging as absent. The open also reclaims a sentinel left by a crashed apply, which previously stayed behind whenever an apply died after deleting its sidecar, and the release is idempotent. The RFC 0040 system-column upgrade keeps the v9 exact protocol, its unmarked staging and its sidecar; the schema-apply sidecar constructor is test-only until step 5 removes the classifier. Lance fence: an Overwrite replayed over its own twin is idempotent (Lance recognises the committed transaction and adds no version), so racing promoters of a schema-apply pin leave no residue. Tests: the schema-apply failpoint cells become no-residue, staging discard, staging promotion, leftover reclaim, same-handle heal and manifest-CAS-loss tests; the sidecar-rollback, foreign-winner and Optimize late-sidecar tests whose producer no longer exists are retired; `detached_commit_matrix` gains the `SchemaApply` writer (52 cells in the full run); `write_cost` gains `schema_apply_writes_no_control_object`; the seam `schema_apply.post_publish_pre_promotion` is new and the DST catalog lists 72 windows.
1 parent c122211 commit 934e7ea

19 files changed

Lines changed: 988 additions & 1629 deletions

‎crates/omnigraph-dst/cost_table.txt‎

Lines changed: 5 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -65,12 +65,12 @@ UpdateV a.list calls=4
6565
UpdateV l.get calls=10
6666
UpdateV l.list calls=20
6767
_audit a.delete calls=1
68-
_audit a.exists calls=79
68+
_audit a.exists calls=81
6969
_audit a.get calls=122
7070
_audit a.list calls=4
7171
_audit a.put_if_absent calls=1
72-
_audit l.get calls=1228
73-
_audit l.list calls=65
72+
_audit l.get calls=1230
73+
_audit l.list calls=66
7474
_close a.exists calls=6
7575
_close a.get calls=10
7676
_close a.list calls=2
@@ -87,9 +87,9 @@ _setup a.list calls=3
8787
_setup a.put calls=2
8888
_setup a.put_if_absent calls=4
8989
_setup l.get calls=139
90-
_setup l.list calls=37
90+
_setup l.list calls=38
9191
_setup l.put calls=35
92-
_verify a.exists calls=376
92+
_verify a.exists calls=389
9393
_verify a.get calls=601
9494
_verify a.list calls=13
9595
_verify l.get calls=2069

‎crates/omnigraph-dst/src/catalog.rs‎

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,13 +1,13 @@
11
//! The full crash-window catalog for the hunt
2-
//! (`dst_hunt_crash_window_sweep`): 71 of the engine's decision seams
2+
//! (`dst_hunt_crash_window_sweep`): 72 of the engine's decision seams
33
//! (`omnigraph::seams::catalog`) at the pinned engine version. A seam added
44
//! to the engine enters here as never-reached until its workload exists.
55
//!
66
//! Kept honest by `catalog_names_are_engine_seams` below: every entry must
77
//! be a name the engine catalog declares, so a typo'd or renamed-away window
88
//! fails the suite instead of compiling and silently never firing.
99
10-
pub const CRASH_WINDOWS: [&str; 71] = [
10+
pub const CRASH_WINDOWS: [&str; 72] = [
1111
"blob_read.post_capture",
1212
"branch_control.post_recovery_barrier",
1313
"branch_create.post_native",
@@ -74,6 +74,7 @@ pub const CRASH_WINDOWS: [&str; 71] = [
7474
"schema_apply.after_staging_write",
7575
"schema_apply.before_staging_write",
7676
"schema_apply.post_sidecar_pre_effect",
77+
"schema_apply.post_publish_pre_promotion",
7778
"schema_apply.post_table_commit",
7879
"schema_reload.before_contract_read",
7980
"storage.local_create_if_absent_probe",

‎crates/omnigraph-dst/tests/scenarios.rs‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2214,6 +2214,11 @@ fn dst_milestone_never_remerges_merged_branch() {
22142214
/// The detached index writer opens each productive table at its pin and
22152215
/// promotes the twin it lands (EnsureIndices l.get 11 -> 15; the closing
22162216
/// pass reads the promoted twins, _close l.get 31 -> 35, _verify 2071 -> 2069).
2217+
/// Detached schema apply: a read-write open lists the branch refs once to
2218+
/// reclaim a stale schema-apply sentinel (_setup/_audit l.list 37/65 ->
2219+
/// 38/66, _audit l.get 1228 -> 1230) and a read-only open probes the staged
2220+
/// schema state once for its coherence proof (_audit a.exists 79 -> 81,
2221+
/// _verify 376 -> 389).
22172222
#[test]
22182223
#[serial]
22192224
fn dst_bench_cost_count_golden() {

‎crates/omnigraph/src/db/manifest.rs‎

Lines changed: 3 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -70,10 +70,9 @@ pub(crate) use recovery::{
7070
HealPendingOutcome, MAX_BRANCH_MERGE_DATA_TRANSACTIONS, RecoveryAuthorityToken,
7171
RecoveryLineageIntent, RecoveryManifestDelta, RecoveryMode, RecoverySchemaApplyEffect,
7272
RecoverySchemaApplyEffectKind, RecoverySidecarHandle, RecoverySystemColumnUpgrade,
73-
RecoveryTableUpdateSlot, SidecarKind, SidecarTablePin, SidecarTableRegistration,
74-
SidecarTableRename, SidecarTombstone, confirm_schema_apply_sidecar_v9, delete_sidecar,
75-
ensure_read_only_schema_coherent, heal_pending_sidecars_roll_forward, list_sidecars,
76-
new_optimize_sidecar_v9, new_schema_apply_sidecar_v9, new_system_column_upgrade_sidecar_v9,
73+
RecoveryTableUpdateSlot, SidecarKind, SidecarTablePin, confirm_schema_apply_sidecar_v9,
74+
delete_sidecar, ensure_read_only_schema_coherent, heal_pending_sidecars_roll_forward,
75+
list_sidecars, new_optimize_sidecar_v9, new_system_column_upgrade_sidecar_v9,
7776
recover_manifest_drift, schema_apply_serial_queue_key, write_sidecar,
7877
};
7978
pub use state::DatasetEntry;

‎crates/omnigraph/src/db/manifest/recovery.rs‎

Lines changed: 28 additions & 33 deletions
Original file line numberDiff line numberDiff line change
@@ -43,8 +43,8 @@ use crate::branch_control::list_branch_contents;
4343
use crate::db::graph_coordinator::GraphCoordinator;
4444
use crate::db::recovery_audit::{RecoveryAudit, RecoveryAuditRecord, RecoveryKind, TableOutcome};
4545
use crate::db::schema_state::{
46-
SchemaStateRecovery, read_schema_state_identity, schema_ir_staging_uri,
47-
schema_source_staging_uri, schema_state_staging_uri,
46+
read_schema_state_identity, schema_ir_staging_uri, schema_source_staging_uri,
47+
schema_state_staging_uri,
4848
};
4949
use crate::error::{OmniError, Result};
5050
use crate::storage::StorageAdapter;
@@ -3265,7 +3265,6 @@ pub(crate) async fn heal_pending_sidecars_roll_forward(
32653265
.iter()
32663266
.map(|pin| (pin.table_key.clone(), pin.table_branch.clone()))
32673267
.collect();
3268-
let is_schema_apply = matches!(sidecar.writer_kind, SidecarKind::SchemaApply);
32693268
let _table_guards = write_queue.acquire_many(&queue_keys).await;
32703269
// Re-read after the wait: the writer we blocked on may have completed
32713270
// Phase C and deleted the sidecar, or may have durably confirmed Phase B
@@ -3283,21 +3282,6 @@ pub(crate) async fn heal_pending_sidecars_roll_forward(
32833282
// It also re-runs per sidecar, so a multi-sidecar pass never
32843283
// classifies against a reconcile result an earlier roll-forward
32853284
// staled. Non-SchemaApply sidecars never consult the value.
3286-
let schema_state_recovery = if is_schema_apply {
3287-
let snapshot = {
3288-
let mut coord = coordinator.write().await;
3289-
coord.refresh().await?;
3290-
coord.snapshot()
3291-
};
3292-
crate::db::schema_state::recover_schema_state_files(
3293-
root_uri,
3294-
std::sync::Arc::clone(&storage),
3295-
&snapshot,
3296-
)
3297-
.await?
3298-
} else {
3299-
SchemaStateRecovery::Noop
3300-
};
33013285
// Fresh per-branch snapshot — same rationale as
33023286
// `recover_manifest_drift`: classify against the branch the
33033287
// sidecar's writer targeted, refreshed after any prior
@@ -3336,7 +3320,6 @@ pub(crate) async fn heal_pending_sidecars_roll_forward(
33363320
&branch_snapshot,
33373321
&sidecar,
33383322
RecoveryMode::RollForwardOnly,
3339-
schema_state_recovery,
33403323
)
33413324
.await?
33423325
{
@@ -3484,7 +3467,6 @@ pub(crate) async fn recover_manifest_drift(
34843467
storage: std::sync::Arc<dyn StorageAdapter>,
34853468
coordinator: &mut GraphCoordinator,
34863469
mode: RecoveryMode,
3487-
schema_state_recovery: SchemaStateRecovery,
34883470
write_queue: &crate::db::write_queue::WriteQueueManager,
34893471
) -> Result<()> {
34903472
let sidecars = list_sidecars(root_uri, storage.as_ref()).await?;
@@ -3554,15 +3536,7 @@ pub(crate) async fn recover_manifest_drift(
35543536
coordinator.snapshot()
35553537
}
35563538
};
3557-
process_sidecar(
3558-
root_uri,
3559-
&storage,
3560-
&branch_snapshot,
3561-
&sidecar,
3562-
mode,
3563-
schema_state_recovery,
3564-
)
3565-
.await?;
3539+
process_sidecar(root_uri, &storage, &branch_snapshot, &sidecar, mode).await?;
35663540
}
35673541
// Final refresh so the caller sees the post-sweep state.
35683542
coordinator.refresh().await?;
@@ -3783,7 +3757,6 @@ async fn process_sidecar(
37833757
snapshot: &Snapshot,
37843758
sidecar: &RecoverySidecar,
37853759
mode: RecoveryMode,
3786-
schema_state_recovery: SchemaStateRecovery,
37873760
) -> Result<bool> {
37883761
// Returns whether durable state changed (roll-forward, roll-back,
37893762
// effect-free retirement, or stale-sidecar audit recovery). `false` =
@@ -4188,9 +4161,7 @@ async fn process_sidecar(
41884161
.map(|()| true)
41894162
}
41904163
SidecarDecision::RollForward => {
4191-
if matches!(sidecar.writer_kind, SidecarKind::SchemaApply)
4192-
&& !schema_state_recovery.completed_schema_apply_sidecar_rename()
4193-
{
4164+
if matches!(sidecar.writer_kind, SidecarKind::SchemaApply) {
41944165
if sidecar.schema_version == SCHEMA_APPLY_CONFIRMATION_SCHEMA_VERSION {
41954166
let target_is_live =
41964167
schema_apply_target_identity_is_live(root_uri, storage.as_ref(), sidecar)
@@ -5613,6 +5584,7 @@ async fn regenerate_system_column_upgrade_staging(
56135584
root_uri,
56145585
storage.as_ref(),
56155586
&target.desired_ir,
5587+
None,
56165588
)
56175589
.await?;
56185590
crate::db::schema_state::validate_exact_schema_staging_target(
@@ -6939,6 +6911,28 @@ pub(crate) async fn ensure_read_only_schema_coherent(
69396911
}
69406912
}
69416913

6914+
// RFC 0067: schema apply arms no sidecar. Its staged contract names the
6915+
// publishing graph commit; once that commit is in lineage the manifest
6916+
// already carries the new registrations, and only a read-write open may
6917+
// install the contract files. An unpublished staging is inert garbage.
6918+
if let crate::db::schema_state::StagedContract::Marked {
6919+
state,
6920+
published: true,
6921+
} = crate::db::schema_state::inspect_staged_contract(root_uri, storage, false).await?
6922+
{
6923+
let graph_commit_id = state
6924+
.publication
6925+
.map(|publication| publication.graph_commit_id)
6926+
.unwrap_or_default();
6927+
return Err(OmniError::recovery_required(
6928+
graph_commit_id.clone(),
6929+
format!(
6930+
"read-only open found SchemaApply manifest outcome for graph commit '{}' but the schema contract promotion is pending; run a read-write open to finish it",
6931+
graph_commit_id
6932+
),
6933+
));
6934+
}
6935+
69426936
for sidecar in sidecars {
69436937
if !matches!(sidecar.writer_kind, SidecarKind::SchemaApply) {
69446938
continue;
@@ -8687,6 +8681,7 @@ pub(crate) async fn confirm_occ_sidecar_v9(
86878681
/// independently durable table operation; metadata/tombstone-only applies may
86888682
/// deliberately carry an empty set because schema staging is confirmed by the
86898683
/// same protocol before publication.
8684+
#[cfg(test)]
86908685
pub(crate) fn new_schema_apply_sidecar_v9(
86918686
actor_id: Option<String>,
86928687
tables: Vec<SidecarTablePin>,

‎crates/omnigraph/src/db/omnigraph.rs‎

Lines changed: 82 additions & 19 deletions
Original file line numberDiff line numberDiff line change
@@ -62,12 +62,13 @@ use super::manifest::{
6262
GenesisManifestAttempt, ManifestChange, Snapshot, TableRegistration, TableTombstone,
6363
};
6464
use super::schema_state::{
65-
SCHEMA_SOURCE_FILENAME, SchemaContractText, load_validated_schema_contract,
66-
load_validated_schema_contract_for_source, read_accepted_schema_ir, read_schema_contract_text,
67-
read_schema_contract_text_for_source, read_schema_state_identity, recover_schema_state_files,
68-
render_schema_contract, schema_ir_uri, schema_source_staging_uri, schema_source_uri,
69-
schema_state_uri, validate_schema_contract, validate_schema_contract_text,
70-
validate_schema_ir_against_snapshot, write_schema_contract, write_schema_contract_staging,
65+
SCHEMA_SOURCE_FILENAME, SchemaContractText, SchemaStagingPolicy, SchemaStateRecovery,
66+
load_validated_schema_contract, load_validated_schema_contract_for_source,
67+
read_accepted_schema_ir, read_schema_contract_text, read_schema_contract_text_for_source,
68+
read_schema_state_identity, recover_schema_state_files, render_schema_contract, schema_ir_uri,
69+
schema_source_staging_uri, schema_source_uri, schema_state_uri, validate_schema_contract,
70+
validate_schema_contract_text, validate_schema_ir_against_snapshot, write_schema_contract,
71+
write_schema_contract_staging,
7172
};
7273
use super::{
7374
ReadTarget, ResolvedTarget, SCHEMA_APPLY_LOCK_BRANCH, SnapshotId, is_internal_system_branch,
@@ -213,6 +214,11 @@ pub struct Omnigraph {
213214
coordinator: Arc<tokio::sync::RwLock<GraphCoordinator>>,
214215
table_store: TableStore,
215216
runtime_cache: RuntimeCache,
217+
/// RFC 0067: this handle's schema apply published its manifest commit but
218+
/// could not install the schema contract. The next write-entry heal on
219+
/// this handle installs it from the staged copy; other handles and
220+
/// processes converge at their next read-write open.
221+
pending_schema_install: std::sync::atomic::AtomicBool,
216222
/// Warm change-feed cut for this handle's bound branch. A cut (head,
217223
/// witness, genesis, lineage projection, forward child index) is a PURE
218224
/// projection of `__manifest`, so it is exactly valid while the manifest
@@ -495,7 +501,7 @@ impl Omnigraph {
495501
let schema_identity_domain = schema_ir.schema_identity_domain.as_str().to_string();
496502
let mut catalog = build_catalog_from_ir(&schema_ir)?;
497503
fixup_physical_schemas(&mut catalog)?;
498-
let (_, ir_json, state_json) = render_schema_contract(&schema_ir)?;
504+
let (_, ir_json, state_json) = render_schema_contract(&schema_ir, None)?;
499505
let contract = SchemaContractText {
500506
source: schema_source.to_string(),
501507
ir_json,
@@ -665,6 +671,7 @@ impl Omnigraph {
665671
// sessions reuse the process-wide object-store registry.
666672
table_store: TableStore::new(&root, session),
667673
runtime_cache: RuntimeCache::default(),
674+
pending_schema_install: std::sync::atomic::AtomicBool::new(false),
668675
feed_cut_cache: tokio::sync::RwLock::new(None),
669676
read_caches,
670677
schema_view: Arc::new(ArcSwap::from_pointee(HandleSchemaView {
@@ -784,12 +791,29 @@ impl Omnigraph {
784791
if matches!(mode, OpenMode::ReadWrite) {
785792
// Schema staging is itself mutable recovery state. Hold the shared
786793
// schema gate across BOTH its file pre-pass and the complete Full
787-
// sidecar sweep, so `schema_state_recovery` cannot go stale in a
794+
// sidecar sweep, so the schema staging decision cannot go stale in a
788795
// release/reacquire gap. The sweep adds branch → sorted table gates
789796
// per sidecar under this outer guard.
790-
let schema_state_recovery =
791-
recover_schema_state_files(&root, Arc::clone(&storage), &coordinator.snapshot())
792-
.await?;
797+
recover_schema_state_files(
798+
&root,
799+
Arc::clone(&storage),
800+
&coordinator.snapshot(),
801+
SchemaStagingPolicy::PromoteOrDiscard,
802+
)
803+
.await?;
804+
// RFC 0067: a crashed schema apply leaves its durable sentinel
805+
// behind with no sidecar to retire it. The open-time pass above
806+
// settled its staging, so the sentinel is stale under the same
807+
// one-mutation-process boundary; reclaim it before the sweep.
808+
if coordinator
809+
.all_branches()
810+
.await?
811+
.iter()
812+
.any(|branch| is_schema_apply_lock_branch(branch))
813+
{
814+
tracing::warn!("reclaiming the schema apply sentinel left by a crashed apply");
815+
coordinator.branch_delete(SCHEMA_APPLY_LOCK_BRANCH).await?;
816+
}
793817
// Recovery sweep: close the Phase B → Phase C residual on
794818
// any sidecar left over from a crashed writer. Long-running
795819
// processes additionally converge in-process: the staged-
@@ -802,7 +826,6 @@ impl Omnigraph {
802826
Arc::clone(&storage),
803827
&mut coordinator,
804828
crate::db::manifest::RecoveryMode::Full,
805-
schema_state_recovery,
806829
write_queue.as_ref(),
807830
)
808831
.await?;
@@ -857,6 +880,7 @@ impl Omnigraph {
857880
// sessions reuse the process-wide object-store registry.
858881
table_store: TableStore::new(&root, session),
859882
runtime_cache: RuntimeCache::default(),
883+
pending_schema_install: std::sync::atomic::AtomicBool::new(false),
860884
feed_cut_cache: tokio::sync::RwLock::new(None),
861885
read_caches,
862886
schema_view: Arc::new(ArcSwap::from_pointee(HandleSchemaView {
@@ -2050,12 +2074,25 @@ impl Omnigraph {
20502074
{
20512075
let mut coord = self.coordinator.write().await;
20522076
coord.refresh().await?;
2053-
recover_schema_state_files(
2077+
let outcome = recover_schema_state_files(
20542078
&self.root_uri,
20552079
Arc::clone(&self.storage),
20562080
&coord.snapshot(),
2081+
SchemaStagingPolicy::PromoteOnly,
20572082
)
20582083
.await?;
2084+
// A promoted staging completes a published apply whose writer
2085+
// died before releasing its sentinel; release it here so the
2086+
// caller's write is not refused until the next open.
2087+
if matches!(outcome, SchemaStateRecovery::Promoted)
2088+
&& coord
2089+
.all_branches()
2090+
.await?
2091+
.iter()
2092+
.any(|branch| is_schema_apply_lock_branch(branch))
2093+
{
2094+
coord.branch_delete(SCHEMA_APPLY_LOCK_BRANCH).await?;
2095+
}
20592096
}
20602097
} // ← guards released before the heal's queue acquisition
20612098
let _outcome = crate::db::manifest::heal_pending_sidecars_roll_forward(
@@ -2330,12 +2367,38 @@ impl Omnigraph {
23302367
&self.write_queue,
23312368
)
23322369
.await?;
2333-
if outcome.processed_any {
2334-
// A rolled-forward SchemaApply sidecar moved disk + manifest
2335-
// to the new schema (staging promoted, registrations
2336-
// published); the in-memory catalog must follow or the very
2337-
// write that triggered the heal validates against the stale
2338-
// schema. Same post-heal step as `refresh`.
2370+
// RFC 0067: finish this handle's own published-but-uninstalled schema
2371+
// apply before the caller's write plans against the manifest. The
2372+
// flag keeps the common path free of any staging probe.
2373+
let mut installed = false;
2374+
if self
2375+
.pending_schema_install
2376+
.swap(false, std::sync::atomic::Ordering::SeqCst)
2377+
{
2378+
let result = {
2379+
let _serial = self
2380+
.write_queue
2381+
.acquire(&crate::db::manifest::schema_apply_serial_queue_key())
2382+
.await;
2383+
let snapshot = self.coordinator.read().await.snapshot();
2384+
recover_schema_state_files(
2385+
&self.root_uri,
2386+
Arc::clone(&self.storage),
2387+
&snapshot,
2388+
SchemaStagingPolicy::PromoteOnly,
2389+
)
2390+
.await
2391+
};
2392+
match result {
2393+
Ok(recovery) => installed = matches!(recovery, SchemaStateRecovery::Promoted),
2394+
Err(error) => {
2395+
self.pending_schema_install
2396+
.store(true, std::sync::atomic::Ordering::SeqCst);
2397+
return Err(error);
2398+
}
2399+
}
2400+
}
2401+
if outcome.processed_any || installed {
23392402
self.reload_schema_if_source_changed().await?;
23402403
self.invalidate_read_caches().await;
23412404
}

0 commit comments

Comments
 (0)