Skip to main content

heddle_api/
v2.rs

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