1pub mod account_metadata;
3pub mod client;
4pub mod custodial_recovery;
5pub mod identity_management;
6pub mod invitation;
7pub mod notifications;
8pub mod passkey_label;
9pub mod read_budget;
10pub mod spool_settings;
11use crate::StreamingShape;
12use crate::heddle::api::common::{
13 AuthorizationAccess, AuthorizationExistence, AuthorizationRole, AuthorizationScopeSource,
14 CallContext, DeploymentTarget, RetryBehavior, RpcEffect, ServiceMaturity, SigningTier,
15};
16use crate::heddle::api::v1alpha2::{StreamDataKind, StreamFrame, stream_frame};
17pub use read_budget::{
18 GUARANTEED_READ_BUDGET, ReadBudgetError, negotiate_read_budget, validate_accepted_read_budget,
19};
20
21include!(concat!(env!("OUT_DIR"), "/heddle_api_v2_methods.rs"));
22
23#[derive(Clone, Copy, Debug, Eq, PartialEq)]
25pub struct AuthorizationTarget {
26 pub path: &'static str,
27 pub role: AuthorizationRole,
28}
29
30#[derive(Clone, Copy, Debug, Eq, PartialEq)]
34pub struct AuthorizationPolicy {
35 pub role: AuthorizationRole,
36 pub scope_source: AuthorizationScopeSource,
37 pub existence: AuthorizationExistence,
38 pub targets: &'static [AuthorizationTarget],
39}
40
41#[derive(Debug)]
43pub struct RoutedCall<'a> {
44 pub method: &'static MethodDescriptor,
45 pub context: &'a CallContext,
46}
47
48impl<'a> RoutedCall<'a> {
49 pub fn new(path: &str, context: &'a CallContext) -> Option<Self> {
50 method_descriptor(path).map(|method| Self { method, context })
51 }
52}
53
54impl MethodDescriptor {
55 pub fn verify_protocol(
57 &self,
58 context: &CallContext,
59 ) -> Result<(), crate::hybrid_codec::Reject> {
60 if self.mandatory_features.is_empty() {
61 return Ok(());
62 }
63 if self.mandatory_features
64 != [crate::heddle::api::common::MandatoryProtocolFeature::ImportAuthorityHostWitnessV1]
65 {
66 return Err(crate::hybrid_codec::Reject::Protocol);
67 }
68 crate::import_authority::require_hybrid_peer(context.protocol.as_ref())
69 }
70
71 pub const fn allows_zero_rtt(&self) -> bool {
73 matches!(self.effect, RpcEffect::ReadOnly)
74 && matches!(self.retry_behavior, RetryBehavior::Safe)
75 }
76 pub fn client_operation_id<'a>(
79 &self,
80 request: &'a [u8],
81 ) -> Result<Option<&'a str>, crate::RequestMetadataError> {
82 self.client_operation_id_field_number
83 .map(|field| crate::transport::protobuf_string_field(request, field))
84 .transpose()
85 .map(Option::flatten)
86 }
87}
88
89pub const MAX_CURSOR_BYTES: usize = 4096;
91
92#[derive(Clone, Copy, Debug, PartialEq, Eq, thiserror::Error)]
94pub enum StreamProtocolError {
95 #[error("non-contiguous stream sequence")]
96 Sequence,
97 #[error("stream query binding mismatch")]
98 Binding,
99 #[error("stream resumed from an unexpected cursor")]
100 Resume,
101 #[error("invalid stream phase")]
102 Phase,
103 #[error("payload does not match the frame kind")]
104 Payload,
105 #[error("invalid or oversized checkpoint cursor")]
106 Cursor,
107 #[error("missing or invalid accepted ReadBudget")]
108 Budget,
109}
110
111#[derive(Clone, Copy, Debug, PartialEq, Eq)]
112enum Phase {
113 Opening,
114 Snapshot,
115 Live,
116 Reset,
117 Complete,
118}
119
120#[derive(Clone, Copy, Debug, PartialEq, Eq)]
123pub enum ObservationAction {
124 BeginSnapshot,
125 Resumed,
126 Stage(StreamDataKind),
127 Commit,
128 Heartbeat,
129 Reset,
130 Complete,
131}
132
133#[derive(Debug, thiserror::Error)]
134pub enum ObservationApplyError<E> {
135 #[error("stream protocol failed: {0}")]
136 Protocol(#[from] StreamProtocolError),
137 #[error("view reducer failed: {0}")]
138 Reducer(E),
139}
140
141#[derive(Clone)]
143pub struct ObservationState {
144 binding_digest: [u8; 32],
145 cursor: Vec<u8>,
146 sequence: u64,
147 phase: Phase,
148 pending: bool,
149}
150
151impl ObservationState {
152 pub async fn apply<E, F, Fut>(
156 &mut self,
157 frame: &StreamFrame,
158 has_payload: bool,
159 reducer: F,
160 ) -> Result<ObservationAction, ObservationApplyError<E>>
161 where
162 F: FnOnce(ObservationAction, Vec<u8>) -> Fut,
163 Fut: std::future::Future<Output = Result<(), E>>,
164 {
165 let mut next = self.clone();
166 let action = next.accept(frame, has_payload)?;
167 reducer(action, next.cursor.clone())
168 .await
169 .map_err(ObservationApplyError::Reducer)?;
170 *self = next;
171 Ok(action)
172 }
173
174 pub fn new(binding_digest: [u8; 32], cursor: Vec<u8>) -> Self {
175 Self {
176 binding_digest,
177 cursor,
178 sequence: 0,
179 phase: Phase::Opening,
180 pending: false,
181 }
182 }
183
184 pub fn cursor(&self) -> &[u8] {
185 &self.cursor
186 }
187
188 pub fn is_complete(&self) -> bool {
189 self.phase == Phase::Complete
190 }
191
192 pub fn accept(
199 &mut self,
200 frame: &StreamFrame,
201 has_payload: bool,
202 ) -> Result<ObservationAction, StreamProtocolError> {
203 use stream_frame::Body;
204 if matches!(self.phase, Phase::Reset | Phase::Complete) {
205 return Err(StreamProtocolError::Phase);
206 }
207 if self.sequence.checked_add(1) != Some(frame.sequence) {
208 return Err(StreamProtocolError::Sequence);
209 }
210 let body = frame.body.as_ref().ok_or(StreamProtocolError::Phase)?;
211 if has_payload != matches!(body, Body::Data(_)) {
212 return Err(StreamProtocolError::Payload);
213 }
214 let action = match body {
215 Body::Open(open) => {
216 if self.phase != Phase::Opening {
217 return Err(StreamProtocolError::Phase);
218 }
219 if open.binding_digest.as_slice() != self.binding_digest {
220 return Err(StreamProtocolError::Binding);
221 }
222 if open.resumed_from != self.cursor {
223 return Err(StreamProtocolError::Resume);
224 }
225 if self.cursor.len() > MAX_CURSOR_BYTES {
226 return Err(StreamProtocolError::Cursor);
227 }
228 validate_accepted_read_budget(&Default::default(), open.accepted_budget.as_ref())
229 .map_err(|_| StreamProtocolError::Budget)?;
230 if self.cursor.is_empty() {
231 self.phase = Phase::Snapshot;
232 ObservationAction::BeginSnapshot
233 } else {
234 self.phase = Phase::Live;
235 ObservationAction::Resumed
236 }
237 }
238 Body::Data(data) => {
239 let kind =
240 StreamDataKind::try_from(data.kind).map_err(|_| StreamProtocolError::Phase)?;
241 let valid = match self.phase {
242 Phase::Snapshot => kind == StreamDataKind::Snapshot,
243 Phase::Live => matches!(kind, StreamDataKind::Upsert | StreamDataKind::Remove),
244 _ => false,
245 };
246 if !valid {
247 return Err(StreamProtocolError::Phase);
248 }
249 self.pending = true;
250 ObservationAction::Stage(kind)
251 }
252 Body::Checkpoint(checkpoint) => {
253 if !matches!(self.phase, Phase::Snapshot | Phase::Live)
254 || checkpoint.snapshot_complete != (self.phase == Phase::Snapshot)
255 {
256 return Err(StreamProtocolError::Phase);
257 }
258 if checkpoint.previous_cursor != self.cursor
259 || checkpoint.cursor.is_empty()
260 || checkpoint.cursor.len() > MAX_CURSOR_BYTES
261 || checkpoint.cursor == self.cursor
262 {
263 return Err(StreamProtocolError::Cursor);
264 }
265 self.cursor.clone_from(&checkpoint.cursor);
266 self.phase = Phase::Live;
267 self.pending = false;
268 ObservationAction::Commit
269 }
270 Body::Reset(_) => {
271 self.phase = Phase::Reset;
272 self.cursor.clear();
273 self.pending = false;
274 ObservationAction::Reset
275 }
276 Body::Complete(complete) => {
277 if self.phase != Phase::Live || self.pending {
278 return Err(StreamProtocolError::Phase);
279 }
280 if complete.cursor != self.cursor {
281 return Err(StreamProtocolError::Cursor);
282 }
283 self.phase = Phase::Complete;
284 ObservationAction::Complete
285 }
286 Body::Heartbeat(_) => {
287 if self.phase == Phase::Opening {
288 return Err(StreamProtocolError::Phase);
289 }
290 ObservationAction::Heartbeat
291 }
292 };
293 self.sequence = frame.sequence;
294 Ok(action)
295 }
296}