1use std::sync::Arc;
8
9use serde::{Deserialize, Serialize};
10
11use crate::identifiers::LogicalRuntimeId;
12use crate::store::{
13 RuntimeDeliveryAuthorityCasOutcome, RuntimeDeliveryAuthorityRecord, RuntimeDeliveryStoreRecord,
14 RuntimeStore, RuntimeStoreError,
15};
16
17pub mod dsl;
18
19const AUTHORITY_ENVELOPE_VERSION: u16 = 1;
20const SUBMISSION_ENVELOPE_VERSION: u16 = 1;
21const MAX_CAS_ATTEMPTS: usize = 32;
22
23#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize)]
24#[serde(transparent)]
25pub struct RuntimeDeliveryId(String);
26
27impl RuntimeDeliveryId {
28 pub fn new(value: impl Into<String>) -> Result<Self, RuntimeDeliveryError> {
29 Ok(Self(validate_component("delivery id", value.into())?))
30 }
31
32 pub fn as_str(&self) -> &str {
33 &self.0
34 }
35}
36
37impl std::fmt::Display for RuntimeDeliveryId {
38 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
39 f.write_str(self.as_str())
40 }
41}
42
43#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
44#[serde(rename_all = "snake_case")]
45#[non_exhaustive]
46pub enum RuntimeDeliveryKind {
47 JobTerminal,
48 JobNotification,
49}
50
51#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
52pub struct RuntimeDeliverySubmission {
53 delivery_id: RuntimeDeliveryId,
54 kind: RuntimeDeliveryKind,
55 source_id: String,
56 source_sequence: u64,
57 interaction_lineage_id: String,
58 payload: Vec<u8>,
59}
60
61impl RuntimeDeliverySubmission {
62 pub fn new(
63 delivery_id: RuntimeDeliveryId,
64 kind: RuntimeDeliveryKind,
65 source_id: impl Into<String>,
66 source_sequence: u64,
67 interaction_lineage_id: impl Into<String>,
68 payload: Vec<u8>,
69 ) -> Result<Self, RuntimeDeliveryError> {
70 let submission = Self {
71 delivery_id,
72 kind,
73 source_id: source_id.into(),
74 source_sequence,
75 interaction_lineage_id: interaction_lineage_id.into(),
76 payload,
77 };
78 submission.validate()?;
79 Ok(submission)
80 }
81
82 pub fn delivery_id(&self) -> &RuntimeDeliveryId {
83 &self.delivery_id
84 }
85
86 pub const fn kind(&self) -> RuntimeDeliveryKind {
87 self.kind
88 }
89
90 pub fn source_id(&self) -> &str {
91 &self.source_id
92 }
93
94 pub const fn source_sequence(&self) -> u64 {
95 self.source_sequence
96 }
97
98 pub fn interaction_lineage_id(&self) -> &str {
99 &self.interaction_lineage_id
100 }
101
102 pub fn payload(&self) -> &[u8] {
103 &self.payload
104 }
105
106 fn validate(&self) -> Result<(), RuntimeDeliveryError> {
107 validate_component_ref("delivery id", self.delivery_id.as_str())?;
108 validate_component_ref("delivery source id", &self.source_id)?;
109 validate_component_ref(
110 "delivery interaction lineage id",
111 &self.interaction_lineage_id,
112 )?;
113 if self.source_sequence == 0 {
114 return Err(RuntimeDeliveryError::InvalidInput(
115 "delivery source sequence must be positive".into(),
116 ));
117 }
118 if self.payload.is_empty() {
119 return Err(RuntimeDeliveryError::InvalidInput(
120 "delivery payload must not be empty".into(),
121 ));
122 }
123 Ok(())
124 }
125}
126
127#[derive(Debug, Clone, PartialEq, Eq)]
128pub struct RuntimeDeliveryReceipt {
129 pub delivery_id: RuntimeDeliveryId,
130 pub sequence: u64,
131 pub deduplicated: bool,
132}
133
134#[derive(Debug, Clone, PartialEq, Eq)]
135pub struct RuntimeDeliveryRecord {
136 pub sequence: u64,
137 pub submission: RuntimeDeliverySubmission,
138}
139
140#[derive(Debug, Clone, thiserror::Error)]
141#[non_exhaustive]
142pub enum RuntimeDeliveryError {
143 #[error("invalid runtime delivery: {0}")]
144 InvalidInput(String),
145 #[error("runtime delivery idempotency conflict for {0}")]
146 IdempotencyConflict(RuntimeDeliveryId),
147 #[error(
148 "runtime delivery {delivery_id} is out of order: next sequence is {expected}, received {actual}"
149 )]
150 OutOfOrder {
151 delivery_id: RuntimeDeliveryId,
152 expected: u64,
153 actual: u64,
154 },
155 #[error("runtime delivery persistence is corrupt: {0}")]
156 Corrupt(String),
157 #[error("runtime delivery authority rejected the transition: {0}")]
158 Authority(String),
159 #[error(transparent)]
160 Store(#[from] RuntimeStoreError),
161}
162
163#[derive(Debug, Clone, Serialize, Deserialize)]
164struct AuthorityEnvelope {
165 version: u16,
166 state: PersistedAuthorityState,
167}
168
169#[derive(Debug, Clone, Serialize, Deserialize)]
170struct PersistedAuthorityState {
171 delivery_ids: std::collections::BTreeSet<String>,
172 delivery_sequences: std::collections::BTreeMap<String, u64>,
173 delivery_source_sequences: std::collections::BTreeMap<String, u64>,
174 committed_sequences: std::collections::BTreeSet<u64>,
175 next_sequence: u64,
176 applied_cursor: u64,
177}
178
179impl From<&dsl::RuntimeDeliveryMachineState> for PersistedAuthorityState {
180 fn from(state: &dsl::RuntimeDeliveryMachineState) -> Self {
181 Self {
182 delivery_ids: state.delivery_ids.clone(),
183 delivery_sequences: state.delivery_sequences.clone(),
184 delivery_source_sequences: state.delivery_source_sequences.clone(),
185 committed_sequences: state.committed_sequences.clone(),
186 next_sequence: state.next_sequence,
187 applied_cursor: state.applied_cursor,
188 }
189 }
190}
191
192impl From<PersistedAuthorityState> for dsl::RuntimeDeliveryMachineState {
193 fn from(state: PersistedAuthorityState) -> Self {
194 Self {
195 lifecycle_phase: dsl::RuntimeDeliveryPhase::Active,
196 delivery_ids: state.delivery_ids,
197 delivery_sequences: state.delivery_sequences,
198 delivery_source_sequences: state.delivery_source_sequences,
199 committed_sequences: state.committed_sequences,
200 next_sequence: state.next_sequence,
201 applied_cursor: state.applied_cursor,
202 }
203 }
204}
205
206impl PersistedAuthorityState {
207 fn validate(&self) -> Result<(), RuntimeDeliveryError> {
208 for delivery_id in &self.delivery_ids {
209 validate_component_ref("persisted delivery id", delivery_id)
210 .map_err(|error| RuntimeDeliveryError::Corrupt(error.to_string()))?;
211 }
212 let sequence_keys = self
213 .delivery_sequences
214 .keys()
215 .cloned()
216 .collect::<std::collections::BTreeSet<_>>();
217 let source_keys = self
218 .delivery_source_sequences
219 .keys()
220 .cloned()
221 .collect::<std::collections::BTreeSet<_>>();
222 if sequence_keys != self.delivery_ids || source_keys != self.delivery_ids {
223 return Err(RuntimeDeliveryError::Corrupt(
224 "runtime delivery authority indexes disagree on delivery identity".into(),
225 ));
226 }
227 if self
228 .delivery_source_sequences
229 .values()
230 .any(|sequence| *sequence == 0)
231 {
232 return Err(RuntimeDeliveryError::Corrupt(
233 "runtime delivery authority contains a zero source sequence".into(),
234 ));
235 }
236 let mapped_sequences = self
237 .delivery_sequences
238 .values()
239 .copied()
240 .collect::<std::collections::BTreeSet<_>>();
241 if mapped_sequences != self.committed_sequences {
242 return Err(RuntimeDeliveryError::Corrupt(
243 "runtime delivery sequence index disagrees with committed sequence authority"
244 .into(),
245 ));
246 }
247 let committed_count = u64::try_from(self.committed_sequences.len()).map_err(|_| {
248 RuntimeDeliveryError::Corrupt(
249 "runtime delivery committed sequence count exceeds u64".into(),
250 )
251 })?;
252 let has_exact_bounds = if self.next_sequence == 0 {
253 self.committed_sequences.is_empty()
254 } else {
255 self.committed_sequences.first() == Some(&1)
256 && self.committed_sequences.last() == Some(&self.next_sequence)
257 };
258 if committed_count != self.next_sequence || !has_exact_bounds {
259 return Err(RuntimeDeliveryError::Corrupt(
260 "runtime delivery committed sequences are not contiguous through the high-water mark"
261 .into(),
262 ));
263 }
264 if self.applied_cursor > self.next_sequence {
265 return Err(RuntimeDeliveryError::Corrupt(
266 "runtime delivery applied cursor exceeds the committed high-water mark".into(),
267 ));
268 }
269 Ok(())
270 }
271}
272
273#[derive(Debug, Clone, Serialize, Deserialize)]
274struct SubmissionEnvelope {
275 version: u16,
276 submission: RuntimeDeliverySubmission,
277}
278
279#[derive(Clone)]
280pub struct RuntimeDeliveryInbox {
281 store: Arc<dyn RuntimeStore>,
282}
283
284impl std::fmt::Debug for RuntimeDeliveryInbox {
285 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
286 f.debug_struct("RuntimeDeliveryInbox")
287 .finish_non_exhaustive()
288 }
289}
290
291impl RuntimeDeliveryInbox {
292 pub fn new(store: Arc<dyn RuntimeStore>) -> Self {
293 Self { store }
294 }
295
296 pub async fn submit(
297 &self,
298 runtime_id: &LogicalRuntimeId,
299 submission: RuntimeDeliverySubmission,
300 ) -> Result<RuntimeDeliveryReceipt, RuntimeDeliveryError> {
301 for _ in 0..MAX_CAS_ATTEMPTS {
302 let observed = self
303 .store
304 .load_runtime_delivery_authority(runtime_id)
305 .await?;
306 let mut authority = decode_or_new_authority(observed.as_ref())?;
307 let delivery_id = submission.delivery_id.clone();
308 let transition = dsl::RuntimeDeliveryMachineMutator::apply(
309 &mut authority,
310 dsl::RuntimeDeliveryInput::CommitDelivery {
311 delivery_id: delivery_id.as_str().to_string(),
312 source_sequence: submission.source_sequence,
313 },
314 )
315 .map_err(|error| RuntimeDeliveryError::Authority(format!("{error:?}")))?;
316
317 let (sequence, deduplicated) = classify_commit_effects(
318 transition.effects(),
319 &delivery_id,
320 submission.source_sequence,
321 )?;
322 if deduplicated {
323 let stored = self
324 .store
325 .load_runtime_delivery_record(runtime_id, delivery_id.as_str())
326 .await?
327 .ok_or_else(|| {
328 RuntimeDeliveryError::Corrupt(format!(
329 "generated authority remembers {delivery_id}, but its inbox row is missing"
330 ))
331 })?;
332 let persisted = decode_submission(&stored)?;
333 if persisted != submission {
334 return Err(RuntimeDeliveryError::IdempotencyConflict(delivery_id));
335 }
336 if stored.sequence() != sequence {
337 return Err(RuntimeDeliveryError::Corrupt(format!(
338 "delivery {delivery_id} row sequence {} disagrees with generated sequence {sequence}",
339 stored.sequence()
340 )));
341 }
342 return Ok(RuntimeDeliveryReceipt {
343 delivery_id,
344 sequence,
345 deduplicated: true,
346 });
347 }
348
349 let next_revision = observed
350 .as_ref()
351 .map_or(Ok(1), |record| next_revision(record.revision()))?;
352 let replacement = RuntimeDeliveryAuthorityRecord::from_parts(
353 next_revision,
354 encode_authority(&authority)?,
355 );
356 let inserted = RuntimeDeliveryStoreRecord::from_parts(
357 delivery_id.as_str(),
358 sequence,
359 encode_submission(&submission)?,
360 );
361 match self
362 .store
363 .compare_and_swap_runtime_delivery_authority(
364 runtime_id,
365 observed
366 .as_ref()
367 .map(RuntimeDeliveryAuthorityRecord::revision),
368 replacement,
369 Some(inserted),
370 )
371 .await?
372 {
373 RuntimeDeliveryAuthorityCasOutcome::Applied(_) => {
374 return Ok(RuntimeDeliveryReceipt {
375 delivery_id,
376 sequence,
377 deduplicated: false,
378 });
379 }
380 RuntimeDeliveryAuthorityCasOutcome::Conflict(_) => continue,
381 }
382 }
383 Err(RuntimeDeliveryError::Store(RuntimeStoreError::WriteFailed(
384 format!("runtime delivery CAS did not converge after {MAX_CAS_ATTEMPTS} attempts"),
385 )))
386 }
387
388 pub async fn list_pending(
389 &self,
390 runtime_id: &LogicalRuntimeId,
391 limit: usize,
392 ) -> Result<Vec<RuntimeDeliveryRecord>, RuntimeDeliveryError> {
393 if limit == 0 {
394 return Ok(Vec::new());
395 }
396 let authority = self
397 .store
398 .load_runtime_delivery_authority(runtime_id)
399 .await?;
400 let authority = decode_or_new_authority(authority.as_ref())?;
401 let cursor = authority.state().applied_cursor;
402 let pending_count = authority.state().next_sequence - cursor;
403 let expected_count = pending_count.min(u64::try_from(limit).unwrap_or(u64::MAX));
404 let rows = self
405 .store
406 .list_runtime_delivery_records(runtime_id, cursor, limit)
407 .await?;
408 let mut expected_sequence = cursor.checked_add(1);
409 let records = rows
410 .into_iter()
411 .map(|row| {
412 let sequence = row.sequence();
413 if expected_sequence != Some(sequence) {
414 return Err(RuntimeDeliveryError::Corrupt(format!(
415 "runtime delivery inbox has a sequence gap after cursor {cursor}: expected {expected_sequence:?}, found {sequence}"
416 )));
417 }
418 expected_sequence = sequence.checked_add(1);
419 let submission = decode_submission(&row)?;
420 let expected = authority
421 .state()
422 .delivery_sequences
423 .get(submission.delivery_id.as_str())
424 .copied()
425 .ok_or_else(|| {
426 RuntimeDeliveryError::Corrupt(format!(
427 "inbox row {} has no generated delivery authority",
428 submission.delivery_id
429 ))
430 })?;
431 if expected != sequence {
432 return Err(RuntimeDeliveryError::Corrupt(format!(
433 "inbox row {} sequence {sequence} disagrees with generated sequence {expected}",
434 submission.delivery_id
435 )));
436 }
437 let expected_source_sequence = authority
438 .state()
439 .delivery_source_sequences
440 .get(submission.delivery_id.as_str())
441 .copied()
442 .ok_or_else(|| {
443 RuntimeDeliveryError::Corrupt(format!(
444 "inbox row {} has no generated source-sequence authority",
445 submission.delivery_id
446 ))
447 })?;
448 if expected_source_sequence != submission.source_sequence {
449 return Err(RuntimeDeliveryError::Corrupt(format!(
450 "inbox row {} source sequence {} disagrees with generated source sequence {expected_source_sequence}",
451 submission.delivery_id, submission.source_sequence
452 )));
453 }
454 Ok(RuntimeDeliveryRecord {
455 sequence,
456 submission,
457 })
458 })
459 .collect::<Result<Vec<_>, _>>()?;
460 if u64::try_from(records.len()).unwrap_or(u64::MAX) != expected_count {
461 return Err(RuntimeDeliveryError::Corrupt(format!(
462 "runtime delivery inbox contains {} rows after cursor {cursor}, but generated authority requires {expected_count}",
463 records.len()
464 )));
465 }
466 Ok(records)
467 }
468
469 pub async fn mark_applied(
470 &self,
471 runtime_id: &LogicalRuntimeId,
472 delivery_id: &RuntimeDeliveryId,
473 sequence: u64,
474 ) -> Result<u64, RuntimeDeliveryError> {
475 for _ in 0..MAX_CAS_ATTEMPTS {
476 let observed = self
477 .store
478 .load_runtime_delivery_authority(runtime_id)
479 .await?
480 .ok_or_else(|| {
481 RuntimeDeliveryError::Corrupt(format!(
482 "runtime {runtime_id} has no delivery authority"
483 ))
484 })?;
485 let mut authority = decode_authority(&observed)?;
486 let current = authority.state().applied_cursor;
487 let expected_sequence = authority
488 .state()
489 .delivery_sequences
490 .get(delivery_id.as_str())
491 .copied()
492 .ok_or_else(|| {
493 RuntimeDeliveryError::Corrupt(format!(
494 "generated authority has no delivery {delivery_id}"
495 ))
496 })?;
497 if expected_sequence != sequence {
498 return Err(RuntimeDeliveryError::Corrupt(format!(
499 "delivery {delivery_id} sequence {sequence} disagrees with generated sequence {expected_sequence}"
500 )));
501 }
502 if sequence > current && sequence - 1 != current {
503 return Err(RuntimeDeliveryError::OutOfOrder {
504 delivery_id: delivery_id.clone(),
505 expected: current.saturating_add(1),
506 actual: sequence,
507 });
508 }
509
510 let transition = dsl::RuntimeDeliveryMachineMutator::apply(
511 &mut authority,
512 dsl::RuntimeDeliveryInput::MarkDeliveryApplied {
513 delivery_id: delivery_id.as_str().to_string(),
514 delivery_sequence: sequence,
515 },
516 )
517 .map_err(|error| RuntimeDeliveryError::Authority(format!("{error:?}")))?;
518 classify_applied_effects(transition.effects(), delivery_id, sequence)?;
519 let applied_cursor = authority.state().applied_cursor;
520 if applied_cursor == current {
521 return Ok(applied_cursor);
522 }
523
524 let replacement = RuntimeDeliveryAuthorityRecord::from_parts(
525 next_revision(observed.revision())?,
526 encode_authority(&authority)?,
527 );
528 match self
529 .store
530 .compare_and_swap_runtime_delivery_authority(
531 runtime_id,
532 Some(observed.revision()),
533 replacement,
534 None,
535 )
536 .await?
537 {
538 RuntimeDeliveryAuthorityCasOutcome::Applied(_) => return Ok(applied_cursor),
539 RuntimeDeliveryAuthorityCasOutcome::Conflict(_) => continue,
540 }
541 }
542 Err(RuntimeDeliveryError::Store(RuntimeStoreError::WriteFailed(
543 format!(
544 "runtime delivery cursor CAS did not converge after {MAX_CAS_ATTEMPTS} attempts"
545 ),
546 )))
547 }
548
549 pub async fn applied_cursor(
550 &self,
551 runtime_id: &LogicalRuntimeId,
552 ) -> Result<u64, RuntimeDeliveryError> {
553 let observed = self
554 .store
555 .load_runtime_delivery_authority(runtime_id)
556 .await?;
557 Ok(decode_or_new_authority(observed.as_ref())?
558 .state()
559 .applied_cursor)
560 }
561}
562
563fn validate_component(label: &str, value: String) -> Result<String, RuntimeDeliveryError> {
564 validate_component_ref(label, &value)?;
565 Ok(value)
566}
567
568fn validate_component_ref(label: &str, value: &str) -> Result<(), RuntimeDeliveryError> {
569 let trimmed = value.trim();
570 if trimmed.is_empty() || trimmed != value || trimmed.chars().any(char::is_control) {
571 return Err(RuntimeDeliveryError::InvalidInput(format!(
572 "{label} must be non-empty, canonical, and contain no control characters"
573 )));
574 }
575 Ok(())
576}
577
578fn next_revision(current: u64) -> Result<u64, RuntimeDeliveryError> {
579 current.checked_add(1).ok_or_else(|| {
580 RuntimeDeliveryError::Store(RuntimeStoreError::WriteFailed(
581 "runtime delivery authority revision exhausted u64".into(),
582 ))
583 })
584}
585
586fn encode_authority(
587 authority: &dsl::RuntimeDeliveryMachineAuthority,
588) -> Result<Vec<u8>, RuntimeDeliveryError> {
589 serde_json::to_vec(&AuthorityEnvelope {
590 version: AUTHORITY_ENVELOPE_VERSION,
591 state: PersistedAuthorityState::from(authority.state()),
592 })
593 .map_err(|error| RuntimeDeliveryError::Corrupt(error.to_string()))
594}
595
596fn decode_authority(
597 record: &RuntimeDeliveryAuthorityRecord,
598) -> Result<dsl::RuntimeDeliveryMachineAuthority, RuntimeDeliveryError> {
599 let envelope: AuthorityEnvelope = serde_json::from_slice(record.state_json())
600 .map_err(|error| RuntimeDeliveryError::Corrupt(error.to_string()))?;
601 if envelope.version != AUTHORITY_ENVELOPE_VERSION {
602 return Err(RuntimeDeliveryError::Corrupt(format!(
603 "unsupported runtime delivery authority envelope version {}",
604 envelope.version
605 )));
606 }
607 envelope.state.validate()?;
608 dsl::RuntimeDeliveryMachineAuthority::recover_from_state(envelope.state.into())
609 .map_err(|error| RuntimeDeliveryError::Corrupt(format!("{error:?}")))
610}
611
612fn decode_or_new_authority(
613 record: Option<&RuntimeDeliveryAuthorityRecord>,
614) -> Result<dsl::RuntimeDeliveryMachineAuthority, RuntimeDeliveryError> {
615 match record {
616 Some(record) => decode_authority(record),
617 None => Ok(dsl::RuntimeDeliveryMachineAuthority::new()),
618 }
619}
620
621fn encode_submission(
622 submission: &RuntimeDeliverySubmission,
623) -> Result<Vec<u8>, RuntimeDeliveryError> {
624 serde_json::to_vec(&SubmissionEnvelope {
625 version: SUBMISSION_ENVELOPE_VERSION,
626 submission: submission.clone(),
627 })
628 .map_err(|error| RuntimeDeliveryError::Corrupt(error.to_string()))
629}
630
631fn decode_submission(
632 record: &RuntimeDeliveryStoreRecord,
633) -> Result<RuntimeDeliverySubmission, RuntimeDeliveryError> {
634 let envelope: SubmissionEnvelope = serde_json::from_slice(record.submission_json())
635 .map_err(|error| RuntimeDeliveryError::Corrupt(error.to_string()))?;
636 if envelope.version != SUBMISSION_ENVELOPE_VERSION {
637 return Err(RuntimeDeliveryError::Corrupt(format!(
638 "unsupported runtime delivery submission envelope version {}",
639 envelope.version
640 )));
641 }
642 if envelope.submission.delivery_id.as_str() != record.delivery_id() {
643 return Err(RuntimeDeliveryError::Corrupt(format!(
644 "runtime delivery row key {} disagrees with payload key {}",
645 record.delivery_id(),
646 envelope.submission.delivery_id
647 )));
648 }
649 envelope
650 .submission
651 .validate()
652 .map_err(|error| RuntimeDeliveryError::Corrupt(error.to_string()))?;
653 Ok(envelope.submission)
654}
655
656fn classify_commit_effects(
657 effects: &[dsl::RuntimeDeliveryEffect],
658 delivery_id: &RuntimeDeliveryId,
659 source_sequence: u64,
660) -> Result<(u64, bool), RuntimeDeliveryError> {
661 let mut matching = effects.iter().filter_map(|effect| match effect {
662 dsl::RuntimeDeliveryEffect::DeliveryCommitted {
663 delivery_id: emitted_id,
664 source_sequence: emitted_source_sequence,
665 delivery_sequence,
666 } if emitted_id == delivery_id.as_str() && *emitted_source_sequence == source_sequence => {
667 Some((*delivery_sequence, false))
668 }
669 dsl::RuntimeDeliveryEffect::DeliveryReused {
670 delivery_id: emitted_id,
671 source_sequence: emitted_source_sequence,
672 delivery_sequence,
673 } if emitted_id == delivery_id.as_str() && *emitted_source_sequence == source_sequence => {
674 Some((*delivery_sequence, true))
675 }
676 _ => None,
677 });
678 let first = matching.next().ok_or_else(|| {
679 RuntimeDeliveryError::Authority(
680 "generated commit emitted no matching delivery acknowledgement".into(),
681 )
682 })?;
683 if matching.next().is_some() {
684 return Err(RuntimeDeliveryError::Authority(
685 "generated commit emitted multiple delivery acknowledgements".into(),
686 ));
687 }
688 Ok(first)
689}
690
691fn classify_applied_effects(
692 effects: &[dsl::RuntimeDeliveryEffect],
693 delivery_id: &RuntimeDeliveryId,
694 sequence: u64,
695) -> Result<(), RuntimeDeliveryError> {
696 let count = effects
697 .iter()
698 .filter(|effect| {
699 matches!(
700 effect,
701 dsl::RuntimeDeliveryEffect::DeliveryApplied {
702 delivery_id: emitted_id,
703 delivery_sequence: emitted_sequence,
704 } if emitted_id == delivery_id.as_str() && *emitted_sequence == sequence
705 )
706 })
707 .count();
708 if count != 1 {
709 return Err(RuntimeDeliveryError::Authority(format!(
710 "generated apply emitted {count} matching acknowledgements"
711 )));
712 }
713 Ok(())
714}