From 1ff4b393c84978a4cf90fa4ceb17ff321b038c1f Mon Sep 17 00:00:00 2001 From: ScriptedAlchemy Date: Wed, 7 Oct 2026 23:45:42 -0700 Subject: [PATCH 1/6] fix(observability): settle retired owners and batch writes --- .../src/observability/producer.rs | 37 +++++ .../src/project_runtime/observability.rs | 22 +-- .../src/observability_outbox_tests.rs | 147 ++++++++++++++++++ .../src/registered_analytics.rs | 73 +++++++-- .../daemon_suite/invocation_observability.rs | 85 ++++++++++ 5 files changed, 346 insertions(+), 18 deletions(-) diff --git a/crates/tracedecay-application/src/observability/producer.rs b/crates/tracedecay-application/src/observability/producer.rs index 735417b8ed..395687739b 100644 --- a/crates/tracedecay-application/src/observability/producer.rs +++ b/crates/tracedecay-application/src/observability/producer.rs @@ -32,6 +32,7 @@ use rollup_rebuild::{RollupAdvanceOutcome, run_one_rollup_maintenance}; const PRODUCER_RUNNING: u8 = 0; const PRODUCER_STOPPING: u8 = 1; const PRODUCER_STOPPED: u8 = 2; +const PRODUCER_SETTLED: u8 = 3; const MAX_PRODUCER_CAPACITY: usize = 1_024; const OBSERVABILITY_WRITE_BATCH: usize = 32; const MAX_PRODUCER_DEADLINE: Duration = Duration::from_secs(60); @@ -428,6 +429,12 @@ impl BoundedObservabilityProducerV1 { self.core.stop(false).await } + /// Confirms worker join and the writer fence completed. Persistence + /// failures remain in the shutdown result and the stream's coverage. + pub fn shutdown_settled(&self) -> bool { + self.core.state.load(Ordering::Acquire) == PRODUCER_SETTLED + } + pub async fn cancel(&self) -> Result { self.core.stop(true).await } @@ -575,6 +582,36 @@ impl ObservabilityProducerCoreV1 { "observability worker join failed: {error}" )) })?; + // A persistence timeout can drop an await while its database + // command still owns the transaction. The canonical writer fence + // proves those commands settled before replacement is permitted. + let fence = timeout_at(shutdown_deadline, async { + let transaction = self.db.begin_write_transaction().await.map_err(|error| { + ApplicationContractError::Domain(format!( + "observability writer fence failed: {error}" + )) + })?; + transaction.rollback().await.map_err(|error| { + ApplicationContractError::Domain(format!( + "observability writer fence failed: {error}" + )) + }) + }) + .await; + match fence { + Ok(Ok(())) => self.state.store(PRODUCER_SETTLED, Ordering::Release), + Ok(Err(error)) => { + tracing::warn!(%error, "observability shutdown writer fence failed"); + return outcome.and(Err(error)); + } + Err(_) => { + let error = ApplicationContractError::Domain( + "observability_shutdown_deadline".to_owned(), + ); + tracing::warn!(%error, "observability shutdown writer fence incomplete"); + return outcome.and(Err(error)); + } + } } outcome } diff --git a/crates/tracedecay-daemon-service/src/project_runtime/observability.rs b/crates/tracedecay-daemon-service/src/project_runtime/observability.rs index 3c03e2c943..b8e77641ca 100644 --- a/crates/tracedecay-daemon-service/src/project_runtime/observability.rs +++ b/crates/tracedecay-daemon-service/src/project_runtime/observability.rs @@ -54,7 +54,7 @@ impl StoreObservabilityCoreV1 { }) } - async fn shutdown(&self) -> Result<(), tracedecay_contracts::ApplicationContractError> { + async fn shutdown(&self) -> (Result<(), ApplicationContractError>, bool) { let mut first_error = None; if let Err(error) = self.work_observations.shutdown().await { tracing::warn!(%error, "registered Work owner-observation recovery was incomplete"); @@ -66,13 +66,17 @@ impl StoreObservabilityCoreV1 { first_error = Some(error); } } + let ancillary_joined = first_error.is_none(); if let Err(error) = self.producer.shutdown().await { tracing::warn!(%error, "registered observability producer shutdown was incomplete"); if first_error.is_none() { first_error = Some(error); } } - first_error.map_or(Ok(()), Err) + ( + first_error.map_or(Ok(()), Err), + ancillary_joined && self.producer.shutdown_settled(), + ) } } @@ -387,11 +391,11 @@ impl StoreObservabilityRegistryV1 { let registry = self.clone(); self.retirement_drains.spawn_on( async move { - let result = core.shutdown().await; + let (result, settled) = core.shutdown().await; if let Err(error) = &result { tracing::warn!(%error, "background observability owner drain was incomplete"); } - if let Err(error) = registry.finish_retirement(&core, result.is_ok()) { + if let Err(error) = registry.finish_retirement(&core, settled) { tracing::warn!(%error, "background observability retirement was incomplete"); } }, @@ -410,9 +414,9 @@ impl StoreObservabilityRegistryV1 { self.retirement_drains.wait().await; } - /// Settles a `Stopping` entry: a releasable retirement removes it so a - /// fresh owner may mount; a genuine shutdown failure is remembered as - /// `Failed` and refuses all future mounts. + /// A confirmed join and writer fence remove the retiring owner even when persistence + /// reported an error. The caller still receives that error; incomplete + /// shutdown remains failed closed to prevent overlapping producers. fn finish_retirement( &self, core: &Arc, @@ -588,8 +592,8 @@ impl RegisteredObservabilityProducerV1 { if !registry.begin_retirement(&core, StoreObservabilityDrainV1::InFlight)? { return Ok(()); } - let result = core.shutdown().await; - let retirement = registry.finish_retirement(&core, result.is_ok()); + let (result, settled) = core.shutdown().await; + let retirement = registry.finish_retirement(&core, settled); match result { Ok(()) => retirement, Err(error) => { diff --git a/crates/tracedecay-global-db/src/observability_outbox_tests.rs b/crates/tracedecay-global-db/src/observability_outbox_tests.rs index a81fdb38c9..5595baa639 100644 --- a/crates/tracedecay-global-db/src/observability_outbox_tests.rs +++ b/crates/tracedecay-global-db/src/observability_outbox_tests.rs @@ -110,6 +110,153 @@ async fn observability_batch_appends_every_event() { ); } +fn ordinary_batch_event(index: usize) -> AnalyticsEventInsert { + AnalyticsEventInsert { + provider: "tracedecay-observability".to_owned(), + project_id: "scope:ordinary-batch".to_owned(), + session_id: None, + timestamp: 1, + event_kind: "retrieval.query.completed.v1".to_owned(), + hook_name: None, + tool_name: None, + tool_category: None, + skill_name: None, + hint_category: None, + hint_id: Some(format!("ordinary:{index}")), + outcome: Some("succeeded".to_owned()), + metadata_json: Some(format!("{{\"index\":{index}}}")), + } +} + +#[tokio::test] +async fn ordinary_observability_batch_matches_single_appends_across_chunks_and_replays() { + let batch = RegisteredGlobalDbHarness::open("ordinary-batch-equivalence").await; + let singles = RegisteredGlobalDbHarness::open("ordinary-single-equivalence").await; + let seeded = ordinary_batch_event(0); + for db in [&batch.registered, &singles.registered] { + db.append_observability_event(&seeded).await.unwrap(); + } + let mut events = vec![ + ordinary_batch_event(1), + seeded.clone(), + ordinary_batch_event(1), + ]; + let mut other_scope = ordinary_batch_event(1); + other_scope.project_id = "scope:other".to_owned(); + events.push(other_scope); + events.extend( + (2..=crate::registered_analytics::ANALYTICS_INSERT_ROWS_PER_STATEMENT) + .map(ordinary_batch_event), + ); + events.extend([ordinary_batch_event(1), seeded]); + let mut expected = Vec::new(); + for event in &events { + expected.push( + singles + .registered + .append_observability_event(event) + .await + .unwrap(), + ); + } + let actual = batch + .registered + .append_observability_events(&events) + .await + .unwrap(); + assert_eq!(actual, expected); + assert_eq!(actual[0], actual[2]); + assert_eq!(actual[0], actual[actual.len() - 2]); + assert_ne!(actual[0], actual[3]); + assert_eq!( + batch + .registered + .append_observability_events(&events) + .await + .unwrap(), + actual + ); + for (event, id) in events.iter().zip(actual) { + let record = batch + .registered + .read_observability_event(&event.project_id, event.hint_id.as_deref().unwrap()) + .await + .unwrap() + .unwrap(); + assert_eq!(record.id, id); + assert_eq!(record.metadata_json, event.metadata_json); + } +} + +#[tokio::test] +async fn ordinary_observability_batch_conflicts_roll_back_all_chunks() { + let harness = RegisteredGlobalDbHarness::open("ordinary-batch-rollback").await; + let seeded = ordinary_batch_event(0); + let seeded_id = harness + .registered + .append_observability_event(&seeded) + .await + .unwrap(); + let count = crate::registered_analytics::ANALYTICS_INSERT_ROWS_PER_STATEMENT; + let mut events = (1..=count).map(ordinary_batch_event).collect::>(); + let mut conflict = events[0].clone(); + conflict.metadata_json = Some("{\"changed\":true}".to_owned()); + events.push(conflict); + let error = harness + .registered + .append_observability_events(&events) + .await + .unwrap_err(); + assert!(error.contains("idempotency conflict"), "{error}"); + assert_eq!( + harness + .registered + .count_analytics_events(Some(&seeded.project_id), 0) + .await + .unwrap(), + 1 + ); + assert_eq!( + harness + .registered + .append_observability_event(&seeded) + .await + .unwrap(), + seeded_id + ); + + // A conflict within the first run must also leave no accepted prefix. + let mut changed = ordinary_batch_event(1); + changed.timestamp += 1; + let error = harness + .registered + .append_observability_events(&[ordinary_batch_event(1), changed]) + .await + .unwrap_err(); + assert!(error.contains("idempotency conflict"), "{error}"); + assert_eq!( + harness + .registered + .count_analytics_events(Some(&seeded.project_id), 0) + .await + .unwrap(), + 1 + ); + + // Preserve input failure precedence: an earlier replay conflict precedes + // an invalid later envelope, even though lookup is performed in a batch. + let mut changed = seeded.clone(); + changed.timestamp += 1; + let mut invalid = ordinary_batch_event(2); + invalid.hint_id = None; + let error = harness + .registered + .append_observability_events(&[changed, invalid]) + .await + .unwrap_err(); + assert!(error.contains("idempotency conflict"), "{error}"); +} + #[tokio::test] async fn observability_outbox_replay_reuses_exact_delivery_and_settles_atomically() { let harness = RegisteredGlobalDbHarness::open("observability-outbox-replay").await; diff --git a/crates/tracedecay-global-db/src/registered_analytics.rs b/crates/tracedecay-global-db/src/registered_analytics.rs index cafcdfd70b..24c98a5ad9 100644 --- a/crates/tracedecay-global-db/src/registered_analytics.rs +++ b/crates/tracedecay-global-db/src/registered_analytics.rs @@ -221,15 +221,7 @@ impl RegisteredGlobalDb { append_analytics_events_in_existing_tx(&transaction, events).await? } AnalyticsAppendKind::Observability => { - let mut ids = Vec::with_capacity(events.len()); - for event in events { - ids.push( - append_observability_event_in_existing_tx(&transaction, event) - .await - .map_err(|error| error.to_string())?, - ); - } - ids + append_observability_events_in_existing_tx(&transaction, events).await? } }; transaction @@ -1419,6 +1411,69 @@ enum ObservabilityAppendError { Failed(String), } +/// Resolve replay and in-batch duplicates before inserting each bounded run. +/// All runs share the caller's transaction, so a later conflict rolls back +/// earlier appends just as the single-event path does. +async fn append_observability_events_in_existing_tx( + transaction: &RegisteredGlobalDbWriteTransaction<'_>, + events: &[AnalyticsEventInsert], +) -> Result, String> { + let mut ids = Vec::with_capacity(events.len()); + for chunk in events.chunks(ANALYTICS_INSERT_ROWS_PER_STATEMENT) { + let keys = chunk + .iter() + .filter_map(|event| { + event + .hint_id + .as_ref() + .map(|hint| (event.project_id.clone(), hint.clone())) + }) + .collect(); + let stored = stored_observability_events(transaction, keys).await?; + let mut pending = BTreeMap::<(String, String), usize>::new(); + let mut appends = Vec::::new(); + let mut ordered_keys = Vec::with_capacity(chunk.len()); + let mut resolved = BTreeMap::new(); + for event in chunk { + validate_observability_event(event)?; + let hint = event + .hint_id + .as_ref() + .ok_or("invalid canonical observability event")?; + let key = (event.project_id.clone(), hint.clone()); + if let Some(record) = stored.get(&key) { + if !analytics_record_matches_insert(record, event) { + return Err(ObservabilityAppendError::Conflict.to_string()); + } + resolved.insert(key.clone(), record.id); + } else if let Some(index) = pending.get(&key) { + if appends[*index] != *event { + return Err(ObservabilityAppendError::Conflict.to_string()); + } + } else { + pending.insert(key.clone(), appends.len()); + appends.push(event.clone()); + } + ordered_keys.push(key); + } + let appended = append_analytics_events_in_existing_tx(transaction, &appends).await?; + for (key, index) in pending { + let id = appended + .get(index) + .ok_or("observability batch lost an appended event id")?; + resolved.insert(key, *id); + } + for key in ordered_keys { + ids.push( + *resolved + .get(&key) + .ok_or("observability batch left an event unresolved")?, + ); + } + } + Ok(ids) +} + async fn append_observability_event_in_existing_tx( transaction: &RegisteredGlobalDbWriteTransaction<'_>, event: &AnalyticsEventInsert, diff --git a/crates/tracedecay/tests/daemon_suite/invocation_observability.rs b/crates/tracedecay/tests/daemon_suite/invocation_observability.rs index 4935d23def..589ad56faf 100644 --- a/crates/tracedecay/tests/daemon_suite/invocation_observability.rs +++ b/crates/tracedecay/tests/daemon_suite/invocation_observability.rs @@ -1308,6 +1308,86 @@ async fn runtimeless_last_alias_drop_keeps_the_store_retiring_until_the_drain_co replacement.shutdown().await.expect("replacement shutdown"); } +#[tokio::test] +async fn registered_persistence_failure_releases_only_after_join_and_writer_fence() { + let (_project, project_id, database, _runtime) = runtime("observability-shutdown-fenced").await; + let identity = ObservabilityProducerIdentityV1 { + authorized_scope_ref: project_id.as_str().to_owned(), + process_boot_id: "daemon:shutdown-failure".to_owned(), + producer_revision: "producer.v1".to_owned(), + configuration_revision: digest('e').as_str().to_owned(), + policy_revision: digest('f').as_str().to_owned(), + }; + let producer = BoundedObservabilityProducerV1::start_with_deadlines( + database.clone(), + identity.clone(), + 1, + ObservabilityProducerDeadlinesV1 { + persistence: Duration::from_millis(50), + shutdown: Duration::from_millis(250), + }, + ) + .expect("producer"); + let registry = StoreObservabilityRegistryV1::default(); + let registered = registry + .acquire_or_start(&database, &store_mount(&identity), || Ok(producer)) + .expect("registered observability producer"); + let blocker = database + .begin_write_transaction() + .await + .expect("hold registered writer"); + registered + .producer() + .try_emit(envelope(&project_id, "shutdown:blocked")) + .expect("enqueue blocked event"); + tokio::task::yield_now().await; + + let retained = registered.producer(); + assert!(!retained.shutdown_settled()); + let shutdown = tokio::spawn(registered.shutdown()); + // Hold the writer past persistence's deadline, then release it within + // shutdown's existing budget so its fence can observe transaction closure. + tokio::time::sleep(Duration::from_millis(100)).await; + assert!(!retained.shutdown_settled()); + blocker.commit().await.expect("release registered writer"); + let error = shutdown + .await + .expect("shutdown task") + .expect_err("successful fencing must preserve the persistence failure"); + assert!(retained.shutdown_settled()); + assert_eq!( + retained + .try_emit(envelope(&project_id, "shutdown:after-fence")) + .unwrap_err(), + "observability_producer_closed", + ); + assert!( + error + .to_string() + .contains("observability_persistence_deadline"), + "unexpected shutdown error: {error}" + ); + let start_called = Arc::new(AtomicBool::new(false)); + let observed_start = Arc::clone(&start_called); + let replacement_identity = ObservabilityProducerIdentityV1 { + process_boot_id: "daemon:shutdown-failure-replacement".to_owned(), + ..identity.clone() + }; + let replacement_mount = store_mount(&replacement_identity); + let replacement = registry.acquire_or_start(&database, &replacement_mount, || { + observed_start.store(true, Ordering::Release); + BoundedObservabilityProducerV1::start(database.clone(), replacement_identity.clone(), 1) + .map_err(StoreObservabilityMountErrorV1::Unavailable) + }); + let replacement = replacement.expect("settled store can reopen without resetting data"); + assert!(start_called.load(Ordering::Acquire)); + replacement + .producer() + .try_emit(envelope(&project_id, "replacement:accepted")) + .unwrap(); + replacement.shutdown().await.expect("replacement drains"); +} + #[tokio::test] async fn registered_shutdown_reports_a_blocked_producer_flush() { let (_project, project_id, database, _runtime) = @@ -1343,10 +1423,15 @@ async fn registered_shutdown_reports_a_blocked_producer_flush() { .expect("enqueue blocked event"); tokio::task::yield_now().await; + let retained = registered.producer(); let error = registered .shutdown() .await .expect_err("blocked flush must fail the registered shutdown"); + assert!( + !retained.shutdown_settled(), + "an unfinished writer fence cannot release the store" + ); blocker.commit().await.expect("release registered writer"); assert!( error From 75ae3c502f8a4a8285f22a4c2506f07b21c40569 Mon Sep 17 00:00:00 2001 From: ScriptedAlchemy Date: Wed, 7 Oct 2026 23:55:12 -0700 Subject: [PATCH 2/6] fix(observability): reuse final write settlement on shutdown --- .../src/observability/producer.rs | 136 +++++++++++++++--- 1 file changed, 113 insertions(+), 23 deletions(-) diff --git a/crates/tracedecay-application/src/observability/producer.rs b/crates/tracedecay-application/src/observability/producer.rs index 395687739b..084379a86b 100644 --- a/crates/tracedecay-application/src/observability/producer.rs +++ b/crates/tracedecay-application/src/observability/producer.rs @@ -129,7 +129,10 @@ impl ObservabilityProducerDeadlinesV1 { enum ProducerControl { Shutdown { cancelled: bool, - reply: oneshot::Sender>, + reply: oneshot::Sender<( + Result, + bool, + )>, }, } @@ -554,7 +557,7 @@ impl ObservabilityProducerCoreV1 { "observability_control_lane_closed".to_owned(), )); } - let outcome = match timeout_at(shutdown_deadline, result).await { + let (outcome, writer_settled) = match timeout_at(shutdown_deadline, result).await { Ok(Ok(outcome)) => outcome, Ok(Err(_)) => { if let Some(worker) = worker.take() { @@ -582,6 +585,10 @@ impl ObservabilityProducerCoreV1 { "observability worker join failed: {error}" )) })?; + if writer_settled { + self.state.store(PRODUCER_SETTLED, Ordering::Release); + return outcome; + } // A persistence timeout can drop an await while its database // command still owns the transaction. The canonical writer fence // proves those commands settled before replacement is permitted. @@ -655,7 +662,7 @@ async fn run_worker( .await; break; }; - let dropped_count = settle_worker( + let (dropped_count, writer_settled) = settle_worker( &db, &identity, &mut data, @@ -674,7 +681,7 @@ async fn run_worker( }), Err, ); - let _ = reply.send(result); + let _ = reply.send((result, writer_settled)); break; } observation = data.recv() => { @@ -844,9 +851,9 @@ async fn record_batch( persisted: &mut u64, first_error: &mut Option, persistence_deadline: Duration, -) { +) -> bool { if envelopes.is_empty() { - return; + return false; } let count = u64::try_from(envelopes.len()).unwrap_or(u64::MAX); match timeout( @@ -855,14 +862,21 @@ async fn record_batch( ) .await { - Ok(Ok(_)) => *persisted = persisted.saturating_add(count), - Ok(Err(error)) if first_error.is_none() => *first_error = Some(error), + Ok(Ok(_)) => { + *persisted = persisted.saturating_add(count); + true + } + Ok(Err(error)) if first_error.is_none() => { + *first_error = Some(error); + false + } Err(_) if first_error.is_none() => { *first_error = Some(ApplicationContractError::Domain( "observability_persistence_deadline".to_owned(), )); + false } - Ok(Err(_)) | Err(_) => {} + Ok(Err(_)) | Err(_) => false, } } @@ -895,8 +909,9 @@ async fn settle_worker( progress: &mut ProducerWorkerProgress, discard_pending: bool, clean_shutdown_observed: bool, -) -> u64 { +) -> (u64, bool) { data.close(); + let mut writer_settled = false; if discard_pending { let mut ranges = Vec::new(); while let Ok(observation) = data.try_recv() { @@ -947,7 +962,7 @@ async fn settle_worker( } for range in ranges { let drop_envelope = telemetry_drop_envelope(range, false); - record( + writer_settled = record( db, drop_envelope, &mut progress.persisted, @@ -977,6 +992,16 @@ async fn settle_worker( state.deadlines.persistence, ) .await; + // Finish optional maintenance before the final carrier. Its successful + // transaction then fences every earlier write, including timed-out + // commands, without acquiring the shared writer again after closure. + let _ = run_one_rollup_maintenance( + db, + identity, + state.deadlines.persistence, + &mut progress.rollup_frontier_initialized, + ) + .await; let pending = match take_pending_drops(state) { Ok(pending) => pending, Err(error) => { @@ -990,7 +1015,7 @@ async fn settle_worker( for pending in pending { let closes_cleanly = clean_shutdown_observed && progress.first_error.is_none(); let drop_envelope = telemetry_drop_envelope(pending, closes_cleanly); - record( + writer_settled = record( db, drop_envelope, &mut progress.persisted, @@ -1012,7 +1037,7 @@ async fn settle_worker( }, progress.first_error.is_none(), ); - record( + writer_settled = record( db, zero_terminal, &mut progress.persisted, @@ -1021,15 +1046,8 @@ async fn settle_worker( ) .await; } - let _ = run_one_rollup_maintenance( - db, - identity, - state.deadlines.persistence, - &mut progress.rollup_frontier_initialized, - ) - .await; } - state.total_dropped.load(Ordering::Acquire) + (state.total_dropped.load(Ordering::Acquire), writer_settled) } fn take_pending_drops( @@ -1154,7 +1172,7 @@ async fn record( persisted: &mut u64, first_error: &mut Option, persistence_deadline: Duration, -) { +) -> bool { record_batch( db, vec![envelope], @@ -1162,7 +1180,7 @@ async fn record( first_error, persistence_deadline, ) - .await; + .await } fn telemetry_drop_envelope( @@ -1245,3 +1263,75 @@ mod cheaper_frontier_tests { assert!(!should_wake_rollup_now(true, true, false)); } } + +#[cfg(test)] +mod shutdown_settlement_tests { + use super::{ + BoundedObservabilityProducerV1, ObservabilityProducerDeadlinesV1, + ObservabilityProducerIdentityV1, + }; + use std::sync::Arc; + use std::time::Duration; + use tokio::sync::oneshot; + use tracedecay_global_db::tests::harness::RegisteredGlobalDbTestRuntime; + + #[tokio::test] + async fn completed_terminal_does_not_wait_for_a_subsequent_unrelated_writer() { + let directory = tempfile::tempdir().unwrap(); + let project = directory.path().join("project"); + std::fs::create_dir(&project).unwrap(); + let project_id = tracedecay_domain::ProjectId::new("project.shutdown-terminal").unwrap(); + let runtime = RegisteredGlobalDbTestRuntime::project( + directory.path().join("profile"), + &project, + project_id.clone(), + ) + .await + .unwrap(); + let db = runtime.project_database_arc().unwrap(); + let producer = Arc::new( + BoundedObservabilityProducerV1::start_with_deadlines( + db.clone(), + ObservabilityProducerIdentityV1 { + authorized_scope_ref: project_id.as_str().to_owned(), + process_boot_id: "boot.shutdown-terminal".to_owned(), + producer_revision: "producer.v1".to_owned(), + configuration_revision: "configuration.v1".to_owned(), + policy_revision: "policy.v1".to_owned(), + }, + 1, + ObservabilityProducerDeadlinesV1 { + persistence: Duration::from_millis(50), + shutdown: Duration::from_millis(250), + }, + ) + .unwrap(), + ); + let worker = producer.core.worker.lock().unwrap().take().unwrap(); + let (release, released) = oneshot::channel(); + let (blocker_ready, blocker_handle) = oneshot::channel(); + // Keep the real worker and its terminal write, but make its join + // observe another writer already owning the database afterward. + *producer.core.worker.lock().unwrap() = Some(tokio::spawn(async move { + worker.await.unwrap(); + let (acquired, acquired_rx) = oneshot::channel(); + let blocker = tokio::spawn(async move { + let transaction = db.begin_write_transaction().await.unwrap(); + acquired.send(()).unwrap(); + released.await.unwrap(); + transaction.rollback().await.unwrap(); + }); + acquired_rx.await.unwrap(); + blocker_ready.send(blocker).unwrap(); + })); + let stopping = Arc::clone(&producer); + let shutdown = tokio::spawn(async move { stopping.shutdown().await }); + let blocker = blocker_handle.await.unwrap(); + let outcome = shutdown.await.unwrap(); + release.send(()).unwrap(); + blocker.await.unwrap(); + let summary = outcome.expect("completed producer never waits on the unrelated writer"); + assert_eq!(summary.persisted, 1); + assert!(producer.shutdown_settled()); + } +} From af903db8e8bcec4a45970bda62683385b8c79906 Mon Sep 17 00:00:00 2001 From: ScriptedAlchemy Date: Thu, 8 Oct 2026 00:02:38 -0700 Subject: [PATCH 3/6] fix(observability): keep terminal writes before maintenance --- .../src/observability/producer.rs | 21 ++++++++++--------- .../observability_producer_shutdown.rs | 7 +++++++ 2 files changed, 18 insertions(+), 10 deletions(-) diff --git a/crates/tracedecay-application/src/observability/producer.rs b/crates/tracedecay-application/src/observability/producer.rs index 084379a86b..a926ebaa15 100644 --- a/crates/tracedecay-application/src/observability/producer.rs +++ b/crates/tracedecay-application/src/observability/producer.rs @@ -992,16 +992,6 @@ async fn settle_worker( state.deadlines.persistence, ) .await; - // Finish optional maintenance before the final carrier. Its successful - // transaction then fences every earlier write, including timed-out - // commands, without acquiring the shared writer again after closure. - let _ = run_one_rollup_maintenance( - db, - identity, - state.deadlines.persistence, - &mut progress.rollup_frontier_initialized, - ) - .await; let pending = match take_pending_drops(state) { Ok(pending) => pending, Err(error) => { @@ -1046,6 +1036,17 @@ async fn settle_worker( ) .await; } + // The mandatory terminal precedes optional maintenance so its day + // can be rebuilt before shutdown. A deferred maintenance operation + // may still own a timed-out database command and requires the fence. + let maintenance = run_one_rollup_maintenance( + db, + identity, + state.deadlines.persistence, + &mut progress.rollup_frontier_initialized, + ) + .await; + writer_settled &= maintenance != RollupAdvanceOutcome::Deferred; } (state.total_dropped.load(Ordering::Acquire), writer_settled) } diff --git a/crates/tracedecay-application/tests/application_suite/observability_producer_shutdown.rs b/crates/tracedecay-application/tests/application_suite/observability_producer_shutdown.rs index f7e30ab709..7780a541d3 100644 --- a/crates/tracedecay-application/tests/application_suite/observability_producer_shutdown.rs +++ b/crates/tracedecay-application/tests/application_suite/observability_producer_shutdown.rs @@ -166,6 +166,13 @@ async fn clean_shutdown_persists_zero_drop_terminal_without_relabeling_cancel() assert_eq!(idle_summary.persisted, 1); assert_eq!(idle_summary.dropped, 0); assert!(!idle_summary.cancelled); + assert!( + db.claim_observability_rollup_dirty_day(scope, "shutdown:rollup-check", 30) + .await + .expect("inspect terminal day") + .is_none(), + "shutdown maintenance must include its final terminal carrier", + ); let active = BoundedObservabilityProducerV1::start(db.clone(), identity(scope, "boot:clean-active"), 4) From adfc5175c53c6d97abfeadff867ec6845f483a93 Mon Sep 17 00:00:00 2001 From: ScriptedAlchemy Date: Thu, 8 Oct 2026 00:16:44 -0700 Subject: [PATCH 4/6] fix(observability): recover joined owners on mount demand --- .../src/observability/producer.rs | 116 ++++++++++++----- .../src/invocation/observability_producer.rs | 8 +- .../src/project_runtime/observability.rs | 86 ++++++++++--- .../src/daemon/project_composition.rs | 7 +- .../src/daemon/project_delivery_mount.rs | 9 +- .../src/daemon/project_open_admission.rs | 6 + .../src/daemon/tests/runtime_identity.rs | 121 ++++++++++++++++++ .../daemon_suite/invocation_observability.rs | 103 +++++++++++++-- 8 files changed, 388 insertions(+), 68 deletions(-) diff --git a/crates/tracedecay-application/src/observability/producer.rs b/crates/tracedecay-application/src/observability/producer.rs index a926ebaa15..af3ce01ad9 100644 --- a/crates/tracedecay-application/src/observability/producer.rs +++ b/crates/tracedecay-application/src/observability/producer.rs @@ -33,6 +33,7 @@ const PRODUCER_RUNNING: u8 = 0; const PRODUCER_STOPPING: u8 = 1; const PRODUCER_STOPPED: u8 = 2; const PRODUCER_SETTLED: u8 = 3; +const PRODUCER_JOINED: u8 = 4; const MAX_PRODUCER_CAPACITY: usize = 1_024; const OBSERVABILITY_WRITE_BATCH: usize = 32; const MAX_PRODUCER_DEADLINE: Duration = Duration::from_secs(60); @@ -438,6 +439,23 @@ impl BoundedObservabilityProducerV1 { self.core.state.load(Ordering::Acquire) == PRODUCER_SETTLED } + /// The worker and admissions are closed, but timed-out database commands + /// may still need the canonical writer fence before the store can reopen. + pub fn shutdown_joined(&self) -> bool { + matches!( + self.core.state.load(Ordering::Acquire), + PRODUCER_JOINED | PRODUCER_SETTLED + ) + } + + /// A later mount demand may retry only the writer fence, never the worker + /// shutdown or accepted observations. Each attempt has the existing bound. + pub async fn finish_shutdown_settlement(&self) -> Result<(), ApplicationContractError> { + self.core + .finish_shutdown_settlement(Instant::now() + self.core.deadlines.shutdown) + .await + } + pub async fn cancel(&self) -> Result { self.core.stop(true).await } @@ -569,11 +587,26 @@ impl ObservabilityProducerCoreV1 { )); } Err(_) => { - if let Some(worker) = worker.take() { + let joined = if let Some(worker) = worker.take() { worker.abort(); - let _ = worker.await; - } - self.state.store(PRODUCER_STOPPED, Ordering::Release); + match worker.await { + Ok(()) => true, + Err(error) => error.is_cancelled(), + } + } else { + false + }; + // Both admission fences passed before this worker wait. An + // aborted, joined worker can leave only database commands, + // which a later bounded writer fence can settle safely. + self.state.store( + if joined { + PRODUCER_JOINED + } else { + PRODUCER_STOPPED + }, + Ordering::Release, + ); return Err(ApplicationContractError::Domain( "observability_shutdown_deadline".to_owned(), )); @@ -589,39 +622,49 @@ impl ObservabilityProducerCoreV1 { self.state.store(PRODUCER_SETTLED, Ordering::Release); return outcome; } - // A persistence timeout can drop an await while its database - // command still owns the transaction. The canonical writer fence - // proves those commands settled before replacement is permitted. - let fence = timeout_at(shutdown_deadline, async { - let transaction = self.db.begin_write_transaction().await.map_err(|error| { - ApplicationContractError::Domain(format!( - "observability writer fence failed: {error}" - )) - })?; - transaction.rollback().await.map_err(|error| { - ApplicationContractError::Domain(format!( - "observability writer fence failed: {error}" - )) - }) - }) - .await; - match fence { - Ok(Ok(())) => self.state.store(PRODUCER_SETTLED, Ordering::Release), - Ok(Err(error)) => { - tracing::warn!(%error, "observability shutdown writer fence failed"); - return outcome.and(Err(error)); - } - Err(_) => { - let error = ApplicationContractError::Domain( - "observability_shutdown_deadline".to_owned(), - ); - tracing::warn!(%error, "observability shutdown writer fence incomplete"); - return outcome.and(Err(error)); - } + self.state.store(PRODUCER_JOINED, Ordering::Release); + if let Err(error) = self.finish_shutdown_settlement(shutdown_deadline).await { + tracing::warn!(%error, "observability shutdown writer fence incomplete"); + return outcome.and(Err(error)); } } outcome } + + async fn finish_shutdown_settlement( + &self, + deadline: Instant, + ) -> Result<(), ApplicationContractError> { + match self.state.load(Ordering::Acquire) { + PRODUCER_SETTLED => return Ok(()), + PRODUCER_JOINED => {} + _ => { + return Err(ApplicationContractError::Domain( + "observability_worker_not_joined".to_owned(), + )); + } + } + // Timed-out statements may retain a transaction after the worker joins. + // The canonical writer fence proves those commands settled before reuse. + timeout_at(deadline, async { + let transaction = self.db.begin_write_transaction().await.map_err(|error| { + ApplicationContractError::Domain(format!( + "observability writer fence failed: {error}" + )) + })?; + transaction.rollback().await.map_err(|error| { + ApplicationContractError::Domain(format!( + "observability writer fence failed: {error}" + )) + }) + }) + .await + .map_err(|_| { + ApplicationContractError::Domain("observability_shutdown_deadline".to_owned()) + })??; + self.state.store(PRODUCER_SETTLED, Ordering::Release); + Ok(()) + } } async fn run_worker( @@ -1308,6 +1351,13 @@ mod shutdown_settlement_tests { ) .unwrap(), ); + assert!(!producer.shutdown_joined()); + let premature = producer.finish_shutdown_settlement().await.unwrap_err(); + assert!( + premature + .to_string() + .contains("observability_worker_not_joined") + ); let worker = producer.core.worker.lock().unwrap().take().unwrap(); let (release, released) = oneshot::channel(); let (blocker_ready, blocker_handle) = oneshot::channel(); diff --git a/crates/tracedecay-daemon-service/src/invocation/observability_producer.rs b/crates/tracedecay-daemon-service/src/invocation/observability_producer.rs index d9ed317a51..b11c362afa 100644 --- a/crates/tracedecay-daemon-service/src/invocation/observability_producer.rs +++ b/crates/tracedecay-daemon-service/src/invocation/observability_producer.rs @@ -118,8 +118,12 @@ impl DaemonInvocationService { "a different observability producer is already mounted for this project store" .to_owned(), }, - error @ (StoreObservabilityMountErrorV1::Retiring - | StoreObservabilityMountErrorV1::ShutdownFailed + StoreObservabilityMountErrorV1::Retiring => TraceDecayError::project_route( + StoreObservabilityMountErrorV1::RETIRING_REASON_CODE, + true, + "The previous observability owner is settling; retry this project shortly", + ), + error @ (StoreObservabilityMountErrorV1::ShutdownFailed | StoreObservabilityMountErrorV1::Unavailable(_)) => { TraceDecayError::Config { message: format!( diff --git a/crates/tracedecay-daemon-service/src/project_runtime/observability.rs b/crates/tracedecay-daemon-service/src/project_runtime/observability.rs index b8e77641ca..436a0ced84 100644 --- a/crates/tracedecay-daemon-service/src/project_runtime/observability.rs +++ b/crates/tracedecay-daemon-service/src/project_runtime/observability.rs @@ -54,7 +54,12 @@ impl StoreObservabilityCoreV1 { }) } - async fn shutdown(&self) -> (Result<(), ApplicationContractError>, bool) { + async fn shutdown( + &self, + ) -> ( + Result<(), ApplicationContractError>, + StoreObservabilityCompletionV1, + ) { let mut first_error = None; if let Err(error) = self.work_observations.shutdown().await { tracing::warn!(%error, "registered Work owner-observation recovery was incomplete"); @@ -75,7 +80,13 @@ impl StoreObservabilityCoreV1 { } ( first_error.map_or(Ok(()), Err), - ancillary_joined && self.producer.shutdown_settled(), + if !ancillary_joined || !self.producer.shutdown_joined() { + StoreObservabilityCompletionV1::Failed + } else if self.producer.shutdown_settled() { + StoreObservabilityCompletionV1::Settled + } else { + StoreObservabilityCompletionV1::AwaitingWriter + }, ) } } @@ -104,6 +115,14 @@ enum StoreObservabilityDrainV1 { /// The last alias was dropped without a tokio runtime, so nothing could /// run the drain. The next mount attempt on a live runtime starts it. Deferred, + /// All workers joined; a later mount may retry their writer fence once. + AwaitingWriter, +} + +enum StoreObservabilityCompletionV1 { + Settled, + AwaitingWriter, + Failed, } enum StoreObservabilityStateV1 { @@ -180,11 +199,15 @@ pub enum StoreObservabilityMountErrorV1 { Unavailable(&'static str), } +impl StoreObservabilityMountErrorV1 { + pub const RETIRING_REASON_CODE: &'static str = "store_observability_retiring"; +} + impl fmt::Display for StoreObservabilityMountErrorV1 { fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { let code = match self { Self::Busy => "store_observability_busy", - Self::Retiring => "store_observability_retiring", + Self::Retiring => Self::RETIRING_REASON_CODE, Self::ShutdownFailed => "store_observability_shutdown_failed", Self::Unavailable(reason) => reason, }; @@ -293,15 +316,18 @@ impl StoreObservabilityRegistryV1 { Ok(registered) } StoreObservabilityStateV1::Stopping { core, drain } => { - // A deferred drain (the last alias was dropped without a - // runtime) starts now that a caller with a live runtime - // has arrived. The mount is still refused: only the - // confirmed drain may vacate the entry. - if matches!(drain, StoreObservabilityDrainV1::Deferred) - && let Ok(runtime) = tokio::runtime::Handle::try_current() + // Demand starts at most one bounded drain or writer recheck. + // The mount stays refused until exact-owner settlement. + if matches!( + drain, + StoreObservabilityDrainV1::Deferred + | StoreObservabilityDrainV1::AwaitingWriter + ) && let Ok(runtime) = tokio::runtime::Handle::try_current() { + let writer_only = + matches!(drain, StoreObservabilityDrainV1::AwaitingWriter); *drain = StoreObservabilityDrainV1::InFlight; - self.spawn_retirement_drain(&runtime, Arc::clone(core)); + self.spawn_retirement_drain(&runtime, Arc::clone(core), writer_only); } Err(StoreObservabilityMountErrorV1::Retiring) } @@ -387,11 +413,22 @@ impl StoreObservabilityRegistryV1 { &self, runtime: &tokio::runtime::Handle, core: Arc, + writer_only: bool, ) { let registry = self.clone(); self.retirement_drains.spawn_on( async move { - let (result, settled) = core.shutdown().await; + let (result, settled) = if writer_only { + let result = core.producer.finish_shutdown_settlement().await; + let completion = if result.is_ok() { + StoreObservabilityCompletionV1::Settled + } else { + StoreObservabilityCompletionV1::AwaitingWriter + }; + (result, completion) + } else { + core.shutdown().await + }; if let Err(error) = &result { tracing::warn!(%error, "background observability owner drain was incomplete"); } @@ -414,13 +451,13 @@ impl StoreObservabilityRegistryV1 { self.retirement_drains.wait().await; } - /// A confirmed join and writer fence remove the retiring owner even when persistence - /// reported an error. The caller still receives that error; incomplete - /// shutdown remains failed closed to prevent overlapping producers. + /// Joined owners retain their exact store while a later mount may retry + /// writer settlement. Only its confirmed fence removes the owner; an + /// unconfirmed worker or ancillary shutdown remains failed closed. fn finish_retirement( &self, core: &Arc, - releasable: bool, + completion: StoreObservabilityCompletionV1, ) -> Result<(), ApplicationContractError> { let mut entries = self.lock_entries() @@ -439,10 +476,19 @@ impl StoreObservabilityRegistryV1 { field: "store_observability_retiring_owner", }); }; - if releasable { - entries.remove(index); - } else { - entries[index].state = StoreObservabilityStateV1::Failed; + match completion { + StoreObservabilityCompletionV1::Settled => { + entries.remove(index); + } + StoreObservabilityCompletionV1::AwaitingWriter => { + entries[index].state = StoreObservabilityStateV1::Stopping { + core: Arc::clone(core), + drain: StoreObservabilityDrainV1::AwaitingWriter, + }; + } + StoreObservabilityCompletionV1::Failed => { + entries[index].state = StoreObservabilityStateV1::Failed; + } } Ok(()) } @@ -637,7 +683,7 @@ impl Drop for RegisteredObservabilityProducerV1 { return; } match runtime { - Some(runtime) => self.registry.spawn_retirement_drain(&runtime, core), + Some(runtime) => self.registry.spawn_retirement_drain(&runtime, core, false), None => tracing::warn!( "observability owners dropped without a runtime; the store stays \ retiring until a deferred drain confirms the close" diff --git a/crates/tracedecay/src/daemon/project_composition.rs b/crates/tracedecay/src/daemon/project_composition.rs index 3104af3051..0be1b85edc 100644 --- a/crates/tracedecay/src/daemon/project_composition.rs +++ b/crates/tracedecay/src/daemon/project_composition.rs @@ -1784,7 +1784,12 @@ impl ProjectOpenInputs<'_> { error: TraceDecayError, ) -> Result<()> { let failed_key = core.current_key.lock().await.clone(); - let retain_core = !self.cancellation.is_cancelled() && failed_key == opened.key; + // A retiring store owner can settle after this attempt. Caching a + // degraded core would bypass composition on every later request and + // make its missing owners permanent instead of using admission retry. + let retain_core = !self.cancellation.is_cancelled() + && failed_key == opened.key + && !project_open_admission::is_observability_retiring(&error); let (core_retained, failed_full_server) = if retain_core { reclaim_core_after_failed_upgrade( self.store_administration, diff --git a/crates/tracedecay/src/daemon/project_delivery_mount.rs b/crates/tracedecay/src/daemon/project_delivery_mount.rs index 178cf996b1..1670437c24 100644 --- a/crates/tracedecay/src/daemon/project_delivery_mount.rs +++ b/crates/tracedecay/src/daemon/project_delivery_mount.rs @@ -32,8 +32,13 @@ pub(super) async fn ensure_project_delivery_settlement( configuration_policy_digest.clone(), ) .await - .map_err(|error| TraceDecayError::Config { - message: format!("project-open observability producer registration failed: {error}"), + .map_err(|error| match error { + error @ TraceDecayError::ProjectRoute { .. } => error, + error => TraceDecayError::Config { + message: format!( + "project-open observability producer registration failed: {error}" + ), + }, })?; Ok(configuration_policy_digest) } diff --git a/crates/tracedecay/src/daemon/project_open_admission.rs b/crates/tracedecay/src/daemon/project_open_admission.rs index eff0fb6730..6701ceebea 100644 --- a/crates/tracedecay/src/daemon/project_open_admission.rs +++ b/crates/tracedecay/src/daemon/project_open_admission.rs @@ -322,10 +322,16 @@ async fn wait_for_project_open_task(mut completion: tokio::sync::watch::Receiver } } +pub(super) fn is_observability_retiring(error: &TraceDecayError) -> bool { + matches!(error, TraceDecayError::ProjectRoute { reason_code, retryable: true, .. } + if reason_code == tracedecay_daemon_service::StoreObservabilityMountErrorV1::RETIRING_REASON_CODE) +} + /// How long a failed project-open route declines reopening, or `None` when the /// failure may clear on its own. pub(super) fn project_open_retry_backoff(error: &TraceDecayError) -> Option { match error { + error if is_observability_retiring(error) => Some(PROJECT_OPEN_RESOURCE_RETRY_BACKOFF), TraceDecayError::ProjectRoute { reason_code, .. } if reason_code == PROJECT_SERVER_CAPACITY_REASON_CODE => { diff --git a/crates/tracedecay/src/daemon/tests/runtime_identity.rs b/crates/tracedecay/src/daemon/tests/runtime_identity.rs index 784bd594cd..900781a5e5 100644 --- a/crates/tracedecay/src/daemon/tests/runtime_identity.rs +++ b/crates/tracedecay/src/daemon/tests/runtime_identity.rs @@ -865,3 +865,124 @@ async fn maintenance_reclaims_a_removed_linked_worktree_and_the_text_artifact_on .await .expect("shutdown after scope retention must remain bounded"); } + +#[tokio::test] +async fn retiring_observability_owner_reopens_same_route_without_cached_degradation() { + let home = TempDir::new().expect("isolated home"); + let root = home.path().canonicalize().unwrap(); + let project = root.join("project"); + std::fs::create_dir_all(project.join("src")).unwrap(); + std::fs::write(project.join("src/main.rs"), "fn main() {}\n").unwrap(); + let profile = root.join("profile"); + let identity = test_client_identity_for(profile.clone()); + initialize_test_project(&project, &identity).await; + let _scope = enter_test_daemon_database_scope(&profile, "retiring observability route"); + let engine = test_daemon_engine_for_profile(&profile); + let handshake = DaemonHandshake { + project_path: Some(project.clone()), + client_identity: identity, + ..test_handshake_defaults() + }; + let server = engine.project_server(&handshake).await.unwrap(); + let graph = server.cg().await; + let key = ProjectServerKey::from_open_project(&graph, &handshake).unwrap(); + let db = engine + .store_administration + .registered_project_session_database(&project, graph.store_layout()) + .await + .unwrap(); + let producer = engine + .invocation + .service + .observability_producer(Some(&project)) + .await + .unwrap(); + let blocker = db.begin_write_transaction().await.unwrap(); + assert!( + producer.cancel().await.is_err(), + "held writer prevents settlement" + ); + assert!( + producer.shutdown_joined(), + "the real producer worker must have joined" + ); + assert!(!producer.shutdown_settled()); + let roots = std::collections::BTreeSet::from([project.clone()]); + assert!( + engine + .invocation + .service + .project_runtimes + .quiesce_roots(&roots) + .await + .is_none() + ); + // Remove the initial published route to model a capacity reopen. From the + // first refusal onward, only production open/retry logic may change it. + assert!( + engine + .store_administration + .project_servers() + .lock() + .await + .remove(&key) + .is_some() + ); + drop(graph); + drop(server); + // Releasing the unrelated writer does not autonomously retry settlement; + // the first real mount must trigger it and still report Retiring. + blocker.rollback().await.unwrap(); + let error = match engine.project_server(&handshake).await { + Ok(_) => panic!("retiring owner must refuse the full project open"), + Err(error) => error, + }; + assert!( + super::super::project_open_admission::is_observability_retiring(&error), + "{error}" + ); + assert_eq!( + super::super::project_open_retry_backoff(&error), + Some(super::super::PROJECT_OPEN_RESOURCE_RETRY_BACKOFF) + ); + assert!( + engine + .store_administration + .project_servers() + .lock() + .await + .get_ready(&key) + .is_none(), + "a transient owner refusal must not cache a permanently degraded core" + ); + let reopened = tokio::time::timeout(std::time::Duration::from_secs(10), async { + loop { + match engine.project_server(&handshake).await { + Ok(server) => break server, + Err(error) + if super::super::project_open_admission::is_observability_retiring(&error) => + { + tokio::time::sleep(std::time::Duration::from_millis(10)).await; + } + Err(error) => panic!("reopen lost its typed recovery state: {error}"), + } + } + }) + .await + .expect("same route must recover after the writer releases"); + let graph = reopened.cg().await; + let reopened_key = ProjectServerKey::from_open_project(&graph, &handshake).unwrap(); + assert_eq!( + key, reopened_key, + "recovery preserves the canonical store identity" + ); + assert_reopened_linked_route_is_not_degraded( + &engine.store_administration.project_servers().lock().await, + &key, + ); + drop(graph); + drop(reopened); + drop(producer); + drop(db); + engine.shutdown_all().await; +} diff --git a/crates/tracedecay/tests/daemon_suite/invocation_observability.rs b/crates/tracedecay/tests/daemon_suite/invocation_observability.rs index 589ad56faf..19c35cb9a6 100644 --- a/crates/tracedecay/tests/daemon_suite/invocation_observability.rs +++ b/crates/tracedecay/tests/daemon_suite/invocation_observability.rs @@ -1389,7 +1389,33 @@ async fn registered_persistence_failure_releases_only_after_join_and_writer_fenc } #[tokio::test] -async fn registered_shutdown_reports_a_blocked_producer_flush() { +async fn registered_blocked_writer_retries_settlement_only_on_mount_demand() { + assert_blocked_writer_reopens( + ObservabilityProducerDeadlinesV1 { + persistence: Duration::from_millis(50), + shutdown: Duration::from_millis(250), + }, + "observability_persistence_deadline", + ) + .await; +} + +#[tokio::test] +async fn registered_aborted_worker_reopens_only_after_writer_settlement() { + assert_blocked_writer_reopens( + ObservabilityProducerDeadlinesV1 { + persistence: Duration::from_millis(250), + shutdown: Duration::from_millis(50), + }, + "observability_shutdown_deadline", + ) + .await; +} + +async fn assert_blocked_writer_reopens( + deadlines: ObservabilityProducerDeadlinesV1, + expected_error: &str, +) { let (_project, project_id, database, _runtime) = runtime("observability-shutdown-failure").await; let identity = ObservabilityProducerIdentityV1 { @@ -1403,10 +1429,7 @@ async fn registered_shutdown_reports_a_blocked_producer_flush() { database.clone(), identity.clone(), 1, - ObservabilityProducerDeadlinesV1 { - persistence: Duration::from_millis(50), - shutdown: Duration::from_millis(250), - }, + deadlines, ) .expect("producer"); let registry = StoreObservabilityRegistryV1::default(); @@ -1432,11 +1455,8 @@ async fn registered_shutdown_reports_a_blocked_producer_flush() { !retained.shutdown_settled(), "an unfinished writer fence cannot release the store" ); - blocker.commit().await.expect("release registered writer"); assert!( - error - .to_string() - .contains("observability_persistence_deadline"), + error.to_string().contains(expected_error), "unexpected shutdown error: {error}" ); let start_called = Arc::new(AtomicBool::new(false)); @@ -1451,9 +1471,72 @@ async fn registered_shutdown_reports_a_blocked_producer_flush() { BoundedObservabilityProducerV1::start(database.clone(), replacement_identity.clone(), 1) .map_err(StoreObservabilityMountErrorV1::Unavailable) }); + assert_eq!( + retained + .try_emit(envelope(&project_id, "closed:while-blocked")) + .unwrap_err(), + "observability_producer_closed" + ); + assert!( + retained.shutdown_joined(), + "worker closed before fence retry" + ); assert!(matches!( failed, - Err(StoreObservabilityMountErrorV1::ShutdownFailed) + Err(StoreObservabilityMountErrorV1::Retiring) )); assert!(!start_called.load(Ordering::Acquire)); + registry.join_retirement_drains().await; + assert!( + !retained.shutdown_settled(), + "blocked recheck must not release ownership" + ); + + // The same logical project in another registered store remains independent. + let (_other_project, other_id, other_database, _other_runtime) = + runtime("observability-shutdown-failure").await; + assert_eq!(project_id, other_id); + let other = registry + .acquire_or_start(&other_database, &store_mount(&identity), || { + BoundedObservabilityProducerV1::start(other_database.clone(), identity.clone(), 1) + .map_err(StoreObservabilityMountErrorV1::Unavailable) + }) + .expect("foreign exact store is not retiring"); + + blocker.commit().await.expect("release registered writer"); + tokio::task::yield_now().await; + assert!( + !retained.shutdown_settled(), + "release alone must not trigger a retry loop" + ); + let retry = registry.acquire_or_start(&database, &replacement_mount, || { + panic!("writer settlement must precede replacement") + }); + assert!(matches!( + retry, + Err(StoreObservabilityMountErrorV1::Retiring) + )); + registry.join_retirement_drains().await; + assert!(retained.shutdown_settled()); + assert_eq!( + retained + .try_emit(envelope(&project_id, "closed:after-retry")) + .unwrap_err(), + "observability_producer_closed" + ); + other + .producer() + .try_emit(envelope(&other_id, "foreign:still-active")) + .unwrap(); + let replacement = registry + .acquire_or_start(&database, &replacement_mount, || { + BoundedObservabilityProducerV1::start(database.clone(), replacement_identity.clone(), 1) + .map_err(StoreObservabilityMountErrorV1::Unavailable) + }) + .expect("settled exact store reopens without reset"); + replacement.shutdown().await.expect("replacement drains"); + other + .shutdown() + .await + .expect("foreign owner drains independently"); } From 8c2d8fcc334bd87f244b84b4e32d760d1f19ad1c Mon Sep 17 00:00:00 2001 From: ScriptedAlchemy Date: Thu, 8 Oct 2026 00:22:07 -0700 Subject: [PATCH 5/6] fix(observability): keep aborted coverage failed closed --- .../src/observability/producer.rs | 30 ++++++------------- .../src/daemon/tests/runtime_identity.rs | 2 +- .../daemon_suite/invocation_observability.rs | 26 +++++++++++++--- 3 files changed, 32 insertions(+), 26 deletions(-) diff --git a/crates/tracedecay-application/src/observability/producer.rs b/crates/tracedecay-application/src/observability/producer.rs index af3ce01ad9..ae50a1942f 100644 --- a/crates/tracedecay-application/src/observability/producer.rs +++ b/crates/tracedecay-application/src/observability/producer.rs @@ -439,8 +439,9 @@ impl BoundedObservabilityProducerV1 { self.core.state.load(Ordering::Acquire) == PRODUCER_SETTLED } - /// The worker and admissions are closed, but timed-out database commands - /// may still need the canonical writer fence before the store can reopen. + /// Admissions closed and the worker completed its coverage settlement + /// attempt before joining. Timed-out database commands may still need the + /// canonical writer fence; a forced abort never establishes this state. pub fn shutdown_joined(&self) -> bool { matches!( self.core.state.load(Ordering::Acquire), @@ -587,26 +588,13 @@ impl ObservabilityProducerCoreV1 { )); } Err(_) => { - let joined = if let Some(worker) = worker.take() { + if let Some(worker) = worker.take() { worker.abort(); - match worker.await { - Ok(()) => true, - Err(error) => error.is_cancelled(), - } - } else { - false - }; - // Both admission fences passed before this worker wait. An - // aborted, joined worker can leave only database commands, - // which a later bounded writer fence can settle safely. - self.state.store( - if joined { - PRODUCER_JOINED - } else { - PRODUCER_STOPPED - }, - Ordering::Release, - ); + let _ = worker.await; + } + // Abort may discard accepted observations before their terminal + // coverage is attempted; writer settlement cannot repair that. + self.state.store(PRODUCER_STOPPED, Ordering::Release); return Err(ApplicationContractError::Domain( "observability_shutdown_deadline".to_owned(), )); diff --git a/crates/tracedecay/src/daemon/tests/runtime_identity.rs b/crates/tracedecay/src/daemon/tests/runtime_identity.rs index 900781a5e5..f57c9d08dc 100644 --- a/crates/tracedecay/src/daemon/tests/runtime_identity.rs +++ b/crates/tracedecay/src/daemon/tests/runtime_identity.rs @@ -977,7 +977,7 @@ async fn retiring_observability_owner_reopens_same_route_without_cached_degradat "recovery preserves the canonical store identity" ); assert_reopened_linked_route_is_not_degraded( - &engine.store_administration.project_servers().lock().await, + &*engine.store_administration.project_servers().lock().await, &key, ); drop(graph); diff --git a/crates/tracedecay/tests/daemon_suite/invocation_observability.rs b/crates/tracedecay/tests/daemon_suite/invocation_observability.rs index 19c35cb9a6..98e6328bdd 100644 --- a/crates/tracedecay/tests/daemon_suite/invocation_observability.rs +++ b/crates/tracedecay/tests/daemon_suite/invocation_observability.rs @@ -1390,31 +1390,34 @@ async fn registered_persistence_failure_releases_only_after_join_and_writer_fenc #[tokio::test] async fn registered_blocked_writer_retries_settlement_only_on_mount_demand() { - assert_blocked_writer_reopens( + assert_blocked_writer_retirement( ObservabilityProducerDeadlinesV1 { persistence: Duration::from_millis(50), shutdown: Duration::from_millis(250), }, "observability_persistence_deadline", + true, ) .await; } #[tokio::test] -async fn registered_aborted_worker_reopens_only_after_writer_settlement() { - assert_blocked_writer_reopens( +async fn registered_aborted_worker_keeps_unsettled_coverage_failed_closed() { + assert_blocked_writer_retirement( ObservabilityProducerDeadlinesV1 { persistence: Duration::from_millis(250), shutdown: Duration::from_millis(50), }, "observability_shutdown_deadline", + false, ) .await; } -async fn assert_blocked_writer_reopens( +async fn assert_blocked_writer_retirement( deadlines: ObservabilityProducerDeadlinesV1, expected_error: &str, + worker_settled: bool, ) { let (_project, project_id, database, _runtime) = runtime("observability-shutdown-failure").await; @@ -1477,6 +1480,21 @@ async fn assert_blocked_writer_reopens( .unwrap_err(), "observability_producer_closed" ); + if !worker_settled { + assert!( + !retained.shutdown_joined(), + "abort never proves coverage settlement" + ); + assert!(matches!( + failed, + Err(StoreObservabilityMountErrorV1::ShutdownFailed) + )); + assert!(!start_called.load(Ordering::Acquire)); + blocker.commit().await.expect("release registered writer"); + assert!(retained.finish_shutdown_settlement().await.is_err()); + assert!(!retained.shutdown_settled()); + return; + } assert!( retained.shutdown_joined(), "worker closed before fence retry" From b4d17d9a4c9f06836a0fe32ac2848a57a1d96e3f Mon Sep 17 00:00:00 2001 From: ScriptedAlchemy Date: Thu, 8 Oct 2026 00:34:20 -0700 Subject: [PATCH 6/6] fix(observability): retain durable terminal coverage proof --- .../src/observability/producer.rs | 156 ++++++++++++++---- .../daemon_suite/invocation_observability.rs | 2 +- 2 files changed, 124 insertions(+), 34 deletions(-) diff --git a/crates/tracedecay-application/src/observability/producer.rs b/crates/tracedecay-application/src/observability/producer.rs index ae50a1942f..db209e077e 100644 --- a/crates/tracedecay-application/src/observability/producer.rs +++ b/crates/tracedecay-application/src/observability/producer.rs @@ -1,6 +1,6 @@ use std::sync::Arc; use std::sync::Mutex; -use std::sync::atomic::{AtomicU8, AtomicU64, Ordering}; +use std::sync::atomic::{AtomicBool, AtomicU8, AtomicU64, Ordering}; use std::time::Duration; use serde::{Deserialize, Serialize}; @@ -173,6 +173,7 @@ struct ProducerWorkerState { total_dropped: Arc, next_sequence: Arc, lifecycle: Arc, + final_coverage_persisted: Arc, durable_emission_lock: Arc>, deadlines: ObservabilityProducerDeadlinesV1, } @@ -196,7 +197,7 @@ struct ObservabilityProducerCoreV1 { identity: ObservabilityProducerIdentityV1, data: mpsc::Sender, control: mpsc::Sender, - // The next five stay `Arc` because the spawned worker shares them. The + // The shared worker fields stay `Arc` because the spawned worker shares them. The // worker must not hold the core itself: the queue senders live in the // core, so a worker-held core would keep its own channels open and the // worker could never wind down when every frontend is dropped. @@ -204,6 +205,7 @@ struct ObservabilityProducerCoreV1 { total_dropped: Arc, next_sequence: Arc, state: Arc, + final_coverage_persisted: Arc, durable_emission_lock: Arc>, deadlines: ObservabilityProducerDeadlinesV1, emission_lock: Mutex<()>, @@ -249,6 +251,7 @@ impl BoundedObservabilityProducerV1 { let total_dropped = Arc::new(AtomicU64::new(0)); let next_sequence = Arc::new(AtomicU64::new(1)); let state = Arc::new(AtomicU8::new(PRODUCER_RUNNING)); + let final_coverage_persisted = Arc::new(AtomicBool::new(false)); let durable_emission_lock = Arc::new(AsyncMutex::new(())); let runtime = tokio::runtime::Handle::try_current() .map_err(|_| "observability_producer_runtime_unavailable")?; @@ -262,6 +265,7 @@ impl BoundedObservabilityProducerV1 { total_dropped: Arc::clone(&total_dropped), next_sequence: Arc::clone(&next_sequence), lifecycle: Arc::clone(&state), + final_coverage_persisted: Arc::clone(&final_coverage_persisted), durable_emission_lock: Arc::clone(&durable_emission_lock), deadlines, }, @@ -275,6 +279,7 @@ impl BoundedObservabilityProducerV1 { total_dropped, next_sequence, state, + final_coverage_persisted, durable_emission_lock, deadlines, emission_lock: Mutex::new(()), @@ -439,9 +444,9 @@ impl BoundedObservabilityProducerV1 { self.core.state.load(Ordering::Acquire) == PRODUCER_SETTLED } - /// Admissions closed and the worker completed its coverage settlement - /// attempt before joining. Timed-out database commands may still need the - /// canonical writer fence; a forced abort never establishes this state. + /// Admissions closed and the worker joined after normal settlement or a + /// proved durable final carrier. Timed-out database commands may still need + /// the writer fence; abort before final coverage never establishes this state. pub fn shutdown_joined(&self) -> bool { matches!( self.core.state.load(Ordering::Acquire), @@ -588,13 +593,27 @@ impl ObservabilityProducerCoreV1 { )); } Err(_) => { - if let Some(worker) = worker.take() { + let joined = if let Some(worker) = worker.take() { worker.abort(); - let _ = worker.await; - } - // Abort may discard accepted observations before their terminal - // coverage is attempted; writer settlement cannot repair that. - self.state.store(PRODUCER_STOPPED, Ordering::Release); + match worker.await { + Ok(()) => true, + Err(error) => error.is_cancelled(), + } + } else { + false + }; + // Only a durable final carrier proves an abort left no queue + // coverage to finish. Optional maintenance may still need the + // writer fence; an earlier abort must remain failed closed. + let covered = self.final_coverage_persisted.load(Ordering::Acquire); + self.state.store( + if joined && covered { + PRODUCER_JOINED + } else { + PRODUCER_STOPPED + }, + Ordering::Release, + ); return Err(ApplicationContractError::Domain( "observability_shutdown_deadline".to_owned(), )); @@ -1067,6 +1086,11 @@ async fn settle_worker( ) .await; } + // Published only after normal drain's final mandatory carrier succeeds. + // Earlier batch success and cancellation do not prove final coverage. + state + .final_coverage_persisted + .store(clean_shutdown_observed && writer_settled, Ordering::Release); // The mandatory terminal precedes optional maintenance so its day // can be rebuilt before shutdown. A deferred maintenance operation // may still own a timed-out database command and requires the fence. @@ -1300,15 +1324,19 @@ mod cheaper_frontier_tests { mod shutdown_settlement_tests { use super::{ BoundedObservabilityProducerV1, ObservabilityProducerDeadlinesV1, - ObservabilityProducerIdentityV1, + ObservabilityProducerIdentityV1, ProducerControl, }; use std::sync::Arc; + use std::sync::atomic::Ordering; use std::time::Duration; - use tokio::sync::oneshot; + use tokio::sync::{mpsc, oneshot}; use tracedecay_global_db::tests::harness::RegisteredGlobalDbTestRuntime; - #[tokio::test] - async fn completed_terminal_does_not_wait_for_a_subsequent_unrelated_writer() { + async fn fixture() -> ( + tempfile::TempDir, + RegisteredGlobalDbTestRuntime, + BoundedObservabilityProducerV1, + ) { let directory = tempfile::tempdir().unwrap(); let project = directory.path().join("project"); std::fs::create_dir(&project).unwrap(); @@ -1321,24 +1349,30 @@ mod shutdown_settlement_tests { .await .unwrap(); let db = runtime.project_database_arc().unwrap(); - let producer = Arc::new( - BoundedObservabilityProducerV1::start_with_deadlines( - db.clone(), - ObservabilityProducerIdentityV1 { - authorized_scope_ref: project_id.as_str().to_owned(), - process_boot_id: "boot.shutdown-terminal".to_owned(), - producer_revision: "producer.v1".to_owned(), - configuration_revision: "configuration.v1".to_owned(), - policy_revision: "policy.v1".to_owned(), - }, - 1, - ObservabilityProducerDeadlinesV1 { - persistence: Duration::from_millis(50), - shutdown: Duration::from_millis(250), - }, - ) - .unwrap(), - ); + let producer = BoundedObservabilityProducerV1::start_with_deadlines( + db.clone(), + ObservabilityProducerIdentityV1 { + authorized_scope_ref: project_id.as_str().to_owned(), + process_boot_id: "boot.shutdown-terminal".to_owned(), + producer_revision: "producer.v1".to_owned(), + configuration_revision: "configuration.v1".to_owned(), + policy_revision: "policy.v1".to_owned(), + }, + 1, + ObservabilityProducerDeadlinesV1 { + persistence: Duration::from_millis(50), + shutdown: Duration::from_millis(250), + }, + ) + .unwrap(); + (directory, runtime, producer) + } + + #[tokio::test] + async fn completed_terminal_does_not_wait_for_a_subsequent_unrelated_writer() { + let (_directory, runtime, producer) = fixture().await; + let db = runtime.project_database_arc().unwrap(); + let producer = Arc::new(producer); assert!(!producer.shutdown_joined()); let premature = producer.finish_shutdown_settlement().await.unwrap_err(); assert!( @@ -1373,4 +1407,60 @@ mod shutdown_settlement_tests { assert_eq!(summary.persisted, 1); assert!(producer.shutdown_settled()); } + + #[tokio::test] + async fn final_coverage_survives_shutdown_reply_deadline_and_fences_before_reuse() { + let (_directory, runtime, mut producer) = fixture().await; + let db = runtime.project_database_arc().unwrap(); + let (control, mut incoming) = mpsc::channel(1); + let core = Arc::get_mut(&mut producer.core).unwrap(); + let actual_control = std::mem::replace(&mut core.control, control); + let (covered, coverage_ready) = oneshot::channel(); + let (release_reply, released_reply) = oneshot::channel(); + // Forward the real shutdown and delay only its reply. The production + // worker must actually persist final coverage; the test never sets it. + let forwarding = tokio::spawn(async move { + let ProducerControl::Shutdown { cancelled, reply } = incoming.recv().await.unwrap(); + let (actual_reply, actual_result) = oneshot::channel(); + assert!( + actual_control + .try_send(ProducerControl::Shutdown { + cancelled, + reply: actual_reply, + }) + .is_ok() + ); + let (result, _) = actual_result.await.unwrap(); + assert_eq!(result.unwrap().persisted, 1); + let blocker = db.begin_write_transaction().await.unwrap(); + covered.send(()).unwrap(); + released_reply.await.unwrap(); + blocker.rollback().await.unwrap(); + drop(reply); + }); + let producer = Arc::new(producer); + let stopping = Arc::clone(&producer); + let shutdown = tokio::spawn(async move { stopping.shutdown().await }); + coverage_ready.await.unwrap(); + assert!( + producer + .core + .final_coverage_persisted + .load(Ordering::Acquire) + ); + let error = shutdown.await.unwrap().unwrap_err(); + assert!( + error + .to_string() + .contains("observability_shutdown_deadline") + ); + assert!(producer.shutdown_joined()); + assert!(!producer.shutdown_settled()); + assert!(producer.finish_shutdown_settlement().await.is_err()); + assert!(!producer.shutdown_settled()); + release_reply.send(()).unwrap(); + forwarding.await.unwrap(); + producer.finish_shutdown_settlement().await.unwrap(); + assert!(producer.shutdown_settled()); + } } diff --git a/crates/tracedecay/tests/daemon_suite/invocation_observability.rs b/crates/tracedecay/tests/daemon_suite/invocation_observability.rs index 98e6328bdd..76cd49c225 100644 --- a/crates/tracedecay/tests/daemon_suite/invocation_observability.rs +++ b/crates/tracedecay/tests/daemon_suite/invocation_observability.rs @@ -1405,7 +1405,7 @@ async fn registered_blocked_writer_retries_settlement_only_on_mount_demand() { async fn registered_aborted_worker_keeps_unsettled_coverage_failed_closed() { assert_blocked_writer_retirement( ObservabilityProducerDeadlinesV1 { - persistence: Duration::from_millis(250), + persistence: Duration::from_millis(50), shutdown: Duration::from_millis(50), }, "observability_shutdown_deadline",