saddle-runtime 0.3.25

Saddle managed asynchronous runtime and lifecycle
Documentation
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>()
    );
}