1use alloc::{string::String, vec::Vec};
9
10use crate::algebra::{ResourceDimension, ResourceVector, WideResourceVector};
11
12use super::codec::CodecError;
13use super::{
14 AttachAttemptToken, AttachSecret, BindingEpoch, ConnectionIncarnation, DetachAttemptToken,
15 EnrollmentToken, Generation, LeaveAttemptToken, ProtocolVersion, RecordAdmissionAttemptToken,
16 closure as c, envelope as e, response as r, tags as t,
17};
18
19const TOKEN_LEN: usize = 16;
20const SECRET_LEN: usize = 32;
21const RECOVERY_STATUS_LEN: u128 = 26;
22
23pub(super) struct TerminalizedWireDecodeAuthority(());
26
27const TERMINALIZED_WIRE_DECODE_AUTHORITY: TerminalizedWireDecodeAuthority =
28 TerminalizedWireDecodeAuthority(());
29
30pub fn encode_server_value_body(
42 value: &r::ServerValue,
43 version: ProtocolVersion,
44) -> Result<(t::ServerDiscriminant, Vec<u8>), CodecError> {
45 require_v1(version)?;
46 let discriminant = value.discriminant();
47 let mut encoder = Encoder::default();
48
49 if carries_origin(discriminant) {
50 let originating_request = value
51 .originating_request()
52 .ok_or(CodecError::InvalidValue)?;
53 encoder.put_u16(originating_request.wire_value());
54 }
55
56 encode_server_suffix(value, &mut encoder)?;
57 Ok((discriminant, encoder.bytes))
58}
59
60pub fn decode_server_value_body(
70 discriminant: t::ServerDiscriminant,
71 version: ProtocolVersion,
72 body: &[u8],
73) -> Result<(r::ServerValue, ProtocolVersion), CodecError> {
74 require_v1(version)?;
75 let mut decoder = Decoder::new(body);
76 let originating_request = if carries_origin(discriminant) {
77 let raw = decoder.take_u16()?;
78 let request = t::ClientDiscriminant::try_from(raw)
79 .map_err(|_| decode_error(t::DecodeClass::InvalidField))?;
80 if !origin_is_valid(discriminant, request) {
81 return Err(decode_error(t::DecodeClass::InvalidField));
82 }
83 Some(request)
84 } else {
85 None
86 };
87
88 let value = decode_server_suffix(discriminant, originating_request, &mut decoder)?;
89 decoder.finish()?;
90 Ok((value, version))
91}
92
93fn require_v1(version: ProtocolVersion) -> Result<(), CodecError> {
94 if version == ProtocolVersion::V1 {
95 Ok(())
96 } else {
97 Err(CodecError::UnsupportedVersion {
98 presented: version,
99 supported: ProtocolVersion::V1,
100 })
101 }
102}
103
104const fn decode_error(class: t::DecodeClass) -> CodecError {
105 CodecError::Decode { class }
106}
107
108#[derive(Default)]
109struct Encoder {
110 bytes: Vec<u8>,
111}
112
113impl Encoder {
114 fn put_u8(&mut self, value: u8) {
115 self.bytes.push(value);
116 }
117
118 fn put_u16(&mut self, value: u16) {
119 self.bytes.extend_from_slice(&value.to_be_bytes());
120 }
121
122 fn put_u32(&mut self, value: u32) {
123 self.bytes.extend_from_slice(&value.to_be_bytes());
124 }
125
126 fn put_u64(&mut self, value: u64) {
127 self.bytes.extend_from_slice(&value.to_be_bytes());
128 }
129
130 fn put_u128(&mut self, value: u128) {
131 self.bytes.extend_from_slice(&value.to_be_bytes());
132 }
133
134 fn put_fixed(&mut self, value: &[u8]) {
135 self.bytes.extend_from_slice(value);
136 }
137
138 fn put_bool(&mut self, value: bool) {
139 self.put_u8(u8::from(value));
140 }
141
142 fn put_option_u64(&mut self, value: Option<u64>) {
143 match value {
144 Some(value) => {
145 self.put_u8(1);
146 self.put_u64(value);
147 }
148 None => self.put_u8(0),
149 }
150 }
151
152 fn put_option_generation(&mut self, value: Option<Generation>) {
153 match value {
154 Some(value) => {
155 self.put_u8(1);
156 self.put_generation(value);
157 }
158 None => self.put_u8(0),
159 }
160 }
161
162 fn put_option_binding_epoch(&mut self, value: Option<BindingEpoch>) {
163 match value {
164 Some(value) => {
165 self.put_u8(1);
166 self.put_binding_epoch(value);
167 }
168 None => self.put_u8(0),
169 }
170 }
171
172 fn put_generation(&mut self, value: Generation) {
173 self.put_u64(value.get());
174 }
175
176 fn put_protocol_version(&mut self, value: ProtocolVersion) {
177 self.put_u16(value.major);
178 self.put_u16(value.minor);
179 }
180
181 fn put_binding_epoch(&mut self, value: BindingEpoch) {
182 self.put_u64(value.connection_incarnation.server_incarnation);
183 self.put_u64(value.connection_incarnation.connection_ordinal);
184 self.put_generation(value.capability_generation);
185 }
186
187 fn put_resource_vector(&mut self, value: ResourceVector) {
188 self.put_u64(value.entries);
189 self.put_u64(value.bytes);
190 }
191
192 fn put_wide_resource_vector(&mut self, value: WideResourceVector) {
193 self.put_u128(value.entries);
194 self.put_u128(value.bytes);
195 }
196
197 fn put_string(&mut self, value: &str) -> Result<(), CodecError> {
198 let length: u32 = value
199 .len()
200 .try_into()
201 .map_err(|_| CodecError::LengthOverflow)?;
202 self.put_u32(length);
203 self.put_fixed(value.as_bytes());
204 Ok(())
205 }
206}
207
208struct Decoder<'a> {
209 input: &'a [u8],
210 position: usize,
211 invalid_field: bool,
212}
213
214impl<'a> Decoder<'a> {
215 const fn new(input: &'a [u8]) -> Self {
216 Self {
217 input,
218 position: 0,
219 invalid_field: false,
220 }
221 }
222
223 const fn remaining(&self) -> usize {
224 self.input.len().saturating_sub(self.position)
225 }
226
227 fn take(&mut self, length: usize) -> Result<&'a [u8], CodecError> {
228 let end = self
229 .position
230 .checked_add(length)
231 .ok_or_else(|| decode_error(t::DecodeClass::MissingRequiredField))?;
232 let value = self
233 .input
234 .get(self.position..end)
235 .ok_or_else(|| decode_error(t::DecodeClass::MissingRequiredField))?;
236 self.position = end;
237 Ok(value)
238 }
239
240 fn take_u8(&mut self) -> Result<u8, CodecError> {
241 let bytes = self.take(1)?;
242 bytes
243 .first()
244 .copied()
245 .ok_or_else(|| decode_error(t::DecodeClass::MissingRequiredField))
246 }
247
248 fn take_u16(&mut self) -> Result<u16, CodecError> {
249 let bytes: [u8; 2] = self
250 .take(2)?
251 .try_into()
252 .map_err(|_| decode_error(t::DecodeClass::MissingRequiredField))?;
253 Ok(u16::from_be_bytes(bytes))
254 }
255
256 fn take_u32(&mut self) -> Result<u32, CodecError> {
257 let bytes: [u8; 4] = self
258 .take(4)?
259 .try_into()
260 .map_err(|_| decode_error(t::DecodeClass::MissingRequiredField))?;
261 Ok(u32::from_be_bytes(bytes))
262 }
263
264 fn take_u64(&mut self) -> Result<u64, CodecError> {
265 let bytes: [u8; 8] = self
266 .take(8)?
267 .try_into()
268 .map_err(|_| decode_error(t::DecodeClass::MissingRequiredField))?;
269 Ok(u64::from_be_bytes(bytes))
270 }
271
272 fn take_u128(&mut self) -> Result<u128, CodecError> {
273 let bytes: [u8; 16] = self
274 .take(16)?
275 .try_into()
276 .map_err(|_| decode_error(t::DecodeClass::MissingRequiredField))?;
277 Ok(u128::from_be_bytes(bytes))
278 }
279
280 fn take_fixed<const N: usize>(&mut self) -> Result<[u8; N], CodecError> {
281 self.take(N)?
282 .try_into()
283 .map_err(|_| decode_error(t::DecodeClass::MissingRequiredField))
284 }
285
286 fn take_bool(&mut self) -> Result<bool, CodecError> {
287 match self.take_u8()? {
288 0 => Ok(false),
289 1 => Ok(true),
290 _ => Err(decode_error(t::DecodeClass::InvalidField)),
291 }
292 }
293
294 fn take_option_u64(&mut self) -> Result<Option<u64>, CodecError> {
295 match self.take_u8()? {
296 0 => Ok(None),
297 1 => self.take_u64().map(Some),
298 _ => Err(decode_error(t::DecodeClass::InvalidField)),
299 }
300 }
301
302 fn take_generation(&mut self) -> Result<Generation, CodecError> {
303 let raw = self.take_u64()?;
304 if let Some(value) = Generation::new(raw) {
305 Ok(value)
306 } else {
307 self.invalid_field = true;
308 generation_one()
309 }
310 }
311
312 fn take_option_generation(&mut self) -> Result<Option<Generation>, CodecError> {
313 match self.take_u8()? {
314 0 => Ok(None),
315 1 => self.take_generation().map(Some),
316 _ => Err(decode_error(t::DecodeClass::InvalidField)),
317 }
318 }
319
320 fn take_binding_epoch(&mut self) -> Result<BindingEpoch, CodecError> {
321 Ok(BindingEpoch::new(
322 ConnectionIncarnation::new(self.take_u64()?, self.take_u64()?),
323 self.take_generation()?,
324 ))
325 }
326
327 fn take_option_binding_epoch(&mut self) -> Result<Option<BindingEpoch>, CodecError> {
328 match self.take_u8()? {
329 0 => Ok(None),
330 1 => self.take_binding_epoch().map(Some),
331 _ => Err(decode_error(t::DecodeClass::InvalidField)),
332 }
333 }
334
335 fn take_protocol_version(&mut self) -> Result<ProtocolVersion, CodecError> {
336 Ok(ProtocolVersion::new(self.take_u16()?, self.take_u16()?))
337 }
338
339 fn take_resource_vector(&mut self) -> Result<ResourceVector, CodecError> {
340 Ok(ResourceVector::new(self.take_u64()?, self.take_u64()?))
341 }
342
343 fn take_wide_resource_vector(&mut self) -> Result<WideResourceVector, CodecError> {
344 Ok(WideResourceVector::new(
345 self.take_u128()?,
346 self.take_u128()?,
347 ))
348 }
349
350 fn take_string(&mut self) -> Result<String, CodecError> {
351 let length: usize = self
352 .take_u32()?
353 .try_into()
354 .map_err(|_| decode_error(t::DecodeClass::MissingRequiredField))?;
355 let bytes = self.take(length)?;
356 if let Ok(value) = core::str::from_utf8(bytes) {
357 Ok(String::from(value))
358 } else {
359 self.invalid_field = true;
360 Ok(String::new())
361 }
362 }
363
364 const fn invalidate(&mut self) {
365 self.invalid_field = true;
366 }
367
368 const fn finish(self) -> Result<(), CodecError> {
369 if self.remaining() != 0 {
370 Err(decode_error(t::DecodeClass::CanonicalEncoding))
371 } else if self.invalid_field {
372 Err(decode_error(t::DecodeClass::InvalidField))
373 } else {
374 Ok(())
375 }
376 }
377}
378
379fn generation_one() -> Result<Generation, CodecError> {
380 Generation::new(1).ok_or(CodecError::InvalidValue)
381}
382
383fn put_enrollment(value: &e::EnrollmentEnvelope, encoder: &mut Encoder) {
384 encoder.put_u64(value.conversation_id);
385 encoder.put_fixed(value.enrollment_token.as_bytes());
386}
387
388fn put_attach(value: &e::AttachEnvelope, encoder: &mut Encoder) {
389 encoder.put_u64(value.conversation_id);
390 encoder.put_u64(value.participant_id);
391 encoder.put_generation(value.capability_generation);
392 encoder.put_fixed(value.attach_attempt_token.as_bytes());
393 encoder.put_option_u64(value.accept_marker_delivery_seq);
394}
395
396fn put_detach(value: &e::DetachEnvelope, encoder: &mut Encoder) {
397 encoder.put_u64(value.conversation_id);
398 encoder.put_u64(value.participant_id);
399 encoder.put_generation(value.capability_generation);
400 encoder.put_fixed(value.detach_attempt_token.as_bytes());
401}
402
403fn put_participant_ack(value: &e::ParticipantAckEnvelope, encoder: &mut Encoder) {
404 encoder.put_u64(value.conversation_id);
405 encoder.put_u64(value.participant_id);
406 encoder.put_generation(value.capability_generation);
407 encoder.put_u64(value.through_seq);
408}
409
410fn put_leave(value: &e::LeaveEnvelope, encoder: &mut Encoder) {
411 encoder.put_u64(value.conversation_id);
412 encoder.put_u64(value.participant_id);
413 encoder.put_generation(value.capability_generation);
414 encoder.put_fixed(value.leave_attempt_token.as_bytes());
415}
416
417fn put_marker_ack(value: &e::MarkerAckEnvelope, encoder: &mut Encoder) {
418 encoder.put_u64(value.conversation_id);
419 encoder.put_u64(value.participant_id);
420 encoder.put_generation(value.capability_generation);
421 encoder.put_u64(value.marker_delivery_seq);
422}
423
424fn put_record_admission(value: &e::RecordAdmissionEnvelope, encoder: &mut Encoder) {
425 encoder.put_u64(value.conversation_id);
426 encoder.put_u64(value.participant_id);
427 encoder.put_generation(value.capability_generation);
428 encoder.put_fixed(value.record_admission_attempt_token.as_bytes());
429}
430
431fn put_response_envelope(value: &e::ResponseEnvelope, encoder: &mut Encoder) {
432 match value {
433 e::ResponseEnvelope::Enrollment(value) => put_enrollment(value, encoder),
434 e::ResponseEnvelope::CredentialAttach(value) => put_attach(value, encoder),
435 e::ResponseEnvelope::Detach(value) => put_detach(value, encoder),
436 e::ResponseEnvelope::ParticipantAck(value) => put_participant_ack(value, encoder),
437 e::ResponseEnvelope::Leave(value) => put_leave(value, encoder),
438 e::ResponseEnvelope::MarkerAck(value) => put_marker_ack(value, encoder),
439 e::ResponseEnvelope::RecordAdmission(value) => put_record_admission(value, encoder),
440 }
441}
442
443fn put_participant_reference(value: &r::ParticipantReferenceEnvelope, encoder: &mut Encoder) {
444 match value {
445 r::ParticipantReferenceEnvelope::CredentialAttach(value) => put_attach(value, encoder),
446 r::ParticipantReferenceEnvelope::Detach(value) => put_detach(value, encoder),
447 r::ParticipantReferenceEnvelope::ParticipantAck(value) => {
448 put_participant_ack(value, encoder);
449 }
450 r::ParticipantReferenceEnvelope::Leave(value) => put_leave(value, encoder),
451 r::ParticipantReferenceEnvelope::MarkerAck(value) => put_marker_ack(value, encoder),
452 r::ParticipantReferenceEnvelope::RecordAdmission(value) => {
453 put_record_admission(value, encoder);
454 }
455 }
456}
457
458fn put_binding_required(value: &r::BindingRequiredEnvelope, encoder: &mut Encoder) {
459 match value {
460 r::BindingRequiredEnvelope::Detach(value) => put_detach(value, encoder),
461 r::BindingRequiredEnvelope::ParticipantAck(value) => put_participant_ack(value, encoder),
462 r::BindingRequiredEnvelope::Leave(value) => put_leave(value, encoder),
463 r::BindingRequiredEnvelope::MarkerAck(value) => put_marker_ack(value, encoder),
464 r::BindingRequiredEnvelope::RecordAdmission(value) => put_record_admission(value, encoder),
465 }
466}
467
468fn put_order_allocating(value: &r::OrderAllocatingEnvelope, encoder: &mut Encoder) {
469 match value {
470 r::OrderAllocatingEnvelope::Enrollment(value) => put_enrollment(value, encoder),
471 r::OrderAllocatingEnvelope::CredentialAttach(value) => put_attach(value, encoder),
472 r::OrderAllocatingEnvelope::RecordAdmission(value) => put_record_admission(value, encoder),
473 }
474}
475
476fn put_closure_checked(value: &c::ClosureCheckedEnvelope, encoder: &mut Encoder) {
477 match value {
478 c::ClosureCheckedEnvelope::Enrollment(value) => put_enrollment(value, encoder),
479 c::ClosureCheckedEnvelope::CredentialAttach(value) => put_attach(value, encoder),
480 c::ClosureCheckedEnvelope::Leave(value) => put_leave(value, encoder),
481 c::ClosureCheckedEnvelope::RecordAdmission(value) => put_record_admission(value, encoder),
482 }
483}
484
485fn put_sequence_allocating(value: &r::SequenceAllocatingEnvelope, encoder: &mut Encoder) {
486 match value {
487 r::SequenceAllocatingEnvelope::Enrollment(value) => put_enrollment(value, encoder),
488 r::SequenceAllocatingEnvelope::CredentialAttach(value) => put_attach(value, encoder),
489 r::SequenceAllocatingEnvelope::RecordAdmission(value) => {
490 put_record_admission(value, encoder);
491 }
492 }
493}
494
495fn put_marker_proof(value: &r::MarkerProofRequest, encoder: &mut Encoder) {
496 match value {
497 r::MarkerProofRequest::CredentialAttach(value) => {
498 encoder.put_u64(value.conversation_id);
499 encoder.put_fixed(value.token.as_bytes());
500 encoder.put_u64(value.participant_id);
501 encoder.put_generation(value.capability_generation);
502 encoder.put_u64(value.requested_marker_delivery_seq);
503 }
504 r::MarkerProofRequest::MarkerAck(value) => {
505 encoder.put_u64(value.conversation_id);
506 encoder.put_u64(value.participant_id);
507 encoder.put_generation(value.capability_generation);
508 encoder.put_u64(value.requested_marker_delivery_seq);
509 }
510 }
511}
512
513fn put_repayment_edge(value: c::RepaymentEdge, encoder: &mut Encoder) {
514 encoder.put_u16(value.tag().wire_value());
515 match value {
516 c::RepaymentEdge::None => {}
517 c::RepaymentEdge::ObserverProjection { through_seq } => encoder.put_u64(through_seq),
518 c::RepaymentEdge::PhysicalCompaction {
519 from_floor,
520 through_seq,
521 } => {
522 encoder.put_u64(from_floor);
523 encoder.put_u64(through_seq);
524 }
525 c::RepaymentEdge::MarkerDelivery {
526 participant_id,
527 binding_epoch,
528 marker_delivery_seq,
529 } => {
530 encoder.put_u64(participant_id);
531 encoder.put_binding_epoch(binding_epoch);
532 encoder.put_u64(marker_delivery_seq);
533 }
534 c::RepaymentEdge::ParticipantCursorProgress(value) => {
535 encoder.put_u64(value.participant_id);
536 encoder.put_binding_epoch(value.binding_epoch);
537 encoder.put_u64(value.through_seq);
538 encoder.put_option_u64(value.marker_delivery_seq);
539 }
540 c::RepaymentEdge::DetachedCredentialRecovery {
541 participant_id,
542 marker_delivery_seq,
543 prior_binding_epoch,
544 } => {
545 encoder.put_u64(participant_id);
546 encoder.put_u64(marker_delivery_seq);
547 encoder.put_binding_epoch(prior_binding_epoch);
548 }
549 c::RepaymentEdge::DetachedMarkerRelease {
550 participant_id,
551 marker_delivery_seq,
552 last_dead_binding_epoch,
553 } => {
554 encoder.put_u64(participant_id);
555 encoder.put_u64(marker_delivery_seq);
556 encoder.put_binding_epoch(last_dead_binding_epoch);
557 }
558 c::RepaymentEdge::DetachedCursorRelease {
559 participant_id,
560 last_dead_binding_epoch,
561 } => {
562 encoder.put_u64(participant_id);
563 encoder.put_binding_epoch(last_dead_binding_epoch);
564 }
565 }
566}
567
568fn put_closure_snapshot(value: c::ClosureSnapshot, encoder: &mut Encoder) {
569 encoder.put_u64(value.marker_capacity_credits);
570 encoder.put_u64(value.marker_anchors);
571 encoder.put_u64(value.entry_debt);
572 encoder.put_u64(value.byte_debt);
573 put_repayment_edge(value.repayment_edge, encoder);
574 encoder.put_u64(value.edge_sequence_claims);
575 encoder.put_u64(value.edge_order_position_claims);
576 encoder.put_resource_vector(value.edge_k_remaining);
577 encoder.put_wide_resource_vector(value.k_headroom);
578 encoder.put_u64(value.episode_churn_used);
579 encoder.put_u64(value.delta_cycles);
580 encoder.put_u64(value.episode_churn_limit);
581}
582
583fn put_sequence_budget(value: super::SequenceBudget, encoder: &mut Encoder) {
584 encoder.put_u64(value.high_watermark);
585 encoder.put_u64(value.remaining);
586 encoder.put_u64(value.e);
587 encoder.put_u64(value.t);
588 encoder.put_u64(value.m);
589 encoder.put_u64(value.rs);
590 encoder.put_u64(value.rt);
591 encoder.put_u128(value.l_times_t);
592 encoder.put_u128(value.l_times_rt);
593 encoder.put_u128(value.l_other_times_e);
594}
595
596#[allow(clippy::too_many_lines)]
597fn encode_server_suffix(value: &r::ServerValue, encoder: &mut Encoder) -> Result<(), CodecError> {
598 match value {
599 r::ServerValue::ParticipantTransportRejected(value) => match &value.reason {
600 r::TransportRejectionReason::FrameTooLarge {
601 complete_frame_bytes,
602 max_frame_bytes,
603 } => {
604 encoder.put_u16(t::TransportReasonTag::FrameTooLarge.wire_value());
605 encoder.put_u64(*complete_frame_bytes);
606 encoder.put_u64(*max_frame_bytes);
607 }
608 r::TransportRejectionReason::DecodeFailed { decode_class } => {
609 encoder.put_u16(t::TransportReasonTag::DecodeFailed.wire_value());
610 encoder.put_u16(decode_class.wire_value());
611 }
612 r::TransportRejectionReason::UnsupportedVersion {
613 presented_version,
614 supported_version,
615 } => {
616 encoder.put_u16(t::TransportReasonTag::UnsupportedVersion.wire_value());
617 encoder.put_protocol_version(*presented_version);
618 encoder.put_protocol_version(*supported_version);
619 }
620 r::TransportRejectionReason::AuthenticationFailed => {
621 encoder.put_u16(t::TransportReasonTag::AuthenticationFailed.wire_value());
622 }
623 r::TransportRejectionReason::ParticipantCapabilityRequired => {
624 encoder.put_u16(t::TransportReasonTag::ParticipantCapabilityRequired.wire_value());
625 encoder.put_string(r::PARTICIPANT_CAPABILITY)?;
626 }
627 },
628 r::ServerValue::AttemptTokenBodyConflict(value) => match value {
629 r::AttemptTokenBodyConflict::CredentialAttach {
630 token,
631 conversation_id,
632 presented_participant_id,
633 presented_generation,
634 presented_marker_delivery_seq,
635 conflict,
636 } => {
637 encoder.put_fixed(token.as_bytes());
638 encoder.put_u16(t::AttemptOperation::CredentialAttachRequest.wire_value());
639 encoder.put_u64(*conversation_id);
640 encoder.put_u64(*presented_participant_id);
641 encoder.put_generation(*presented_generation);
642 encoder.put_option_u64(*presented_marker_delivery_seq);
643 encoder.put_u16(conflict.wire_value());
644 }
645 r::AttemptTokenBodyConflict::Leave {
646 token,
647 conversation_id,
648 presented_participant_id,
649 presented_generation,
650 } => {
651 encoder.put_fixed(token.as_bytes());
652 encoder.put_u16(t::AttemptOperation::LeaveRequest.wire_value());
653 encoder.put_u64(*conversation_id);
654 encoder.put_u64(*presented_participant_id);
655 encoder.put_generation(*presented_generation);
656 encoder.put_u16(t::AttemptConflict::Generation.wire_value());
657 }
658 r::AttemptTokenBodyConflict::RecordAdmission {
659 token,
660 conversation_id,
661 presented_participant_id,
662 presented_generation,
663 } => {
664 encoder.put_fixed(token.as_bytes());
665 encoder.put_u16(t::AttemptOperation::RecordAdmission.wire_value());
666 encoder.put_u64(*conversation_id);
667 encoder.put_u64(*presented_participant_id);
668 encoder.put_generation(*presented_generation);
669 }
670 },
671 r::ServerValue::ConnectionConversationCapacityExceeded(value) => match value {
672 r::ConnectionConversationCapacityExceeded::SemanticRequest { request, limit } => {
673 put_response_envelope(request, encoder);
674 encoder.put_u64(*limit);
675 }
676 r::ConnectionConversationCapacityExceeded::ObserverRecovery {
677 conversation_id,
678 limit,
679 } => {
680 encoder.put_u64(0);
681 encoder.put_u64(*conversation_id);
682 encoder.put_u64(*limit);
683 }
684 },
685 r::ServerValue::ConnectionConversationBindingOccupied(value) => match value {
686 r::ConnectionConversationBindingOccupied::Enrollment {
687 conversation_id,
688 enrollment_token,
689 } => {
690 encoder.put_u64(*conversation_id);
691 encoder.put_fixed(enrollment_token.as_bytes());
692 encoder.put_option_u64(None);
693 }
694 r::ConnectionConversationBindingOccupied::CredentialAttach {
695 conversation_id,
696 participant_id,
697 capability_generation,
698 attach_attempt_token,
699 accept_marker_delivery_seq,
700 } => {
701 encoder.put_u64(*conversation_id);
702 encoder.put_u64(*participant_id);
703 encoder.put_generation(*capability_generation);
704 encoder.put_fixed(attach_attempt_token.as_bytes());
705 encoder.put_option_u64(*accept_marker_delivery_seq);
706 encoder.put_option_u64(Some(*participant_id));
707 }
708 },
709 r::ServerValue::ConversationOrderExhausted(value) => {
710 put_order_allocating(value.request(), encoder);
711 encoder.put_u16(value.counter().wire_value());
712 encoder.put_u64(value.high());
713 encoder.put_option_u64(value.next_value());
714 encoder.put_u128(value.order_remaining());
715 encoder.put_u128(value.reserved_claims());
716 encoder.put_u64(r::ConversationOrderExhausted::REQUIRED_MAJORS);
717 encoder.put_u128(value.resulting_order_remaining());
718 encoder.put_u128(value.resulting_reserved_claims());
719 }
720 r::ServerValue::ParticipantUnknown(value) => {
721 put_participant_reference(&value.request, encoder);
722 }
723 r::ServerValue::NoBinding(value) => put_binding_required(&value.request, encoder),
724 r::ServerValue::StaleAuthority(value) => put_stale_authority(value, encoder),
725 r::ServerValue::Retired(value) => match value {
726 r::Retired::Enrollment {
727 request,
728 participant_id,
729 retired_generation,
730 } => {
731 put_enrollment(request, encoder);
732 encoder.put_u64(*participant_id);
733 encoder.put_generation(*retired_generation);
734 }
735 r::Retired::Participant {
736 request,
737 retired_generation,
738 } => {
739 put_participant_reference(request, encoder);
740 encoder.put_generation(*retired_generation);
741 }
742 },
743 r::ServerValue::MarkerClosureCapacityExceeded(value) => {
744 put_closure_checked(&value.request, encoder);
745 let scope = match value.reason {
746 c::ClosureRefusalReason::Capacity(_) => t::ClosureScope::Capacity,
747 c::ClosureRefusalReason::RecoveryFence => t::ClosureScope::RecoveryFence,
748 c::ClosureRefusalReason::DeliveredMarkerAwaitingAck => {
749 t::ClosureScope::DeliveredMarkerAwaitingAck
750 }
751 c::ClosureRefusalReason::EpisodeChurnLimit => t::ClosureScope::EpisodeChurnLimit,
752 };
753 encoder.put_u16(scope.wire_value());
754 put_closure_snapshot(value.snapshot, encoder);
755 if let c::ClosureRefusalReason::Capacity(reason) = value.reason {
756 encoder.put_u16(t::ResourceDimensionTag::from(reason.dimension).wire_value());
757 encoder.put_u128(reason.required);
758 encoder.put_u128(reason.limit);
759 }
760 }
761 r::ServerValue::EnrollBound(value) => {
762 encoder.put_u64(value.conversation_id());
763 encoder.put_fixed(value.token().as_bytes());
764 encoder.put_u64(value.participant_id());
765 encoder.put_option_generation(value.request_generation());
766 encoder.put_generation(value.capability_generation());
767 encoder.put_fixed(value.attach_secret().as_bytes());
768 encoder.put_binding_epoch(value.origin_binding_epoch());
769 encoder.put_u64(value.persisted_cursor());
770 encoder.put_option_u64(value.accepted_marker_delivery_seq());
771 encoder.put_u128(value.receipt_expires_at());
772 encoder.put_u128(value.provenance_expires_at());
773 }
774 r::ServerValue::EnrollmentKnown(value) => {
775 encoder.put_u64(value.conversation_id);
776 encoder.put_fixed(value.token.as_bytes());
777 encoder.put_u64(value.participant_id);
778 encoder.put_generation(value.current_generation);
779 }
780 r::ServerValue::ReceiptExpired(value) => put_receipt_expired(value, encoder),
781 r::ServerValue::ReceiptCapacityExceeded(value) => {
782 put_receipt_capacity(value, encoder);
783 }
784 r::ServerValue::IdentityCapacityExceeded(value) => {
785 put_enrollment(&value.request, encoder);
786 encoder.put_u16(value.scope.wire_value());
787 encoder.put_u64(value.limit);
788 encoder.put_u64(value.occupied);
789 encoder.put_u64(r::IdentityCapacityExceeded::REQUESTED);
790 }
791 r::ServerValue::ObserverBackpressure(value) => put_observer_backpressure(value, encoder),
792 r::ServerValue::ConversationSequenceExhausted(value) => {
793 put_sequence_allocating(&value.request, encoder);
794 put_sequence_budget(value.sequence_budget, encoder);
795 }
796 r::ServerValue::AttachBound(value) => {
797 encoder.put_u64(value.conversation_id());
798 encoder.put_fixed(value.token().as_bytes());
799 encoder.put_u64(value.participant_id());
800 encoder.put_option_generation(Some(value.request_generation()));
801 encoder.put_generation(value.capability_generation());
802 encoder.put_fixed(value.attach_secret().as_bytes());
803 encoder.put_binding_epoch(value.origin_binding_epoch());
804 encoder.put_u64(value.persisted_cursor());
805 encoder.put_option_u64(value.accepted_marker_delivery_seq());
806 encoder.put_u128(value.receipt_expires_at());
807 encoder.put_u128(value.provenance_expires_at());
808 }
809 r::ServerValue::StaleOrUnknownReceipt(value) => {
810 encoder.put_u64(value.conversation_id);
811 encoder.put_fixed(value.token.as_bytes());
812 encoder.put_u64(value.participant_id);
813 encoder.put_generation(value.presented_generation);
814 encoder.put_option_u64(value.presented_marker_delivery_seq);
815 encoder.put_generation(value.current_generation);
816 }
817 r::ServerValue::MarkerNotDelivered(value) => {
818 put_marker_proof(&value.request, encoder);
819 encoder.put_u16(value.reason.wire_value());
820 encoder.put_u64(value.expected_marker_delivery_seq);
821 }
822 r::ServerValue::MarkerMismatch(value) => {
823 put_marker_proof(&value.request, encoder);
824 encoder.put_u16(value.mismatch.reason().wire_value());
825 match value.mismatch {
826 r::MarkerMismatchBody::BelowCursor { current_cursor } => {
827 encoder.put_u64(current_cursor);
828 }
829 r::MarkerMismatchBody::NoMarkerExpected => {}
830 r::MarkerMismatchBody::ExpectedDifferentMarker {
831 expected_marker_delivery_seq,
832 } => encoder.put_u64(expected_marker_delivery_seq),
833 }
834 }
835 r::ServerValue::Bound(value) | r::ServerValue::UnboundReceipt(value) => {
836 put_receipt_replay(value, encoder);
837 }
838 r::ServerValue::DetachCommitted(value) => {
839 encoder.put_u64(value.conversation_id());
840 encoder.put_u64(value.participant_id());
841 encoder.put_generation(value.capability_generation());
842 encoder.put_fixed(value.detach_attempt_token().as_bytes());
843 encoder.put_binding_epoch(value.committed_binding_epoch());
844 encoder.put_u64(value.detached_delivery_seq());
845 }
846 r::ServerValue::DetachInProgress(value) => {
847 encoder.put_u64(value.conversation_id);
848 encoder.put_u64(value.participant_id);
849 encoder.put_fixed(value.presented_token.as_bytes());
850 encoder.put_generation(value.presented_generation);
851 encoder.put_binding_epoch(value.committed_binding_epoch);
852 }
853 r::ServerValue::AckCommitted(value) => {
854 put_participant_ack(value.request(), encoder);
855 encoder.put_u64(value.current_cursor());
856 }
857 r::ServerValue::AckNoOp(value) => match value {
858 r::AckNoOp::ParticipantAck(request) => {
859 put_participant_ack(request, encoder);
860 encoder.put_u64(value.current_cursor());
861 }
862 r::AckNoOp::MarkerAck(request) => {
863 put_marker_ack(request, encoder);
864 encoder.put_u64(value.current_cursor());
865 }
866 },
867 r::ServerValue::AckGap(value) => {
868 put_participant_ack(value.request(), encoder);
869 encoder.put_u64(value.current_cursor());
870 encoder.put_u16(value.reason().wire_value());
871 }
872 r::ServerValue::AckRegression(value) => {
873 put_participant_ack(value.request(), encoder);
874 encoder.put_u64(value.current_cursor());
875 encoder.put_u16(value.reason().wire_value());
876 }
877 r::ServerValue::LeaveCommitted(value) => {
878 encoder.put_u64(value.conversation_id());
879 encoder.put_fixed(value.leave_attempt_token().as_bytes());
880 encoder.put_u64(value.participant_id());
881 encoder.put_generation(value.presented_generation());
882 encoder.put_generation(value.retired_generation());
883 encoder.put_option_binding_epoch(value.ended_binding_epoch());
884 encoder.put_option_u64(value.prior_terminal_delivery_seq());
885 encoder.put_u64(value.left_delivery_seq());
886 }
887 r::ServerValue::MarkerAckCommitted(value) => {
888 put_marker_ack(value.request(), encoder);
889 encoder.put_u64(value.current_cursor());
890 }
891 r::ServerValue::RecordCommitted(value) => {
892 put_record_admission(value.request(), encoder);
893 encoder.put_u64(value.sender_participant_id());
894 encoder.put_u64(value.delivery_seq());
895 }
896 r::ServerValue::RecordTooLarge(value) => {
897 put_record_admission(&value.request, encoder);
898 encoder.put_u16(t::ResourceDimensionTag::from(value.dimension).wire_value());
899 encoder.put_resource_vector(value.encoded_record_charge);
900 encoder.put_resource_vector(value.max_ordinary_record_charge);
901 }
902 r::ServerValue::ObserverRecoveryAccepted(value) => {
903 let count: u64 = value
904 .statuses
905 .len()
906 .try_into()
907 .map_err(|_| CodecError::LengthOverflow)?;
908 encoder.put_u64(count);
909 for status in &value.statuses {
910 encoder.put_u64(status.conversation_id);
911 encoder.put_u64(status.refused_epoch);
912 encoder.put_u64(status.current_observer_progress);
913 encoder.put_bool(status.armed);
914 encoder.put_bool(status.progressed);
915 }
916 }
917 r::ServerValue::InvalidObserverEpoch(value) => {
918 encoder.put_u64(0);
919 encoder.put_u16(value.reason().wire_value());
920 match value {
921 r::InvalidObserverEpoch::ConversationUnknown {
922 conversation_id,
923 presented_epoch,
924 } => {
925 encoder.put_u64(*conversation_id);
926 encoder.put_u64(*presented_epoch);
927 encoder.put_option_u64(None);
928 }
929 r::InvalidObserverEpoch::EpochAhead {
930 conversation_id,
931 presented_epoch,
932 current_observer_progress,
933 } => {
934 encoder.put_u64(*conversation_id);
935 encoder.put_u64(*presented_epoch);
936 encoder.put_option_u64(Some(*current_observer_progress));
937 }
938 }
939 }
940 r::ServerValue::InvalidObserverEpochList(value) => {
941 encoder.put_u64(0);
942 encoder.put_u16(value.reason().wire_value());
943 match value {
944 r::InvalidObserverEpochList::TooManyEntries {
945 presented_entries,
946 max_entries,
947 } => {
948 encoder.put_u64(*presented_entries);
949 encoder.put_u64(*max_entries);
950 }
951 r::InvalidObserverEpochList::DuplicateConversation {
952 conversation_id,
953 first_index,
954 duplicate_index,
955 } => {
956 encoder.put_u64(*conversation_id);
957 encoder.put_u64(*first_index);
958 encoder.put_u64(*duplicate_index);
959 }
960 }
961 }
962 r::ServerValue::MarkerSettlementBackpressure(value) => match value {
963 r::MarkerSettlementBackpressure::CredentialAttach {
964 conversation_id,
965 refused_epoch,
966 }
967 | r::MarkerSettlementBackpressure::Detach {
968 conversation_id,
969 refused_epoch,
970 } => {
971 encoder.put_u64(*conversation_id);
972 encoder.put_u64(*refused_epoch);
973 }
974 },
975 r::ServerValue::EnrollmentSettlementBackpressure(value) => {
976 encoder.put_u64(value.conversation_id);
977 }
978 r::ServerValue::RecordAdmissionProtocolFault(value) => {
979 put_record_admission(&value.request, encoder);
980 encoder.put_u16(value.class.tag().wire_value());
981 }
982 }
983 Ok(())
984}
985
986fn put_stale_authority(value: &r::StaleAuthority, encoder: &mut Encoder) {
987 match value {
988 r::StaleAuthority::Live {
989 request,
990 current_generation,
991 } => {
992 match request {
993 r::CommonStaleAuthorityEnvelope::CredentialAttach(value) => {
994 put_attach(value, encoder);
995 }
996 r::CommonStaleAuthorityEnvelope::ParticipantAck(value) => {
997 put_participant_ack(value, encoder);
998 }
999 r::CommonStaleAuthorityEnvelope::MarkerAck(value) => {
1000 put_marker_ack(value, encoder);
1001 }
1002 r::CommonStaleAuthorityEnvelope::RecordAdmission(value) => {
1003 put_record_admission(value, encoder);
1004 }
1005 }
1006 encoder.put_generation(*current_generation);
1007 }
1008 r::StaleAuthority::Detach(value) => {
1009 encoder.put_u16(value.authority_state_tag().wire_value());
1010 match value {
1011 r::DetachStaleAuthority::Live {
1012 conversation_id,
1013 participant_id,
1014 capability_generation,
1015 detach_attempt_token,
1016 current_generation,
1017 } => {
1018 encoder.put_u64(*conversation_id);
1019 encoder.put_u64(*participant_id);
1020 encoder.put_generation(*capability_generation);
1021 encoder.put_fixed(detach_attempt_token.as_bytes());
1022 encoder.put_generation(*current_generation);
1023 }
1024 r::DetachStaleAuthority::TerminalizedDetachCell(value) => {
1025 encoder.put_u64(value.conversation_id());
1026 encoder.put_u64(value.participant_id());
1027 encoder.put_generation(value.capability_generation());
1028 encoder.put_fixed(value.detach_attempt_token().as_bytes());
1029 encoder.put_generation(value.current_generation());
1030 encoder.put_binding_epoch(value.committed_binding_epoch());
1031 encoder.put_u16(value.binding_state().tag().wire_value());
1032 if let r::BindingStateView::Bound {
1033 current_binding_epoch,
1034 } = value.binding_state()
1035 {
1036 encoder.put_binding_epoch(current_binding_epoch);
1037 }
1038 }
1039 }
1040 }
1041 r::StaleAuthority::Leave(value) => {
1042 encoder.put_u16(value.authority_state_tag().wire_value());
1043 match value {
1044 r::LeaveStaleAuthority::Live {
1045 conversation_id,
1046 participant_id,
1047 presented_generation,
1048 leave_attempt_token,
1049 current_generation,
1050 } => {
1051 encoder.put_u64(*conversation_id);
1052 encoder.put_u64(*participant_id);
1053 encoder.put_generation(*presented_generation);
1054 encoder.put_fixed(leave_attempt_token.as_bytes());
1055 encoder.put_generation(*current_generation);
1056 }
1057 r::LeaveStaleAuthority::CommittedLeaveTombstone {
1058 conversation_id,
1059 participant_id,
1060 presented_generation,
1061 leave_attempt_token,
1062 retired_generation,
1063 } => {
1064 encoder.put_u64(*conversation_id);
1065 encoder.put_u64(*participant_id);
1066 encoder.put_generation(*presented_generation);
1067 encoder.put_fixed(leave_attempt_token.as_bytes());
1068 encoder.put_generation(*retired_generation);
1069 }
1070 }
1071 }
1072 }
1073}
1074
1075fn put_receipt_expired(value: &r::ReceiptExpired, encoder: &mut Encoder) {
1076 match value {
1077 r::ReceiptExpired::Enrollment {
1078 conversation_id,
1079 token,
1080 participant_id,
1081 result_generation,
1082 current_generation,
1083 reason,
1084 } => {
1085 encoder.put_u64(*conversation_id);
1086 encoder.put_fixed(token.as_bytes());
1087 encoder.put_u64(*participant_id);
1088 encoder.put_option_generation(None);
1089 encoder.put_generation(*result_generation);
1090 encoder.put_generation(*current_generation);
1091 encoder.put_u16(reason.wire_value());
1092 }
1093 r::ReceiptExpired::CredentialAttach {
1094 conversation_id,
1095 token,
1096 participant_id,
1097 presented_generation,
1098 presented_marker_delivery_seq,
1099 result_generation,
1100 current_generation,
1101 reason,
1102 } => {
1103 encoder.put_u64(*conversation_id);
1104 encoder.put_fixed(token.as_bytes());
1105 encoder.put_u64(*participant_id);
1106 encoder.put_option_generation(Some(*presented_generation));
1107 encoder.put_option_u64(*presented_marker_delivery_seq);
1108 encoder.put_generation(*result_generation);
1109 encoder.put_generation(*current_generation);
1110 encoder.put_u16(reason.wire_value());
1111 }
1112 }
1113}
1114
1115fn put_receipt_capacity(value: &r::ReceiptCapacityExceeded, encoder: &mut Encoder) {
1116 match value {
1117 r::ReceiptCapacityExceeded::Enrollment {
1118 request,
1119 scope,
1120 limit,
1121 occupied,
1122 } => {
1123 put_enrollment(request, encoder);
1124 encoder.put_u16(scope.wire_scope().wire_value());
1125 encoder.put_u64(*limit);
1126 encoder.put_u64(*occupied);
1127 encoder.put_u64(r::ReceiptCapacityExceeded::REQUESTED);
1128 }
1129 r::ReceiptCapacityExceeded::CredentialAttach {
1130 request,
1131 scope,
1132 limit,
1133 occupied,
1134 } => {
1135 put_attach(request, encoder);
1136 encoder.put_u16(scope.wire_value());
1137 encoder.put_u64(*limit);
1138 encoder.put_u64(*occupied);
1139 encoder.put_u64(r::ReceiptCapacityExceeded::REQUESTED);
1140 }
1141 }
1142}
1143
1144fn put_observer_backpressure(value: &r::ObserverBackpressure, encoder: &mut Encoder) {
1145 match value {
1146 r::ObserverBackpressure::Enrollment { request, state } => {
1147 put_enrollment(request, encoder);
1148 put_backpressure_state(*state, encoder);
1149 }
1150 r::ObserverBackpressure::CredentialAttach { request, state } => {
1151 put_attach(request, encoder);
1152 put_backpressure_state(*state, encoder);
1153 }
1154 r::ObserverBackpressure::Detach {
1155 request,
1156 committed_binding_epoch,
1157 state,
1158 } => {
1159 put_detach(request, encoder);
1160 encoder.put_binding_epoch(*committed_binding_epoch);
1161 put_backpressure_state(*state, encoder);
1162 }
1163 r::ObserverBackpressure::Leave {
1164 request,
1165 state,
1166 prior_terminal_cell_exists,
1167 } => {
1168 put_leave(request, encoder);
1169 put_backpressure_state(*state, encoder);
1170 encoder.put_bool(*prior_terminal_cell_exists);
1171 }
1172 r::ObserverBackpressure::RecordAdmission { request, state } => {
1173 put_record_admission(request, encoder);
1174 put_backpressure_state(*state, encoder);
1175 }
1176 }
1177}
1178
1179fn put_backpressure_state(value: r::ObserverBackpressureState, encoder: &mut Encoder) {
1180 encoder.put_u64(value.backpressure_epoch());
1181 encoder.put_u64(value.observer_progress());
1182}
1183
1184fn put_receipt_replay(value: &r::ReceiptReplay, encoder: &mut Encoder) {
1185 match value {
1186 r::ReceiptReplay::Enrollment(value) => {
1187 encoder.put_u64(value.conversation_id());
1188 encoder.put_fixed(value.token().as_bytes());
1189 encoder.put_u64(value.participant_id());
1190 encoder.put_option_generation(None);
1191 encoder.put_generation(value.capability_generation());
1192 encoder.put_fixed(value.attach_secret().as_bytes());
1193 encoder.put_binding_epoch(value.origin_binding_epoch());
1194 encoder.put_u64(value.persisted_cursor());
1195 encoder.put_option_u64(None);
1196 encoder.put_u128(value.receipt_expires_at());
1197 encoder.put_u128(value.provenance_expires_at());
1198 }
1199 r::ReceiptReplay::CredentialAttach(value) => {
1200 encoder.put_u64(value.conversation_id());
1201 encoder.put_fixed(value.token().as_bytes());
1202 encoder.put_u64(value.participant_id());
1203 encoder.put_option_generation(Some(value.request_generation()));
1204 encoder.put_generation(value.capability_generation());
1205 encoder.put_fixed(value.attach_secret().as_bytes());
1206 encoder.put_binding_epoch(value.origin_binding_epoch());
1207 encoder.put_u64(value.persisted_cursor());
1208 encoder.put_option_u64(value.accepted_marker_delivery_seq());
1209 encoder.put_u128(value.receipt_expires_at());
1210 encoder.put_u128(value.provenance_expires_at());
1211 }
1212 }
1213}
1214
1215const fn carries_origin(discriminant: t::ServerDiscriminant) -> bool {
1223 use t::ServerDiscriminant as D;
1224
1225 !matches!(
1226 discriminant,
1227 D::ParticipantTransportRejected
1228 | D::ObserverRecoveryAccepted
1229 | D::InvalidObserverEpoch
1230 | D::InvalidObserverEpochList
1231 | D::ObserverRecoveryConnectionCapacityExceeded
1232 )
1233}
1234
1235const fn origin_is_valid(
1236 discriminant: t::ServerDiscriminant,
1237 origin: t::ClientDiscriminant,
1238) -> bool {
1239 use t::ClientDiscriminant as O;
1240 use t::ServerDiscriminant as D;
1241
1242 match discriminant {
1243 D::AttemptTokenBodyConflict => matches!(
1244 origin,
1245 O::CredentialAttachRequest | O::LeaveRequest | O::RecordAdmission
1246 ),
1247 D::ConnectionConversationCapacityExceeded => {
1248 !matches!(origin, O::ObserverRecoveryHandshake)
1249 }
1250 D::ConnectionConversationBindingOccupied => {
1251 matches!(origin, O::EnrollmentRequest | O::CredentialAttachRequest)
1252 }
1253 D::ConversationOrderExhausted => matches!(
1254 origin,
1255 O::EnrollmentRequest | O::CredentialAttachRequest | O::RecordAdmission
1256 ),
1257 D::ParticipantUnknown => {
1258 !matches!(origin, O::EnrollmentRequest | O::ObserverRecoveryHandshake)
1259 }
1260 D::NoBinding => !matches!(
1261 origin,
1262 O::EnrollmentRequest | O::CredentialAttachRequest | O::ObserverRecoveryHandshake
1263 ),
1264 D::StaleAuthority => !matches!(origin, O::EnrollmentRequest | O::ObserverRecoveryHandshake),
1265 D::Retired => !matches!(origin, O::ObserverRecoveryHandshake),
1266 D::MarkerClosureCapacityExceeded => matches!(
1267 origin,
1268 O::EnrollmentRequest
1269 | O::CredentialAttachRequest
1270 | O::LeaveRequest
1271 | O::RecordAdmission
1272 ),
1273 D::EnrollBound | D::EnrollmentKnown | D::IdentityCapacityExceeded => {
1274 matches!(origin, O::EnrollmentRequest)
1275 }
1276 D::ReceiptExpired | D::ReceiptCapacityExceeded | D::Bound | D::UnboundReceipt => {
1277 matches!(origin, O::EnrollmentRequest | O::CredentialAttachRequest)
1278 }
1279 D::ObserverBackpressure => matches!(
1280 origin,
1281 O::EnrollmentRequest
1282 | O::CredentialAttachRequest
1283 | O::DetachRequest
1284 | O::LeaveRequest
1285 | O::RecordAdmission
1286 ),
1287 D::ConversationSequenceExhausted => matches!(
1288 origin,
1289 O::EnrollmentRequest | O::CredentialAttachRequest | O::RecordAdmission
1290 ),
1291 D::AttachBound | D::StaleOrUnknownReceipt => {
1292 matches!(origin, O::CredentialAttachRequest)
1293 }
1294 D::MarkerNotDelivered | D::MarkerMismatch => {
1295 matches!(origin, O::CredentialAttachRequest | O::MarkerAck)
1296 }
1297 D::DetachCommitted | D::DetachInProgress => matches!(origin, O::DetachRequest),
1298 D::AckCommitted | D::AckGap | D::AckRegression => matches!(origin, O::ParticipantAck),
1299 D::AckNoOp => matches!(origin, O::ParticipantAck | O::MarkerAck),
1300 D::LeaveCommitted => matches!(origin, O::LeaveRequest),
1301 D::MarkerAckCommitted => matches!(origin, O::MarkerAck),
1302 D::RecordCommitted | D::RecordTooLarge | D::RecordAdmissionProtocolFault => {
1303 matches!(origin, O::RecordAdmission)
1304 }
1305 D::MarkerSettlementBackpressure => {
1306 matches!(origin, O::CredentialAttachRequest | O::DetachRequest)
1307 }
1308 D::EnrollmentSettlementBackpressure => matches!(origin, O::EnrollmentRequest),
1309 D::ParticipantTransportRejected
1310 | D::ObserverRecoveryAccepted
1311 | D::InvalidObserverEpoch
1312 | D::InvalidObserverEpochList
1313 | D::ObserverRecoveryConnectionCapacityExceeded => false,
1314 }
1315}
1316
1317fn take_enrollment(decoder: &mut Decoder<'_>) -> Result<e::EnrollmentEnvelope, CodecError> {
1321 Ok(e::EnrollmentEnvelope {
1322 conversation_id: decoder.take_u64()?,
1323 enrollment_token: EnrollmentToken::new(decoder.take_fixed::<TOKEN_LEN>()?),
1324 })
1325}
1326
1327fn take_attach(decoder: &mut Decoder<'_>) -> Result<e::AttachEnvelope, CodecError> {
1328 Ok(e::AttachEnvelope {
1329 conversation_id: decoder.take_u64()?,
1330 participant_id: decoder.take_u64()?,
1331 capability_generation: decoder.take_generation()?,
1332 attach_attempt_token: AttachAttemptToken::new(decoder.take_fixed::<TOKEN_LEN>()?),
1333 accept_marker_delivery_seq: decoder.take_option_u64()?,
1334 })
1335}
1336
1337fn take_detach(decoder: &mut Decoder<'_>) -> Result<e::DetachEnvelope, CodecError> {
1338 Ok(e::DetachEnvelope {
1339 conversation_id: decoder.take_u64()?,
1340 participant_id: decoder.take_u64()?,
1341 capability_generation: decoder.take_generation()?,
1342 detach_attempt_token: DetachAttemptToken::new(decoder.take_fixed::<TOKEN_LEN>()?),
1343 })
1344}
1345
1346fn take_participant_ack(
1347 decoder: &mut Decoder<'_>,
1348) -> Result<e::ParticipantAckEnvelope, CodecError> {
1349 Ok(e::ParticipantAckEnvelope {
1350 conversation_id: decoder.take_u64()?,
1351 participant_id: decoder.take_u64()?,
1352 capability_generation: decoder.take_generation()?,
1353 through_seq: decoder.take_u64()?,
1354 })
1355}
1356
1357fn take_leave(decoder: &mut Decoder<'_>) -> Result<e::LeaveEnvelope, CodecError> {
1358 Ok(e::LeaveEnvelope {
1359 conversation_id: decoder.take_u64()?,
1360 participant_id: decoder.take_u64()?,
1361 capability_generation: decoder.take_generation()?,
1362 leave_attempt_token: LeaveAttemptToken::new(decoder.take_fixed::<TOKEN_LEN>()?),
1363 })
1364}
1365
1366fn take_marker_ack(decoder: &mut Decoder<'_>) -> Result<e::MarkerAckEnvelope, CodecError> {
1367 Ok(e::MarkerAckEnvelope {
1368 conversation_id: decoder.take_u64()?,
1369 participant_id: decoder.take_u64()?,
1370 capability_generation: decoder.take_generation()?,
1371 marker_delivery_seq: decoder.take_u64()?,
1372 })
1373}
1374
1375fn take_record_admission(
1376 decoder: &mut Decoder<'_>,
1377) -> Result<e::RecordAdmissionEnvelope, CodecError> {
1378 Ok(e::RecordAdmissionEnvelope {
1379 conversation_id: decoder.take_u64()?,
1380 participant_id: decoder.take_u64()?,
1381 capability_generation: decoder.take_generation()?,
1382 record_admission_attempt_token: RecordAdmissionAttemptToken::new(
1383 decoder.take_fixed::<TOKEN_LEN>()?,
1384 ),
1385 })
1386}
1387
1388fn take_response_envelope(
1389 origin: t::ClientDiscriminant,
1390 decoder: &mut Decoder<'_>,
1391) -> Result<e::ResponseEnvelope, CodecError> {
1392 match origin {
1393 t::ClientDiscriminant::EnrollmentRequest => {
1394 take_enrollment(decoder).map(e::ResponseEnvelope::Enrollment)
1395 }
1396 t::ClientDiscriminant::CredentialAttachRequest => {
1397 take_attach(decoder).map(e::ResponseEnvelope::CredentialAttach)
1398 }
1399 t::ClientDiscriminant::DetachRequest => {
1400 take_detach(decoder).map(e::ResponseEnvelope::Detach)
1401 }
1402 t::ClientDiscriminant::ParticipantAck => {
1403 take_participant_ack(decoder).map(e::ResponseEnvelope::ParticipantAck)
1404 }
1405 t::ClientDiscriminant::LeaveRequest => take_leave(decoder).map(e::ResponseEnvelope::Leave),
1406 t::ClientDiscriminant::MarkerAck => {
1407 take_marker_ack(decoder).map(e::ResponseEnvelope::MarkerAck)
1408 }
1409 t::ClientDiscriminant::RecordAdmission => {
1410 take_record_admission(decoder).map(e::ResponseEnvelope::RecordAdmission)
1411 }
1412 t::ClientDiscriminant::ObserverRecoveryHandshake => {
1413 Err(decode_error(t::DecodeClass::InvalidField))
1414 }
1415 }
1416}
1417
1418fn take_participant_reference(
1419 origin: t::ClientDiscriminant,
1420 decoder: &mut Decoder<'_>,
1421) -> Result<r::ParticipantReferenceEnvelope, CodecError> {
1422 match origin {
1423 t::ClientDiscriminant::CredentialAttachRequest => {
1424 take_attach(decoder).map(r::ParticipantReferenceEnvelope::CredentialAttach)
1425 }
1426 t::ClientDiscriminant::DetachRequest => {
1427 take_detach(decoder).map(r::ParticipantReferenceEnvelope::Detach)
1428 }
1429 t::ClientDiscriminant::ParticipantAck => {
1430 take_participant_ack(decoder).map(r::ParticipantReferenceEnvelope::ParticipantAck)
1431 }
1432 t::ClientDiscriminant::LeaveRequest => {
1433 take_leave(decoder).map(r::ParticipantReferenceEnvelope::Leave)
1434 }
1435 t::ClientDiscriminant::MarkerAck => {
1436 take_marker_ack(decoder).map(r::ParticipantReferenceEnvelope::MarkerAck)
1437 }
1438 t::ClientDiscriminant::RecordAdmission => {
1439 take_record_admission(decoder).map(r::ParticipantReferenceEnvelope::RecordAdmission)
1440 }
1441 t::ClientDiscriminant::EnrollmentRequest
1442 | t::ClientDiscriminant::ObserverRecoveryHandshake => {
1443 Err(decode_error(t::DecodeClass::InvalidField))
1444 }
1445 }
1446}
1447
1448fn take_binding_required(
1449 origin: t::ClientDiscriminant,
1450 decoder: &mut Decoder<'_>,
1451) -> Result<r::BindingRequiredEnvelope, CodecError> {
1452 match origin {
1453 t::ClientDiscriminant::DetachRequest => {
1454 take_detach(decoder).map(r::BindingRequiredEnvelope::Detach)
1455 }
1456 t::ClientDiscriminant::ParticipantAck => {
1457 take_participant_ack(decoder).map(r::BindingRequiredEnvelope::ParticipantAck)
1458 }
1459 t::ClientDiscriminant::LeaveRequest => {
1460 take_leave(decoder).map(r::BindingRequiredEnvelope::Leave)
1461 }
1462 t::ClientDiscriminant::MarkerAck => {
1463 take_marker_ack(decoder).map(r::BindingRequiredEnvelope::MarkerAck)
1464 }
1465 t::ClientDiscriminant::RecordAdmission => {
1466 take_record_admission(decoder).map(r::BindingRequiredEnvelope::RecordAdmission)
1467 }
1468 t::ClientDiscriminant::EnrollmentRequest
1469 | t::ClientDiscriminant::CredentialAttachRequest
1470 | t::ClientDiscriminant::ObserverRecoveryHandshake => {
1471 Err(decode_error(t::DecodeClass::InvalidField))
1472 }
1473 }
1474}
1475
1476fn take_order_allocating(
1477 origin: t::ClientDiscriminant,
1478 decoder: &mut Decoder<'_>,
1479) -> Result<r::OrderAllocatingEnvelope, CodecError> {
1480 match origin {
1481 t::ClientDiscriminant::EnrollmentRequest => {
1482 take_enrollment(decoder).map(r::OrderAllocatingEnvelope::Enrollment)
1483 }
1484 t::ClientDiscriminant::CredentialAttachRequest => {
1485 take_attach(decoder).map(r::OrderAllocatingEnvelope::CredentialAttach)
1486 }
1487 t::ClientDiscriminant::RecordAdmission => {
1488 take_record_admission(decoder).map(r::OrderAllocatingEnvelope::RecordAdmission)
1489 }
1490 _ => Err(decode_error(t::DecodeClass::InvalidField)),
1491 }
1492}
1493
1494fn take_closure_checked(
1495 origin: t::ClientDiscriminant,
1496 decoder: &mut Decoder<'_>,
1497) -> Result<c::ClosureCheckedEnvelope, CodecError> {
1498 match origin {
1499 t::ClientDiscriminant::EnrollmentRequest => {
1500 take_enrollment(decoder).map(c::ClosureCheckedEnvelope::Enrollment)
1501 }
1502 t::ClientDiscriminant::CredentialAttachRequest => {
1503 take_attach(decoder).map(c::ClosureCheckedEnvelope::CredentialAttach)
1504 }
1505 t::ClientDiscriminant::LeaveRequest => {
1506 take_leave(decoder).map(c::ClosureCheckedEnvelope::Leave)
1507 }
1508 t::ClientDiscriminant::RecordAdmission => {
1509 take_record_admission(decoder).map(c::ClosureCheckedEnvelope::RecordAdmission)
1510 }
1511 _ => Err(decode_error(t::DecodeClass::InvalidField)),
1512 }
1513}
1514
1515fn take_sequence_allocating(
1516 origin: t::ClientDiscriminant,
1517 decoder: &mut Decoder<'_>,
1518) -> Result<r::SequenceAllocatingEnvelope, CodecError> {
1519 match origin {
1520 t::ClientDiscriminant::EnrollmentRequest => {
1521 take_enrollment(decoder).map(r::SequenceAllocatingEnvelope::Enrollment)
1522 }
1523 t::ClientDiscriminant::CredentialAttachRequest => {
1524 take_attach(decoder).map(r::SequenceAllocatingEnvelope::CredentialAttach)
1525 }
1526 t::ClientDiscriminant::RecordAdmission => {
1527 take_record_admission(decoder).map(r::SequenceAllocatingEnvelope::RecordAdmission)
1528 }
1529 _ => Err(decode_error(t::DecodeClass::InvalidField)),
1530 }
1531}
1532
1533fn take_marker_proof(
1534 origin: t::ClientDiscriminant,
1535 decoder: &mut Decoder<'_>,
1536) -> Result<r::MarkerProofRequest, CodecError> {
1537 match origin {
1538 t::ClientDiscriminant::CredentialAttachRequest => Ok(
1539 r::MarkerProofRequest::CredentialAttach(r::AttachMarkerProof {
1540 conversation_id: decoder.take_u64()?,
1541 token: AttachAttemptToken::new(decoder.take_fixed::<TOKEN_LEN>()?),
1542 participant_id: decoder.take_u64()?,
1543 capability_generation: decoder.take_generation()?,
1544 requested_marker_delivery_seq: decoder.take_u64()?,
1545 }),
1546 ),
1547 t::ClientDiscriminant::MarkerAck => {
1548 Ok(r::MarkerProofRequest::MarkerAck(r::MarkerAckProof {
1549 conversation_id: decoder.take_u64()?,
1550 participant_id: decoder.take_u64()?,
1551 capability_generation: decoder.take_generation()?,
1552 requested_marker_delivery_seq: decoder.take_u64()?,
1553 }))
1554 }
1555 _ => Err(decode_error(t::DecodeClass::InvalidField)),
1556 }
1557}
1558
1559fn take_repayment_edge(decoder: &mut Decoder<'_>) -> Result<c::RepaymentEdge, CodecError> {
1560 let tag = t::RepaymentEdgeTag::try_from(decoder.take_u16()?)
1561 .map_err(|_| decode_error(t::DecodeClass::InvalidField))?;
1562 match tag {
1563 t::RepaymentEdgeTag::None => Ok(c::RepaymentEdge::None),
1564 t::RepaymentEdgeTag::ObserverProjection => Ok(c::RepaymentEdge::ObserverProjection {
1565 through_seq: decoder.take_u64()?,
1566 }),
1567 t::RepaymentEdgeTag::PhysicalCompaction => Ok(c::RepaymentEdge::PhysicalCompaction {
1568 from_floor: decoder.take_u64()?,
1569 through_seq: decoder.take_u64()?,
1570 }),
1571 t::RepaymentEdgeTag::MarkerDelivery => Ok(c::RepaymentEdge::MarkerDelivery {
1572 participant_id: decoder.take_u64()?,
1573 binding_epoch: decoder.take_binding_epoch()?,
1574 marker_delivery_seq: decoder.take_u64()?,
1575 }),
1576 t::RepaymentEdgeTag::ParticipantCursorProgress => Ok(
1577 c::RepaymentEdge::ParticipantCursorProgress(c::ParticipantCursorProgressEdge {
1578 participant_id: decoder.take_u64()?,
1579 binding_epoch: decoder.take_binding_epoch()?,
1580 through_seq: decoder.take_u64()?,
1581 marker_delivery_seq: decoder.take_option_u64()?,
1582 }),
1583 ),
1584 t::RepaymentEdgeTag::DetachedCredentialRecovery => {
1585 Ok(c::RepaymentEdge::DetachedCredentialRecovery {
1586 participant_id: decoder.take_u64()?,
1587 marker_delivery_seq: decoder.take_u64()?,
1588 prior_binding_epoch: decoder.take_binding_epoch()?,
1589 })
1590 }
1591 t::RepaymentEdgeTag::DetachedMarkerRelease => Ok(c::RepaymentEdge::DetachedMarkerRelease {
1592 participant_id: decoder.take_u64()?,
1593 marker_delivery_seq: decoder.take_u64()?,
1594 last_dead_binding_epoch: decoder.take_binding_epoch()?,
1595 }),
1596 t::RepaymentEdgeTag::DetachedCursorRelease => Ok(c::RepaymentEdge::DetachedCursorRelease {
1597 participant_id: decoder.take_u64()?,
1598 last_dead_binding_epoch: decoder.take_binding_epoch()?,
1599 }),
1600 }
1601}
1602
1603fn take_closure_snapshot(decoder: &mut Decoder<'_>) -> Result<c::ClosureSnapshot, CodecError> {
1604 Ok(c::ClosureSnapshot {
1605 marker_capacity_credits: decoder.take_u64()?,
1606 marker_anchors: decoder.take_u64()?,
1607 entry_debt: decoder.take_u64()?,
1608 byte_debt: decoder.take_u64()?,
1609 repayment_edge: take_repayment_edge(decoder)?,
1610 edge_sequence_claims: decoder.take_u64()?,
1611 edge_order_position_claims: decoder.take_u64()?,
1612 edge_k_remaining: decoder.take_resource_vector()?,
1613 k_headroom: decoder.take_wide_resource_vector()?,
1614 episode_churn_used: decoder.take_u64()?,
1615 delta_cycles: decoder.take_u64()?,
1616 episode_churn_limit: decoder.take_u64()?,
1617 })
1618}
1619
1620fn take_sequence_budget(decoder: &mut Decoder<'_>) -> Result<super::SequenceBudget, CodecError> {
1621 Ok(super::SequenceBudget {
1622 high_watermark: decoder.take_u64()?,
1623 remaining: decoder.take_u64()?,
1624 e: decoder.take_u64()?,
1625 t: decoder.take_u64()?,
1626 m: decoder.take_u64()?,
1627 rs: decoder.take_u64()?,
1628 rt: decoder.take_u64()?,
1629 l_times_t: decoder.take_u128()?,
1630 l_times_rt: decoder.take_u128()?,
1631 l_other_times_e: decoder.take_u128()?,
1632 })
1633}
1634
1635fn required_origin(
1636 origin: Option<t::ClientDiscriminant>,
1637) -> Result<t::ClientDiscriminant, CodecError> {
1638 origin.ok_or_else(|| decode_error(t::DecodeClass::InvalidField))
1639}
1640
1641#[allow(clippy::too_many_lines)]
1642fn decode_server_suffix(
1643 discriminant: t::ServerDiscriminant,
1644 origin: Option<t::ClientDiscriminant>,
1645 decoder: &mut Decoder<'_>,
1646) -> Result<r::ServerValue, CodecError> {
1647 use t::ServerDiscriminant as D;
1648
1649 match discriminant {
1650 D::ParticipantTransportRejected => {
1651 let reason = match take_tag::<t::TransportReasonTag>(decoder)? {
1652 t::TransportReasonTag::FrameTooLarge => {
1653 r::TransportRejectionReason::FrameTooLarge {
1654 complete_frame_bytes: decoder.take_u64()?,
1655 max_frame_bytes: decoder.take_u64()?,
1656 }
1657 }
1658 t::TransportReasonTag::DecodeFailed => {
1659 let decode_class = take_tag::<t::DecodeClass>(decoder)?;
1660 r::TransportRejectionReason::DecodeFailed { decode_class }
1661 }
1662 t::TransportReasonTag::UnsupportedVersion => {
1663 r::TransportRejectionReason::UnsupportedVersion {
1664 presented_version: decoder.take_protocol_version()?,
1665 supported_version: decoder.take_protocol_version()?,
1666 }
1667 }
1668 t::TransportReasonTag::AuthenticationFailed => {
1669 r::TransportRejectionReason::AuthenticationFailed
1670 }
1671 t::TransportReasonTag::ParticipantCapabilityRequired => {
1672 let capability = decoder.take_string()?;
1673 if capability != r::PARTICIPANT_CAPABILITY {
1674 decoder.invalidate();
1675 }
1676 r::TransportRejectionReason::ParticipantCapabilityRequired
1677 }
1678 };
1679 Ok(r::ServerValue::ParticipantTransportRejected(
1680 r::ParticipantTransportRejected { reason },
1681 ))
1682 }
1683 D::AttemptTokenBodyConflict => {
1684 let origin = required_origin(origin)?;
1685 let token = decoder.take_fixed::<TOKEN_LEN>()?;
1686 let operation = take_tag::<t::AttemptOperation>(decoder)?;
1687 match origin {
1688 t::ClientDiscriminant::CredentialAttachRequest => {
1689 if operation != t::AttemptOperation::CredentialAttachRequest {
1690 return Err(decode_error(t::DecodeClass::InvalidField));
1691 }
1692 Ok(r::ServerValue::AttemptTokenBodyConflict(
1693 r::AttemptTokenBodyConflict::CredentialAttach {
1694 token: AttachAttemptToken::new(token),
1695 conversation_id: decoder.take_u64()?,
1696 presented_participant_id: decoder.take_u64()?,
1697 presented_generation: decoder.take_generation()?,
1698 presented_marker_delivery_seq: decoder.take_option_u64()?,
1699 conflict: take_tag::<t::AttemptConflict>(decoder)?,
1700 },
1701 ))
1702 }
1703 t::ClientDiscriminant::LeaveRequest => {
1704 if operation != t::AttemptOperation::LeaveRequest {
1705 return Err(decode_error(t::DecodeClass::InvalidField));
1706 }
1707 let value = r::AttemptTokenBodyConflict::Leave {
1708 token: LeaveAttemptToken::new(token),
1709 conversation_id: decoder.take_u64()?,
1710 presented_participant_id: decoder.take_u64()?,
1711 presented_generation: decoder.take_generation()?,
1712 };
1713 if take_tag::<t::AttemptConflict>(decoder)? != t::AttemptConflict::Generation {
1714 decoder.invalidate();
1715 }
1716 Ok(r::ServerValue::AttemptTokenBodyConflict(value))
1717 }
1718 t::ClientDiscriminant::RecordAdmission => {
1719 if operation != t::AttemptOperation::RecordAdmission {
1720 return Err(decode_error(t::DecodeClass::InvalidField));
1721 }
1722 Ok(r::ServerValue::AttemptTokenBodyConflict(
1723 r::AttemptTokenBodyConflict::RecordAdmission {
1724 token: RecordAdmissionAttemptToken::new(token),
1725 conversation_id: decoder.take_u64()?,
1726 presented_participant_id: decoder.take_u64()?,
1727 presented_generation: decoder.take_generation()?,
1728 },
1729 ))
1730 }
1731 _ => Err(decode_error(t::DecodeClass::InvalidField)),
1732 }
1733 }
1734 D::ConnectionConversationCapacityExceeded => {
1735 let origin = required_origin(origin)?;
1736 Ok(r::ServerValue::ConnectionConversationCapacityExceeded(
1737 r::ConnectionConversationCapacityExceeded::SemanticRequest {
1738 request: take_response_envelope(origin, decoder)?,
1739 limit: decoder.take_u64()?,
1740 },
1741 ))
1742 }
1743 D::ConnectionConversationBindingOccupied => {
1744 let origin = required_origin(origin)?;
1745 let value = match origin {
1746 t::ClientDiscriminant::EnrollmentRequest => {
1747 let conversation_id = decoder.take_u64()?;
1748 let enrollment_token = EnrollmentToken::new(decoder.take_fixed::<TOKEN_LEN>()?);
1749 if decoder.take_option_u64()?.is_some() {
1750 decoder.invalidate();
1751 }
1752 r::ConnectionConversationBindingOccupied::Enrollment {
1753 conversation_id,
1754 enrollment_token,
1755 }
1756 }
1757 t::ClientDiscriminant::CredentialAttachRequest => {
1758 let conversation_id = decoder.take_u64()?;
1759 let participant_id = decoder.take_u64()?;
1760 let capability_generation = decoder.take_generation()?;
1761 let attach_attempt_token =
1762 AttachAttemptToken::new(decoder.take_fixed::<TOKEN_LEN>()?);
1763 let accept_marker_delivery_seq = decoder.take_option_u64()?;
1764 if decoder.take_option_u64()? != Some(participant_id) {
1765 decoder.invalidate();
1766 }
1767 r::ConnectionConversationBindingOccupied::CredentialAttach {
1768 conversation_id,
1769 participant_id,
1770 capability_generation,
1771 attach_attempt_token,
1772 accept_marker_delivery_seq,
1773 }
1774 }
1775 _ => return Err(decode_error(t::DecodeClass::InvalidField)),
1776 };
1777 Ok(r::ServerValue::ConnectionConversationBindingOccupied(value))
1778 }
1779 D::ConversationOrderExhausted => {
1780 let origin = required_origin(origin)?;
1781 let request = take_order_allocating(origin, decoder)?;
1782 let _counter = take_tag::<t::Counter>(decoder)?;
1783 let high = decoder.take_u64()?;
1784 let next_value = decoder.take_option_u64()?;
1785 let order_remaining = decoder.take_u128()?;
1786 let reserved_claims = decoder.take_u128()?;
1787 if decoder.take_u64()? != r::ConversationOrderExhausted::REQUIRED_MAJORS {
1788 decoder.invalidate();
1789 }
1790 let resulting_order_remaining = decoder.take_u128()?;
1791 let resulting_reserved_claims = decoder.take_u128()?;
1792 let value = r::ConversationOrderExhausted::new(
1793 request,
1794 high,
1795 order_remaining,
1796 reserved_claims,
1797 resulting_order_remaining,
1798 resulting_reserved_claims,
1799 );
1800 if next_value != value.next_value() {
1801 decoder.invalidate();
1802 }
1803 Ok(r::ServerValue::ConversationOrderExhausted(
1804 alloc::boxed::Box::new(value),
1805 ))
1806 }
1807 D::ParticipantUnknown => {
1808 let origin = required_origin(origin)?;
1809 Ok(r::ServerValue::ParticipantUnknown(r::ParticipantUnknown {
1810 request: take_participant_reference(origin, decoder)?,
1811 }))
1812 }
1813 D::NoBinding => {
1814 let origin = required_origin(origin)?;
1815 Ok(r::ServerValue::NoBinding(r::NoBinding {
1816 request: take_binding_required(origin, decoder)?,
1817 }))
1818 }
1819 D::StaleAuthority => take_stale_authority(required_origin(origin)?, decoder)
1820 .map(r::ServerValue::StaleAuthority),
1821 D::Retired => {
1822 let origin = required_origin(origin)?;
1823 let value = if origin == t::ClientDiscriminant::EnrollmentRequest {
1824 r::Retired::Enrollment {
1825 request: take_enrollment(decoder)?,
1826 participant_id: decoder.take_u64()?,
1827 retired_generation: decoder.take_generation()?,
1828 }
1829 } else {
1830 r::Retired::Participant {
1831 request: take_participant_reference(origin, decoder)?,
1832 retired_generation: decoder.take_generation()?,
1833 }
1834 };
1835 Ok(r::ServerValue::Retired(value))
1836 }
1837 D::MarkerClosureCapacityExceeded => {
1838 let origin = required_origin(origin)?;
1839 let request = take_closure_checked(origin, decoder)?;
1840 let scope = take_tag::<t::ClosureScope>(decoder)?;
1841 let snapshot = take_closure_snapshot(decoder)?;
1842 let reason = match scope {
1843 t::ClosureScope::Capacity => {
1844 let dimension =
1845 ResourceDimension::from(take_tag::<t::ResourceDimensionTag>(decoder)?);
1846 c::ClosureRefusalReason::Capacity(c::ClosureCapacityReason {
1847 dimension,
1848 required: decoder.take_u128()?,
1849 limit: decoder.take_u128()?,
1850 })
1851 }
1852 t::ClosureScope::RecoveryFence => c::ClosureRefusalReason::RecoveryFence,
1853 t::ClosureScope::DeliveredMarkerAwaitingAck => {
1854 c::ClosureRefusalReason::DeliveredMarkerAwaitingAck
1855 }
1856 t::ClosureScope::EpisodeChurnLimit => c::ClosureRefusalReason::EpisodeChurnLimit,
1857 };
1858 Ok(r::ServerValue::MarkerClosureCapacityExceeded(
1859 alloc::boxed::Box::new(c::MarkerClosureCapacityExceeded {
1860 request,
1861 snapshot,
1862 reason,
1863 }),
1864 ))
1865 }
1866 D::EnrollBound => take_enroll_bound(decoder).map(r::ServerValue::EnrollBound),
1867 D::EnrollmentKnown => Ok(r::ServerValue::EnrollmentKnown(r::EnrollmentKnown {
1868 conversation_id: decoder.take_u64()?,
1869 token: EnrollmentToken::new(decoder.take_fixed::<TOKEN_LEN>()?),
1870 participant_id: decoder.take_u64()?,
1871 current_generation: decoder.take_generation()?,
1872 })),
1873 D::ReceiptExpired => take_receipt_expired(required_origin(origin)?, decoder)
1874 .map(r::ServerValue::ReceiptExpired),
1875 D::ReceiptCapacityExceeded => take_receipt_capacity(required_origin(origin)?, decoder)
1876 .map(r::ServerValue::ReceiptCapacityExceeded),
1877 D::IdentityCapacityExceeded => {
1878 let request = take_enrollment(decoder)?;
1879 let scope = take_tag::<t::IdentityCapacityScope>(decoder)?;
1880 let limit = decoder.take_u64()?;
1881 let occupied = decoder.take_u64()?;
1882 if decoder.take_u64()? != r::IdentityCapacityExceeded::REQUESTED {
1883 decoder.invalidate();
1884 }
1885 Ok(r::ServerValue::IdentityCapacityExceeded(
1886 r::IdentityCapacityExceeded {
1887 request,
1888 scope,
1889 limit,
1890 occupied,
1891 },
1892 ))
1893 }
1894 D::ObserverBackpressure => take_observer_backpressure(required_origin(origin)?, decoder)
1895 .map(r::ServerValue::ObserverBackpressure),
1896 D::ConversationSequenceExhausted => {
1897 let request = take_sequence_allocating(required_origin(origin)?, decoder)?;
1898 Ok(r::ServerValue::ConversationSequenceExhausted(
1899 alloc::boxed::Box::new(r::ConversationSequenceExhausted {
1900 request,
1901 sequence_budget: take_sequence_budget(decoder)?,
1902 }),
1903 ))
1904 }
1905 D::AttachBound => take_attach_bound(decoder).map(r::ServerValue::AttachBound),
1906 D::StaleOrUnknownReceipt => Ok(r::ServerValue::StaleOrUnknownReceipt(
1907 r::StaleOrUnknownReceipt {
1908 conversation_id: decoder.take_u64()?,
1909 token: AttachAttemptToken::new(decoder.take_fixed::<TOKEN_LEN>()?),
1910 participant_id: decoder.take_u64()?,
1911 presented_generation: decoder.take_generation()?,
1912 presented_marker_delivery_seq: decoder.take_option_u64()?,
1913 current_generation: decoder.take_generation()?,
1914 },
1915 )),
1916 D::MarkerNotDelivered => {
1917 let request = take_marker_proof(required_origin(origin)?, decoder)?;
1918 let reason = take_tag::<t::MarkerNotDeliveredReason>(decoder)?;
1919 Ok(r::ServerValue::MarkerNotDelivered(r::MarkerNotDelivered {
1920 request,
1921 reason,
1922 expected_marker_delivery_seq: decoder.take_u64()?,
1923 }))
1924 }
1925 D::MarkerMismatch => {
1926 let request = take_marker_proof(required_origin(origin)?, decoder)?;
1927 let reason = take_tag::<t::MarkerMismatchReason>(decoder)?;
1928 let mismatch = match reason {
1929 t::MarkerMismatchReason::BelowCursor => r::MarkerMismatchBody::BelowCursor {
1930 current_cursor: decoder.take_u64()?,
1931 },
1932 t::MarkerMismatchReason::NoMarkerExpected => {
1933 r::MarkerMismatchBody::NoMarkerExpected
1934 }
1935 t::MarkerMismatchReason::ExpectedDifferentMarker => {
1936 r::MarkerMismatchBody::ExpectedDifferentMarker {
1937 expected_marker_delivery_seq: decoder.take_u64()?,
1938 }
1939 }
1940 };
1941 Ok(r::ServerValue::MarkerMismatch(r::MarkerMismatch {
1942 request,
1943 mismatch,
1944 }))
1945 }
1946 D::Bound => {
1947 take_receipt_replay(required_origin(origin)?, decoder).map(r::ServerValue::Bound)
1948 }
1949 D::UnboundReceipt => take_receipt_replay(required_origin(origin)?, decoder)
1950 .map(r::ServerValue::UnboundReceipt),
1951 D::DetachCommitted => {
1952 let conversation_id = decoder.take_u64()?;
1953 let participant_id = decoder.take_u64()?;
1954 let capability_generation = decoder.take_generation()?;
1955 let detach_attempt_token = DetachAttemptToken::new(decoder.take_fixed::<TOKEN_LEN>()?);
1956 let committed_binding_epoch = decoder.take_binding_epoch()?;
1957 let detached_delivery_seq = decoder.take_u64()?;
1958 if capability_generation != committed_binding_epoch.capability_generation {
1959 decoder.invalidate();
1960 }
1961 Ok(r::ServerValue::DetachCommitted(r::DetachCommitted::new(
1962 conversation_id,
1963 participant_id,
1964 detach_attempt_token,
1965 committed_binding_epoch,
1966 detached_delivery_seq,
1967 )))
1968 }
1969 D::DetachInProgress => Ok(r::ServerValue::DetachInProgress(r::DetachInProgress {
1970 conversation_id: decoder.take_u64()?,
1971 participant_id: decoder.take_u64()?,
1972 presented_token: DetachAttemptToken::new(decoder.take_fixed::<TOKEN_LEN>()?),
1973 presented_generation: decoder.take_generation()?,
1974 committed_binding_epoch: decoder.take_binding_epoch()?,
1975 })),
1976 D::AckCommitted => {
1977 let request = take_participant_ack(decoder)?;
1978 if decoder.take_u64()? != request.through_seq {
1979 decoder.invalidate();
1980 }
1981 Ok(r::ServerValue::AckCommitted(r::AckCommitted::new(request)))
1982 }
1983 D::AckNoOp => {
1984 take_ack_no_op(required_origin(origin)?, decoder).map(r::ServerValue::AckNoOp)
1985 }
1986 D::AckGap => {
1987 let request = take_participant_ack(decoder)?;
1988 let current_cursor = decoder.take_u64()?;
1989 let _reason = take_tag::<t::AckGapReason>(decoder)?;
1990 r::AckGap::new(request, current_cursor)
1991 .map(r::ServerValue::AckGap)
1992 .ok_or_else(|| decode_error(t::DecodeClass::InvalidField))
1993 }
1994 D::AckRegression => {
1995 let request = take_participant_ack(decoder)?;
1996 let current_cursor = decoder.take_u64()?;
1997 let _reason = take_tag::<t::AckRegressionReason>(decoder)?;
1998 r::AckRegression::new(request, current_cursor)
1999 .map(r::ServerValue::AckRegression)
2000 .ok_or_else(|| decode_error(t::DecodeClass::InvalidField))
2001 }
2002 D::LeaveCommitted => {
2003 let conversation_id = decoder.take_u64()?;
2004 let leave_attempt_token = LeaveAttemptToken::new(decoder.take_fixed::<TOKEN_LEN>()?);
2005 let participant_id = decoder.take_u64()?;
2006 let presented_generation = decoder.take_generation()?;
2007 let retired_generation = decoder.take_generation()?;
2008 let ended_binding_epoch = decoder.take_option_binding_epoch()?;
2009 let prior_terminal_delivery_seq = decoder.take_option_u64()?;
2010 let left_delivery_seq = decoder.take_u64()?;
2011 if presented_generation != retired_generation {
2012 return Err(decode_error(t::DecodeClass::InvalidField));
2013 }
2014 r::LeaveCommitted::new(
2015 conversation_id,
2016 leave_attempt_token,
2017 participant_id,
2018 retired_generation,
2019 ended_binding_epoch,
2020 prior_terminal_delivery_seq,
2021 left_delivery_seq,
2022 )
2023 .map(r::ServerValue::LeaveCommitted)
2024 .ok_or_else(|| decode_error(t::DecodeClass::InvalidField))
2025 }
2026 D::MarkerAckCommitted => {
2027 let request = take_marker_ack(decoder)?;
2028 if decoder.take_u64()? != request.marker_delivery_seq {
2029 decoder.invalidate();
2030 }
2031 Ok(r::ServerValue::MarkerAckCommitted(
2032 r::MarkerAckCommitted::new(request),
2033 ))
2034 }
2035 D::RecordCommitted => {
2036 let request = take_record_admission(decoder)?;
2037 if decoder.take_u64()? != request.participant_id {
2038 decoder.invalidate();
2039 }
2040 let delivery_seq = decoder.take_u64()?;
2041 Ok(r::ServerValue::RecordCommitted(r::RecordCommitted::new(
2042 request,
2043 delivery_seq,
2044 )))
2045 }
2046 D::RecordTooLarge => Ok(r::ServerValue::RecordTooLarge(r::RecordTooLarge {
2047 request: take_record_admission(decoder)?,
2048 dimension: ResourceDimension::from(take_tag::<t::ResourceDimensionTag>(decoder)?),
2049 encoded_record_charge: decoder.take_resource_vector()?,
2050 max_ordinary_record_charge: decoder.take_resource_vector()?,
2051 })),
2052 D::ObserverRecoveryAccepted => {
2053 take_recovery_accepted(decoder).map(r::ServerValue::ObserverRecoveryAccepted)
2054 }
2055 D::InvalidObserverEpoch => {
2056 require_zero_recovery_count(decoder)?;
2057 take_invalid_observer_epoch(decoder).map(r::ServerValue::InvalidObserverEpoch)
2058 }
2059 D::InvalidObserverEpochList => {
2060 require_zero_recovery_count(decoder)?;
2061 take_invalid_observer_epoch_list(decoder).map(r::ServerValue::InvalidObserverEpochList)
2062 }
2063 D::ObserverRecoveryConnectionCapacityExceeded => {
2064 require_zero_recovery_count(decoder)?;
2065 Ok(r::ServerValue::ConnectionConversationCapacityExceeded(
2066 r::ConnectionConversationCapacityExceeded::ObserverRecovery {
2067 conversation_id: decoder.take_u64()?,
2068 limit: decoder.take_u64()?,
2069 },
2070 ))
2071 }
2072 D::MarkerSettlementBackpressure => {
2073 let origin = required_origin(origin)?;
2074 let conversation_id = decoder.take_u64()?;
2075 let refused_epoch = decoder.take_u64()?;
2076 let value = match origin {
2077 t::ClientDiscriminant::CredentialAttachRequest => {
2078 r::MarkerSettlementBackpressure::CredentialAttach {
2079 conversation_id,
2080 refused_epoch,
2081 }
2082 }
2083 t::ClientDiscriminant::DetachRequest => r::MarkerSettlementBackpressure::Detach {
2084 conversation_id,
2085 refused_epoch,
2086 },
2087 _ => return Err(decode_error(t::DecodeClass::InvalidField)),
2088 };
2089 Ok(r::ServerValue::MarkerSettlementBackpressure(value))
2090 }
2091 D::EnrollmentSettlementBackpressure => {
2092 required_origin(origin)?;
2093 Ok(r::ServerValue::EnrollmentSettlementBackpressure(
2094 r::EnrollmentSettlementBackpressure {
2095 conversation_id: decoder.take_u64()?,
2096 },
2097 ))
2098 }
2099 D::RecordAdmissionProtocolFault => Ok(r::ServerValue::RecordAdmissionProtocolFault(
2100 r::RecordAdmissionProtocolFault {
2101 request: take_record_admission(decoder)?,
2102 class: r::RecordAdmissionFaultClass::from(take_tag::<
2103 t::RecordAdmissionFaultClassTag,
2104 >(decoder)?),
2105 },
2106 )),
2107 }
2108}
2109
2110trait WireTag: TryFrom<u16> {}
2111
2112impl WireTag for t::TransportReasonTag {}
2113impl WireTag for t::DecodeClass {}
2114impl WireTag for t::AttemptOperation {}
2115impl WireTag for t::AttemptConflict {}
2116impl WireTag for t::Counter {}
2117impl WireTag for t::ClosureScope {}
2118impl WireTag for t::ResourceDimensionTag {}
2119impl WireTag for t::IdentityCapacityScope {}
2120impl WireTag for t::MarkerNotDeliveredReason {}
2121impl WireTag for t::MarkerMismatchReason {}
2122impl WireTag for t::AckGapReason {}
2123impl WireTag for t::AckRegressionReason {}
2124impl WireTag for t::InvalidObserverEpochReason {}
2125impl WireTag for t::InvalidObserverEpochListReason {}
2126impl WireTag for t::RecordAdmissionFaultClassTag {}
2127
2128fn take_tag<T>(decoder: &mut Decoder<'_>) -> Result<T, CodecError>
2129where
2130 T: WireTag,
2131{
2132 T::try_from(decoder.take_u16()?).map_err(|_| decode_error(t::DecodeClass::InvalidField))
2133}
2134
2135fn take_stale_authority(
2136 origin: t::ClientDiscriminant,
2137 decoder: &mut Decoder<'_>,
2138) -> Result<r::StaleAuthority, CodecError> {
2139 match origin {
2140 t::ClientDiscriminant::CredentialAttachRequest => Ok(r::StaleAuthority::Live {
2141 request: r::CommonStaleAuthorityEnvelope::CredentialAttach(take_attach(decoder)?),
2142 current_generation: decoder.take_generation()?,
2143 }),
2144 t::ClientDiscriminant::ParticipantAck => Ok(r::StaleAuthority::Live {
2145 request: r::CommonStaleAuthorityEnvelope::ParticipantAck(take_participant_ack(
2146 decoder,
2147 )?),
2148 current_generation: decoder.take_generation()?,
2149 }),
2150 t::ClientDiscriminant::MarkerAck => Ok(r::StaleAuthority::Live {
2151 request: r::CommonStaleAuthorityEnvelope::MarkerAck(take_marker_ack(decoder)?),
2152 current_generation: decoder.take_generation()?,
2153 }),
2154 t::ClientDiscriminant::RecordAdmission => Ok(r::StaleAuthority::Live {
2155 request: r::CommonStaleAuthorityEnvelope::RecordAdmission(take_record_admission(
2156 decoder,
2157 )?),
2158 current_generation: decoder.take_generation()?,
2159 }),
2160 t::ClientDiscriminant::DetachRequest => {
2161 let authority = t::DetachAuthorityStateTag::try_from(decoder.take_u16()?)
2162 .map_err(|_| decode_error(t::DecodeClass::InvalidField))?;
2163 let value = match authority {
2164 t::DetachAuthorityStateTag::Live => r::DetachStaleAuthority::Live {
2165 conversation_id: decoder.take_u64()?,
2166 participant_id: decoder.take_u64()?,
2167 capability_generation: decoder.take_generation()?,
2168 detach_attempt_token: DetachAttemptToken::new(
2169 decoder.take_fixed::<TOKEN_LEN>()?,
2170 ),
2171 current_generation: decoder.take_generation()?,
2172 },
2173 t::DetachAuthorityStateTag::TerminalizedDetachCell => {
2174 let conversation_id = decoder.take_u64()?;
2175 let participant_id = decoder.take_u64()?;
2176 let capability_generation = decoder.take_generation()?;
2177 let detach_attempt_token =
2178 DetachAttemptToken::new(decoder.take_fixed::<TOKEN_LEN>()?);
2179 let current_generation = decoder.take_generation()?;
2180 let committed_binding_epoch = decoder.take_binding_epoch()?;
2181 let binding_state = match t::BindingStateTag::try_from(decoder.take_u16()?)
2182 .map_err(|_| decode_error(t::DecodeClass::InvalidField))?
2183 {
2184 t::BindingStateTag::Bound => r::BindingStateView::Bound {
2185 current_binding_epoch: decoder.take_binding_epoch()?,
2186 },
2187 t::BindingStateTag::Detached => r::BindingStateView::Detached,
2188 };
2189 r::DetachStaleAuthority::TerminalizedDetachCell(
2190 r::TerminalizedDetachCell::from_wire_decode(
2191 TERMINALIZED_WIRE_DECODE_AUTHORITY,
2192 conversation_id,
2193 participant_id,
2194 capability_generation,
2195 detach_attempt_token,
2196 current_generation,
2197 committed_binding_epoch,
2198 binding_state,
2199 ),
2200 )
2201 }
2202 };
2203 Ok(r::StaleAuthority::Detach(value))
2204 }
2205 t::ClientDiscriminant::LeaveRequest => {
2206 let authority = t::LeaveAuthorityStateTag::try_from(decoder.take_u16()?)
2207 .map_err(|_| decode_error(t::DecodeClass::InvalidField))?;
2208 let conversation_id = decoder.take_u64()?;
2209 let participant_id = decoder.take_u64()?;
2210 let presented_generation = decoder.take_generation()?;
2211 let leave_attempt_token = LeaveAttemptToken::new(decoder.take_fixed::<TOKEN_LEN>()?);
2212 let value = match authority {
2213 t::LeaveAuthorityStateTag::Live => r::LeaveStaleAuthority::Live {
2214 conversation_id,
2215 participant_id,
2216 presented_generation,
2217 leave_attempt_token,
2218 current_generation: decoder.take_generation()?,
2219 },
2220 t::LeaveAuthorityStateTag::CommittedLeaveTombstone => {
2221 r::LeaveStaleAuthority::CommittedLeaveTombstone {
2222 conversation_id,
2223 participant_id,
2224 presented_generation,
2225 leave_attempt_token,
2226 retired_generation: decoder.take_generation()?,
2227 }
2228 }
2229 };
2230 Ok(r::StaleAuthority::Leave(value))
2231 }
2232 t::ClientDiscriminant::EnrollmentRequest
2233 | t::ClientDiscriminant::ObserverRecoveryHandshake => {
2234 Err(decode_error(t::DecodeClass::InvalidField))
2235 }
2236 }
2237}
2238
2239fn take_enroll_bound(decoder: &mut Decoder<'_>) -> Result<r::EnrollBound, CodecError> {
2240 let conversation_id = decoder.take_u64()?;
2241 let token = EnrollmentToken::new(decoder.take_fixed::<TOKEN_LEN>()?);
2242 let participant_id = decoder.take_u64()?;
2243 if decoder.take_option_generation()?.is_some() {
2244 decoder.invalidate();
2245 }
2246 let capability_generation = decoder.take_generation()?;
2247 if capability_generation.get() != 1 {
2248 decoder.invalidate();
2249 }
2250 let attach_secret = AttachSecret::new(decoder.take_fixed::<SECRET_LEN>()?);
2251 let origin_binding_epoch = decoder.take_binding_epoch()?;
2252 if origin_binding_epoch.capability_generation.get() != 1 {
2253 decoder.invalidate();
2254 }
2255 if decoder.take_u64()? != 0 {
2256 decoder.invalidate();
2257 }
2258 if decoder.take_option_u64()?.is_some() {
2259 decoder.invalidate();
2260 }
2261 let receipt_expires_at = decoder.take_u128()?;
2262 let provenance_expires_at = decoder.take_u128()?;
2263
2264 let valid_epoch = if origin_binding_epoch.capability_generation.get() == 1 {
2265 origin_binding_epoch
2266 } else {
2267 BindingEpoch::new(
2268 origin_binding_epoch.connection_incarnation,
2269 generation_one()?,
2270 )
2271 };
2272 r::EnrollBound::new(
2273 conversation_id,
2274 token,
2275 participant_id,
2276 attach_secret,
2277 valid_epoch,
2278 receipt_expires_at,
2279 provenance_expires_at,
2280 )
2281 .ok_or(CodecError::InvalidValue)
2282}
2283
2284fn take_receipt_expired(
2285 origin: t::ClientDiscriminant,
2286 decoder: &mut Decoder<'_>,
2287) -> Result<r::ReceiptExpired, CodecError> {
2288 let conversation_id = decoder.take_u64()?;
2289 let token = decoder.take_fixed::<TOKEN_LEN>()?;
2290 let participant_id = decoder.take_u64()?;
2291 match origin {
2292 t::ClientDiscriminant::EnrollmentRequest => {
2293 if decoder.take_option_generation()?.is_some() {
2294 decoder.invalidate();
2295 }
2296 Ok(r::ReceiptExpired::Enrollment {
2297 conversation_id,
2298 token: EnrollmentToken::new(token),
2299 participant_id,
2300 result_generation: decoder.take_generation()?,
2301 current_generation: decoder.take_generation()?,
2302 reason: take_tag::<t::ReceiptExpiryReason>(decoder)?,
2303 })
2304 }
2305 t::ClientDiscriminant::CredentialAttachRequest => {
2306 let presented_generation = required_generation_option(decoder)?;
2307 Ok(r::ReceiptExpired::CredentialAttach {
2308 conversation_id,
2309 token: AttachAttemptToken::new(token),
2310 participant_id,
2311 presented_generation,
2312 presented_marker_delivery_seq: decoder.take_option_u64()?,
2313 result_generation: decoder.take_generation()?,
2314 current_generation: decoder.take_generation()?,
2315 reason: take_tag::<t::ReceiptExpiryReason>(decoder)?,
2316 })
2317 }
2318 _ => Err(decode_error(t::DecodeClass::InvalidField)),
2319 }
2320}
2321
2322impl WireTag for t::ReceiptExpiryReason {}
2323impl WireTag for t::ReceiptCapacityScope {}
2324
2325fn take_receipt_capacity(
2326 origin: t::ClientDiscriminant,
2327 decoder: &mut Decoder<'_>,
2328) -> Result<r::ReceiptCapacityExceeded, CodecError> {
2329 match origin {
2330 t::ClientDiscriminant::EnrollmentRequest => {
2331 let request = take_enrollment(decoder)?;
2332 let scope = match take_tag::<t::ReceiptCapacityScope>(decoder)? {
2333 t::ReceiptCapacityScope::LiveReceiptServer => {
2334 r::EnrollmentReceiptCapacityScope::LiveReceiptServer
2335 }
2336 t::ReceiptCapacityScope::ProvenanceServer => {
2337 r::EnrollmentReceiptCapacityScope::ProvenanceServer
2338 }
2339 t::ReceiptCapacityScope::ProvenanceConversation => {
2340 r::EnrollmentReceiptCapacityScope::ProvenanceConversation
2341 }
2342 t::ReceiptCapacityScope::LiveReceiptParticipant
2343 | t::ReceiptCapacityScope::ProvenanceParticipant => {
2344 return Err(decode_error(t::DecodeClass::InvalidField));
2345 }
2346 };
2347 let limit = decoder.take_u64()?;
2348 let occupied = decoder.take_u64()?;
2349 require_one(decoder, r::ReceiptCapacityExceeded::REQUESTED)?;
2350 Ok(r::ReceiptCapacityExceeded::Enrollment {
2351 request,
2352 scope,
2353 limit,
2354 occupied,
2355 })
2356 }
2357 t::ClientDiscriminant::CredentialAttachRequest => {
2358 let request = take_attach(decoder)?;
2359 let scope = take_tag::<t::ReceiptCapacityScope>(decoder)?;
2360 let limit = decoder.take_u64()?;
2361 let occupied = decoder.take_u64()?;
2362 require_one(decoder, r::ReceiptCapacityExceeded::REQUESTED)?;
2363 Ok(r::ReceiptCapacityExceeded::CredentialAttach {
2364 request,
2365 scope,
2366 limit,
2367 occupied,
2368 })
2369 }
2370 _ => Err(decode_error(t::DecodeClass::InvalidField)),
2371 }
2372}
2373
2374fn require_one(decoder: &mut Decoder<'_>, expected: u64) -> Result<(), CodecError> {
2375 if decoder.take_u64()? != expected {
2376 decoder.invalidate();
2377 }
2378 Ok(())
2379}
2380
2381fn take_observer_backpressure(
2382 origin: t::ClientDiscriminant,
2383 decoder: &mut Decoder<'_>,
2384) -> Result<r::ObserverBackpressure, CodecError> {
2385 match origin {
2386 t::ClientDiscriminant::EnrollmentRequest => Ok(r::ObserverBackpressure::Enrollment {
2387 request: take_enrollment(decoder)?,
2388 state: take_backpressure_state(decoder)?,
2389 }),
2390 t::ClientDiscriminant::CredentialAttachRequest => {
2391 Ok(r::ObserverBackpressure::CredentialAttach {
2392 request: take_attach(decoder)?,
2393 state: take_backpressure_state(decoder)?,
2394 })
2395 }
2396 t::ClientDiscriminant::DetachRequest => Ok(r::ObserverBackpressure::Detach {
2397 request: take_detach(decoder)?,
2398 committed_binding_epoch: decoder.take_binding_epoch()?,
2399 state: take_backpressure_state(decoder)?,
2400 }),
2401 t::ClientDiscriminant::LeaveRequest => Ok(r::ObserverBackpressure::Leave {
2402 request: take_leave(decoder)?,
2403 state: take_backpressure_state(decoder)?,
2404 prior_terminal_cell_exists: decoder.take_bool()?,
2405 }),
2406 t::ClientDiscriminant::RecordAdmission => Ok(r::ObserverBackpressure::RecordAdmission {
2407 request: take_record_admission(decoder)?,
2408 state: take_backpressure_state(decoder)?,
2409 }),
2410 _ => Err(decode_error(t::DecodeClass::InvalidField)),
2411 }
2412}
2413
2414fn take_backpressure_state(
2415 decoder: &mut Decoder<'_>,
2416) -> Result<r::ObserverBackpressureState, CodecError> {
2417 let backpressure_epoch = decoder.take_u64()?;
2418 let observer_progress = decoder.take_u64()?;
2419 r::ObserverBackpressureState::replay(backpressure_epoch, observer_progress)
2420 .ok_or_else(|| decode_error(t::DecodeClass::InvalidField))
2421}
2422
2423fn required_generation_option(decoder: &mut Decoder<'_>) -> Result<Generation, CodecError> {
2424 decoder.take_option_generation()?.map_or_else(
2425 || {
2426 decoder.invalidate();
2427 generation_one()
2428 },
2429 Ok,
2430 )
2431}
2432
2433fn take_attach_bound(decoder: &mut Decoder<'_>) -> Result<r::AttachBound, CodecError> {
2434 let conversation_id = decoder.take_u64()?;
2435 let token = AttachAttemptToken::new(decoder.take_fixed::<TOKEN_LEN>()?);
2436 let participant_id = decoder.take_u64()?;
2437 let request_generation = required_generation_option(decoder)?;
2438 let capability_generation = decoder.take_generation()?;
2439 let attach_secret = AttachSecret::new(decoder.take_fixed::<SECRET_LEN>()?);
2440 let origin_binding_epoch = decoder.take_binding_epoch()?;
2441 let persisted_cursor = decoder.take_u64()?;
2442 let accepted_marker_delivery_seq = decoder.take_option_u64()?;
2443 let receipt_expires_at = decoder.take_u128()?;
2444 let provenance_expires_at = decoder.take_u128()?;
2445
2446 let value = match accepted_marker_delivery_seq {
2447 Some(marker) if persisted_cursor == marker => r::AttachBound::fenced(
2448 conversation_id,
2449 token,
2450 participant_id,
2451 request_generation,
2452 attach_secret,
2453 origin_binding_epoch,
2454 marker,
2455 receipt_expires_at,
2456 provenance_expires_at,
2457 ),
2458 Some(_) => None,
2459 None => r::AttachBound::ordinary(
2460 conversation_id,
2461 token,
2462 participant_id,
2463 request_generation,
2464 attach_secret,
2465 origin_binding_epoch,
2466 persisted_cursor,
2467 receipt_expires_at,
2468 provenance_expires_at,
2469 ),
2470 }
2471 .ok_or_else(|| decode_error(t::DecodeClass::InvalidField))?;
2472 if value.capability_generation() == capability_generation {
2473 Ok(value)
2474 } else {
2475 Err(decode_error(t::DecodeClass::InvalidField))
2476 }
2477}
2478
2479fn take_receipt_replay(
2480 origin: t::ClientDiscriminant,
2481 decoder: &mut Decoder<'_>,
2482) -> Result<r::ReceiptReplay, CodecError> {
2483 let conversation_id = decoder.take_u64()?;
2484 let token = decoder.take_fixed::<TOKEN_LEN>()?;
2485 let participant_id = decoder.take_u64()?;
2486 match origin {
2487 t::ClientDiscriminant::EnrollmentRequest => {
2488 let request_generation = decoder.take_option_generation()?;
2489 let capability_generation = decoder.take_generation()?;
2490 let attach_secret = AttachSecret::new(decoder.take_fixed::<SECRET_LEN>()?);
2491 let origin_binding_epoch = decoder.take_binding_epoch()?;
2492 let persisted_cursor = decoder.take_u64()?;
2493 let accepted_marker_delivery_seq = decoder.take_option_u64()?;
2494 let receipt_expires_at = decoder.take_u128()?;
2495 let provenance_expires_at = decoder.take_u128()?;
2496 let value = r::EnrollBound::new(
2497 conversation_id,
2498 EnrollmentToken::new(token),
2499 participant_id,
2500 attach_secret,
2501 origin_binding_epoch,
2502 receipt_expires_at,
2503 provenance_expires_at,
2504 )
2505 .ok_or_else(|| decode_error(t::DecodeClass::InvalidField))?;
2506 if request_generation.is_none()
2507 && capability_generation == Generation::ONE
2508 && persisted_cursor == 0
2509 && accepted_marker_delivery_seq.is_none()
2510 {
2511 Ok(r::ReceiptReplay::Enrollment(value))
2512 } else {
2513 Err(decode_error(t::DecodeClass::InvalidField))
2514 }
2515 }
2516 t::ClientDiscriminant::CredentialAttachRequest => {
2517 let request_generation = required_generation_option(decoder)?;
2518 let capability_generation = decoder.take_generation()?;
2519 let attach_secret = AttachSecret::new(decoder.take_fixed::<SECRET_LEN>()?);
2520 let origin_binding_epoch = decoder.take_binding_epoch()?;
2521 let persisted_cursor = decoder.take_u64()?;
2522 let accepted_marker_delivery_seq = decoder.take_option_u64()?;
2523 let receipt_expires_at = decoder.take_u128()?;
2524 let provenance_expires_at = decoder.take_u128()?;
2525 let attach_token = AttachAttemptToken::new(token);
2526 let value = match accepted_marker_delivery_seq {
2527 Some(marker) if persisted_cursor == marker => r::AttachBound::fenced(
2528 conversation_id,
2529 attach_token,
2530 participant_id,
2531 request_generation,
2532 attach_secret,
2533 origin_binding_epoch,
2534 marker,
2535 receipt_expires_at,
2536 provenance_expires_at,
2537 ),
2538 Some(_) => None,
2539 None => r::AttachBound::ordinary(
2540 conversation_id,
2541 attach_token,
2542 participant_id,
2543 request_generation,
2544 attach_secret,
2545 origin_binding_epoch,
2546 persisted_cursor,
2547 receipt_expires_at,
2548 provenance_expires_at,
2549 ),
2550 }
2551 .ok_or_else(|| decode_error(t::DecodeClass::InvalidField))?;
2552 if value.capability_generation() == capability_generation {
2553 Ok(r::ReceiptReplay::CredentialAttach(value))
2554 } else {
2555 Err(decode_error(t::DecodeClass::InvalidField))
2556 }
2557 }
2558 _ => Err(decode_error(t::DecodeClass::InvalidField)),
2559 }
2560}
2561
2562fn take_ack_no_op(
2563 origin: t::ClientDiscriminant,
2564 decoder: &mut Decoder<'_>,
2565) -> Result<r::AckNoOp, CodecError> {
2566 match origin {
2567 t::ClientDiscriminant::ParticipantAck => {
2568 let request = take_participant_ack(decoder)?;
2569 if decoder.take_u64()? != request.through_seq {
2570 decoder.invalidate();
2571 }
2572 Ok(r::AckNoOp::participant_ack(request))
2573 }
2574 t::ClientDiscriminant::MarkerAck => {
2575 let request = take_marker_ack(decoder)?;
2576 if decoder.take_u64()? != request.marker_delivery_seq {
2577 decoder.invalidate();
2578 }
2579 Ok(r::AckNoOp::marker_ack(request))
2580 }
2581 _ => Err(decode_error(t::DecodeClass::InvalidField)),
2582 }
2583}
2584
2585fn take_recovery_accepted(
2586 decoder: &mut Decoder<'_>,
2587) -> Result<r::ObserverRecoveryAccepted, CodecError> {
2588 let count = decoder.take_u64()?;
2589 let required = u128::from(count) * RECOVERY_STATUS_LEN;
2590 let remaining = u64::try_from(decoder.remaining())
2591 .map_err(|_| decode_error(t::DecodeClass::MissingRequiredField))?;
2592 if required > u128::from(remaining) {
2593 return Err(decode_error(t::DecodeClass::MissingRequiredField));
2594 }
2595 let count_usize: usize = count
2596 .try_into()
2597 .map_err(|_| decode_error(t::DecodeClass::MissingRequiredField))?;
2598
2599 let mut statuses = Vec::new();
2600 for _ in 0..count_usize {
2601 statuses.push(r::ObserverProgressStatus {
2602 conversation_id: decoder.take_u64()?,
2603 refused_epoch: decoder.take_u64()?,
2604 current_observer_progress: decoder.take_u64()?,
2605 armed: decoder.take_bool()?,
2606 progressed: decoder.take_bool()?,
2607 });
2608 }
2609 Ok(r::ObserverRecoveryAccepted { statuses })
2610}
2611
2612fn require_zero_recovery_count(decoder: &mut Decoder<'_>) -> Result<(), CodecError> {
2613 if decoder.take_u64()? == 0 {
2614 Ok(())
2615 } else {
2616 Err(decode_error(t::DecodeClass::InvalidField))
2617 }
2618}
2619
2620fn take_invalid_observer_epoch(
2621 decoder: &mut Decoder<'_>,
2622) -> Result<r::InvalidObserverEpoch, CodecError> {
2623 let reason = take_tag::<t::InvalidObserverEpochReason>(decoder)?;
2624 let conversation_id = decoder.take_u64()?;
2625 let presented_epoch = decoder.take_u64()?;
2626 let current = decoder.take_option_u64()?;
2627 match reason {
2628 t::InvalidObserverEpochReason::ConversationUnknown => {
2629 if current.is_some() {
2630 decoder.invalidate();
2631 }
2632 Ok(r::InvalidObserverEpoch::ConversationUnknown {
2633 conversation_id,
2634 presented_epoch,
2635 })
2636 }
2637 t::InvalidObserverEpochReason::EpochAhead => {
2638 let current_observer_progress = current.unwrap_or_else(|| {
2639 decoder.invalidate();
2640 0
2641 });
2642 Ok(r::InvalidObserverEpoch::EpochAhead {
2643 conversation_id,
2644 presented_epoch,
2645 current_observer_progress,
2646 })
2647 }
2648 }
2649}
2650
2651fn take_invalid_observer_epoch_list(
2652 decoder: &mut Decoder<'_>,
2653) -> Result<r::InvalidObserverEpochList, CodecError> {
2654 match take_tag::<t::InvalidObserverEpochListReason>(decoder)? {
2655 t::InvalidObserverEpochListReason::TooManyEntries => {
2656 Ok(r::InvalidObserverEpochList::TooManyEntries {
2657 presented_entries: decoder.take_u64()?,
2658 max_entries: decoder.take_u64()?,
2659 })
2660 }
2661 t::InvalidObserverEpochListReason::DuplicateConversation => {
2662 Ok(r::InvalidObserverEpochList::DuplicateConversation {
2663 conversation_id: decoder.take_u64()?,
2664 first_index: decoder.take_u64()?,
2665 duplicate_index: decoder.take_u64()?,
2666 })
2667 }
2668 }
2669}