Skip to main content

meerkat_runtime/
delivery_inbox.rs

1//! Runtime-owned durable delivery inbox.
2//!
3//! `RuntimeDeliveryMachine` is the semantic authority for idempotent sequence
4//! assignment and ordered application. `RuntimeStore` implementations retain
5//! its exact state with CAS and atomically insert opaque inbox rows.
6
7use 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}