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
302 changes: 279 additions & 23 deletions crates/tracedecay-application/src/observability/producer.rs

Large diffs are not rendered by default.

Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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!(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -54,7 +54,12 @@ impl StoreObservabilityCoreV1 {
})
}

async fn shutdown(&self) -> Result<(), tracedecay_contracts::ApplicationContractError> {
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");
Expand All @@ -66,13 +71,23 @@ 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),
if !ancillary_joined || !self.producer.shutdown_joined() {
StoreObservabilityCompletionV1::Failed
} else if self.producer.shutdown_settled() {
StoreObservabilityCompletionV1::Settled
} else {
StoreObservabilityCompletionV1::AwaitingWriter
},
)
}
}

Expand Down Expand Up @@ -100,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 {
Expand Down Expand Up @@ -176,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,
};
Expand Down Expand Up @@ -289,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)
}
Expand Down Expand Up @@ -383,15 +413,26 @@ impl StoreObservabilityRegistryV1 {
&self,
runtime: &tokio::runtime::Handle,
core: Arc<StoreObservabilityCoreV1>,
writer_only: bool,
) {
let registry = self.clone();
self.retirement_drains.spawn_on(
async move {
let result = 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");
}
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");
}
},
Expand All @@ -410,13 +451,13 @@ 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.
/// 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<StoreObservabilityCoreV1>,
releasable: bool,
completion: StoreObservabilityCompletionV1,
) -> Result<(), ApplicationContractError> {
let mut entries =
self.lock_entries()
Expand All @@ -435,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(())
}
Expand Down Expand Up @@ -588,8 +638,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) => {
Expand Down Expand Up @@ -633,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"
Expand Down
147 changes: 147 additions & 0 deletions crates/tracedecay-global-db/src/observability_outbox_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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::<Vec<_>>();
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;
Expand Down
Loading
Loading