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