1pub mod client;
3pub mod custodial_recovery;
4pub mod identity_management;
5pub mod invitation;
6pub mod passkey_label;
7use crate::StreamingShape;
8use crate::heddle::api::common::{
9 AuthorizationAccess, AuthorizationExistence, AuthorizationRole, AuthorizationScopeSource,
10 CallContext, DeploymentTarget, RetryBehavior, RpcEffect, ServiceMaturity, SigningTier,
11};
12use crate::heddle::api::v1alpha2::{StreamDataKind, StreamFrame, stream_frame};
13
14include!(concat!(env!("OUT_DIR"), "/heddle_api_v2_methods.rs"));
15
16#[derive(Clone, Copy, Debug, Eq, PartialEq)]
18pub struct AuthorizationTarget {
19 pub path: &'static str,
20 pub role: AuthorizationRole,
21}
22
23#[derive(Clone, Copy, Debug, Eq, PartialEq)]
27pub struct AuthorizationPolicy {
28 pub role: AuthorizationRole,
29 pub scope_source: AuthorizationScopeSource,
30 pub existence: AuthorizationExistence,
31 pub targets: &'static [AuthorizationTarget],
32}
33
34#[derive(Debug)]
36pub struct RoutedCall<'a> {
37 pub method: &'static MethodDescriptor,
38 pub context: &'a CallContext,
39}
40
41impl<'a> RoutedCall<'a> {
42 pub fn new(path: &str, context: &'a CallContext) -> Option<Self> {
43 method_descriptor(path).map(|method| Self { method, context })
44 }
45}
46
47impl MethodDescriptor {
48 pub fn verify_protocol(
50 &self,
51 context: &CallContext,
52 ) -> Result<(), crate::hybrid_codec::Reject> {
53 if self.mandatory_features.is_empty() {
54 return Ok(());
55 }
56 if self.mandatory_features
57 != [crate::heddle::api::common::MandatoryProtocolFeature::ImportAuthorityHostWitnessV1]
58 {
59 return Err(crate::hybrid_codec::Reject::Protocol);
60 }
61 crate::import_authority::require_hybrid_peer(context.protocol.as_ref())
62 }
63
64 pub const fn allows_zero_rtt(&self) -> bool {
66 matches!(self.effect, RpcEffect::ReadOnly)
67 && matches!(self.retry_behavior, RetryBehavior::Safe)
68 }
69 pub fn client_operation_id<'a>(
72 &self,
73 request: &'a [u8],
74 ) -> Result<Option<&'a str>, crate::RequestMetadataError> {
75 self.client_operation_id_field_number
76 .map(|field| crate::transport::protobuf_string_field(request, field))
77 .transpose()
78 .map(Option::flatten)
79 }
80}
81
82pub const MAX_CURSOR_BYTES: usize = 4096;
84
85#[derive(Clone, Copy, Debug, PartialEq, Eq, thiserror::Error)]
87pub enum StreamProtocolError {
88 #[error("non-contiguous stream sequence")]
89 Sequence,
90 #[error("stream query binding mismatch")]
91 Binding,
92 #[error("stream resumed from an unexpected cursor")]
93 Resume,
94 #[error("invalid stream phase")]
95 Phase,
96 #[error("payload does not match the frame kind")]
97 Payload,
98 #[error("invalid or oversized checkpoint cursor")]
99 Cursor,
100}
101
102#[derive(Clone, Copy, Debug, PartialEq, Eq)]
103enum Phase {
104 Opening,
105 Snapshot,
106 Live,
107 Reset,
108 Complete,
109}
110
111#[derive(Clone, Copy, Debug, PartialEq, Eq)]
114pub enum ObservationAction {
115 BeginSnapshot,
116 Resumed,
117 Stage(StreamDataKind),
118 Commit,
119 Heartbeat,
120 Reset,
121 Complete,
122}
123
124#[derive(Debug, thiserror::Error)]
125pub enum ObservationApplyError<E> {
126 #[error("stream protocol failed: {0}")]
127 Protocol(#[from] StreamProtocolError),
128 #[error("view reducer failed: {0}")]
129 Reducer(E),
130}
131
132#[derive(Clone)]
134pub struct ObservationState {
135 binding_digest: [u8; 32],
136 cursor: Vec<u8>,
137 sequence: u64,
138 phase: Phase,
139 pending: bool,
140}
141
142impl ObservationState {
143 pub async fn apply<E, F, Fut>(
147 &mut self,
148 frame: &StreamFrame,
149 has_payload: bool,
150 reducer: F,
151 ) -> Result<ObservationAction, ObservationApplyError<E>>
152 where
153 F: FnOnce(ObservationAction, Vec<u8>) -> Fut,
154 Fut: std::future::Future<Output = Result<(), E>>,
155 {
156 let mut next = self.clone();
157 let action = next.accept(frame, has_payload)?;
158 reducer(action, next.cursor.clone())
159 .await
160 .map_err(ObservationApplyError::Reducer)?;
161 *self = next;
162 Ok(action)
163 }
164
165 pub fn new(binding_digest: [u8; 32], cursor: Vec<u8>) -> Self {
166 Self {
167 binding_digest,
168 cursor,
169 sequence: 0,
170 phase: Phase::Opening,
171 pending: false,
172 }
173 }
174
175 pub fn cursor(&self) -> &[u8] {
176 &self.cursor
177 }
178
179 pub fn is_complete(&self) -> bool {
180 self.phase == Phase::Complete
181 }
182
183 pub fn accept(
190 &mut self,
191 frame: &StreamFrame,
192 has_payload: bool,
193 ) -> Result<ObservationAction, StreamProtocolError> {
194 use stream_frame::Body;
195 if matches!(self.phase, Phase::Reset | Phase::Complete) {
196 return Err(StreamProtocolError::Phase);
197 }
198 if self.sequence.checked_add(1) != Some(frame.sequence) {
199 return Err(StreamProtocolError::Sequence);
200 }
201 let body = frame.body.as_ref().ok_or(StreamProtocolError::Phase)?;
202 if has_payload != matches!(body, Body::Data(_)) {
203 return Err(StreamProtocolError::Payload);
204 }
205 let action = match body {
206 Body::Open(open) => {
207 if self.phase != Phase::Opening {
208 return Err(StreamProtocolError::Phase);
209 }
210 if open.binding_digest.as_slice() != self.binding_digest {
211 return Err(StreamProtocolError::Binding);
212 }
213 if open.resumed_from != self.cursor {
214 return Err(StreamProtocolError::Resume);
215 }
216 if self.cursor.len() > MAX_CURSOR_BYTES {
217 return Err(StreamProtocolError::Cursor);
218 }
219 if self.cursor.is_empty() {
220 self.phase = Phase::Snapshot;
221 ObservationAction::BeginSnapshot
222 } else {
223 self.phase = Phase::Live;
224 ObservationAction::Resumed
225 }
226 }
227 Body::Data(data) => {
228 let kind =
229 StreamDataKind::try_from(data.kind).map_err(|_| StreamProtocolError::Phase)?;
230 let valid = match self.phase {
231 Phase::Snapshot => kind == StreamDataKind::Snapshot,
232 Phase::Live => matches!(kind, StreamDataKind::Upsert | StreamDataKind::Remove),
233 _ => false,
234 };
235 if !valid {
236 return Err(StreamProtocolError::Phase);
237 }
238 self.pending = true;
239 ObservationAction::Stage(kind)
240 }
241 Body::Checkpoint(checkpoint) => {
242 if !matches!(self.phase, Phase::Snapshot | Phase::Live)
243 || checkpoint.snapshot_complete != (self.phase == Phase::Snapshot)
244 {
245 return Err(StreamProtocolError::Phase);
246 }
247 if checkpoint.previous_cursor != self.cursor
248 || checkpoint.cursor.is_empty()
249 || checkpoint.cursor.len() > MAX_CURSOR_BYTES
250 || checkpoint.cursor == self.cursor
251 {
252 return Err(StreamProtocolError::Cursor);
253 }
254 self.cursor.clone_from(&checkpoint.cursor);
255 self.phase = Phase::Live;
256 self.pending = false;
257 ObservationAction::Commit
258 }
259 Body::Reset(_) => {
260 self.phase = Phase::Reset;
261 self.cursor.clear();
262 self.pending = false;
263 ObservationAction::Reset
264 }
265 Body::Complete(complete) => {
266 if self.phase != Phase::Live || self.pending {
267 return Err(StreamProtocolError::Phase);
268 }
269 if complete.cursor != self.cursor {
270 return Err(StreamProtocolError::Cursor);
271 }
272 self.phase = Phase::Complete;
273 ObservationAction::Complete
274 }
275 Body::Heartbeat(_) => {
276 if self.phase == Phase::Opening {
277 return Err(StreamProtocolError::Phase);
278 }
279 ObservationAction::Heartbeat
280 }
281 };
282 self.sequence = frame.sequence;
283 Ok(action)
284 }
285}