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