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