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;
9use crate::StreamingShape;
10use crate::heddle::api::common::{
11    AuthorizationAccess, AuthorizationExistence, AuthorizationRole, AuthorizationScopeSource,
12    CallContext, DeploymentTarget, RetryBehavior, RpcEffect, ServiceMaturity, SigningTier,
13};
14use crate::heddle::api::v1alpha2::{StreamDataKind, StreamFrame, stream_frame};
15
16include!(concat!(env!("OUT_DIR"), "/heddle_api_v2_methods.rs"));
17
18/// Generated resource guard selected by a protobuf request field path.
19#[derive(Clone, Copy, Debug, Eq, PartialEq)]
20pub struct AuthorizationTarget {
21    pub path: &'static str,
22    pub role: AuthorizationRole,
23}
24
25/// Complete method authorization declaration. Hosts must resolve every
26/// populated target and enforce the role against current authority; metadata
27/// describes that obligation and is never itself a capability.
28#[derive(Clone, Copy, Debug, Eq, PartialEq)]
29pub struct AuthorizationPolicy {
30    pub role: AuthorizationRole,
31    pub scope_source: AuthorizationScopeSource,
32    pub existence: AuthorizationExistence,
33    pub targets: &'static [AuthorizationTarget],
34}
35
36/// Decoded transport context bound to a v2 route, before body interpretation.
37#[derive(Debug)]
38pub struct RoutedCall<'a> {
39    pub method: &'static MethodDescriptor,
40    pub context: &'a CallContext,
41}
42
43impl<'a> RoutedCall<'a> {
44    pub fn new(path: &str, context: &'a CallContext) -> Option<Self> {
45        method_descriptor(path).map(|method| Self { method, context })
46    }
47}
48
49impl MethodDescriptor {
50    /// Mandatory semantic gate before decoding/staging an authenticated route.
51    pub fn verify_protocol(
52        &self,
53        context: &CallContext,
54    ) -> Result<(), crate::hybrid_codec::Reject> {
55        if self.mandatory_features.is_empty() {
56            return Ok(());
57        }
58        if self.mandatory_features
59            != [crate::heddle::api::common::MandatoryProtocolFeature::ImportAuthorityHostWitnessV1]
60        {
61            return Err(crate::hybrid_codec::Reject::Protocol);
62        }
63        crate::import_authority::require_hybrid_peer(context.protocol.as_ref())
64    }
65
66    /// Only safe reads can be sent on replayable 0-RTT connections.
67    pub const fn allows_zero_rtt(&self) -> bool {
68        matches!(self.effect, RpcEffect::ReadOnly)
69            && matches!(self.retry_behavior, RetryBehavior::Safe)
70    }
71    /// Extract the operation ID for the transport context from the exact request
72    /// bytes. Adapters must not maintain a second route/field-number inventory.
73    pub fn client_operation_id<'a>(
74        &self,
75        request: &'a [u8],
76    ) -> Result<Option<&'a str>, crate::RequestMetadataError> {
77        self.client_operation_id_field_number
78            .map(|field| crate::transport::protobuf_string_field(request, field))
79            .transpose()
80            .map(Option::flatten)
81    }
82}
83
84/// Upper bound for an opaque observation cursor; independent of payload limits.
85pub const MAX_CURSOR_BYTES: usize = 4096;
86
87/// Invalid stream framing must never advance a client's durable observation.
88#[derive(Clone, Copy, Debug, PartialEq, Eq, thiserror::Error)]
89pub enum StreamProtocolError {
90    #[error("non-contiguous stream sequence")]
91    Sequence,
92    #[error("stream query binding mismatch")]
93    Binding,
94    #[error("stream resumed from an unexpected cursor")]
95    Resume,
96    #[error("invalid stream phase")]
97    Phase,
98    #[error("payload does not match the frame kind")]
99    Payload,
100    #[error("invalid or oversized checkpoint cursor")]
101    Cursor,
102}
103
104#[derive(Clone, Copy, Debug, PartialEq, Eq)]
105enum Phase {
106    Opening,
107    Snapshot,
108    Live,
109    Reset,
110    Complete,
111}
112
113/// Instructions for the typed view reducer. Apply staged data atomically at
114/// Commit, and persist the cursor only after that reducer succeeds.
115#[derive(Clone, Copy, Debug, PartialEq, Eq)]
116pub enum ObservationAction {
117    BeginSnapshot,
118    Resumed,
119    Stage(StreamDataKind),
120    Commit,
121    Heartbeat,
122    Reset,
123    Complete,
124}
125
126#[derive(Debug, thiserror::Error)]
127pub enum ObservationApplyError<E> {
128    #[error("stream protocol failed: {0}")]
129    Protocol(#[from] StreamProtocolError),
130    #[error("view reducer failed: {0}")]
131    Reducer(E),
132}
133
134/// Cursor and lifecycle tracker; typed consumers own staged view data.
135#[derive(Clone)]
136pub struct ObservationState {
137    binding_digest: [u8; 32],
138    cursor: Vec<u8>,
139    sequence: u64,
140    phase: Phase,
141    pending: bool,
142}
143
144impl ObservationState {
145    /// Advances this tracker only after the view reducer succeeds. At Commit,
146    /// the reducer must atomically apply staged data and persist the supplied
147    /// cursor. A failed reducer leaves the tracker at its previous checkpoint.
148    pub async fn apply<E, F, Fut>(
149        &mut self,
150        frame: &StreamFrame,
151        has_payload: bool,
152        reducer: F,
153    ) -> Result<ObservationAction, ObservationApplyError<E>>
154    where
155        F: FnOnce(ObservationAction, Vec<u8>) -> Fut,
156        Fut: std::future::Future<Output = Result<(), E>>,
157    {
158        let mut next = self.clone();
159        let action = next.accept(frame, has_payload)?;
160        reducer(action, next.cursor.clone())
161            .await
162            .map_err(ObservationApplyError::Reducer)?;
163        *self = next;
164        Ok(action)
165    }
166
167    pub fn new(binding_digest: [u8; 32], cursor: Vec<u8>) -> Self {
168        Self {
169            binding_digest,
170            cursor,
171            sequence: 0,
172            phase: Phase::Opening,
173            pending: false,
174        }
175    }
176
177    pub fn cursor(&self) -> &[u8] {
178        &self.cursor
179    }
180
181    pub fn is_complete(&self) -> bool {
182        self.phase == Phase::Complete
183    }
184
185    /// Validate a frame before changing lifecycle state. A rejected frame leaves
186    /// the state unchanged. Transport adapters enforce negotiated byte budgets
187    /// before decoding; callers validate the method-specific typed payload.
188    ///
189    /// For reducers that can fail, run this on a clone, apply the returned action,
190    /// then replace the original tracker only if the reducer succeeds.
191    pub fn accept(
192        &mut self,
193        frame: &StreamFrame,
194        has_payload: bool,
195    ) -> Result<ObservationAction, StreamProtocolError> {
196        use stream_frame::Body;
197        if matches!(self.phase, Phase::Reset | Phase::Complete) {
198            return Err(StreamProtocolError::Phase);
199        }
200        if self.sequence.checked_add(1) != Some(frame.sequence) {
201            return Err(StreamProtocolError::Sequence);
202        }
203        let body = frame.body.as_ref().ok_or(StreamProtocolError::Phase)?;
204        if has_payload != matches!(body, Body::Data(_)) {
205            return Err(StreamProtocolError::Payload);
206        }
207        let action = match body {
208            Body::Open(open) => {
209                if self.phase != Phase::Opening {
210                    return Err(StreamProtocolError::Phase);
211                }
212                if open.binding_digest.as_slice() != self.binding_digest {
213                    return Err(StreamProtocolError::Binding);
214                }
215                if open.resumed_from != self.cursor {
216                    return Err(StreamProtocolError::Resume);
217                }
218                if self.cursor.len() > MAX_CURSOR_BYTES {
219                    return Err(StreamProtocolError::Cursor);
220                }
221                if self.cursor.is_empty() {
222                    self.phase = Phase::Snapshot;
223                    ObservationAction::BeginSnapshot
224                } else {
225                    self.phase = Phase::Live;
226                    ObservationAction::Resumed
227                }
228            }
229            Body::Data(data) => {
230                let kind =
231                    StreamDataKind::try_from(data.kind).map_err(|_| StreamProtocolError::Phase)?;
232                let valid = match self.phase {
233                    Phase::Snapshot => kind == StreamDataKind::Snapshot,
234                    Phase::Live => matches!(kind, StreamDataKind::Upsert | StreamDataKind::Remove),
235                    _ => false,
236                };
237                if !valid {
238                    return Err(StreamProtocolError::Phase);
239                }
240                self.pending = true;
241                ObservationAction::Stage(kind)
242            }
243            Body::Checkpoint(checkpoint) => {
244                if !matches!(self.phase, Phase::Snapshot | Phase::Live)
245                    || checkpoint.snapshot_complete != (self.phase == Phase::Snapshot)
246                {
247                    return Err(StreamProtocolError::Phase);
248                }
249                if checkpoint.previous_cursor != self.cursor
250                    || checkpoint.cursor.is_empty()
251                    || checkpoint.cursor.len() > MAX_CURSOR_BYTES
252                    || checkpoint.cursor == self.cursor
253                {
254                    return Err(StreamProtocolError::Cursor);
255                }
256                self.cursor.clone_from(&checkpoint.cursor);
257                self.phase = Phase::Live;
258                self.pending = false;
259                ObservationAction::Commit
260            }
261            Body::Reset(_) => {
262                self.phase = Phase::Reset;
263                self.cursor.clear();
264                self.pending = false;
265                ObservationAction::Reset
266            }
267            Body::Complete(complete) => {
268                if self.phase != Phase::Live || self.pending {
269                    return Err(StreamProtocolError::Phase);
270                }
271                if complete.cursor != self.cursor {
272                    return Err(StreamProtocolError::Cursor);
273                }
274                self.phase = Phase::Complete;
275                ObservationAction::Complete
276            }
277            Body::Heartbeat(_) => {
278                if self.phase == Phase::Opening {
279                    return Err(StreamProtocolError::Phase);
280                }
281                ObservationAction::Heartbeat
282            }
283        };
284        self.sequence = frame.sequence;
285        Ok(action)
286    }
287}