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 }
979 Ok(())
980}
981
982fn put_stale_authority(value: &r::StaleAuthority, encoder: &mut Encoder) {
983 match value {
984 r::StaleAuthority::Live {
985 request,
986 current_generation,
987 } => {
988 match request {
989 r::CommonStaleAuthorityEnvelope::CredentialAttach(value) => {
990 put_attach(value, encoder);
991 }
992 r::CommonStaleAuthorityEnvelope::ParticipantAck(value) => {
993 put_participant_ack(value, encoder);
994 }
995 r::CommonStaleAuthorityEnvelope::MarkerAck(value) => {
996 put_marker_ack(value, encoder);
997 }
998 r::CommonStaleAuthorityEnvelope::RecordAdmission(value) => {
999 put_record_admission(value, encoder);
1000 }
1001 }
1002 encoder.put_generation(*current_generation);
1003 }
1004 r::StaleAuthority::Detach(value) => {
1005 encoder.put_u16(value.authority_state_tag().wire_value());
1006 match value {
1007 r::DetachStaleAuthority::Live {
1008 conversation_id,
1009 participant_id,
1010 capability_generation,
1011 detach_attempt_token,
1012 current_generation,
1013 } => {
1014 encoder.put_u64(*conversation_id);
1015 encoder.put_u64(*participant_id);
1016 encoder.put_generation(*capability_generation);
1017 encoder.put_fixed(detach_attempt_token.as_bytes());
1018 encoder.put_generation(*current_generation);
1019 }
1020 r::DetachStaleAuthority::TerminalizedDetachCell(value) => {
1021 encoder.put_u64(value.conversation_id());
1022 encoder.put_u64(value.participant_id());
1023 encoder.put_generation(value.capability_generation());
1024 encoder.put_fixed(value.detach_attempt_token().as_bytes());
1025 encoder.put_generation(value.current_generation());
1026 encoder.put_binding_epoch(value.committed_binding_epoch());
1027 encoder.put_u16(value.binding_state().tag().wire_value());
1028 if let r::BindingStateView::Bound {
1029 current_binding_epoch,
1030 } = value.binding_state()
1031 {
1032 encoder.put_binding_epoch(current_binding_epoch);
1033 }
1034 }
1035 }
1036 }
1037 r::StaleAuthority::Leave(value) => {
1038 encoder.put_u16(value.authority_state_tag().wire_value());
1039 match value {
1040 r::LeaveStaleAuthority::Live {
1041 conversation_id,
1042 participant_id,
1043 presented_generation,
1044 leave_attempt_token,
1045 current_generation,
1046 } => {
1047 encoder.put_u64(*conversation_id);
1048 encoder.put_u64(*participant_id);
1049 encoder.put_generation(*presented_generation);
1050 encoder.put_fixed(leave_attempt_token.as_bytes());
1051 encoder.put_generation(*current_generation);
1052 }
1053 r::LeaveStaleAuthority::CommittedLeaveTombstone {
1054 conversation_id,
1055 participant_id,
1056 presented_generation,
1057 leave_attempt_token,
1058 retired_generation,
1059 } => {
1060 encoder.put_u64(*conversation_id);
1061 encoder.put_u64(*participant_id);
1062 encoder.put_generation(*presented_generation);
1063 encoder.put_fixed(leave_attempt_token.as_bytes());
1064 encoder.put_generation(*retired_generation);
1065 }
1066 }
1067 }
1068 }
1069}
1070
1071fn put_receipt_expired(value: &r::ReceiptExpired, encoder: &mut Encoder) {
1072 match value {
1073 r::ReceiptExpired::Enrollment {
1074 conversation_id,
1075 token,
1076 participant_id,
1077 result_generation,
1078 current_generation,
1079 reason,
1080 } => {
1081 encoder.put_u64(*conversation_id);
1082 encoder.put_fixed(token.as_bytes());
1083 encoder.put_u64(*participant_id);
1084 encoder.put_option_generation(None);
1085 encoder.put_generation(*result_generation);
1086 encoder.put_generation(*current_generation);
1087 encoder.put_u16(reason.wire_value());
1088 }
1089 r::ReceiptExpired::CredentialAttach {
1090 conversation_id,
1091 token,
1092 participant_id,
1093 presented_generation,
1094 presented_marker_delivery_seq,
1095 result_generation,
1096 current_generation,
1097 reason,
1098 } => {
1099 encoder.put_u64(*conversation_id);
1100 encoder.put_fixed(token.as_bytes());
1101 encoder.put_u64(*participant_id);
1102 encoder.put_option_generation(Some(*presented_generation));
1103 encoder.put_option_u64(*presented_marker_delivery_seq);
1104 encoder.put_generation(*result_generation);
1105 encoder.put_generation(*current_generation);
1106 encoder.put_u16(reason.wire_value());
1107 }
1108 }
1109}
1110
1111fn put_receipt_capacity(value: &r::ReceiptCapacityExceeded, encoder: &mut Encoder) {
1112 match value {
1113 r::ReceiptCapacityExceeded::Enrollment {
1114 request,
1115 scope,
1116 limit,
1117 occupied,
1118 } => {
1119 put_enrollment(request, encoder);
1120 encoder.put_u16(scope.wire_scope().wire_value());
1121 encoder.put_u64(*limit);
1122 encoder.put_u64(*occupied);
1123 encoder.put_u64(r::ReceiptCapacityExceeded::REQUESTED);
1124 }
1125 r::ReceiptCapacityExceeded::CredentialAttach {
1126 request,
1127 scope,
1128 limit,
1129 occupied,
1130 } => {
1131 put_attach(request, encoder);
1132 encoder.put_u16(scope.wire_value());
1133 encoder.put_u64(*limit);
1134 encoder.put_u64(*occupied);
1135 encoder.put_u64(r::ReceiptCapacityExceeded::REQUESTED);
1136 }
1137 }
1138}
1139
1140fn put_observer_backpressure(value: &r::ObserverBackpressure, encoder: &mut Encoder) {
1141 match value {
1142 r::ObserverBackpressure::Enrollment { request, state } => {
1143 put_enrollment(request, encoder);
1144 put_backpressure_state(*state, encoder);
1145 }
1146 r::ObserverBackpressure::CredentialAttach { request, state } => {
1147 put_attach(request, encoder);
1148 put_backpressure_state(*state, encoder);
1149 }
1150 r::ObserverBackpressure::Detach {
1151 request,
1152 committed_binding_epoch,
1153 state,
1154 } => {
1155 put_detach(request, encoder);
1156 encoder.put_binding_epoch(*committed_binding_epoch);
1157 put_backpressure_state(*state, encoder);
1158 }
1159 r::ObserverBackpressure::Leave {
1160 request,
1161 state,
1162 prior_terminal_cell_exists,
1163 } => {
1164 put_leave(request, encoder);
1165 put_backpressure_state(*state, encoder);
1166 encoder.put_bool(*prior_terminal_cell_exists);
1167 }
1168 r::ObserverBackpressure::RecordAdmission { request, state } => {
1169 put_record_admission(request, encoder);
1170 put_backpressure_state(*state, encoder);
1171 }
1172 }
1173}
1174
1175fn put_backpressure_state(value: r::ObserverBackpressureState, encoder: &mut Encoder) {
1176 encoder.put_u64(value.backpressure_epoch());
1177 encoder.put_u64(value.observer_progress());
1178}
1179
1180fn put_receipt_replay(value: &r::ReceiptReplay, encoder: &mut Encoder) {
1181 match value {
1182 r::ReceiptReplay::Enrollment(value) => {
1183 encoder.put_u64(value.conversation_id());
1184 encoder.put_fixed(value.token().as_bytes());
1185 encoder.put_u64(value.participant_id());
1186 encoder.put_option_generation(None);
1187 encoder.put_generation(value.capability_generation());
1188 encoder.put_fixed(value.attach_secret().as_bytes());
1189 encoder.put_binding_epoch(value.origin_binding_epoch());
1190 encoder.put_u64(value.persisted_cursor());
1191 encoder.put_option_u64(None);
1192 encoder.put_u128(value.receipt_expires_at());
1193 encoder.put_u128(value.provenance_expires_at());
1194 }
1195 r::ReceiptReplay::CredentialAttach(value) => {
1196 encoder.put_u64(value.conversation_id());
1197 encoder.put_fixed(value.token().as_bytes());
1198 encoder.put_u64(value.participant_id());
1199 encoder.put_option_generation(Some(value.request_generation()));
1200 encoder.put_generation(value.capability_generation());
1201 encoder.put_fixed(value.attach_secret().as_bytes());
1202 encoder.put_binding_epoch(value.origin_binding_epoch());
1203 encoder.put_u64(value.persisted_cursor());
1204 encoder.put_option_u64(value.accepted_marker_delivery_seq());
1205 encoder.put_u128(value.receipt_expires_at());
1206 encoder.put_u128(value.provenance_expires_at());
1207 }
1208 }
1209}
1210
1211const fn carries_origin(discriminant: t::ServerDiscriminant) -> bool {
1219 use t::ServerDiscriminant as D;
1220
1221 !matches!(
1222 discriminant,
1223 D::ParticipantTransportRejected
1224 | D::ObserverRecoveryAccepted
1225 | D::InvalidObserverEpoch
1226 | D::InvalidObserverEpochList
1227 | D::ObserverRecoveryConnectionCapacityExceeded
1228 )
1229}
1230
1231const fn origin_is_valid(
1232 discriminant: t::ServerDiscriminant,
1233 origin: t::ClientDiscriminant,
1234) -> bool {
1235 use t::ClientDiscriminant as O;
1236 use t::ServerDiscriminant as D;
1237
1238 match discriminant {
1239 D::AttemptTokenBodyConflict => matches!(
1240 origin,
1241 O::CredentialAttachRequest | O::LeaveRequest | O::RecordAdmission
1242 ),
1243 D::ConnectionConversationCapacityExceeded => {
1244 !matches!(origin, O::ObserverRecoveryHandshake)
1245 }
1246 D::ConnectionConversationBindingOccupied => {
1247 matches!(origin, O::EnrollmentRequest | O::CredentialAttachRequest)
1248 }
1249 D::ConversationOrderExhausted => matches!(
1250 origin,
1251 O::EnrollmentRequest | O::CredentialAttachRequest | O::RecordAdmission
1252 ),
1253 D::ParticipantUnknown => {
1254 !matches!(origin, O::EnrollmentRequest | O::ObserverRecoveryHandshake)
1255 }
1256 D::NoBinding => !matches!(
1257 origin,
1258 O::EnrollmentRequest | O::CredentialAttachRequest | O::ObserverRecoveryHandshake
1259 ),
1260 D::StaleAuthority => !matches!(origin, O::EnrollmentRequest | O::ObserverRecoveryHandshake),
1261 D::Retired => !matches!(origin, O::ObserverRecoveryHandshake),
1262 D::MarkerClosureCapacityExceeded => matches!(
1263 origin,
1264 O::EnrollmentRequest
1265 | O::CredentialAttachRequest
1266 | O::LeaveRequest
1267 | O::RecordAdmission
1268 ),
1269 D::EnrollBound | D::EnrollmentKnown | D::IdentityCapacityExceeded => {
1270 matches!(origin, O::EnrollmentRequest)
1271 }
1272 D::ReceiptExpired | D::ReceiptCapacityExceeded | D::Bound | D::UnboundReceipt => {
1273 matches!(origin, O::EnrollmentRequest | O::CredentialAttachRequest)
1274 }
1275 D::ObserverBackpressure => matches!(
1276 origin,
1277 O::EnrollmentRequest
1278 | O::CredentialAttachRequest
1279 | O::DetachRequest
1280 | O::LeaveRequest
1281 | O::RecordAdmission
1282 ),
1283 D::ConversationSequenceExhausted => matches!(
1284 origin,
1285 O::EnrollmentRequest | O::CredentialAttachRequest | O::RecordAdmission
1286 ),
1287 D::AttachBound | D::StaleOrUnknownReceipt => {
1288 matches!(origin, O::CredentialAttachRequest)
1289 }
1290 D::MarkerNotDelivered | D::MarkerMismatch => {
1291 matches!(origin, O::CredentialAttachRequest | O::MarkerAck)
1292 }
1293 D::DetachCommitted | D::DetachInProgress => matches!(origin, O::DetachRequest),
1294 D::AckCommitted | D::AckGap | D::AckRegression => matches!(origin, O::ParticipantAck),
1295 D::AckNoOp => matches!(origin, O::ParticipantAck | O::MarkerAck),
1296 D::LeaveCommitted => matches!(origin, O::LeaveRequest),
1297 D::MarkerAckCommitted => matches!(origin, O::MarkerAck),
1298 D::RecordCommitted | D::RecordTooLarge => matches!(origin, O::RecordAdmission),
1299 D::MarkerSettlementBackpressure => {
1300 matches!(origin, O::CredentialAttachRequest | O::DetachRequest)
1301 }
1302 D::EnrollmentSettlementBackpressure => matches!(origin, O::EnrollmentRequest),
1303 D::ParticipantTransportRejected
1304 | D::ObserverRecoveryAccepted
1305 | D::InvalidObserverEpoch
1306 | D::InvalidObserverEpochList
1307 | D::ObserverRecoveryConnectionCapacityExceeded => false,
1308 }
1309}
1310
1311fn take_enrollment(decoder: &mut Decoder<'_>) -> Result<e::EnrollmentEnvelope, CodecError> {
1315 Ok(e::EnrollmentEnvelope {
1316 conversation_id: decoder.take_u64()?,
1317 enrollment_token: EnrollmentToken::new(decoder.take_fixed::<TOKEN_LEN>()?),
1318 })
1319}
1320
1321fn take_attach(decoder: &mut Decoder<'_>) -> Result<e::AttachEnvelope, CodecError> {
1322 Ok(e::AttachEnvelope {
1323 conversation_id: decoder.take_u64()?,
1324 participant_id: decoder.take_u64()?,
1325 capability_generation: decoder.take_generation()?,
1326 attach_attempt_token: AttachAttemptToken::new(decoder.take_fixed::<TOKEN_LEN>()?),
1327 accept_marker_delivery_seq: decoder.take_option_u64()?,
1328 })
1329}
1330
1331fn take_detach(decoder: &mut Decoder<'_>) -> Result<e::DetachEnvelope, CodecError> {
1332 Ok(e::DetachEnvelope {
1333 conversation_id: decoder.take_u64()?,
1334 participant_id: decoder.take_u64()?,
1335 capability_generation: decoder.take_generation()?,
1336 detach_attempt_token: DetachAttemptToken::new(decoder.take_fixed::<TOKEN_LEN>()?),
1337 })
1338}
1339
1340fn take_participant_ack(
1341 decoder: &mut Decoder<'_>,
1342) -> Result<e::ParticipantAckEnvelope, CodecError> {
1343 Ok(e::ParticipantAckEnvelope {
1344 conversation_id: decoder.take_u64()?,
1345 participant_id: decoder.take_u64()?,
1346 capability_generation: decoder.take_generation()?,
1347 through_seq: decoder.take_u64()?,
1348 })
1349}
1350
1351fn take_leave(decoder: &mut Decoder<'_>) -> Result<e::LeaveEnvelope, CodecError> {
1352 Ok(e::LeaveEnvelope {
1353 conversation_id: decoder.take_u64()?,
1354 participant_id: decoder.take_u64()?,
1355 capability_generation: decoder.take_generation()?,
1356 leave_attempt_token: LeaveAttemptToken::new(decoder.take_fixed::<TOKEN_LEN>()?),
1357 })
1358}
1359
1360fn take_marker_ack(decoder: &mut Decoder<'_>) -> Result<e::MarkerAckEnvelope, CodecError> {
1361 Ok(e::MarkerAckEnvelope {
1362 conversation_id: decoder.take_u64()?,
1363 participant_id: decoder.take_u64()?,
1364 capability_generation: decoder.take_generation()?,
1365 marker_delivery_seq: decoder.take_u64()?,
1366 })
1367}
1368
1369fn take_record_admission(
1370 decoder: &mut Decoder<'_>,
1371) -> Result<e::RecordAdmissionEnvelope, CodecError> {
1372 Ok(e::RecordAdmissionEnvelope {
1373 conversation_id: decoder.take_u64()?,
1374 participant_id: decoder.take_u64()?,
1375 capability_generation: decoder.take_generation()?,
1376 record_admission_attempt_token: RecordAdmissionAttemptToken::new(
1377 decoder.take_fixed::<TOKEN_LEN>()?,
1378 ),
1379 })
1380}
1381
1382fn take_response_envelope(
1383 origin: t::ClientDiscriminant,
1384 decoder: &mut Decoder<'_>,
1385) -> Result<e::ResponseEnvelope, CodecError> {
1386 match origin {
1387 t::ClientDiscriminant::EnrollmentRequest => {
1388 take_enrollment(decoder).map(e::ResponseEnvelope::Enrollment)
1389 }
1390 t::ClientDiscriminant::CredentialAttachRequest => {
1391 take_attach(decoder).map(e::ResponseEnvelope::CredentialAttach)
1392 }
1393 t::ClientDiscriminant::DetachRequest => {
1394 take_detach(decoder).map(e::ResponseEnvelope::Detach)
1395 }
1396 t::ClientDiscriminant::ParticipantAck => {
1397 take_participant_ack(decoder).map(e::ResponseEnvelope::ParticipantAck)
1398 }
1399 t::ClientDiscriminant::LeaveRequest => take_leave(decoder).map(e::ResponseEnvelope::Leave),
1400 t::ClientDiscriminant::MarkerAck => {
1401 take_marker_ack(decoder).map(e::ResponseEnvelope::MarkerAck)
1402 }
1403 t::ClientDiscriminant::RecordAdmission => {
1404 take_record_admission(decoder).map(e::ResponseEnvelope::RecordAdmission)
1405 }
1406 t::ClientDiscriminant::ObserverRecoveryHandshake => {
1407 Err(decode_error(t::DecodeClass::InvalidField))
1408 }
1409 }
1410}
1411
1412fn take_participant_reference(
1413 origin: t::ClientDiscriminant,
1414 decoder: &mut Decoder<'_>,
1415) -> Result<r::ParticipantReferenceEnvelope, CodecError> {
1416 match origin {
1417 t::ClientDiscriminant::CredentialAttachRequest => {
1418 take_attach(decoder).map(r::ParticipantReferenceEnvelope::CredentialAttach)
1419 }
1420 t::ClientDiscriminant::DetachRequest => {
1421 take_detach(decoder).map(r::ParticipantReferenceEnvelope::Detach)
1422 }
1423 t::ClientDiscriminant::ParticipantAck => {
1424 take_participant_ack(decoder).map(r::ParticipantReferenceEnvelope::ParticipantAck)
1425 }
1426 t::ClientDiscriminant::LeaveRequest => {
1427 take_leave(decoder).map(r::ParticipantReferenceEnvelope::Leave)
1428 }
1429 t::ClientDiscriminant::MarkerAck => {
1430 take_marker_ack(decoder).map(r::ParticipantReferenceEnvelope::MarkerAck)
1431 }
1432 t::ClientDiscriminant::RecordAdmission => {
1433 take_record_admission(decoder).map(r::ParticipantReferenceEnvelope::RecordAdmission)
1434 }
1435 t::ClientDiscriminant::EnrollmentRequest
1436 | t::ClientDiscriminant::ObserverRecoveryHandshake => {
1437 Err(decode_error(t::DecodeClass::InvalidField))
1438 }
1439 }
1440}
1441
1442fn take_binding_required(
1443 origin: t::ClientDiscriminant,
1444 decoder: &mut Decoder<'_>,
1445) -> Result<r::BindingRequiredEnvelope, CodecError> {
1446 match origin {
1447 t::ClientDiscriminant::DetachRequest => {
1448 take_detach(decoder).map(r::BindingRequiredEnvelope::Detach)
1449 }
1450 t::ClientDiscriminant::ParticipantAck => {
1451 take_participant_ack(decoder).map(r::BindingRequiredEnvelope::ParticipantAck)
1452 }
1453 t::ClientDiscriminant::LeaveRequest => {
1454 take_leave(decoder).map(r::BindingRequiredEnvelope::Leave)
1455 }
1456 t::ClientDiscriminant::MarkerAck => {
1457 take_marker_ack(decoder).map(r::BindingRequiredEnvelope::MarkerAck)
1458 }
1459 t::ClientDiscriminant::RecordAdmission => {
1460 take_record_admission(decoder).map(r::BindingRequiredEnvelope::RecordAdmission)
1461 }
1462 t::ClientDiscriminant::EnrollmentRequest
1463 | t::ClientDiscriminant::CredentialAttachRequest
1464 | t::ClientDiscriminant::ObserverRecoveryHandshake => {
1465 Err(decode_error(t::DecodeClass::InvalidField))
1466 }
1467 }
1468}
1469
1470fn take_order_allocating(
1471 origin: t::ClientDiscriminant,
1472 decoder: &mut Decoder<'_>,
1473) -> Result<r::OrderAllocatingEnvelope, CodecError> {
1474 match origin {
1475 t::ClientDiscriminant::EnrollmentRequest => {
1476 take_enrollment(decoder).map(r::OrderAllocatingEnvelope::Enrollment)
1477 }
1478 t::ClientDiscriminant::CredentialAttachRequest => {
1479 take_attach(decoder).map(r::OrderAllocatingEnvelope::CredentialAttach)
1480 }
1481 t::ClientDiscriminant::RecordAdmission => {
1482 take_record_admission(decoder).map(r::OrderAllocatingEnvelope::RecordAdmission)
1483 }
1484 _ => Err(decode_error(t::DecodeClass::InvalidField)),
1485 }
1486}
1487
1488fn take_closure_checked(
1489 origin: t::ClientDiscriminant,
1490 decoder: &mut Decoder<'_>,
1491) -> Result<c::ClosureCheckedEnvelope, CodecError> {
1492 match origin {
1493 t::ClientDiscriminant::EnrollmentRequest => {
1494 take_enrollment(decoder).map(c::ClosureCheckedEnvelope::Enrollment)
1495 }
1496 t::ClientDiscriminant::CredentialAttachRequest => {
1497 take_attach(decoder).map(c::ClosureCheckedEnvelope::CredentialAttach)
1498 }
1499 t::ClientDiscriminant::LeaveRequest => {
1500 take_leave(decoder).map(c::ClosureCheckedEnvelope::Leave)
1501 }
1502 t::ClientDiscriminant::RecordAdmission => {
1503 take_record_admission(decoder).map(c::ClosureCheckedEnvelope::RecordAdmission)
1504 }
1505 _ => Err(decode_error(t::DecodeClass::InvalidField)),
1506 }
1507}
1508
1509fn take_sequence_allocating(
1510 origin: t::ClientDiscriminant,
1511 decoder: &mut Decoder<'_>,
1512) -> Result<r::SequenceAllocatingEnvelope, CodecError> {
1513 match origin {
1514 t::ClientDiscriminant::EnrollmentRequest => {
1515 take_enrollment(decoder).map(r::SequenceAllocatingEnvelope::Enrollment)
1516 }
1517 t::ClientDiscriminant::CredentialAttachRequest => {
1518 take_attach(decoder).map(r::SequenceAllocatingEnvelope::CredentialAttach)
1519 }
1520 t::ClientDiscriminant::RecordAdmission => {
1521 take_record_admission(decoder).map(r::SequenceAllocatingEnvelope::RecordAdmission)
1522 }
1523 _ => Err(decode_error(t::DecodeClass::InvalidField)),
1524 }
1525}
1526
1527fn take_marker_proof(
1528 origin: t::ClientDiscriminant,
1529 decoder: &mut Decoder<'_>,
1530) -> Result<r::MarkerProofRequest, CodecError> {
1531 match origin {
1532 t::ClientDiscriminant::CredentialAttachRequest => Ok(
1533 r::MarkerProofRequest::CredentialAttach(r::AttachMarkerProof {
1534 conversation_id: decoder.take_u64()?,
1535 token: AttachAttemptToken::new(decoder.take_fixed::<TOKEN_LEN>()?),
1536 participant_id: decoder.take_u64()?,
1537 capability_generation: decoder.take_generation()?,
1538 requested_marker_delivery_seq: decoder.take_u64()?,
1539 }),
1540 ),
1541 t::ClientDiscriminant::MarkerAck => {
1542 Ok(r::MarkerProofRequest::MarkerAck(r::MarkerAckProof {
1543 conversation_id: decoder.take_u64()?,
1544 participant_id: decoder.take_u64()?,
1545 capability_generation: decoder.take_generation()?,
1546 requested_marker_delivery_seq: decoder.take_u64()?,
1547 }))
1548 }
1549 _ => Err(decode_error(t::DecodeClass::InvalidField)),
1550 }
1551}
1552
1553fn take_repayment_edge(decoder: &mut Decoder<'_>) -> Result<c::RepaymentEdge, CodecError> {
1554 let tag = t::RepaymentEdgeTag::try_from(decoder.take_u16()?)
1555 .map_err(|_| decode_error(t::DecodeClass::InvalidField))?;
1556 match tag {
1557 t::RepaymentEdgeTag::None => Ok(c::RepaymentEdge::None),
1558 t::RepaymentEdgeTag::ObserverProjection => Ok(c::RepaymentEdge::ObserverProjection {
1559 through_seq: decoder.take_u64()?,
1560 }),
1561 t::RepaymentEdgeTag::PhysicalCompaction => Ok(c::RepaymentEdge::PhysicalCompaction {
1562 from_floor: decoder.take_u64()?,
1563 through_seq: decoder.take_u64()?,
1564 }),
1565 t::RepaymentEdgeTag::MarkerDelivery => Ok(c::RepaymentEdge::MarkerDelivery {
1566 participant_id: decoder.take_u64()?,
1567 binding_epoch: decoder.take_binding_epoch()?,
1568 marker_delivery_seq: decoder.take_u64()?,
1569 }),
1570 t::RepaymentEdgeTag::ParticipantCursorProgress => Ok(
1571 c::RepaymentEdge::ParticipantCursorProgress(c::ParticipantCursorProgressEdge {
1572 participant_id: decoder.take_u64()?,
1573 binding_epoch: decoder.take_binding_epoch()?,
1574 through_seq: decoder.take_u64()?,
1575 marker_delivery_seq: decoder.take_option_u64()?,
1576 }),
1577 ),
1578 t::RepaymentEdgeTag::DetachedCredentialRecovery => {
1579 Ok(c::RepaymentEdge::DetachedCredentialRecovery {
1580 participant_id: decoder.take_u64()?,
1581 marker_delivery_seq: decoder.take_u64()?,
1582 prior_binding_epoch: decoder.take_binding_epoch()?,
1583 })
1584 }
1585 t::RepaymentEdgeTag::DetachedMarkerRelease => Ok(c::RepaymentEdge::DetachedMarkerRelease {
1586 participant_id: decoder.take_u64()?,
1587 marker_delivery_seq: decoder.take_u64()?,
1588 last_dead_binding_epoch: decoder.take_binding_epoch()?,
1589 }),
1590 t::RepaymentEdgeTag::DetachedCursorRelease => Ok(c::RepaymentEdge::DetachedCursorRelease {
1591 participant_id: decoder.take_u64()?,
1592 last_dead_binding_epoch: decoder.take_binding_epoch()?,
1593 }),
1594 }
1595}
1596
1597fn take_closure_snapshot(decoder: &mut Decoder<'_>) -> Result<c::ClosureSnapshot, CodecError> {
1598 Ok(c::ClosureSnapshot {
1599 marker_capacity_credits: decoder.take_u64()?,
1600 marker_anchors: decoder.take_u64()?,
1601 entry_debt: decoder.take_u64()?,
1602 byte_debt: decoder.take_u64()?,
1603 repayment_edge: take_repayment_edge(decoder)?,
1604 edge_sequence_claims: decoder.take_u64()?,
1605 edge_order_position_claims: decoder.take_u64()?,
1606 edge_k_remaining: decoder.take_resource_vector()?,
1607 k_headroom: decoder.take_wide_resource_vector()?,
1608 episode_churn_used: decoder.take_u64()?,
1609 delta_cycles: decoder.take_u64()?,
1610 episode_churn_limit: decoder.take_u64()?,
1611 })
1612}
1613
1614fn take_sequence_budget(decoder: &mut Decoder<'_>) -> Result<super::SequenceBudget, CodecError> {
1615 Ok(super::SequenceBudget {
1616 high_watermark: decoder.take_u64()?,
1617 remaining: decoder.take_u64()?,
1618 e: decoder.take_u64()?,
1619 t: decoder.take_u64()?,
1620 m: decoder.take_u64()?,
1621 rs: decoder.take_u64()?,
1622 rt: decoder.take_u64()?,
1623 l_times_t: decoder.take_u128()?,
1624 l_times_rt: decoder.take_u128()?,
1625 l_other_times_e: decoder.take_u128()?,
1626 })
1627}
1628
1629fn required_origin(
1630 origin: Option<t::ClientDiscriminant>,
1631) -> Result<t::ClientDiscriminant, CodecError> {
1632 origin.ok_or_else(|| decode_error(t::DecodeClass::InvalidField))
1633}
1634
1635#[allow(clippy::too_many_lines)]
1636fn decode_server_suffix(
1637 discriminant: t::ServerDiscriminant,
1638 origin: Option<t::ClientDiscriminant>,
1639 decoder: &mut Decoder<'_>,
1640) -> Result<r::ServerValue, CodecError> {
1641 use t::ServerDiscriminant as D;
1642
1643 match discriminant {
1644 D::ParticipantTransportRejected => {
1645 let reason = match take_tag::<t::TransportReasonTag>(decoder)? {
1646 t::TransportReasonTag::FrameTooLarge => {
1647 r::TransportRejectionReason::FrameTooLarge {
1648 complete_frame_bytes: decoder.take_u64()?,
1649 max_frame_bytes: decoder.take_u64()?,
1650 }
1651 }
1652 t::TransportReasonTag::DecodeFailed => {
1653 let decode_class = take_tag::<t::DecodeClass>(decoder)?;
1654 r::TransportRejectionReason::DecodeFailed { decode_class }
1655 }
1656 t::TransportReasonTag::UnsupportedVersion => {
1657 r::TransportRejectionReason::UnsupportedVersion {
1658 presented_version: decoder.take_protocol_version()?,
1659 supported_version: decoder.take_protocol_version()?,
1660 }
1661 }
1662 t::TransportReasonTag::AuthenticationFailed => {
1663 r::TransportRejectionReason::AuthenticationFailed
1664 }
1665 t::TransportReasonTag::ParticipantCapabilityRequired => {
1666 let capability = decoder.take_string()?;
1667 if capability != r::PARTICIPANT_CAPABILITY {
1668 decoder.invalidate();
1669 }
1670 r::TransportRejectionReason::ParticipantCapabilityRequired
1671 }
1672 };
1673 Ok(r::ServerValue::ParticipantTransportRejected(
1674 r::ParticipantTransportRejected { reason },
1675 ))
1676 }
1677 D::AttemptTokenBodyConflict => {
1678 let origin = required_origin(origin)?;
1679 let token = decoder.take_fixed::<TOKEN_LEN>()?;
1680 let operation = take_tag::<t::AttemptOperation>(decoder)?;
1681 match origin {
1682 t::ClientDiscriminant::CredentialAttachRequest => {
1683 if operation != t::AttemptOperation::CredentialAttachRequest {
1684 return Err(decode_error(t::DecodeClass::InvalidField));
1685 }
1686 Ok(r::ServerValue::AttemptTokenBodyConflict(
1687 r::AttemptTokenBodyConflict::CredentialAttach {
1688 token: AttachAttemptToken::new(token),
1689 conversation_id: decoder.take_u64()?,
1690 presented_participant_id: decoder.take_u64()?,
1691 presented_generation: decoder.take_generation()?,
1692 presented_marker_delivery_seq: decoder.take_option_u64()?,
1693 conflict: take_tag::<t::AttemptConflict>(decoder)?,
1694 },
1695 ))
1696 }
1697 t::ClientDiscriminant::LeaveRequest => {
1698 if operation != t::AttemptOperation::LeaveRequest {
1699 return Err(decode_error(t::DecodeClass::InvalidField));
1700 }
1701 let value = r::AttemptTokenBodyConflict::Leave {
1702 token: LeaveAttemptToken::new(token),
1703 conversation_id: decoder.take_u64()?,
1704 presented_participant_id: decoder.take_u64()?,
1705 presented_generation: decoder.take_generation()?,
1706 };
1707 if take_tag::<t::AttemptConflict>(decoder)? != t::AttemptConflict::Generation {
1708 decoder.invalidate();
1709 }
1710 Ok(r::ServerValue::AttemptTokenBodyConflict(value))
1711 }
1712 t::ClientDiscriminant::RecordAdmission => {
1713 if operation != t::AttemptOperation::RecordAdmission {
1714 return Err(decode_error(t::DecodeClass::InvalidField));
1715 }
1716 Ok(r::ServerValue::AttemptTokenBodyConflict(
1717 r::AttemptTokenBodyConflict::RecordAdmission {
1718 token: RecordAdmissionAttemptToken::new(token),
1719 conversation_id: decoder.take_u64()?,
1720 presented_participant_id: decoder.take_u64()?,
1721 presented_generation: decoder.take_generation()?,
1722 },
1723 ))
1724 }
1725 _ => Err(decode_error(t::DecodeClass::InvalidField)),
1726 }
1727 }
1728 D::ConnectionConversationCapacityExceeded => {
1729 let origin = required_origin(origin)?;
1730 Ok(r::ServerValue::ConnectionConversationCapacityExceeded(
1731 r::ConnectionConversationCapacityExceeded::SemanticRequest {
1732 request: take_response_envelope(origin, decoder)?,
1733 limit: decoder.take_u64()?,
1734 },
1735 ))
1736 }
1737 D::ConnectionConversationBindingOccupied => {
1738 let origin = required_origin(origin)?;
1739 let value = match origin {
1740 t::ClientDiscriminant::EnrollmentRequest => {
1741 let conversation_id = decoder.take_u64()?;
1742 let enrollment_token = EnrollmentToken::new(decoder.take_fixed::<TOKEN_LEN>()?);
1743 if decoder.take_option_u64()?.is_some() {
1744 decoder.invalidate();
1745 }
1746 r::ConnectionConversationBindingOccupied::Enrollment {
1747 conversation_id,
1748 enrollment_token,
1749 }
1750 }
1751 t::ClientDiscriminant::CredentialAttachRequest => {
1752 let conversation_id = decoder.take_u64()?;
1753 let participant_id = decoder.take_u64()?;
1754 let capability_generation = decoder.take_generation()?;
1755 let attach_attempt_token =
1756 AttachAttemptToken::new(decoder.take_fixed::<TOKEN_LEN>()?);
1757 let accept_marker_delivery_seq = decoder.take_option_u64()?;
1758 if decoder.take_option_u64()? != Some(participant_id) {
1759 decoder.invalidate();
1760 }
1761 r::ConnectionConversationBindingOccupied::CredentialAttach {
1762 conversation_id,
1763 participant_id,
1764 capability_generation,
1765 attach_attempt_token,
1766 accept_marker_delivery_seq,
1767 }
1768 }
1769 _ => return Err(decode_error(t::DecodeClass::InvalidField)),
1770 };
1771 Ok(r::ServerValue::ConnectionConversationBindingOccupied(value))
1772 }
1773 D::ConversationOrderExhausted => {
1774 let origin = required_origin(origin)?;
1775 let request = take_order_allocating(origin, decoder)?;
1776 let _counter = take_tag::<t::Counter>(decoder)?;
1777 let high = decoder.take_u64()?;
1778 let next_value = decoder.take_option_u64()?;
1779 let order_remaining = decoder.take_u128()?;
1780 let reserved_claims = decoder.take_u128()?;
1781 if decoder.take_u64()? != r::ConversationOrderExhausted::REQUIRED_MAJORS {
1782 decoder.invalidate();
1783 }
1784 let resulting_order_remaining = decoder.take_u128()?;
1785 let resulting_reserved_claims = decoder.take_u128()?;
1786 let value = r::ConversationOrderExhausted::new(
1787 request,
1788 high,
1789 order_remaining,
1790 reserved_claims,
1791 resulting_order_remaining,
1792 resulting_reserved_claims,
1793 );
1794 if next_value != value.next_value() {
1795 decoder.invalidate();
1796 }
1797 Ok(r::ServerValue::ConversationOrderExhausted(
1798 alloc::boxed::Box::new(value),
1799 ))
1800 }
1801 D::ParticipantUnknown => {
1802 let origin = required_origin(origin)?;
1803 Ok(r::ServerValue::ParticipantUnknown(r::ParticipantUnknown {
1804 request: take_participant_reference(origin, decoder)?,
1805 }))
1806 }
1807 D::NoBinding => {
1808 let origin = required_origin(origin)?;
1809 Ok(r::ServerValue::NoBinding(r::NoBinding {
1810 request: take_binding_required(origin, decoder)?,
1811 }))
1812 }
1813 D::StaleAuthority => take_stale_authority(required_origin(origin)?, decoder)
1814 .map(r::ServerValue::StaleAuthority),
1815 D::Retired => {
1816 let origin = required_origin(origin)?;
1817 let value = if origin == t::ClientDiscriminant::EnrollmentRequest {
1818 r::Retired::Enrollment {
1819 request: take_enrollment(decoder)?,
1820 participant_id: decoder.take_u64()?,
1821 retired_generation: decoder.take_generation()?,
1822 }
1823 } else {
1824 r::Retired::Participant {
1825 request: take_participant_reference(origin, decoder)?,
1826 retired_generation: decoder.take_generation()?,
1827 }
1828 };
1829 Ok(r::ServerValue::Retired(value))
1830 }
1831 D::MarkerClosureCapacityExceeded => {
1832 let origin = required_origin(origin)?;
1833 let request = take_closure_checked(origin, decoder)?;
1834 let scope = take_tag::<t::ClosureScope>(decoder)?;
1835 let snapshot = take_closure_snapshot(decoder)?;
1836 let reason = match scope {
1837 t::ClosureScope::Capacity => {
1838 let dimension =
1839 ResourceDimension::from(take_tag::<t::ResourceDimensionTag>(decoder)?);
1840 c::ClosureRefusalReason::Capacity(c::ClosureCapacityReason {
1841 dimension,
1842 required: decoder.take_u128()?,
1843 limit: decoder.take_u128()?,
1844 })
1845 }
1846 t::ClosureScope::RecoveryFence => c::ClosureRefusalReason::RecoveryFence,
1847 t::ClosureScope::DeliveredMarkerAwaitingAck => {
1848 c::ClosureRefusalReason::DeliveredMarkerAwaitingAck
1849 }
1850 t::ClosureScope::EpisodeChurnLimit => c::ClosureRefusalReason::EpisodeChurnLimit,
1851 };
1852 Ok(r::ServerValue::MarkerClosureCapacityExceeded(
1853 alloc::boxed::Box::new(c::MarkerClosureCapacityExceeded {
1854 request,
1855 snapshot,
1856 reason,
1857 }),
1858 ))
1859 }
1860 D::EnrollBound => take_enroll_bound(decoder).map(r::ServerValue::EnrollBound),
1861 D::EnrollmentKnown => Ok(r::ServerValue::EnrollmentKnown(r::EnrollmentKnown {
1862 conversation_id: decoder.take_u64()?,
1863 token: EnrollmentToken::new(decoder.take_fixed::<TOKEN_LEN>()?),
1864 participant_id: decoder.take_u64()?,
1865 current_generation: decoder.take_generation()?,
1866 })),
1867 D::ReceiptExpired => take_receipt_expired(required_origin(origin)?, decoder)
1868 .map(r::ServerValue::ReceiptExpired),
1869 D::ReceiptCapacityExceeded => take_receipt_capacity(required_origin(origin)?, decoder)
1870 .map(r::ServerValue::ReceiptCapacityExceeded),
1871 D::IdentityCapacityExceeded => {
1872 let request = take_enrollment(decoder)?;
1873 let scope = take_tag::<t::IdentityCapacityScope>(decoder)?;
1874 let limit = decoder.take_u64()?;
1875 let occupied = decoder.take_u64()?;
1876 if decoder.take_u64()? != r::IdentityCapacityExceeded::REQUESTED {
1877 decoder.invalidate();
1878 }
1879 Ok(r::ServerValue::IdentityCapacityExceeded(
1880 r::IdentityCapacityExceeded {
1881 request,
1882 scope,
1883 limit,
1884 occupied,
1885 },
1886 ))
1887 }
1888 D::ObserverBackpressure => take_observer_backpressure(required_origin(origin)?, decoder)
1889 .map(r::ServerValue::ObserverBackpressure),
1890 D::ConversationSequenceExhausted => {
1891 let request = take_sequence_allocating(required_origin(origin)?, decoder)?;
1892 Ok(r::ServerValue::ConversationSequenceExhausted(
1893 alloc::boxed::Box::new(r::ConversationSequenceExhausted {
1894 request,
1895 sequence_budget: take_sequence_budget(decoder)?,
1896 }),
1897 ))
1898 }
1899 D::AttachBound => take_attach_bound(decoder).map(r::ServerValue::AttachBound),
1900 D::StaleOrUnknownReceipt => Ok(r::ServerValue::StaleOrUnknownReceipt(
1901 r::StaleOrUnknownReceipt {
1902 conversation_id: decoder.take_u64()?,
1903 token: AttachAttemptToken::new(decoder.take_fixed::<TOKEN_LEN>()?),
1904 participant_id: decoder.take_u64()?,
1905 presented_generation: decoder.take_generation()?,
1906 presented_marker_delivery_seq: decoder.take_option_u64()?,
1907 current_generation: decoder.take_generation()?,
1908 },
1909 )),
1910 D::MarkerNotDelivered => {
1911 let request = take_marker_proof(required_origin(origin)?, decoder)?;
1912 let reason = take_tag::<t::MarkerNotDeliveredReason>(decoder)?;
1913 Ok(r::ServerValue::MarkerNotDelivered(r::MarkerNotDelivered {
1914 request,
1915 reason,
1916 expected_marker_delivery_seq: decoder.take_u64()?,
1917 }))
1918 }
1919 D::MarkerMismatch => {
1920 let request = take_marker_proof(required_origin(origin)?, decoder)?;
1921 let reason = take_tag::<t::MarkerMismatchReason>(decoder)?;
1922 let mismatch = match reason {
1923 t::MarkerMismatchReason::BelowCursor => r::MarkerMismatchBody::BelowCursor {
1924 current_cursor: decoder.take_u64()?,
1925 },
1926 t::MarkerMismatchReason::NoMarkerExpected => {
1927 r::MarkerMismatchBody::NoMarkerExpected
1928 }
1929 t::MarkerMismatchReason::ExpectedDifferentMarker => {
1930 r::MarkerMismatchBody::ExpectedDifferentMarker {
1931 expected_marker_delivery_seq: decoder.take_u64()?,
1932 }
1933 }
1934 };
1935 Ok(r::ServerValue::MarkerMismatch(r::MarkerMismatch {
1936 request,
1937 mismatch,
1938 }))
1939 }
1940 D::Bound => {
1941 take_receipt_replay(required_origin(origin)?, decoder).map(r::ServerValue::Bound)
1942 }
1943 D::UnboundReceipt => take_receipt_replay(required_origin(origin)?, decoder)
1944 .map(r::ServerValue::UnboundReceipt),
1945 D::DetachCommitted => {
1946 let conversation_id = decoder.take_u64()?;
1947 let participant_id = decoder.take_u64()?;
1948 let capability_generation = decoder.take_generation()?;
1949 let detach_attempt_token = DetachAttemptToken::new(decoder.take_fixed::<TOKEN_LEN>()?);
1950 let committed_binding_epoch = decoder.take_binding_epoch()?;
1951 let detached_delivery_seq = decoder.take_u64()?;
1952 if capability_generation != committed_binding_epoch.capability_generation {
1953 decoder.invalidate();
1954 }
1955 Ok(r::ServerValue::DetachCommitted(r::DetachCommitted::new(
1956 conversation_id,
1957 participant_id,
1958 detach_attempt_token,
1959 committed_binding_epoch,
1960 detached_delivery_seq,
1961 )))
1962 }
1963 D::DetachInProgress => Ok(r::ServerValue::DetachInProgress(r::DetachInProgress {
1964 conversation_id: decoder.take_u64()?,
1965 participant_id: decoder.take_u64()?,
1966 presented_token: DetachAttemptToken::new(decoder.take_fixed::<TOKEN_LEN>()?),
1967 presented_generation: decoder.take_generation()?,
1968 committed_binding_epoch: decoder.take_binding_epoch()?,
1969 })),
1970 D::AckCommitted => {
1971 let request = take_participant_ack(decoder)?;
1972 if decoder.take_u64()? != request.through_seq {
1973 decoder.invalidate();
1974 }
1975 Ok(r::ServerValue::AckCommitted(r::AckCommitted::new(request)))
1976 }
1977 D::AckNoOp => {
1978 take_ack_no_op(required_origin(origin)?, decoder).map(r::ServerValue::AckNoOp)
1979 }
1980 D::AckGap => {
1981 let request = take_participant_ack(decoder)?;
1982 let current_cursor = decoder.take_u64()?;
1983 let _reason = take_tag::<t::AckGapReason>(decoder)?;
1984 r::AckGap::new(request, current_cursor)
1985 .map(r::ServerValue::AckGap)
1986 .ok_or_else(|| decode_error(t::DecodeClass::InvalidField))
1987 }
1988 D::AckRegression => {
1989 let request = take_participant_ack(decoder)?;
1990 let current_cursor = decoder.take_u64()?;
1991 let _reason = take_tag::<t::AckRegressionReason>(decoder)?;
1992 r::AckRegression::new(request, current_cursor)
1993 .map(r::ServerValue::AckRegression)
1994 .ok_or_else(|| decode_error(t::DecodeClass::InvalidField))
1995 }
1996 D::LeaveCommitted => {
1997 let conversation_id = decoder.take_u64()?;
1998 let leave_attempt_token = LeaveAttemptToken::new(decoder.take_fixed::<TOKEN_LEN>()?);
1999 let participant_id = decoder.take_u64()?;
2000 let presented_generation = decoder.take_generation()?;
2001 let retired_generation = decoder.take_generation()?;
2002 let ended_binding_epoch = decoder.take_option_binding_epoch()?;
2003 let prior_terminal_delivery_seq = decoder.take_option_u64()?;
2004 let left_delivery_seq = decoder.take_u64()?;
2005 if presented_generation != retired_generation {
2006 return Err(decode_error(t::DecodeClass::InvalidField));
2007 }
2008 r::LeaveCommitted::new(
2009 conversation_id,
2010 leave_attempt_token,
2011 participant_id,
2012 retired_generation,
2013 ended_binding_epoch,
2014 prior_terminal_delivery_seq,
2015 left_delivery_seq,
2016 )
2017 .map(r::ServerValue::LeaveCommitted)
2018 .ok_or_else(|| decode_error(t::DecodeClass::InvalidField))
2019 }
2020 D::MarkerAckCommitted => {
2021 let request = take_marker_ack(decoder)?;
2022 if decoder.take_u64()? != request.marker_delivery_seq {
2023 decoder.invalidate();
2024 }
2025 Ok(r::ServerValue::MarkerAckCommitted(
2026 r::MarkerAckCommitted::new(request),
2027 ))
2028 }
2029 D::RecordCommitted => {
2030 let request = take_record_admission(decoder)?;
2031 if decoder.take_u64()? != request.participant_id {
2032 decoder.invalidate();
2033 }
2034 let delivery_seq = decoder.take_u64()?;
2035 Ok(r::ServerValue::RecordCommitted(r::RecordCommitted::new(
2036 request,
2037 delivery_seq,
2038 )))
2039 }
2040 D::RecordTooLarge => Ok(r::ServerValue::RecordTooLarge(r::RecordTooLarge {
2041 request: take_record_admission(decoder)?,
2042 dimension: ResourceDimension::from(take_tag::<t::ResourceDimensionTag>(decoder)?),
2043 encoded_record_charge: decoder.take_resource_vector()?,
2044 max_ordinary_record_charge: decoder.take_resource_vector()?,
2045 })),
2046 D::ObserverRecoveryAccepted => {
2047 take_recovery_accepted(decoder).map(r::ServerValue::ObserverRecoveryAccepted)
2048 }
2049 D::InvalidObserverEpoch => {
2050 require_zero_recovery_count(decoder)?;
2051 take_invalid_observer_epoch(decoder).map(r::ServerValue::InvalidObserverEpoch)
2052 }
2053 D::InvalidObserverEpochList => {
2054 require_zero_recovery_count(decoder)?;
2055 take_invalid_observer_epoch_list(decoder).map(r::ServerValue::InvalidObserverEpochList)
2056 }
2057 D::ObserverRecoveryConnectionCapacityExceeded => {
2058 require_zero_recovery_count(decoder)?;
2059 Ok(r::ServerValue::ConnectionConversationCapacityExceeded(
2060 r::ConnectionConversationCapacityExceeded::ObserverRecovery {
2061 conversation_id: decoder.take_u64()?,
2062 limit: decoder.take_u64()?,
2063 },
2064 ))
2065 }
2066 D::MarkerSettlementBackpressure => {
2067 let origin = required_origin(origin)?;
2068 let conversation_id = decoder.take_u64()?;
2069 let refused_epoch = decoder.take_u64()?;
2070 let value = match origin {
2071 t::ClientDiscriminant::CredentialAttachRequest => {
2072 r::MarkerSettlementBackpressure::CredentialAttach {
2073 conversation_id,
2074 refused_epoch,
2075 }
2076 }
2077 t::ClientDiscriminant::DetachRequest => r::MarkerSettlementBackpressure::Detach {
2078 conversation_id,
2079 refused_epoch,
2080 },
2081 _ => return Err(decode_error(t::DecodeClass::InvalidField)),
2082 };
2083 Ok(r::ServerValue::MarkerSettlementBackpressure(value))
2084 }
2085 D::EnrollmentSettlementBackpressure => {
2086 required_origin(origin)?;
2087 Ok(r::ServerValue::EnrollmentSettlementBackpressure(
2088 r::EnrollmentSettlementBackpressure {
2089 conversation_id: decoder.take_u64()?,
2090 },
2091 ))
2092 }
2093 }
2094}
2095
2096trait WireTag: TryFrom<u16> {}
2097
2098impl WireTag for t::TransportReasonTag {}
2099impl WireTag for t::DecodeClass {}
2100impl WireTag for t::AttemptOperation {}
2101impl WireTag for t::AttemptConflict {}
2102impl WireTag for t::Counter {}
2103impl WireTag for t::ClosureScope {}
2104impl WireTag for t::ResourceDimensionTag {}
2105impl WireTag for t::IdentityCapacityScope {}
2106impl WireTag for t::MarkerNotDeliveredReason {}
2107impl WireTag for t::MarkerMismatchReason {}
2108impl WireTag for t::AckGapReason {}
2109impl WireTag for t::AckRegressionReason {}
2110impl WireTag for t::InvalidObserverEpochReason {}
2111impl WireTag for t::InvalidObserverEpochListReason {}
2112
2113fn take_tag<T>(decoder: &mut Decoder<'_>) -> Result<T, CodecError>
2114where
2115 T: WireTag,
2116{
2117 T::try_from(decoder.take_u16()?).map_err(|_| decode_error(t::DecodeClass::InvalidField))
2118}
2119
2120fn take_stale_authority(
2121 origin: t::ClientDiscriminant,
2122 decoder: &mut Decoder<'_>,
2123) -> Result<r::StaleAuthority, CodecError> {
2124 match origin {
2125 t::ClientDiscriminant::CredentialAttachRequest => Ok(r::StaleAuthority::Live {
2126 request: r::CommonStaleAuthorityEnvelope::CredentialAttach(take_attach(decoder)?),
2127 current_generation: decoder.take_generation()?,
2128 }),
2129 t::ClientDiscriminant::ParticipantAck => Ok(r::StaleAuthority::Live {
2130 request: r::CommonStaleAuthorityEnvelope::ParticipantAck(take_participant_ack(
2131 decoder,
2132 )?),
2133 current_generation: decoder.take_generation()?,
2134 }),
2135 t::ClientDiscriminant::MarkerAck => Ok(r::StaleAuthority::Live {
2136 request: r::CommonStaleAuthorityEnvelope::MarkerAck(take_marker_ack(decoder)?),
2137 current_generation: decoder.take_generation()?,
2138 }),
2139 t::ClientDiscriminant::RecordAdmission => Ok(r::StaleAuthority::Live {
2140 request: r::CommonStaleAuthorityEnvelope::RecordAdmission(take_record_admission(
2141 decoder,
2142 )?),
2143 current_generation: decoder.take_generation()?,
2144 }),
2145 t::ClientDiscriminant::DetachRequest => {
2146 let authority = t::DetachAuthorityStateTag::try_from(decoder.take_u16()?)
2147 .map_err(|_| decode_error(t::DecodeClass::InvalidField))?;
2148 let value = match authority {
2149 t::DetachAuthorityStateTag::Live => r::DetachStaleAuthority::Live {
2150 conversation_id: decoder.take_u64()?,
2151 participant_id: decoder.take_u64()?,
2152 capability_generation: decoder.take_generation()?,
2153 detach_attempt_token: DetachAttemptToken::new(
2154 decoder.take_fixed::<TOKEN_LEN>()?,
2155 ),
2156 current_generation: decoder.take_generation()?,
2157 },
2158 t::DetachAuthorityStateTag::TerminalizedDetachCell => {
2159 let conversation_id = decoder.take_u64()?;
2160 let participant_id = decoder.take_u64()?;
2161 let capability_generation = decoder.take_generation()?;
2162 let detach_attempt_token =
2163 DetachAttemptToken::new(decoder.take_fixed::<TOKEN_LEN>()?);
2164 let current_generation = decoder.take_generation()?;
2165 let committed_binding_epoch = decoder.take_binding_epoch()?;
2166 let binding_state = match t::BindingStateTag::try_from(decoder.take_u16()?)
2167 .map_err(|_| decode_error(t::DecodeClass::InvalidField))?
2168 {
2169 t::BindingStateTag::Bound => r::BindingStateView::Bound {
2170 current_binding_epoch: decoder.take_binding_epoch()?,
2171 },
2172 t::BindingStateTag::Detached => r::BindingStateView::Detached,
2173 };
2174 r::DetachStaleAuthority::TerminalizedDetachCell(
2175 r::TerminalizedDetachCell::from_wire_decode(
2176 TERMINALIZED_WIRE_DECODE_AUTHORITY,
2177 conversation_id,
2178 participant_id,
2179 capability_generation,
2180 detach_attempt_token,
2181 current_generation,
2182 committed_binding_epoch,
2183 binding_state,
2184 ),
2185 )
2186 }
2187 };
2188 Ok(r::StaleAuthority::Detach(value))
2189 }
2190 t::ClientDiscriminant::LeaveRequest => {
2191 let authority = t::LeaveAuthorityStateTag::try_from(decoder.take_u16()?)
2192 .map_err(|_| decode_error(t::DecodeClass::InvalidField))?;
2193 let conversation_id = decoder.take_u64()?;
2194 let participant_id = decoder.take_u64()?;
2195 let presented_generation = decoder.take_generation()?;
2196 let leave_attempt_token = LeaveAttemptToken::new(decoder.take_fixed::<TOKEN_LEN>()?);
2197 let value = match authority {
2198 t::LeaveAuthorityStateTag::Live => r::LeaveStaleAuthority::Live {
2199 conversation_id,
2200 participant_id,
2201 presented_generation,
2202 leave_attempt_token,
2203 current_generation: decoder.take_generation()?,
2204 },
2205 t::LeaveAuthorityStateTag::CommittedLeaveTombstone => {
2206 r::LeaveStaleAuthority::CommittedLeaveTombstone {
2207 conversation_id,
2208 participant_id,
2209 presented_generation,
2210 leave_attempt_token,
2211 retired_generation: decoder.take_generation()?,
2212 }
2213 }
2214 };
2215 Ok(r::StaleAuthority::Leave(value))
2216 }
2217 t::ClientDiscriminant::EnrollmentRequest
2218 | t::ClientDiscriminant::ObserverRecoveryHandshake => {
2219 Err(decode_error(t::DecodeClass::InvalidField))
2220 }
2221 }
2222}
2223
2224fn take_enroll_bound(decoder: &mut Decoder<'_>) -> Result<r::EnrollBound, CodecError> {
2225 let conversation_id = decoder.take_u64()?;
2226 let token = EnrollmentToken::new(decoder.take_fixed::<TOKEN_LEN>()?);
2227 let participant_id = decoder.take_u64()?;
2228 if decoder.take_option_generation()?.is_some() {
2229 decoder.invalidate();
2230 }
2231 let capability_generation = decoder.take_generation()?;
2232 if capability_generation.get() != 1 {
2233 decoder.invalidate();
2234 }
2235 let attach_secret = AttachSecret::new(decoder.take_fixed::<SECRET_LEN>()?);
2236 let origin_binding_epoch = decoder.take_binding_epoch()?;
2237 if origin_binding_epoch.capability_generation.get() != 1 {
2238 decoder.invalidate();
2239 }
2240 if decoder.take_u64()? != 0 {
2241 decoder.invalidate();
2242 }
2243 if decoder.take_option_u64()?.is_some() {
2244 decoder.invalidate();
2245 }
2246 let receipt_expires_at = decoder.take_u128()?;
2247 let provenance_expires_at = decoder.take_u128()?;
2248
2249 let valid_epoch = if origin_binding_epoch.capability_generation.get() == 1 {
2250 origin_binding_epoch
2251 } else {
2252 BindingEpoch::new(
2253 origin_binding_epoch.connection_incarnation,
2254 generation_one()?,
2255 )
2256 };
2257 r::EnrollBound::new(
2258 conversation_id,
2259 token,
2260 participant_id,
2261 attach_secret,
2262 valid_epoch,
2263 receipt_expires_at,
2264 provenance_expires_at,
2265 )
2266 .ok_or(CodecError::InvalidValue)
2267}
2268
2269fn take_receipt_expired(
2270 origin: t::ClientDiscriminant,
2271 decoder: &mut Decoder<'_>,
2272) -> Result<r::ReceiptExpired, CodecError> {
2273 let conversation_id = decoder.take_u64()?;
2274 let token = decoder.take_fixed::<TOKEN_LEN>()?;
2275 let participant_id = decoder.take_u64()?;
2276 match origin {
2277 t::ClientDiscriminant::EnrollmentRequest => {
2278 if decoder.take_option_generation()?.is_some() {
2279 decoder.invalidate();
2280 }
2281 Ok(r::ReceiptExpired::Enrollment {
2282 conversation_id,
2283 token: EnrollmentToken::new(token),
2284 participant_id,
2285 result_generation: decoder.take_generation()?,
2286 current_generation: decoder.take_generation()?,
2287 reason: take_tag::<t::ReceiptExpiryReason>(decoder)?,
2288 })
2289 }
2290 t::ClientDiscriminant::CredentialAttachRequest => {
2291 let presented_generation = required_generation_option(decoder)?;
2292 Ok(r::ReceiptExpired::CredentialAttach {
2293 conversation_id,
2294 token: AttachAttemptToken::new(token),
2295 participant_id,
2296 presented_generation,
2297 presented_marker_delivery_seq: decoder.take_option_u64()?,
2298 result_generation: decoder.take_generation()?,
2299 current_generation: decoder.take_generation()?,
2300 reason: take_tag::<t::ReceiptExpiryReason>(decoder)?,
2301 })
2302 }
2303 _ => Err(decode_error(t::DecodeClass::InvalidField)),
2304 }
2305}
2306
2307impl WireTag for t::ReceiptExpiryReason {}
2308impl WireTag for t::ReceiptCapacityScope {}
2309
2310fn take_receipt_capacity(
2311 origin: t::ClientDiscriminant,
2312 decoder: &mut Decoder<'_>,
2313) -> Result<r::ReceiptCapacityExceeded, CodecError> {
2314 match origin {
2315 t::ClientDiscriminant::EnrollmentRequest => {
2316 let request = take_enrollment(decoder)?;
2317 let scope = match take_tag::<t::ReceiptCapacityScope>(decoder)? {
2318 t::ReceiptCapacityScope::LiveReceiptServer => {
2319 r::EnrollmentReceiptCapacityScope::LiveReceiptServer
2320 }
2321 t::ReceiptCapacityScope::ProvenanceServer => {
2322 r::EnrollmentReceiptCapacityScope::ProvenanceServer
2323 }
2324 t::ReceiptCapacityScope::ProvenanceConversation => {
2325 r::EnrollmentReceiptCapacityScope::ProvenanceConversation
2326 }
2327 t::ReceiptCapacityScope::LiveReceiptParticipant
2328 | t::ReceiptCapacityScope::ProvenanceParticipant => {
2329 return Err(decode_error(t::DecodeClass::InvalidField));
2330 }
2331 };
2332 let limit = decoder.take_u64()?;
2333 let occupied = decoder.take_u64()?;
2334 require_one(decoder, r::ReceiptCapacityExceeded::REQUESTED)?;
2335 Ok(r::ReceiptCapacityExceeded::Enrollment {
2336 request,
2337 scope,
2338 limit,
2339 occupied,
2340 })
2341 }
2342 t::ClientDiscriminant::CredentialAttachRequest => {
2343 let request = take_attach(decoder)?;
2344 let scope = take_tag::<t::ReceiptCapacityScope>(decoder)?;
2345 let limit = decoder.take_u64()?;
2346 let occupied = decoder.take_u64()?;
2347 require_one(decoder, r::ReceiptCapacityExceeded::REQUESTED)?;
2348 Ok(r::ReceiptCapacityExceeded::CredentialAttach {
2349 request,
2350 scope,
2351 limit,
2352 occupied,
2353 })
2354 }
2355 _ => Err(decode_error(t::DecodeClass::InvalidField)),
2356 }
2357}
2358
2359fn require_one(decoder: &mut Decoder<'_>, expected: u64) -> Result<(), CodecError> {
2360 if decoder.take_u64()? != expected {
2361 decoder.invalidate();
2362 }
2363 Ok(())
2364}
2365
2366fn take_observer_backpressure(
2367 origin: t::ClientDiscriminant,
2368 decoder: &mut Decoder<'_>,
2369) -> Result<r::ObserverBackpressure, CodecError> {
2370 match origin {
2371 t::ClientDiscriminant::EnrollmentRequest => Ok(r::ObserverBackpressure::Enrollment {
2372 request: take_enrollment(decoder)?,
2373 state: take_backpressure_state(decoder)?,
2374 }),
2375 t::ClientDiscriminant::CredentialAttachRequest => {
2376 Ok(r::ObserverBackpressure::CredentialAttach {
2377 request: take_attach(decoder)?,
2378 state: take_backpressure_state(decoder)?,
2379 })
2380 }
2381 t::ClientDiscriminant::DetachRequest => Ok(r::ObserverBackpressure::Detach {
2382 request: take_detach(decoder)?,
2383 committed_binding_epoch: decoder.take_binding_epoch()?,
2384 state: take_backpressure_state(decoder)?,
2385 }),
2386 t::ClientDiscriminant::LeaveRequest => Ok(r::ObserverBackpressure::Leave {
2387 request: take_leave(decoder)?,
2388 state: take_backpressure_state(decoder)?,
2389 prior_terminal_cell_exists: decoder.take_bool()?,
2390 }),
2391 t::ClientDiscriminant::RecordAdmission => Ok(r::ObserverBackpressure::RecordAdmission {
2392 request: take_record_admission(decoder)?,
2393 state: take_backpressure_state(decoder)?,
2394 }),
2395 _ => Err(decode_error(t::DecodeClass::InvalidField)),
2396 }
2397}
2398
2399fn take_backpressure_state(
2400 decoder: &mut Decoder<'_>,
2401) -> Result<r::ObserverBackpressureState, CodecError> {
2402 let backpressure_epoch = decoder.take_u64()?;
2403 let observer_progress = decoder.take_u64()?;
2404 r::ObserverBackpressureState::replay(backpressure_epoch, observer_progress)
2405 .ok_or_else(|| decode_error(t::DecodeClass::InvalidField))
2406}
2407
2408fn required_generation_option(decoder: &mut Decoder<'_>) -> Result<Generation, CodecError> {
2409 decoder.take_option_generation()?.map_or_else(
2410 || {
2411 decoder.invalidate();
2412 generation_one()
2413 },
2414 Ok,
2415 )
2416}
2417
2418fn take_attach_bound(decoder: &mut Decoder<'_>) -> Result<r::AttachBound, CodecError> {
2419 let conversation_id = decoder.take_u64()?;
2420 let token = AttachAttemptToken::new(decoder.take_fixed::<TOKEN_LEN>()?);
2421 let participant_id = decoder.take_u64()?;
2422 let request_generation = required_generation_option(decoder)?;
2423 let capability_generation = decoder.take_generation()?;
2424 let attach_secret = AttachSecret::new(decoder.take_fixed::<SECRET_LEN>()?);
2425 let origin_binding_epoch = decoder.take_binding_epoch()?;
2426 let persisted_cursor = decoder.take_u64()?;
2427 let accepted_marker_delivery_seq = decoder.take_option_u64()?;
2428 let receipt_expires_at = decoder.take_u128()?;
2429 let provenance_expires_at = decoder.take_u128()?;
2430
2431 let value = match accepted_marker_delivery_seq {
2432 Some(marker) if persisted_cursor == marker => r::AttachBound::fenced(
2433 conversation_id,
2434 token,
2435 participant_id,
2436 request_generation,
2437 attach_secret,
2438 origin_binding_epoch,
2439 marker,
2440 receipt_expires_at,
2441 provenance_expires_at,
2442 ),
2443 Some(_) => None,
2444 None => r::AttachBound::ordinary(
2445 conversation_id,
2446 token,
2447 participant_id,
2448 request_generation,
2449 attach_secret,
2450 origin_binding_epoch,
2451 persisted_cursor,
2452 receipt_expires_at,
2453 provenance_expires_at,
2454 ),
2455 }
2456 .ok_or_else(|| decode_error(t::DecodeClass::InvalidField))?;
2457 if value.capability_generation() == capability_generation {
2458 Ok(value)
2459 } else {
2460 Err(decode_error(t::DecodeClass::InvalidField))
2461 }
2462}
2463
2464fn take_receipt_replay(
2465 origin: t::ClientDiscriminant,
2466 decoder: &mut Decoder<'_>,
2467) -> Result<r::ReceiptReplay, CodecError> {
2468 let conversation_id = decoder.take_u64()?;
2469 let token = decoder.take_fixed::<TOKEN_LEN>()?;
2470 let participant_id = decoder.take_u64()?;
2471 match origin {
2472 t::ClientDiscriminant::EnrollmentRequest => {
2473 let request_generation = decoder.take_option_generation()?;
2474 let capability_generation = decoder.take_generation()?;
2475 let attach_secret = AttachSecret::new(decoder.take_fixed::<SECRET_LEN>()?);
2476 let origin_binding_epoch = decoder.take_binding_epoch()?;
2477 let persisted_cursor = decoder.take_u64()?;
2478 let accepted_marker_delivery_seq = decoder.take_option_u64()?;
2479 let receipt_expires_at = decoder.take_u128()?;
2480 let provenance_expires_at = decoder.take_u128()?;
2481 let value = r::EnrollBound::new(
2482 conversation_id,
2483 EnrollmentToken::new(token),
2484 participant_id,
2485 attach_secret,
2486 origin_binding_epoch,
2487 receipt_expires_at,
2488 provenance_expires_at,
2489 )
2490 .ok_or_else(|| decode_error(t::DecodeClass::InvalidField))?;
2491 if request_generation.is_none()
2492 && capability_generation == Generation::ONE
2493 && persisted_cursor == 0
2494 && accepted_marker_delivery_seq.is_none()
2495 {
2496 Ok(r::ReceiptReplay::Enrollment(value))
2497 } else {
2498 Err(decode_error(t::DecodeClass::InvalidField))
2499 }
2500 }
2501 t::ClientDiscriminant::CredentialAttachRequest => {
2502 let request_generation = required_generation_option(decoder)?;
2503 let capability_generation = decoder.take_generation()?;
2504 let attach_secret = AttachSecret::new(decoder.take_fixed::<SECRET_LEN>()?);
2505 let origin_binding_epoch = decoder.take_binding_epoch()?;
2506 let persisted_cursor = decoder.take_u64()?;
2507 let accepted_marker_delivery_seq = decoder.take_option_u64()?;
2508 let receipt_expires_at = decoder.take_u128()?;
2509 let provenance_expires_at = decoder.take_u128()?;
2510 let attach_token = AttachAttemptToken::new(token);
2511 let value = match accepted_marker_delivery_seq {
2512 Some(marker) if persisted_cursor == marker => r::AttachBound::fenced(
2513 conversation_id,
2514 attach_token,
2515 participant_id,
2516 request_generation,
2517 attach_secret,
2518 origin_binding_epoch,
2519 marker,
2520 receipt_expires_at,
2521 provenance_expires_at,
2522 ),
2523 Some(_) => None,
2524 None => r::AttachBound::ordinary(
2525 conversation_id,
2526 attach_token,
2527 participant_id,
2528 request_generation,
2529 attach_secret,
2530 origin_binding_epoch,
2531 persisted_cursor,
2532 receipt_expires_at,
2533 provenance_expires_at,
2534 ),
2535 }
2536 .ok_or_else(|| decode_error(t::DecodeClass::InvalidField))?;
2537 if value.capability_generation() == capability_generation {
2538 Ok(r::ReceiptReplay::CredentialAttach(value))
2539 } else {
2540 Err(decode_error(t::DecodeClass::InvalidField))
2541 }
2542 }
2543 _ => Err(decode_error(t::DecodeClass::InvalidField)),
2544 }
2545}
2546
2547fn take_ack_no_op(
2548 origin: t::ClientDiscriminant,
2549 decoder: &mut Decoder<'_>,
2550) -> Result<r::AckNoOp, CodecError> {
2551 match origin {
2552 t::ClientDiscriminant::ParticipantAck => {
2553 let request = take_participant_ack(decoder)?;
2554 if decoder.take_u64()? != request.through_seq {
2555 decoder.invalidate();
2556 }
2557 Ok(r::AckNoOp::participant_ack(request))
2558 }
2559 t::ClientDiscriminant::MarkerAck => {
2560 let request = take_marker_ack(decoder)?;
2561 if decoder.take_u64()? != request.marker_delivery_seq {
2562 decoder.invalidate();
2563 }
2564 Ok(r::AckNoOp::marker_ack(request))
2565 }
2566 _ => Err(decode_error(t::DecodeClass::InvalidField)),
2567 }
2568}
2569
2570fn take_recovery_accepted(
2571 decoder: &mut Decoder<'_>,
2572) -> Result<r::ObserverRecoveryAccepted, CodecError> {
2573 let count = decoder.take_u64()?;
2574 let required = u128::from(count) * RECOVERY_STATUS_LEN;
2575 let remaining = u64::try_from(decoder.remaining())
2576 .map_err(|_| decode_error(t::DecodeClass::MissingRequiredField))?;
2577 if required > u128::from(remaining) {
2578 return Err(decode_error(t::DecodeClass::MissingRequiredField));
2579 }
2580 let count_usize: usize = count
2581 .try_into()
2582 .map_err(|_| decode_error(t::DecodeClass::MissingRequiredField))?;
2583
2584 let mut statuses = Vec::new();
2585 for _ in 0..count_usize {
2586 statuses.push(r::ObserverProgressStatus {
2587 conversation_id: decoder.take_u64()?,
2588 refused_epoch: decoder.take_u64()?,
2589 current_observer_progress: decoder.take_u64()?,
2590 armed: decoder.take_bool()?,
2591 progressed: decoder.take_bool()?,
2592 });
2593 }
2594 Ok(r::ObserverRecoveryAccepted { statuses })
2595}
2596
2597fn require_zero_recovery_count(decoder: &mut Decoder<'_>) -> Result<(), CodecError> {
2598 if decoder.take_u64()? == 0 {
2599 Ok(())
2600 } else {
2601 Err(decode_error(t::DecodeClass::InvalidField))
2602 }
2603}
2604
2605fn take_invalid_observer_epoch(
2606 decoder: &mut Decoder<'_>,
2607) -> Result<r::InvalidObserverEpoch, CodecError> {
2608 let reason = take_tag::<t::InvalidObserverEpochReason>(decoder)?;
2609 let conversation_id = decoder.take_u64()?;
2610 let presented_epoch = decoder.take_u64()?;
2611 let current = decoder.take_option_u64()?;
2612 match reason {
2613 t::InvalidObserverEpochReason::ConversationUnknown => {
2614 if current.is_some() {
2615 decoder.invalidate();
2616 }
2617 Ok(r::InvalidObserverEpoch::ConversationUnknown {
2618 conversation_id,
2619 presented_epoch,
2620 })
2621 }
2622 t::InvalidObserverEpochReason::EpochAhead => {
2623 let current_observer_progress = current.unwrap_or_else(|| {
2624 decoder.invalidate();
2625 0
2626 });
2627 Ok(r::InvalidObserverEpoch::EpochAhead {
2628 conversation_id,
2629 presented_epoch,
2630 current_observer_progress,
2631 })
2632 }
2633 }
2634}
2635
2636fn take_invalid_observer_epoch_list(
2637 decoder: &mut Decoder<'_>,
2638) -> Result<r::InvalidObserverEpochList, CodecError> {
2639 match take_tag::<t::InvalidObserverEpochListReason>(decoder)? {
2640 t::InvalidObserverEpochListReason::TooManyEntries => {
2641 Ok(r::InvalidObserverEpochList::TooManyEntries {
2642 presented_entries: decoder.take_u64()?,
2643 max_entries: decoder.take_u64()?,
2644 })
2645 }
2646 t::InvalidObserverEpochListReason::DuplicateConversation => {
2647 Ok(r::InvalidObserverEpochList::DuplicateConversation {
2648 conversation_id: decoder.take_u64()?,
2649 first_index: decoder.take_u64()?,
2650 duplicate_index: decoder.take_u64()?,
2651 })
2652 }
2653 }
2654}