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