1pub 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#[derive(Clone, Copy, Debug, Eq, PartialEq)]
20pub struct AuthorizationTarget {
21 pub path: &'static str,
22 pub role: AuthorizationRole,
23}
24
25#[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#[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 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 pub const fn allows_zero_rtt(&self) -> bool {
68 matches!(self.effect, RpcEffect::ReadOnly)
69 && matches!(self.retry_behavior, RetryBehavior::Safe)
70 }
71 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
84pub const MAX_CURSOR_BYTES: usize = 4096;
86
87#[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#[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#[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 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 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}