Skip to main content

turnframe_runtime/orchestrator/
mod.rs

1//! The public facade: one turn, start to finish (spec §23, §29).
2//!
3//! [`Orchestrator::handle_turn`] runs the steps of §23 in order, and the order is the
4//! contract: accept (A, B), load and project cases (D, E), judge a card answer (C),
5//! issue tokens (F), understand (G), resolve and reduce (J, K), persist cards and
6//! pending commands (L), execute and commit once (M, N), re-project and raise the
7//! cards the cases now need (O, P), then answer and persist the turn as returned
8//! (Q..V). Every step that can fail leaves a phase marker, so [`crate::recover`]
9//! can tell an interrupted understanding from an interrupted commit.
10//!
11//! A turn that carries only a card answer is understood without a model call: its
12//! meaning is the stored option.
13
14mod builder;
15mod directory;
16mod session;
17
18use std::fmt;
19use std::sync::Arc;
20
21use async_trait::async_trait;
22use chrono::{DateTime, Utc};
23use indexmap::IndexMap;
24use turnframe_core::case::CaseKey;
25use turnframe_core::command::CommandBatch;
26use turnframe_core::error::{OrchestratorError, StoreError};
27use turnframe_core::flow::WorkflowRegistry;
28use turnframe_core::ids::{AccountId, TurnId};
29use turnframe_core::observe::{Observer, Signal, SignalLabels};
30use turnframe_core::policy::PolicySnapshot;
31use turnframe_core::replay::TurnPhase;
32use turnframe_core::response::{AssistantTurn, ResponseBlock};
33use turnframe_core::turn::{ActorContext, AttachmentSource, TurnInput};
34use turnframe_store::stores::Stores;
35use turnframe_understand::TurnUnderstander;
36
37pub use self::builder::{BuildError, OrchestratorBuilder};
38pub use self::directory::{CaseCandidate, CaseDirectory, StaticCaseDirectory};
39use crate::attachments::AttachmentCopy;
40use crate::compose::{Composer, CompositionInput};
41use crate::config::OrchestratorConfig;
42use crate::execute::CommandExecutor;
43use crate::interactions::InteractionEngine;
44use crate::planning::{PlannedTurn, SeededCase, SeededTurnPlanner, SharedPlanning, TurnPlanner};
45use crate::policy::PolicyEngine;
46use crate::recover::{Recovery, RecoveryAction};
47use crate::reduce::NoticeCopy;
48use crate::resolve::CaseIdFactory;
49use crate::stream::{TurnPublisher, TurnStream};
50
51/// Source of the runtime's own clock, so a turn is testable without waiting.
52pub trait TurnClock: Send + Sync + fmt::Debug {
53    /// Now.
54    fn now(&self) -> DateTime<Utc>;
55}
56
57/// The system clock.
58#[derive(Debug, Clone, Copy, Default)]
59pub struct SystemTurnClock;
60
61impl TurnClock for SystemTurnClock {
62    fn now(&self) -> DateTime<Utc> {
63        Utc::now()
64    }
65}
66
67/// A clock stopped at one instant, for tests and replays.
68#[derive(Debug, Clone, Copy)]
69pub struct FixedTurnClock(pub DateTime<Utc>);
70
71impl TurnClock for FixedTurnClock {
72    fn now(&self) -> DateTime<Utc> {
73        self.0
74    }
75}
76
77/// What a turn's writes imply on OTHER cases (spec §23 step M).
78///
79/// Every hook that compiles a command sees one case, so «a fact that becomes true of
80/// case A settles a question on case B» has only the application to say it. Asked
81/// once per turn, after reduction, with the batches the turn carries; what it returns
82/// executes in the same turn under the same journal and idempotency rules. It sees
83/// writes, never words, so a consequence holds whether the turn came from a sentence,
84/// a click or a replay; and it is not asked again about its own answer.
85#[async_trait]
86pub trait TurnConsequences: Send + Sync {
87    /// The commands `batches` imply on other cases, or none.
88    ///
89    /// # Errors
90    ///
91    /// [`StoreError`] when the application could not read what it needed. The turn
92    /// fails rather than executing half a rule.
93    async fn following(
94        &self,
95        actor: &ActorContext,
96        batches: &[CommandBatch<serde_json::Value>],
97    ) -> Result<Vec<CommandBatch<serde_json::Value>>, StoreError>;
98}
99
100/// The runtime that answers a turn (spec §29).
101pub struct Orchestrator {
102    workflows: Arc<WorkflowRegistry>,
103    stores: Stores,
104    directory: Arc<dyn CaseDirectory>,
105    consequences: Option<Arc<dyn TurnConsequences>>,
106    understander: Arc<dyn TurnUnderstander>,
107    composer: Composer,
108    executor: CommandExecutor,
109    interactions: InteractionEngine,
110    recovery: Recovery,
111    policy_engine: PolicyEngine,
112    policy: PolicySnapshot,
113    observer: Arc<dyn Observer>,
114    attachment_source: Option<Arc<dyn AttachmentSource>>,
115    clock: Arc<dyn TurnClock>,
116    case_ids: Arc<dyn CaseIdFactory>,
117    config: OrchestratorConfig,
118    notice_copy: NoticeCopy,
119    attachment_copy: AttachmentCopy,
120    trace: Option<Arc<dyn crate::trace::TurnTrace>>,
121}
122
123impl fmt::Debug for Orchestrator {
124    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
125        f.debug_struct("Orchestrator")
126            .field("workflows", &self.workflows.len())
127            .field("mode", &self.config.mode)
128            .finish_non_exhaustive()
129    }
130}
131
132impl Orchestrator {
133    /// Starts a builder.
134    #[must_use]
135    pub fn builder() -> OrchestratorBuilder {
136        OrchestratorBuilder::new()
137    }
138
139    /// The configuration in force.
140    #[must_use]
141    pub const fn config(&self) -> &OrchestratorConfig {
142        &self.config
143    }
144
145    /// The persistence layer, for a caller that inspects what a turn wrote.
146    #[must_use]
147    pub const fn stores(&self) -> &Stores {
148        &self.stores
149    }
150
151    /// The workflows this runtime was built with. An evaluation reads revisions back
152    /// through them; a second registry would be a second answer to «which executor
153    /// owns this workflow».
154    #[must_use]
155    pub const fn workflows(&self) -> &Arc<WorkflowRegistry> {
156        &self.workflows
157    }
158
159    /// The crash-recovery reader (spec §23.1).
160    #[must_use]
161    pub const fn recovery(&self) -> &Recovery {
162        &self.recovery
163    }
164
165    /// The interaction engine, for a caller that settles a card out of band.
166    #[must_use]
167    pub const fn interactions(&self) -> &InteractionEngine {
168        &self.interactions
169    }
170
171    /// The same runtime with everything that writes taken away: it runs a turn's
172    /// decision pipeline and cannot persist anything. See [`crate::planning`].
173    #[must_use]
174    pub fn planner(&self) -> TurnPlanner {
175        TurnPlanner::assemble(
176            self.workflows.read_only(),
177            self.stores.read_only(),
178            Arc::clone(&self.directory),
179            self.shared_planning(),
180        )
181    }
182
183    /// The same runtime planning from the state it is handed, with no persistence at
184    /// all: what turns a recorded corpus into a deterministic shadow corpus.
185    #[must_use]
186    pub fn seeded_planner(&self) -> SeededTurnPlanner {
187        crate::planning::seeded_from_parts(self.workflows.definitions(), self.shared_planning())
188    }
189
190    fn shared_planning(&self) -> SharedPlanning {
191        SharedPlanning {
192            understander: Arc::clone(&self.understander),
193            policy_engine: self.policy_engine.clone(),
194            policy: self.policy.clone(),
195            config: self.config.clone(),
196            clock: Arc::clone(&self.clock),
197            case_ids: Arc::clone(&self.case_ids),
198            observer: Arc::clone(&self.observer),
199        }
200    }
201
202    /// Runs a turn through resolution, reduction and policy, and returns before the
203    /// first side effect. Shorthand for `self.planner().plan(input)`.
204    ///
205    /// # Errors
206    ///
207    /// See [`TurnPlanner::plan`].
208    pub async fn plan_turn(&self, input: TurnInput) -> Result<PlannedTurn, OrchestratorError> {
209        self.planner().plan(input).await
210    }
211
212    /// The same, against the cases it is handed. Shorthand for
213    /// `self.seeded_planner().plan(input, cases)`.
214    ///
215    /// # Errors
216    ///
217    /// See [`SeededTurnPlanner::plan`].
218    pub async fn plan_turn_from(
219        &self,
220        input: TurnInput,
221        cases: Vec<SeededCase>,
222    ) -> Result<PlannedTurn, OrchestratorError> {
223        self.seeded_planner().plan(input, cases).await
224    }
225
226    /// Handles one turn end to end (spec §23).
227    ///
228    /// # Errors
229    ///
230    /// The [`OrchestratorError`] family. The phase marker says how far the turn got,
231    /// and a command that may have taken effect is in the journal.
232    pub async fn handle_turn(&self, input: TurnInput) -> Result<AssistantTurn, OrchestratorError> {
233        let publisher = TurnPublisher::null();
234        self.run(input, &publisher).await
235    }
236
237    /// Handles one turn, publishing its events through `sink` under the §18.5 gate:
238    /// understanding's steps as they happen, and nothing that states an outcome
239    /// before the commit.
240    ///
241    /// # Errors
242    ///
243    /// See [`Self::handle_turn`].
244    pub async fn handle_turn_streaming(
245        &self,
246        input: TurnInput,
247        sink: Arc<dyn crate::stream::TurnSink>,
248    ) -> Result<AssistantTurn, OrchestratorError> {
249        let publisher = TurnPublisher::new(sink);
250        self.run(input, &publisher).await
251    }
252
253    /// Handles one turn on a background task and streams its events (spec §18.5).
254    #[must_use]
255    pub fn stream_turn(self: Arc<Self>, input: TurnInput) -> TurnStream {
256        let (stream, sink) = TurnStream::channel();
257        let sink: Arc<dyn crate::stream::TurnSink> = Arc::new(sink);
258        tokio::spawn(async move {
259            let publisher = TurnPublisher::new(sink);
260            match self.run(input, &publisher).await {
261                Ok(turn) => publisher.completed(&turn),
262                Err(error) => publisher.failed(error_code(&error)),
263            }
264        });
265        stream
266    }
267
268    /// Decides what an unfinished turn needs, without acting on it (spec §23.1).
269    ///
270    /// # Errors
271    ///
272    /// [`OrchestratorError::Store`] when the phase marker or the journal could not be
273    /// read.
274    pub async fn plan_recovery(
275        &self,
276        account: &AccountId,
277        turn_id: &TurnId,
278    ) -> Result<RecoveryAction, OrchestratorError> {
279        self.recovery.decide(account, turn_id).await
280    }
281
282    /// Acts on that decision, and finishes the turn where it can (spec §23.1): pending
283    /// commands resume by idempotency key (I14), a committed turn regenerates only its
284    /// answer, an unknown external outcome is handed back (§16.5).
285    ///
286    /// # Errors
287    ///
288    /// [`OrchestratorError::Store`] when the journal, the conversation or the commit
289    /// could not be reached, and whatever composition returns.
290    pub async fn resume_turn(
291        &self,
292        account: &AccountId,
293        turn_id: &TurnId,
294    ) -> Result<ResumeOutcome, OrchestratorError> {
295        match self.recovery.decide(account, turn_id).await? {
296            RecoveryAction::Nothing { phase } => Ok(ResumeOutcome::Nothing { phase }),
297            RecoveryAction::RestartInterpretation => Ok(ResumeOutcome::Restartable),
298            RecoveryAction::ReconcileExternal { attempts, .. } => {
299                Ok(ResumeOutcome::AwaitingReconciliation { attempts })
300            }
301            RecoveryAction::RegenerateResponse {
302                events,
303                answer_tasks,
304            } => {
305                let turn = self
306                    .regenerate(account, turn_id, &events, &answer_tasks)
307                    .await?;
308                Ok(match turn {
309                    Some(turn) => ResumeOutcome::Regenerated {
310                        turn: Box::new(turn),
311                    },
312                    None => ResumeOutcome::Nothing {
313                        phase: TurnPhase::Delivered,
314                    },
315                })
316            }
317            RecoveryAction::ResumeCommands { entries } => {
318                let stored = self.recovery.stored_turn(account, turn_id).await?;
319                let batches = resume_batches(&stored.user.input.actor, *turn_id, &entries);
320                let execution = self
321                    .executor
322                    .execute(account, &batches, self.clock.now())
323                    .await?;
324                let bundle = execution
325                    .bundle()
326                    .with_turn_phase(*turn_id, TurnPhase::Committed);
327                self.executor.commit(account, bundle).await?;
328                let events =
329                    self.recovery.decide(account, turn_id).await.ok().and_then(
330                        |action| match action {
331                            RecoveryAction::RegenerateResponse {
332                                events,
333                                answer_tasks,
334                            } => Some((events, answer_tasks)),
335                            _ => None,
336                        },
337                    );
338                let turn = match events {
339                    Some((events, answer_tasks)) => {
340                        self.regenerate(account, turn_id, &events, &answer_tasks)
341                            .await?
342                    }
343                    None => None,
344                };
345                Ok(ResumeOutcome::Resumed {
346                    outcomes: execution.outcomes.clone(),
347                    turn: turn.map(Box::new),
348                })
349            }
350        }
351    }
352
353    /// Rebuilds the answer of a committed turn from the ledger, and persists it.
354    /// `None` when the turn already has one.
355    async fn regenerate(
356        &self,
357        account: &AccountId,
358        turn_id: &TurnId,
359        events: &[turnframe_store::events::StoredEvent],
360        answer_tasks: &[turnframe_core::reduce::AnswerTask],
361    ) -> Result<Option<AssistantTurn>, OrchestratorError> {
362        let stored = self.recovery.stored_turn(account, turn_id).await?;
363        if stored.assistant.is_some() {
364            return Ok(None);
365        }
366        // Grouped by the store, which carries a redaction through as one.
367        let groups = turnframe_store::events::group_for_receipts(events);
368        let interactions = self
369            .interactions
370            .open_for_conversation(account, &stored.user.input.conversation_id)
371            .await
372            .unwrap_or_default();
373        let input = CompositionInput::new(&stored.user.input)
374            .with_answer_tasks(answer_tasks)
375            .with_ledger(&groups)
376            .with_interactions(&interactions);
377        let composition = self.composer.compose(input).await?;
378        self.stores
379            .conversations()
380            .append_assistant_turn(account, composition.turn.clone())
381            .await
382            .map_err(OrchestratorError::Store)?;
383        self.stores
384            .conversations()
385            .set_turn_phase(account, turn_id, TurnPhase::Delivered)
386            .await
387            .map_err(OrchestratorError::Store)?;
388        Ok(Some(composition.turn))
389    }
390
391    async fn run(
392        &self,
393        input: TurnInput,
394        publisher: &TurnPublisher,
395    ) -> Result<AssistantTurn, OrchestratorError> {
396        let labels =
397            SignalLabels::none().with_effort(input.effort.unwrap_or(self.config.effort.default));
398        self.observer
399            .observe_labeled(&Signal::TurnReceived, &labels);
400        let stage = crate::signals::Stage::enter();
401        // A turn's state is large, and every caller awaits it: it lives on the heap.
402        let outcome = Box::pin(session::Session::new(self, input, publisher).run()).await;
403        stage.observe(self.observer.as_ref(), Signal::TurnDuration, &labels);
404        match outcome {
405            Ok(turn) => {
406                self.observer
407                    .observe_labeled(&Signal::TurnCompleted, &labels);
408                Ok(turn)
409            }
410            Err(error) => {
411                self.observer.observe_labeled(
412                    &Signal::TurnFailed,
413                    &labels.with_error_code(error_code(&error)),
414                );
415                Err(error)
416            }
417        }
418    }
419}
420
421/// What acting on a recovery decision achieved (spec §23.1).
422#[derive(Debug, Clone)]
423#[non_exhaustive]
424pub enum ResumeOutcome {
425    /// The turn was already finished.
426    Nothing {
427        /// The phase it is in.
428        phase: TurnPhase,
429    },
430    /// Nothing was ever admitted, so the turn may simply be submitted again.
431    Restartable,
432    /// Pending commands were resumed by idempotency key.
433    Resumed {
434        /// What each of them ended up doing.
435        outcomes: Vec<turnframe_core::replay::CommandOutcomeRecord>,
436        /// The answer, when the resumed turn could be finished as well.
437        turn: Option<Box<AssistantTurn>>,
438    },
439    /// The effects were already committed; only the answer was rebuilt.
440    Regenerated {
441        /// The answer, now persisted.
442        turn: Box<AssistantTurn>,
443    },
444    /// An external effect may or may not have happened; only the application can
445    /// settle it (§16.5, I15).
446    AwaitingReconciliation {
447        /// The attempts to settle.
448        attempts: Vec<turnframe_core::ids::AttemptId>,
449    },
450}
451
452/// Rebuilds the batches of a set of journal entries, one per case, so the executor
453/// sees the very commands that were admitted (spec §23.1).
454fn resume_batches(
455    actor: &ActorContext,
456    turn_id: TurnId,
457    entries: &[turnframe_store::journal::CommandJournalEntry],
458) -> Vec<CommandBatch<serde_json::Value>> {
459    let mut batches: IndexMap<CaseKey, CommandBatch<serde_json::Value>> = IndexMap::new();
460    for entry in entries {
461        let key = entry.case_ref.key();
462        let scope = turnframe_core::command::AtomicityScope::PerCase;
463        let batch_id = turnframe_core::ids::BatchId::derive(&turn_id, &key, &scope);
464        batches
465            .entry(key)
466            .or_insert_with(|| CommandBatch {
467                batch_id,
468                scope,
469                envelopes: Vec::new(),
470            })
471            .envelopes
472            .push(turnframe_core::command::CommandEnvelope {
473                command_id: entry.command_id,
474                turn_id,
475                actor: actor.clone(),
476                case_ref: entry.case_ref.clone(),
477                idempotency_key: entry.idempotency_key.clone(),
478                origin: entry.origin.clone(),
479                command: entry.command_payload.clone(),
480            });
481    }
482    batches.into_values().collect()
483}
484
485/// A stable code for an orchestrator failure, safe to put on the wire.
486#[must_use]
487pub fn error_code(error: &OrchestratorError) -> String {
488    use turnframe_core::error::ErrorClassification;
489    error.user_message_key().to_owned()
490}
491
492/// The model-authored text of a stored assistant turn, for the transcript.
493pub(crate) fn assistant_text(turn: &AssistantTurn) -> String {
494    reply_text(&turn.blocks)
495}
496
497/// The reply as the user read it: the transition, which carries the turn's answers, or the
498/// answers alone when no transition was written.
499fn reply_text(blocks: &[ResponseBlock]) -> String {
500    let said = |transitions: bool| {
501        blocks
502            .iter()
503            .filter_map(|block| match block {
504                ResponseBlock::Transition(transition) if transitions => {
505                    Some(transition.text.as_str())
506                }
507                ResponseBlock::Answer(answer) if !transitions => Some(answer.text.as_str()),
508                _ => None,
509            })
510            .collect::<Vec<_>>()
511            .join(" ")
512    };
513    let reply = said(true);
514    if reply.is_empty() { said(false) } else { reply }
515}
516
517#[cfg(test)]
518mod tests {
519    use super::*;
520    use turnframe_core::ids::BlockId;
521    use turnframe_core::plan::AnswerBasis;
522    use turnframe_core::response::{AnswerStatus, GeneratedAnswer, GeneratedTransition};
523
524    fn answer(text: &str) -> ResponseBlock {
525        ResponseBlock::Answer(GeneratedAnswer {
526            block_id: BlockId::from("answer:0"),
527            question_id: None,
528            text: text.to_owned(),
529            basis: AnswerBasis::CurrentCommittedState,
530            status: AnswerStatus::Answered,
531            facts_used: Vec::new(),
532            citations: Vec::new(),
533            enumerations: Vec::new(),
534        })
535    }
536
537    fn transition(text: &str) -> ResponseBlock {
538        ResponseBlock::Transition(GeneratedTransition {
539            block_id: BlockId::from("transition:0"),
540            text: text.to_owned(),
541            facts_used: Vec::new(),
542        })
543    }
544
545    #[test]
546    fn the_reply_is_read_back_once() {
547        let said = "The airline pays for it.";
548        assert_eq!(reply_text(&[answer(said), transition(said)]), said);
549        assert_eq!(reply_text(&[answer(said)]), said);
550    }
551}