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