Skip to main content

turnframe_runtime/
recover.rs

1//! Crash recovery (spec §23.1).
2//!
3//! A turn writes a phase marker as it goes, and a crash leaves that marker
4//! behind. Recovery reads it, together with the command journal, and decides
5//! which of four things to do — and the order the decision is taken in is the
6//! safety property, because two of the four would be wrong if taken first:
7//!
8//! 1. **An unknown external outcome is reconciled, never retried.** A command
9//!    that timed out after transmission may have taken effect. Repeating it
10//!    would be the duplicate the whole library exists to prevent, so this case
11//!    is checked before anything else (§16.5, I15).
12//! 2. **Pending commands are resumed by idempotency key.** An entry in
13//!    `Pending` or `Executing` was admitted and may or may not have reached the
14//!    domain. Running it again is safe *only* through the key, which is why
15//!    recovery hands back the entries rather than the plan.
16//! 3. **A committed turn regenerates its response.** The effects are done; what
17//!    is missing is the answer. It is rebuilt from the committed events and the
18//!    stored plan, and nothing is executed.
19//! 4. **A turn with no journal entry restarts interpretation.** Nothing was
20//!    admitted, so nothing can have happened, and the turn may simply be
21//!    interpreted again.
22//!
23//! # Why the answer tasks come back from the plan
24//!
25//! §23.1 asks for the response to be regenerated "from events and stored answer
26//! tasks". The reduction is pure, so its questions are exactly the questions of
27//! the accepted plan the replay record already stores — no separate table is
28//! needed, and re-deriving them cannot drift from what the turn actually
29//! planned. Their basis is normalized to
30//! [`AnswerBasis::CurrentCommittedState`]: after commit, "what will be true
31//! after this turn" and "what is true now" are the same state, and the second
32//! is the one that can be read.
33
34use std::fmt;
35use std::sync::Arc;
36
37use turnframe_core::error::OrchestratorError;
38use turnframe_core::ids::{AccountId, AttemptId, EventId, TurnId};
39use turnframe_core::plan::AnswerBasis;
40use turnframe_core::reduce::{AnswerTask, SourcePolicy};
41use turnframe_core::replay::{ReplayRecord, TurnPhase};
42use turnframe_store::conversation::{
43    ConversationStore, RecoveryScope, StoredTurn, TurnPhaseMarker,
44};
45use turnframe_store::error::StoreError;
46use turnframe_store::events::{EventJournal, StoredEvent};
47use turnframe_store::journal::{CommandJournal, CommandJournalEntry};
48use turnframe_store::replay::ReplayStore;
49
50/// Maximum events one regeneration reads back per case.
51const MAX_REPLAYED_EVENTS: usize = 1024;
52
53/// What recovery decided to do with a turn (spec §23.1).
54///
55/// "Unfinished" is not quite the right word for what this covers, and the
56/// difference matters: a turn can be [`TurnPhase::Delivered`] — the user has a
57/// truthful answer, saying the request is being verified — while an external
58/// effect it started is still unsettled. Recovery therefore looks at the
59/// journal before it looks at the phase.
60#[derive(Debug, Clone, PartialEq)]
61#[non_exhaustive]
62pub enum RecoveryAction {
63    /// The turn is finished; there is nothing to do.
64    Nothing {
65        /// The phase it finished in.
66        phase: TurnPhase,
67    },
68    /// No command was ever admitted, so nothing can have happened and the turn
69    /// may be interpreted again from the start.
70    RestartInterpretation,
71    /// Commands were admitted and not settled. They are resumed **by
72    /// idempotency key**: handing the executor these entries lets it recognise
73    /// what already ran instead of running it twice (I14).
74    ResumeCommands {
75        /// The entries to resume, in admission order.
76        entries: Vec<CommandJournalEntry>,
77    },
78    /// Everything that was going to commit has committed. The response is
79    /// rebuilt from the ledger and the stored plan, and **nothing is
80    /// executed**.
81    RegenerateResponse {
82        /// The events the turn committed, in append order.
83        events: Vec<StoredEvent>,
84        /// The questions the turn planned, as answer tasks.
85        answer_tasks: Vec<AnswerTask>,
86    },
87    /// An external effect may or may not have happened. It is settled against
88    /// the remote system by attempt identifier, never repeated blindly (I15).
89    ReconcileExternal {
90        /// The attempts to settle, in the order they were made.
91        attempts: Vec<AttemptId>,
92        /// The journal entries that carry them.
93        entries: Vec<CommandJournalEntry>,
94    },
95}
96
97impl RecoveryAction {
98    /// Returns `true` when acting on this decision may cause a domain effect.
99    ///
100    /// Only [`Self::ResumeCommands`] can, and even then only through the
101    /// idempotency key.
102    #[must_use]
103    pub const fn may_cause_effects(&self) -> bool {
104        matches!(self, Self::ResumeCommands { .. })
105    }
106
107    /// Stable snake-case label, for metrics and logs.
108    #[must_use]
109    pub const fn as_str(&self) -> &'static str {
110        match self {
111            Self::Nothing { .. } => "nothing",
112            Self::RestartInterpretation => "restart_interpretation",
113            Self::ResumeCommands { .. } => "resume_commands",
114            Self::RegenerateResponse { .. } => "regenerate_response",
115            Self::ReconcileExternal { .. } => "reconcile_external",
116        }
117    }
118}
119
120/// Reads what a crashed turn left behind and decides what to do about it.
121#[derive(Clone)]
122pub struct Recovery {
123    conversations: Arc<dyn ConversationStore>,
124    journal: Arc<dyn CommandJournal>,
125    events: Arc<dyn EventJournal>,
126    replay: Arc<dyn ReplayStore>,
127}
128
129impl fmt::Debug for Recovery {
130    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
131        f.debug_struct("Recovery").finish_non_exhaustive()
132    }
133}
134
135impl Recovery {
136    /// Builds a recovery reader over the stores a turn writes to.
137    #[must_use]
138    pub fn new(
139        conversations: Arc<dyn ConversationStore>,
140        journal: Arc<dyn CommandJournal>,
141        events: Arc<dyn EventJournal>,
142        replay: Arc<dyn ReplayStore>,
143    ) -> Self {
144        Self {
145            conversations,
146            journal,
147            events,
148            replay,
149        }
150    }
151
152    /// The unfinished turns, oldest first, for a sweep to work through.
153    ///
154    /// # Errors
155    ///
156    /// [`OrchestratorError::Store`] when the sweep could not read.
157    pub async fn unfinished(
158        &self,
159        scope: RecoveryScope,
160        limit: usize,
161    ) -> Result<Vec<TurnPhaseMarker>, OrchestratorError> {
162        self.conversations
163            .list_unfinished_turns(scope, limit)
164            .await
165            .map_err(OrchestratorError::Store)
166    }
167
168    /// Decides what one turn needs (spec §23.1).
169    ///
170    /// # Errors
171    ///
172    /// [`OrchestratorError::Store`] when the phase marker or the journal could
173    /// not be read. A turn that does not exist for `account` is
174    /// [`StoreError::NotFound`], indistinguishable from another tenant's.
175    pub async fn decide(
176        &self,
177        account: &AccountId,
178        turn_id: &TurnId,
179    ) -> Result<RecoveryAction, OrchestratorError> {
180        let marker = self
181            .conversations
182            .turn_phase(account, turn_id)
183            .await
184            .map_err(OrchestratorError::Store)?;
185        let entries = self
186            .journal
187            .for_turn(account, turn_id)
188            .await
189            .map_err(OrchestratorError::Store)?;
190
191        // 1. Uncertainty first — before the phase is even consulted. A turn can
192        //    be delivered, and truthfully so ("it is being verified"), while an
193        //    effect it started is still unsettled. Letting the terminal phase
194        //    answer first would drop that attempt on the floor (§16.5, I15).
195        let unknown: Vec<CommandJournalEntry> = entries
196            .iter()
197            .filter(|entry| {
198                entry.status == turnframe_store::journal::CommandJournalStatus::OutcomeUnknown
199            })
200            .cloned()
201            .collect();
202        if !unknown.is_empty() {
203            let attempts = unknown
204                .iter()
205                .filter_map(|entry| match entry.result.as_ref() {
206                    Some(turnframe_store::journal::JournalOutcome::OutcomeUnknown {
207                        attempt_id,
208                        ..
209                    }) => Some(attempt_id.clone()),
210                    _ => None,
211                })
212                .collect();
213            return Ok(RecoveryAction::ReconcileExternal {
214                attempts,
215                entries: unknown,
216            });
217        }
218
219        if marker.phase.is_terminal() {
220            // Nothing outstanding, and the turn is finished.
221            return Ok(RecoveryAction::Nothing {
222                phase: marker.phase,
223            });
224        }
225
226        // 2. Admitted and unsettled: resume by key.
227        let pending: Vec<CommandJournalEntry> = entries
228            .iter()
229            .filter(|entry| entry.status.is_pending())
230            .cloned()
231            .collect();
232        if !pending.is_empty() {
233            return Ok(RecoveryAction::ResumeCommands { entries: pending });
234        }
235
236        // 3. Settled: the effects are done, the answer is not.
237        if !entries.is_empty() {
238            let record = self.replay.get(account, turn_id).await.ok();
239            let text = self
240                .conversations
241                .load_turn(account, turn_id)
242                .await
243                .ok()
244                .and_then(|turn| turn.user.input.text);
245            let event_ids = committed_event_ids(&entries);
246            let events = self
247                .events
248                .get_by_ids(account, &event_ids)
249                .await
250                .map_err(OrchestratorError::Store)?;
251            return Ok(RecoveryAction::RegenerateResponse {
252                events,
253                answer_tasks: answer_tasks_of(record.as_ref(), text.as_deref()),
254            });
255        }
256
257        // 4. Nothing was ever admitted.
258        Ok(RecoveryAction::RestartInterpretation)
259    }
260
261    /// The stored turn, for a caller that regenerates a response.
262    ///
263    /// # Errors
264    ///
265    /// [`OrchestratorError::Store`]; a turn of another tenant is
266    /// [`StoreError::NotFound`].
267    pub async fn stored_turn(
268        &self,
269        account: &AccountId,
270        turn_id: &TurnId,
271    ) -> Result<StoredTurn, OrchestratorError> {
272        self.conversations
273            .load_turn(account, turn_id)
274            .await
275            .map_err(OrchestratorError::Store)
276    }
277
278    /// The replay record of a turn, when one was written.
279    ///
280    /// # Errors
281    ///
282    /// [`OrchestratorError::Store`] for anything but a missing record, which is
283    /// reported as `Ok(None)`: a turn that crashed before its first record is a
284    /// normal thing to find.
285    pub async fn record(
286        &self,
287        account: &AccountId,
288        turn_id: &TurnId,
289    ) -> Result<Option<ReplayRecord>, OrchestratorError> {
290        match self.replay.get(account, turn_id).await {
291            Ok(record) => Ok(Some(record)),
292            Err(StoreError::NotFound) => Ok(None),
293            Err(error) => Err(OrchestratorError::Store(error)),
294        }
295    }
296
297    /// Every event of one case since a revision, for a regeneration that needs
298    /// more than the turn's own events.
299    ///
300    /// # Errors
301    ///
302    /// [`OrchestratorError::Store`].
303    pub async fn case_events(
304        &self,
305        account: &AccountId,
306        case_key: &turnframe_core::case::CaseKey,
307        since: turnframe_core::ids::CaseRevision,
308    ) -> Result<Vec<StoredEvent>, OrchestratorError> {
309        self.events
310            .list_since(account, case_key, since, MAX_REPLAYED_EVENTS)
311            .await
312            .map_err(OrchestratorError::Store)
313    }
314}
315
316/// Every event identifier a turn's settled journal entries committed.
317#[must_use]
318pub fn committed_event_ids(entries: &[CommandJournalEntry]) -> Vec<EventId> {
319    entries
320        .iter()
321        .filter_map(|entry| match entry.result.as_ref() {
322            Some(turnframe_store::journal::JournalOutcome::Committed { event_ids, .. }) => {
323                Some(event_ids.clone())
324            }
325            _ => None,
326        })
327        .flatten()
328        .collect()
329}
330
331/// The answer tasks a stored replay record implies (spec §23.1), with the question words
332/// read from the turn's `text`. The basis becomes the committed state, because the
333/// turn's commands have already committed by the time this runs.
334#[must_use]
335pub fn answer_tasks_of(record: Option<&ReplayRecord>, text: Option<&str>) -> Vec<AnswerTask> {
336    let (Some(record), Some(text)) = (record, text) else {
337        return Vec::new();
338    };
339    let Some(understanding) = record.understanding.as_ref() else {
340        return Vec::new();
341    };
342    understanding
343        .questions
344        .iter()
345        .map(|question| AnswerTask {
346            question_id: turnframe_core::ids::QuestionId::from(question.unit.to_string()),
347            question: text
348                .get(question.words.start..question.words.end)
349                .unwrap_or_default()
350                .to_owned(),
351            basis: match question.basis {
352                AnswerBasis::GeneralDomainKnowledge => AnswerBasis::GeneralDomainKnowledge,
353                _ => AnswerBasis::CurrentCommittedState,
354            },
355            case_refs: record.loaded_cases.clone(),
356            proposed_diff_ref: None,
357            required_sources: if question.basis == AnswerBasis::GeneralDomainKnowledge {
358                SourcePolicy::AnySource
359            } else {
360                SourcePolicy::AuthoritativeOnly
361            },
362            // What the workflow accepted and offered then is not on the record.
363            enumerations: Vec::new(),
364            capabilities: Vec::new(),
365            continues_previous: question.continues_previous,
366            asked_at: Some(turnframe_core::reduce::TextSpan {
367                start_byte: question.words.start,
368                end_byte: question.words.end,
369            }),
370        })
371        .collect()
372}
373
374#[cfg(test)]
375mod tests {
376    use turnframe_core::ids::{AccountId, ConversationId};
377
378    use super::*;
379
380    #[test]
381    fn answer_tasks_come_back_from_the_stored_plan() {
382        let mut record = ReplayRecord::received(
383            TurnId::nil(),
384            ConversationId::nil(),
385            AccountId::from("acct"),
386            chrono::Utc::now(),
387        );
388        let text = "done. why the loyalty number?";
389        record.understanding = Some(turnframe_core::understanding::Understanding {
390            questions: vec![turnframe_core::understanding::UnderstoodQuestion {
391                unit: turnframe_core::understanding::UnitId(2),
392                words: turnframe_core::understanding::WordRange {
393                    first: 1,
394                    last: 4,
395                    start: 6,
396                    end: 29,
397                },
398                workflow: None,
399                record: None,
400                subjects: Vec::new(),
401                basis: AnswerBasis::CommittedStateAfterTurn,
402                topic: turnframe_core::understanding::QuestionTopic::default(),
403                continues_previous: false,
404            }],
405            ..turnframe_core::understanding::Understanding::default()
406        });
407        let tasks = answer_tasks_of(Some(&record), Some(text));
408        assert_eq!(tasks.len(), 1);
409        assert_eq!(tasks[0].question, "why the loyalty number?");
410        assert_eq!(
411            tasks[0].basis,
412            AnswerBasis::CurrentCommittedState,
413            "after commit, the state after the turn is the state now"
414        );
415    }
416
417    #[test]
418    fn a_record_that_was_never_written_implies_no_questions() {
419        assert!(answer_tasks_of(None, Some("x")).is_empty());
420    }
421
422    #[test]
423    fn only_resuming_commands_can_cause_an_effect() {
424        assert!(
425            RecoveryAction::ResumeCommands {
426                entries: Vec::new()
427            }
428            .may_cause_effects()
429        );
430        assert!(!RecoveryAction::RestartInterpretation.may_cause_effects());
431        assert!(
432            !RecoveryAction::RegenerateResponse {
433                events: Vec::new(),
434                answer_tasks: Vec::new(),
435            }
436            .may_cause_effects()
437        );
438        assert!(
439            !RecoveryAction::ReconcileExternal {
440                attempts: Vec::new(),
441                entries: Vec::new(),
442            }
443            .may_cause_effects()
444        );
445    }
446}