1pub 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#[derive(Clone, Copy, Debug, Eq, PartialEq)]
17pub struct AuthorizationTarget {
18 pub path: &'static str,
19 pub role: AuthorizationRole,
20}
21
22#[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#[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 pub const fn allows_zero_rtt(&self) -> bool {
49 matches!(self.effect, RpcEffect::ReadOnly)
50 && matches!(self.retry_behavior, RetryBehavior::Safe)
51 }
52 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
65pub const MAX_CURSOR_BYTES: usize = 4096;
67
68#[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#[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#[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 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 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}