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