Skip to main content

mkit_server/pipeline/
durable_outcome.rs

1//! Public, transport-neutral terminal outcomes for the delivery sink.
2
3use core::time::Duration;
4
5use crate::error::Redacted;
6use crate::store::codec::{AbortReason, OutcomeRef, ReservationV1};
7
8/// One durable reservation result. Delivery is at least once.
9#[derive(Debug, Clone, PartialEq, Eq)]
10#[non_exhaustive]
11pub struct Outcome {
12    /// Idempotency key; sinks deduplicate by this id.
13    pub reservation_id: String,
14    /// Canonical origin of the server that performed the operation.
15    pub audience: String,
16    /// Full repository identity.
17    pub repository: String,
18    /// Time the terminal result occurred.
19    pub occurred_unix_ms: i64,
20    /// Terminal result payload.
21    pub kind: OutcomeKind,
22}
23
24/// Hooks.v1 terminal variants, field for field.
25#[derive(Debug, Clone, PartialEq, Eq)]
26#[non_exhaustive]
27pub enum OutcomeKind {
28    /// A write committed durably.
29    Committed {
30        bytes_stored: u64,
31        new_to_repo: u64,
32        new_to_store: u64,
33        refs: Vec<OutcomeRef>,
34    },
35    /// A write or read was aborted.
36    Aborted { reason: AbortReason, detail: String },
37    /// An unused ticket expired.
38    Expired,
39    /// An HTTP read completed, possibly after partial transmission.
40    ReadServed { object: [u8; 32], bytes_served: u64 },
41}
42
43impl Outcome {
44    /// Map a terminal stored row to the delivery contract.
45    ///
46    /// # Errors
47    /// A pending or ticketed row is not deliverable.
48    pub fn from_reservation(
49        reservation_id: String,
50        audience: String,
51        row: ReservationV1,
52    ) -> Result<Self, &'static str> {
53        let (repository, occurred, kind) = match row {
54            ReservationV1::Committed {
55                repository,
56                occurred_at_ms,
57                bytes_stored,
58                new_to_repo,
59                new_to_store,
60                refs,
61            } => (
62                repository,
63                occurred_at_ms,
64                OutcomeKind::Committed {
65                    bytes_stored,
66                    new_to_repo,
67                    new_to_store,
68                    refs,
69                },
70            ),
71            ReservationV1::Aborted {
72                repository,
73                occurred_at_ms,
74                reason,
75                detail,
76            } => (
77                repository,
78                occurred_at_ms,
79                OutcomeKind::Aborted { reason, detail },
80            ),
81            ReservationV1::Expired {
82                repository,
83                occurred_at_ms,
84            } => (repository, occurred_at_ms, OutcomeKind::Expired),
85            ReservationV1::ReadServed {
86                repository,
87                occurred_at_ms,
88                object,
89                bytes_served,
90            } => (
91                repository,
92                occurred_at_ms,
93                OutcomeKind::ReadServed {
94                    object,
95                    bytes_served,
96                },
97            ),
98            ReservationV1::Pending { .. } | ReservationV1::Ticketed { .. } => {
99                return Err("outcome row is not terminal");
100            }
101        };
102        Ok(Self {
103            reservation_id,
104            audience,
105            repository,
106            occurred_unix_ms: i64::try_from(occurred).map_err(|_| "outcome timestamp overflow")?,
107            kind,
108        })
109    }
110}
111
112/// A failed sink call always means retry. Its reason is never printed.
113#[derive(Debug, Clone)]
114#[non_exhaustive]
115pub struct DeliveryError {
116    /// Operator-only diagnostic text.
117    pub reason: Redacted,
118    /// Optional sink retry hint.
119    pub retry_after: Option<Duration>,
120}
121
122impl core::fmt::Display for DeliveryError {
123    fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
124        f.write_str("outcome delivery failed")
125    }
126}
127
128impl std::error::Error for DeliveryError {}
129
130impl DeliveryError {
131    /// Construct a retryable delivery error.
132    #[must_use]
133    pub fn new(reason: impl Into<String>, retry_after: Option<Duration>) -> Self {
134        Self {
135            reason: Redacted::new(reason),
136            retry_after,
137        }
138    }
139}