1use super::{
2 ClientBindingState, ClientParticipantAggregate, ClientResponseCorrelation, correlation,
3};
4use crate::wire::{AttachBound, ReceiptReplay, ServerValue};
5
6#[derive(Clone, Copy, Debug, PartialEq, Eq)]
8pub enum ClientInboundRefusalReason {
9 AlreadyDead,
11 ForeignResponse,
13 DelayedResponse,
15 AmbiguousResponse,
17 MissingResponseAuthority,
19 LostAuthorityPending,
22}
23
24#[derive(Debug, PartialEq, Eq)]
26pub struct ClientInboundApplied {
27 aggregate: ClientParticipantAggregate,
28 value: ServerValue,
29}
30
31impl ClientInboundApplied {
32 #[must_use]
34 pub fn into_parts(self) -> (ClientParticipantAggregate, ServerValue) {
35 (self.aggregate, self.value)
36 }
37}
38
39#[derive(Debug, PartialEq, Eq)]
41pub struct ClientInboundRefusal {
42 aggregate: ClientParticipantAggregate,
43 value: ServerValue,
44 reason: ClientInboundRefusalReason,
45}
46
47impl ClientInboundRefusal {
48 #[must_use]
50 pub const fn reason(&self) -> ClientInboundRefusalReason {
51 self.reason
52 }
53
54 #[must_use]
56 pub fn into_parts(self) -> (ClientParticipantAggregate, ServerValue) {
57 (self.aggregate, self.value)
58 }
59}
60
61#[derive(Debug, PartialEq, Eq)]
63pub enum ClientInboundDecision {
64 Applied(ClientInboundApplied),
66 Refused(ClientInboundRefusal),
68}
69
70#[derive(Debug, PartialEq, Eq)]
72pub struct ClientCorrelatedInboundRefusal {
73 aggregate: ClientParticipantAggregate,
74 value: ServerValue,
75 correlation: ClientResponseCorrelation,
76 reason: ClientInboundRefusalReason,
77}
78
79impl ClientCorrelatedInboundRefusal {
80 #[must_use]
82 pub const fn reason(&self) -> ClientInboundRefusalReason {
83 self.reason
84 }
85
86 #[must_use]
88 pub fn into_parts(
89 self,
90 ) -> (
91 ClientParticipantAggregate,
92 ServerValue,
93 ClientResponseCorrelation,
94 ) {
95 (self.aggregate, self.value, self.correlation)
96 }
97}
98
99#[derive(Debug, PartialEq, Eq)]
101pub enum ClientCorrelatedInboundDecision {
102 Applied(ClientInboundApplied),
104 Refused(ClientCorrelatedInboundRefusal),
106}
107
108#[must_use]
110pub fn decide_inbound(
111 aggregate: ClientParticipantAggregate,
112 value: ServerValue,
113) -> ClientInboundDecision {
114 decide_inbound_inner(aggregate, value, false)
115}
116
117#[must_use]
119pub fn decide_correlated_inbound(
120 aggregate: ClientParticipantAggregate,
121 value: ServerValue,
122 correlation: ClientResponseCorrelation,
123) -> ClientCorrelatedInboundDecision {
124 let current_authority = aggregate.expected.as_ref().is_some_and(|expected| {
125 expected.issued && expected.authorization == correlation.authorization
126 });
127 if !current_authority {
128 return ClientCorrelatedInboundDecision::Refused(ClientCorrelatedInboundRefusal {
129 aggregate,
130 value,
131 correlation,
132 reason: ClientInboundRefusalReason::DelayedResponse,
133 });
134 }
135 match decide_inbound_inner(aggregate, value, true) {
136 ClientInboundDecision::Applied(applied) => {
137 ClientCorrelatedInboundDecision::Applied(applied)
138 }
139 ClientInboundDecision::Refused(refusal) => {
140 let reason = refusal.reason();
141 let (aggregate, value) = refusal.into_parts();
142 ClientCorrelatedInboundDecision::Refused(ClientCorrelatedInboundRefusal {
143 aggregate,
144 value,
145 correlation,
146 reason,
147 })
148 }
149 }
150}
151
152fn decide_inbound_inner(
153 mut aggregate: ClientParticipantAggregate,
154 value: ServerValue,
155 has_response_authority: bool,
156) -> ClientInboundDecision {
157 if aggregate.binding.is_left() {
158 return inbound_refusal(aggregate, value, ClientInboundRefusalReason::AlreadyDead);
159 }
160
161 if let Some(request) = correlation::participant_ack_request(&value) {
162 if aggregate.binding.matches_ack(request) {
163 return ClientInboundDecision::Applied(ClientInboundApplied { aggregate, value });
164 }
165 return inbound_refusal(
166 aggregate,
167 value,
168 ClientInboundRefusalReason::ForeignResponse,
169 );
170 }
171
172 if matches!(value, ServerValue::ParticipantTransportRejected(_)) {
173 return ClientInboundDecision::Applied(ClientInboundApplied { aggregate, value });
174 }
175
176 let Some(expected) = aggregate.expected.as_ref() else {
177 return inbound_refusal(
178 aggregate,
179 value,
180 ClientInboundRefusalReason::DelayedResponse,
181 );
182 };
183
184 if expected.lost.is_some() {
185 return inbound_refusal(
186 aggregate,
187 value,
188 ClientInboundRefusalReason::LostAuthorityPending,
189 );
190 }
191
192 if !has_response_authority {
193 return inbound_refusal(
194 aggregate,
195 value,
196 ClientInboundRefusalReason::MissingResponseAuthority,
197 );
198 }
199
200 if !aggregate.binding.accepts_request(&expected.request) {
201 return inbound_refusal(
202 aggregate,
203 value,
204 ClientInboundRefusalReason::ForeignResponse,
205 );
206 }
207
208 if !correlation::matches_request(&value, &expected.request) {
209 let same_request_class = value.originating_request()
210 == Some(expected.request.discriminant())
211 || matches!(
212 (&value, &expected.request),
213 (
214 ServerValue::ObserverRecoveryAccepted(_)
215 | ServerValue::InvalidObserverEpoch(_)
216 | ServerValue::InvalidObserverEpochList(_),
217 crate::wire::ClientRequest::ObserverRecovery(_)
218 )
219 );
220 let same_identity = correlation::same_identity(&value, &expected.request);
221 let reason = if same_request_class && same_identity {
222 if matches!(
223 expected.request,
224 crate::wire::ClientRequest::RecordAdmission(_)
225 ) {
226 ClientInboundRefusalReason::AmbiguousResponse
227 } else {
228 ClientInboundRefusalReason::DelayedResponse
229 }
230 } else {
231 ClientInboundRefusalReason::ForeignResponse
232 };
233 return inbound_refusal(aggregate, value, reason);
234 }
235
236 retire_expected_operation(&mut aggregate, &value);
237 ClientInboundDecision::Applied(ClientInboundApplied { aggregate, value })
238}
239
240fn retire_expected_operation(aggregate: &mut ClientParticipantAggregate, value: &ServerValue) {
267 let expected_detach = match aggregate
268 .expected
269 .as_ref()
270 .map(|expected| &expected.request)
271 {
272 Some(crate::wire::ClientRequest::Detach(request)) => Some(request.clone()),
273 _ => None,
274 };
275 aggregate.expected = None;
276 apply_correlated_value(aggregate, value);
277 if let Some(request) = expected_detach {
278 let envelope = crate::wire::DetachEnvelope {
279 conversation_id: request.conversation_id,
280 participant_id: request.participant_id,
281 capability_generation: request.capability_generation,
282 detach_attempt_token: request.detach_attempt_token,
283 };
284 aggregate
285 .detach_replay
286 .settle_refused_authority(&envelope, value);
287 debug_assert!(
288 !aggregate.detach_replay.is_active()
289 || aggregate.detach_replay.request() != Some(&envelope),
290 "retiring an expected detach must leave its own replay settled"
291 );
292 }
293 debug_assert!(
294 aggregate.expected.is_none(),
295 "the expected slot must be retired by this statement"
296 );
297}
298
299const fn inbound_refusal(
300 aggregate: ClientParticipantAggregate,
301 value: ServerValue,
302 reason: ClientInboundRefusalReason,
303) -> ClientInboundDecision {
304 ClientInboundDecision::Refused(ClientInboundRefusal {
305 aggregate,
306 value,
307 reason,
308 })
309}
310
311fn apply_correlated_value(aggregate: &mut ClientParticipantAggregate, value: &ServerValue) {
312 match value {
313 ServerValue::EnrollBound(value) => apply_enroll_bound(aggregate, value),
314 ServerValue::Bound(ReceiptReplay::Enrollment(value)) => {
315 apply_enroll_bound(aggregate, value);
316 }
317 ServerValue::AttachBound(value)
318 | ServerValue::Bound(ReceiptReplay::CredentialAttach(value)) => {
319 apply_attach_bound(aggregate, value);
320 aggregate.detach_replay.apply_attach(value);
321 }
322 ServerValue::UnboundReceipt(ReceiptReplay::CredentialAttach(value)) => {
323 apply_unbound_attach_receipt(aggregate, value);
324 aggregate.detach_replay.apply_attach(value);
325 }
326 ServerValue::DetachCommitted(value) => {
327 let attach_secret = match aggregate.binding {
328 ClientBindingState::Bound { attach_secret, .. }
329 | ClientBindingState::Detached { attach_secret, .. } => attach_secret,
330 ClientBindingState::Unbound | ClientBindingState::Left { .. } => return,
331 };
332 aggregate.binding = ClientBindingState::Detached {
333 conversation_id: value.conversation_id(),
334 participant_id: value.participant_id(),
335 generation: value.capability_generation(),
336 attach_secret,
337 };
338 aggregate.detach_replay.apply_detach_committed(value);
339 }
340 ServerValue::DetachInProgress(value) => {
341 aggregate.detach_replay.apply_detach_in_progress(value);
342 }
343 ServerValue::StaleAuthority(crate::wire::StaleAuthority::Detach(
344 crate::wire::DetachStaleAuthority::TerminalizedDetachCell(value),
345 )) => {
346 aggregate
347 .detach_replay
348 .apply_terminalized_detach_cell(value);
349 }
350 ServerValue::LeaveCommitted(value) => {
351 aggregate.binding = ClientBindingState::Left {
352 conversation_id: value.conversation_id(),
353 participant_id: value.participant_id(),
354 generation: value.retired_generation(),
355 };
356 aggregate.detach_replay.apply_leave(value);
357 }
358 ServerValue::Retired(value) => {
359 apply_retired(aggregate, value);
360 }
361 ServerValue::ParticipantTransportRejected(_)
362 | ServerValue::AttemptTokenBodyConflict(_)
363 | ServerValue::ConnectionConversationCapacityExceeded(_)
364 | ServerValue::ConnectionConversationBindingOccupied(_)
365 | ServerValue::ConversationOrderExhausted(_)
366 | ServerValue::ParticipantUnknown(_)
367 | ServerValue::NoBinding(_)
368 | ServerValue::StaleAuthority(_)
369 | ServerValue::MarkerClosureCapacityExceeded(_)
370 | ServerValue::EnrollmentKnown(_)
371 | ServerValue::ReceiptExpired(_)
372 | ServerValue::ReceiptCapacityExceeded(_)
373 | ServerValue::IdentityCapacityExceeded(_)
374 | ServerValue::ObserverBackpressure(_)
375 | ServerValue::ConversationSequenceExhausted(_)
376 | ServerValue::StaleOrUnknownReceipt(_)
377 | ServerValue::MarkerNotDelivered(_)
378 | ServerValue::MarkerMismatch(_)
379 | ServerValue::UnboundReceipt(ReceiptReplay::Enrollment(_))
384 | ServerValue::AckCommitted(_)
385 | ServerValue::AckNoOp(_)
386 | ServerValue::AckGap(_)
387 | ServerValue::AckRegression(_)
388 | ServerValue::MarkerAckCommitted(_)
389 | ServerValue::RecordCommitted(_)
390 | ServerValue::RecordTooLarge(_)
391 | ServerValue::ObserverRecoveryAccepted(_)
392 | ServerValue::InvalidObserverEpoch(_)
393 | ServerValue::InvalidObserverEpochList(_) => {}
394 }
395}
396
397const fn apply_enroll_bound(
398 aggregate: &mut ClientParticipantAggregate,
399 value: &crate::wire::EnrollBound,
400) {
401 aggregate.binding = ClientBindingState::Bound {
402 conversation_id: value.conversation_id(),
403 participant_id: value.participant_id(),
404 generation: value.capability_generation(),
405 attach_secret: value.attach_secret(),
406 binding_epoch: value.origin_binding_epoch(),
407 };
408}
409
410const fn apply_unbound_attach_receipt(
438 aggregate: &mut ClientParticipantAggregate,
439 value: &AttachBound,
440) {
441 aggregate.binding = ClientBindingState::Detached {
442 conversation_id: value.conversation_id(),
443 participant_id: value.participant_id(),
444 generation: value.capability_generation(),
445 attach_secret: value.attach_secret(),
446 };
447}
448
449const fn apply_attach_bound(aggregate: &mut ClientParticipantAggregate, value: &AttachBound) {
450 aggregate.binding = ClientBindingState::Bound {
451 conversation_id: value.conversation_id(),
452 participant_id: value.participant_id(),
453 generation: value.capability_generation(),
454 attach_secret: value.attach_secret(),
455 binding_epoch: value.origin_binding_epoch(),
456 };
457}
458
459fn apply_retired(aggregate: &mut ClientParticipantAggregate, value: &crate::wire::Retired) {
460 let (conversation_id, participant_id, generation) = match value {
461 crate::wire::Retired::Enrollment {
462 request,
463 participant_id,
464 retired_generation,
465 } => (
466 request.conversation_id,
467 *participant_id,
468 *retired_generation,
469 ),
470 crate::wire::Retired::Participant {
471 request,
472 retired_generation,
473 } => {
474 let (conversation_id, participant_id) = participant_reference_identity(request);
475 (conversation_id, participant_id, *retired_generation)
476 }
477 };
478 aggregate.binding = ClientBindingState::Left {
479 conversation_id,
480 participant_id,
481 generation,
482 };
483 aggregate
484 .detach_replay
485 .apply_retired(conversation_id, participant_id, generation);
486}
487
488const fn participant_reference_identity(
489 request: &crate::wire::ParticipantReferenceEnvelope,
490) -> (u64, u64) {
491 match request {
492 crate::wire::ParticipantReferenceEnvelope::CredentialAttach(value) => {
493 (value.conversation_id, value.participant_id)
494 }
495 crate::wire::ParticipantReferenceEnvelope::Detach(value) => {
496 (value.conversation_id, value.participant_id)
497 }
498 crate::wire::ParticipantReferenceEnvelope::ParticipantAck(value) => {
499 (value.conversation_id, value.participant_id)
500 }
501 crate::wire::ParticipantReferenceEnvelope::Leave(value) => {
502 (value.conversation_id, value.participant_id)
503 }
504 crate::wire::ParticipantReferenceEnvelope::MarkerAck(value) => {
505 (value.conversation_id, value.participant_id)
506 }
507 crate::wire::ParticipantReferenceEnvelope::RecordAdmission(value) => {
508 (value.conversation_id, value.participant_id)
509 }
510 }
511}