Skip to main content

made_core/entities/ceremony_instance/
step_claims.rs

1use time::OffsetDateTime;
2
3use crate::entities::{CeremonyDefinition, CeremonyInstance};
4use crate::error::DomainError;
5use crate::value_objects::{MaxParallel, StateExecution, StepAttempt, StepId, StepStatus};
6
7impl CeremonyInstance {
8    pub(super) fn resolved_step_is_claimable_at(
9        &self,
10        step_id: &StepId,
11        definition: &CeremonyDefinition,
12        now: OffsetDateTime,
13        host_ceiling: MaxParallel,
14    ) -> Result<bool, DomainError> {
15        self.require_definition(definition)?;
16        let state = definition
17            .state(&self.current_state)
18            .ok_or(DomainError::NotFound {
19                what: "ceremony_instance.current_state",
20            })?;
21        let steps = definition
22            .steps_for_state(&self.current_state)
23            .collect::<Vec<_>>();
24        if state.execution() == StateExecution::Sequential {
25            let first = steps.into_iter().find(|step| {
26                self.step_records
27                    .get(step.id())
28                    .is_some_and(|record| !record.status().is_success())
29            });
30            if first.is_none_or(|step| step.id() != step_id) {
31                return Ok(false);
32            }
33        } else {
34            let capacity =
35                usize::from(definition.max_parallel().effective_with(host_ceiling).get());
36            let live = steps
37                .iter()
38                .filter(|step| {
39                    self.step_records
40                        .get(step.id())
41                        .is_some_and(|record| record.has_live_lease_at(now))
42                })
43                .count();
44            if live >= capacity {
45                return Ok(false);
46            }
47        }
48        self.step_is_ready_without_role_at(step_id, definition, now)
49    }
50
51    /// Steps a caller may claim from the current state at one observed instant.
52    ///
53    /// Concurrent results contain every eligible alternative while capacity
54    /// remains. They are not truncated to the number of free slots: whichever
55    /// caller wins is persisted first and every optimistic retry recomputes the
56    /// set against that new state.
57    pub fn claimable_step_ids_at<'a>(
58        &self,
59        definition: &'a CeremonyDefinition,
60        now: OffsetDateTime,
61        host_ceiling: MaxParallel,
62    ) -> Result<Vec<&'a StepId>, DomainError> {
63        self.require_definition(definition)?;
64        let state = definition
65            .state(&self.current_state)
66            .ok_or(DomainError::NotFound {
67                what: "ceremony_instance.current_state",
68            })?;
69        let steps = definition
70            .steps_for_state(&self.current_state)
71            .collect::<Vec<_>>();
72        if state.execution() == StateExecution::Sequential {
73            let Some(step) = steps.into_iter().find(|step| {
74                self.step_records
75                    .get(step.id())
76                    .is_some_and(|record| !record.status().is_success())
77            }) else {
78                return Ok(Vec::new());
79            };
80            return Ok(if self.step_is_claimable_at(step.id(), definition, now)? {
81                vec![step.id()]
82            } else {
83                Vec::new()
84            });
85        }
86
87        let capacity = usize::from(definition.max_parallel().effective_with(host_ceiling).get());
88        let live = steps
89            .iter()
90            .filter(|step| {
91                self.step_records
92                    .get(step.id())
93                    .is_some_and(|record| record.has_live_lease_at(now))
94            })
95            .count();
96        if live >= capacity {
97            return Ok(Vec::new());
98        }
99        steps
100            .into_iter()
101            .filter_map(
102                |step| match self.step_is_claimable_at(step.id(), definition, now) {
103                    Ok(true) => Some(Ok(step.id())),
104                    Ok(false) => None,
105                    Err(error) => Some(Err(error)),
106                },
107            )
108            .collect()
109    }
110
111    fn step_is_claimable_at(
112        &self,
113        step_id: &StepId,
114        definition: &CeremonyDefinition,
115        now: OffsetDateTime,
116    ) -> Result<bool, DomainError> {
117        if self.resolve_step_role(definition, step_id, None).is_err() {
118            return Ok(false);
119        }
120        self.step_is_ready_without_role_at(step_id, definition, now)
121    }
122
123    fn step_is_ready_without_role_at(
124        &self,
125        step_id: &StepId,
126        definition: &CeremonyDefinition,
127        now: OffsetDateTime,
128    ) -> Result<bool, DomainError> {
129        let step = definition.step(step_id).ok_or(DomainError::NotFound {
130            what: "ceremony_instance.step",
131        })?;
132        let record = self
133            .step_records
134            .get(step_id)
135            .ok_or(DomainError::NotFound {
136                what: "ceremony_instance.step_record",
137            })?;
138        if !record.can_be_started_at(now) {
139            return Ok(false);
140        }
141        let attempt = if matches!(record.status(), StepStatus::Failed | StepStatus::InProgress) {
142            record.attempt().next()?
143        } else {
144            StepAttempt::new(record.attempt().get())?
145        };
146        Ok(step.retry_policy().allows_attempt(attempt))
147    }
148
149    #[must_use]
150    pub fn has_live_step_leases_at(
151        &self,
152        definition: &CeremonyDefinition,
153        now: OffsetDateTime,
154    ) -> bool {
155        definition.steps_for_state(&self.current_state).any(|step| {
156            self.step_records
157                .get(step.id())
158                .is_some_and(|record| record.has_live_lease_at(now))
159        })
160    }
161}