1use std::sync::Arc;
13
14use aion_core::{Event, RunId, WorkflowId};
15use aion_store::EventStore;
16use aion_store::visibility::VisibilityStore;
17use aion_store::workloop::WorkloopStore;
18use chrono::Utc;
19
20use super::iteration::{self, WorkloopIterationClose};
21use super::service::{WorkloopService, window_context};
22use crate::durability::Recorder;
23use crate::error::EngineError;
24use crate::registry::{Registry, TerminalOutcome};
25
26#[derive(Clone)]
32pub struct IterationCloseContext {
33 pub workloop_store: Arc<dyn WorkloopStore>,
35 pub service: Arc<WorkloopService>,
37 pub store: Arc<dyn EventStore>,
39 pub visibility_store: Arc<dyn VisibilityStore>,
41 pub registry: Arc<Registry>,
43}
44
45pub async fn close_iteration(
59 context: &IterationCloseContext,
60 loop_id: &WorkflowId,
61 close: WorkloopIterationClose,
62) -> Result<RunId, EngineError> {
63 let record = context
64 .workloop_store
65 .get_workloop(loop_id)
66 .await
67 .map_err(EngineError::from)?
68 .ok_or_else(|| EngineError::InvalidState {
69 reason: format!("workflow {loop_id} is not a registered workloop"),
70 })?;
71 let window_seq = window_context(&record);
72 let spec = record.spec.clone();
73 let retention = record.spec.retention();
74
75 let carry_for_notify = close.carry.clone();
76 let close_for_recorder = close;
77 let outcome = with_loop_recorder(context, loop_id, move |recorder, history| {
78 let spec = spec.clone();
79 let close = close_for_recorder.clone();
80 let loop_id = loop_id.clone();
81 Box::pin(async move {
82 let run_id = active_run_id(history).ok_or_else(|| {
83 crate::durability::DurabilityError::HistoryShape {
84 reason: format!("workloop {loop_id} has no recorded generation"),
85 }
86 })?;
87 iteration::close_iteration(
88 recorder,
89 iteration::IterationContext {
90 history,
91 run_id: &run_id,
92 loop_id: &loop_id,
93 spec: &spec,
94 window_seq,
95 recorded_at: Utc::now(),
96 },
97 close,
98 )
99 .await
100 .map_err(|error| crate::durability::DurabilityError::HistoryShape {
101 reason: error.to_string(),
102 })
103 })
104 })
105 .await?;
106
107 if let Some(handle) = registry_handle(context, loop_id)? {
113 let closed_run = handle.run_id().clone();
114 handle.completion().notify(TerminalOutcome::ContinuedAsNew {
115 input: carry_for_notify,
116 workflow_type: None,
117 parent_run_id: closed_run.clone(),
118 });
119 context.registry.remove(loop_id, &closed_run)?;
120 }
121
122 let cutoff = retention_cutoff(Utc::now(), retention)?;
128 for state_record in outcome.records.clone() {
129 context
130 .workloop_store
131 .put_invariant_record(state_record, cutoff)
132 .await
133 .map_err(EngineError::from)?;
134 }
135
136 context
137 .service
138 .note_iteration_closed(loop_id, &outcome.samples)
139 .await
140 .map_err(EngineError::from)?;
141
142 Ok(outcome.next_run_id)
143}
144
145fn retention_cutoff(
164 now: chrono::DateTime<Utc>,
165 retention: std::time::Duration,
166) -> Result<chrono::DateTime<Utc>, EngineError> {
167 let window =
168 chrono::Duration::from_std(retention).map_err(|error| EngineError::InvalidState {
169 reason: format!(
170 "workloop retention window of {seconds}s cannot be expressed as a calendar \
171 duration ({error}), so no retention cutoff exists; refusing rather than \
172 pruning against a fallback that would delete every prior generation",
173 seconds = retention.as_secs()
174 ),
175 })?;
176 now.checked_sub_signed(window)
177 .ok_or_else(|| EngineError::InvalidState {
178 reason: format!(
179 "subtracting the declared workloop retention window of {seconds}s from \
180 {now} left the representable calendar range, so no retention cutoff exists",
181 seconds = retention.as_secs()
182 ),
183 })
184}
185
186fn registry_handle(
187 context: &IterationCloseContext,
188 loop_id: &WorkflowId,
189) -> Result<Option<crate::registry::WorkflowHandle>, EngineError> {
190 context.registry.sole_handle(loop_id)
195}
196
197async fn with_loop_recorder<T>(
202 context: &IterationCloseContext,
203 loop_id: &WorkflowId,
204 record: impl for<'a> FnOnce(
205 &'a mut Recorder,
206 &'a [Event],
207 ) -> std::pin::Pin<
208 Box<
209 dyn std::future::Future<Output = Result<T, crate::durability::DurabilityError>>
210 + Send
211 + 'a,
212 >,
213 >,
214) -> Result<T, EngineError> {
215 if let Some(handle) = registry_handle(context, loop_id)? {
216 let recorder = handle.recorder();
217 let mut recorder = recorder.lock().await;
218 let history = context.store.read_history(loop_id).await?;
219 let value = record(&mut recorder, &history).await?;
220 return Ok(value);
221 }
222 let history = context.store.read_history(loop_id).await?;
223 let head = history.iter().map(Event::seq).max().unwrap_or_default();
224 let mut recorder = Recorder::resume_at(loop_id.clone(), Arc::clone(&context.store), head);
225 if let Some(run_id) = active_run_id(&history) {
226 recorder = recorder.with_visibility(run_id, Arc::clone(&context.visibility_store));
227 }
228 let value = record(&mut recorder, &history).await?;
229 Ok(value)
230}
231
232pub(crate) fn active_run_id(history: &[Event]) -> Option<RunId> {
234 history.iter().rev().find_map(|event| match event {
235 Event::WorkflowStarted { run_id, .. } => Some(run_id.clone()),
236 _ => None,
237 })
238}
239
240#[cfg(test)]
241mod tests {
242 use std::time::Duration;
243
244 use chrono::TimeZone;
245
246 use super::retention_cutoff;
247 use crate::error::EngineError;
248
249 fn now() -> Result<chrono::DateTime<chrono::Utc>, Box<dyn std::error::Error>> {
250 chrono::Utc
251 .with_ymd_and_hms(2026, 8, 26, 12, 0, 0)
252 .single()
253 .ok_or_else(|| "test instant must be valid".into())
254 }
255
256 #[test]
259 fn a_declared_window_subtracts_to_its_own_past() -> Result<(), Box<dyn std::error::Error>> {
260 let now = now()?;
261 let cutoff = retention_cutoff(now, Duration::from_secs(14 * 86_400))?;
262 assert_eq!(cutoff, now - chrono::Duration::days(14));
263 Ok(())
264 }
265
266 #[test]
275 fn an_unrepresentable_window_refuses_instead_of_pruning_everything()
276 -> Result<(), Box<dyn std::error::Error>> {
277 let now = now()?;
278 let outcome = retention_cutoff(now, Duration::from_secs(u64::MAX / 1_000));
279 assert!(
280 !matches!(&outcome, Ok(cutoff) if *cutoff == now),
281 "an unrepresentable retention must never yield a cutoff of `now`: that prunes \
282 every prior generation, which is the opposite of what it declares"
283 );
284 let refusal = outcome
285 .err()
286 .ok_or("an unrepresentable retention window must be refused")?;
287 assert!(
288 matches!(&refusal, EngineError::InvalidState { reason }
289 if reason.contains("retention")),
290 "the refusal must name the retention window: {refusal}"
291 );
292 Ok(())
293 }
294
295 #[test]
299 fn a_window_that_underflows_the_calendar_refuses() {
300 let early = chrono::DateTime::<chrono::Utc>::MIN_UTC + chrono::Duration::days(1);
301 let outcome = retention_cutoff(early, Duration::from_secs(1_000 * 365 * 86_400));
302 assert!(
303 outcome.is_err(),
304 "subtracting past the representable range must refuse, not saturate: {outcome:?}"
305 );
306 }
307}