use super::*;
use saddle_core::*;
use saddle_observability::{
EmergencyDiagnostics, FileLoggingConfig, ObserverConfig, Rotation, Stage, StageOutcome,
};
fn outcome(operation: OperationOutcome) -> RootOutcomeFacts {
RootOutcomeFacts {
axes: DiagnosticOutcomeAxes {
operation,
..Default::default()
},
..Default::default()
}
}
fn code() -> DiagnosticCode {
DiagnosticCode::new("rg.original").unwrap()
}
fn count(observer: &Observer, stage: Stage) -> u64 {
observer
.metrics_snapshot()
.stage_latency(stage)
.iter()
.sum()
}
#[test]
fn guarded_sources_stages_foreign_and_last_ref() {
use saddle_admission::ProfuseGwLightweightAdmissionOutcome as Admission;
let path = std::env::var_os("RG_EVIDENCE")
.map(std::path::PathBuf::from)
.unwrap_or_else(|| std::env::temp_dir().join(format!("rg-{}", std::process::id())));
std::fs::create_dir_all(&path).unwrap();
let output =
EmergencyDiagnostics::start(&FileLoggingConfig::new(path.clone(), Rotation::Daily))
.unwrap();
let handle = output.handle();
let observer = Observer::with_writer(ObserverConfig::default(), std::io::sink()).unwrap();
let process = crate::request_task::reserved::tests::process();
let Admission::Ready(permit) = process.verified_profile().try_admit() else {
panic!("actual admission")
};
let (permit, root) =
ReservedRequestRoot::try_admitted(permit, ContextLabel::checked("rg").unwrap())
.ok()
.unwrap();
let view = root.view(RequestViewPhase::Reading).unwrap();
let foreign_root = ReservedRequestRoot::create(
process
.verified_profile()
.try_rejection_storage(root_demand().unwrap())
.unwrap(),
ContextLabel::checked("other").unwrap(),
)
.unwrap();
let foreign = foreign_root.view(RequestViewPhase::Reading).unwrap();
let primary = BoundedDiagnostic::capture(
DiagnosticCategory::UnexpectedError,
CaptureSite::FirstObserved,
BoundedDiagnosticCause::new(DiagnosticStage::RequestDb, code()),
);
let facts = BoundedDiagnostic::capture(
DiagnosticCategory::UnexpectedError,
CaptureSite::FirstObserved,
BoundedDiagnosticCause::new(DiagnosticStage::RequestDb, code())
.with_column(Some(2), Some(4))
.with_types(
Some(InlineDiagnosticText::metadata("u64")),
Some(InlineDiagnosticText::metadata("VARCHAR")),
None,
),
)
.during_cleanup_of(&primary);
let expected = serde_json::to_value(&facts).unwrap();
let error = std::io::Error::new(std::io::ErrorKind::PermissionDenied, "RG_ORIGINAL_中文");
let failure = view.source_error_with_facts(
&error,
facts,
code(),
Some(&handle),
RootRequestEvent::Database,
outcome(OperationOutcome::Failed),
);
assert_eq!(
serde_json::to_value(failure.occurrence()).unwrap()["diagnostic_id"],
expected["diagnostic_id"]
);
drop(error);
let stage = foreign.start_stage(&observer, Some(&handle), ReservedObservationStage::Database);
let (stage, failure) = stage
.finish_failure(failure, outcome(OperationOutcome::Failed))
.err()
.unwrap();
assert_eq!(
count(&observer, Stage::Database),
0,
"foreign does not finish or cancel interval"
);
let own = foreign.source_description(
&"foreign own",
BoundedDiagnostic::capture(
DiagnosticCategory::ExpectedRejection,
CaptureSite::FirstObserved,
BoundedDiagnosticCause::new(DiagnosticStage::RequestDb, code()),
),
code(),
None,
RootRequestEvent::Database,
outcome(OperationOutcome::Rejected),
);
let (own, _) = stage
.finish_failure(own, outcome(OperationOutcome::Rejected))
.ok()
.unwrap();
let _ = own.finish(&foreign, None, Default::default()).ok().unwrap();
let (failure, _) = view
.start_stage(&observer, Some(&handle), ReservedObservationStage::Database)
.finish_failure(failure, outcome(OperationOutcome::Failed))
.ok()
.unwrap();
assert_eq!(count(&observer, Stage::Database), 2);
let public = failure
.finish(&view, Some(&handle), Default::default())
.ok()
.unwrap();
assert_eq!(
count(&observer, Stage::Database),
2,
"final projection must not repeat stage metric"
);
let facts = BoundedDiagnostic::capture(
DiagnosticCategory::ExpectedRejection,
CaptureSite::FirstObserved,
BoundedDiagnosticCause::new(DiagnosticStage::RequestDb, code()),
)
.during_cleanup_occurrence(primary.occurrence());
let expected_description = serde_json::to_value(&facts).unwrap();
let failure = view.source_description(
&"RG_NO_ERROR",
facts,
code(),
Some(&handle),
RootRequestEvent::Outbound,
outcome(OperationOutcome::Cancelled),
);
let (failure, submission) = view
.start_stage(&observer, None, ReservedObservationStage::Outbound)
.finish_failure(failure, outcome(OperationOutcome::Cancelled))
.ok()
.unwrap();
assert_eq!(submission, DiagnosticSubmission::OutputUnavailable);
let _ = failure
.finish(&view, None, Default::default())
.ok()
.unwrap();
assert_eq!(count(&observer, Stage::ProfuseContract), 1);
let legacy = Diagnostic::capture(
DiagnosticCategory::ExpectedRejection,
CaptureSite::FirstObserved,
DiagnosticCause::new(DiagnosticStage::RequestDb, code()),
)
.during_cleanup_of_occurrence(&primary.occurrence());
let expected_legacy = serde_json::to_value(&legacy).unwrap();
let failure = view.source_existing_description(
&"RG_LEGACY_CANCEL",
legacy,
code(),
Some(&handle),
RootRequestEvent::Outbound,
outcome(OperationOutcome::Cancelled),
);
let _ = failure
.finish(&view, None, Default::default())
.ok()
.unwrap();
let legacy = Diagnostic::capture(
DiagnosticCategory::UnexpectedError,
CaptureSite::FirstObserved,
DiagnosticCause::new(DiagnosticStage::RequestDb, code()),
);
let expected_legacy_error = serde_json::to_value(&legacy).unwrap();
let failure = view.source_existing_error(
&std::io::Error::other("RG_LEGACY_ERROR"),
legacy,
code(),
Some(&handle),
RootRequestEvent::Database,
outcome(OperationOutcome::Failed),
);
let _ = failure
.finish(&view, None, Default::default())
.ok()
.unwrap();
let stage = view.start_stage(&observer, None, ReservedObservationStage::Response);
let stage = stage
.finish_nonfailure(outcome(OperationOutcome::Failed))
.err()
.unwrap();
assert_eq!(count(&observer, Stage::Response), 0);
assert_eq!(
stage
.finish_nonfailure(outcome(OperationOutcome::Rejected))
.ok()
.unwrap(),
DiagnosticSubmission::OutputUnavailable
);
assert_eq!(
observer.metrics_snapshot().requests(StageOutcome::Rejected),
1
);
assert_eq!(
observer
.metrics_snapshot()
.requests(StageOutcome::Cancelled),
0
);
view.start_stage(&observer, None, ReservedObservationStage::Handler)
.finish_nonfailure(outcome(OperationOutcome::Succeeded))
.ok()
.unwrap();
assert_eq!(count(&observer, Stage::Handler), 1);
drop(view.start_stage(&observer, None, ReservedObservationStage::Handler));
assert_eq!(
count(&observer, Stage::Handler),
2,
"unfinished Drop records exactly one cancellation"
);
drop((root, foreign, foreign_root));
permit.into_execution().cancel_observed();
assert!(matches!(
process.verified_profile().try_admit(),
Admission::CapacityRejected
));
drop(view);
let Admission::Ready(next) = process.verified_profile().try_admit() else {
panic!("last guard refunded")
};
drop(next);
assert!(
serde_json::to_value(public.occurrence()).unwrap()["diagnostic_id"]
.as_u64()
.unwrap()
> 0,
"public failure remains alive after refund"
);
process.finish().unwrap();
let exit = crate::diagnostics::close_output(
output,
Some(std::time::Instant::now() + std::time::Duration::from_secs(2)),
);
assert_eq!(exit.snapshot.enqueued, exit.snapshot.written);
assert_eq!(exit.snapshot.dropped, 0);
let raw = std::fs::read_to_string(path.join("saddle.emergency.log")).unwrap();
let rows: Vec<serde_json::Value> = raw
.lines()
.map(|s| serde_json::from_str(s).unwrap())
.collect();
let headers: Vec<serde_json::Value> = rows
.iter()
.filter(|r| r["channel"] == "context")
.map(|r| serde_json::from_str(r["payload"].as_str().unwrap()).unwrap())
.collect();
for expected in [
&expected,
&expected_description,
&expected_legacy,
&expected_legacy_error,
] {
assert!(
headers.iter().any(|h| h["facts"] == *expected),
"exact original ID/location/primary/facts preserved"
);
}
assert!(headers.iter().any(
|h| h["source_interface"] == "unavailable_description_debug_only"
&& h["facts"] == expected_description
));
assert!(
raw.contains("RG_ORIGINAL_中文")
&& raw.contains("PermissionDenied")
&& raw.contains("RG_NO_ERROR")
);
let boundaries: Vec<_> = rows
.iter()
.filter(|r| r["event"] == "request_failure_boundary")
.collect();
assert_eq!(
boundaries
.iter()
.filter(|r| r["stage"] == "database")
.count(),
2
);
assert_eq!(
boundaries
.iter()
.filter(|r| r["stage"] == "finalization")
.count(),
1
);
println!(
"RG PASS stage={} align={} guard_ref={} root_free_public={} originals=4 foreign_retry=PASS no_output_metrics=PASS last_guard_refund=PASS",
std::mem::size_of::<ReservedActiveStage>(),
std::mem::align_of::<ReservedActiveStage>(),
std::mem::size_of::<&ReservedRequestView>(),
std::mem::size_of::<PublicRequestFailure>()
);
}