Skip to main content

heddle_api/
v2.rs

1//! Shared v2 client behavior. Transport adapters retain key and connection ownership.
2pub 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/// Generated resource guard selected by a protobuf request field path.
24#[derive(Clone, Copy, Debug, Eq, PartialEq)]
25pub struct AuthorizationTarget {
26    pub path: &'static str,
27    pub role: AuthorizationRole,
28}
29
30/// Complete method authorization declaration. Hosts must resolve every
31/// populated target and enforce the role against current authority; metadata
32/// describes that obligation and is never itself a capability.
33#[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/// Decoded transport context bound to a v2 route, before body interpretation.
42#[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    /// Mandatory semantic gate before decoding/staging an authenticated route.
56    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    /// Only safe reads can be sent on replayable 0-RTT connections.
72    pub const fn allows_zero_rtt(&self) -> bool {
73        matches!(self.effect, RpcEffect::ReadOnly)
74            && matches!(self.retry_behavior, RetryBehavior::Safe)
75    }
76    /// Extract the operation ID for the transport context from the exact request
77    /// bytes. Adapters must not maintain a second route/field-number inventory.
78    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
89/// Upper bound for an opaque observation cursor; independent of payload limits.
90pub const MAX_CURSOR_BYTES: usize = 4096;
91
92/// Invalid stream framing must never advance a client's durable observation.
93#[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/// Instructions for the typed view reducer. Apply staged data atomically at
121/// Commit, and persist the cursor only after that reducer succeeds.
122#[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/// Cursor and lifecycle tracker; typed consumers own staged view data.
142#[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    /// Advances this tracker only after the view reducer succeeds. At Commit,
153    /// the reducer must atomically apply staged data and persist the supplied
154    /// cursor. A failed reducer leaves the tracker at its previous checkpoint.
155    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    /// Validate a frame before changing lifecycle state. A rejected frame leaves
193    /// the state unchanged. Transport adapters enforce negotiated byte budgets
194    /// before decoding; callers validate the method-specific typed payload.
195    ///
196    /// For reducers that can fail, run this on a clone, apply the returned action,
197    /// then replace the original tracker only if the reducer succeeds.
198    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}