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