1pub 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#[derive(Clone, Copy, Debug, Eq, PartialEq)]
28pub struct AuthorizationTarget {
29 pub path: &'static str,
30 pub role: AuthorizationRole,
31}
32
33#[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#[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 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 pub const fn allows_zero_rtt(&self) -> bool {
76 matches!(self.effect, RpcEffect::ReadOnly)
77 && matches!(self.retry_behavior, RetryBehavior::Safe)
78 }
79 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
92pub const MAX_CURSOR_BYTES: usize = 4096;
94
95#[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#[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#[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 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 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}