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            knowledge: self.composer.has_knowledge(),
200        }
201    }
202
203    /// Runs a turn through resolution, reduction and policy, and returns before the
204    /// first side effect. Shorthand for `self.planner().plan(input)`.
205    ///
206    /// # Errors
207    ///
208    /// See [`TurnPlanner::plan`].
209    pub async fn plan_turn(&self, input: TurnInput) -> Result<PlannedTurn, OrchestratorError> {
210        self.planner().plan(input).await
211    }
212
213    /// The same, against the cases it is handed. Shorthand for
214    /// `self.seeded_planner().plan(input, cases)`.
215    ///
216    /// # Errors
217    ///
218    /// See [`SeededTurnPlanner::plan`].
219    pub async fn plan_turn_from(
220        &self,
221        input: TurnInput,
222        cases: Vec<SeededCase>,
223    ) -> Result<PlannedTurn, OrchestratorError> {
224        self.seeded_planner().plan(input, cases).await
225    }
226
227    /// Handles one turn end to end (spec §23).
228    ///
229    /// # Errors
230    ///
231    /// The [`OrchestratorError`] family. The phase marker says how far the turn got,
232    /// and a command that may have taken effect is in the journal.
233    pub async fn handle_turn(&self, input: TurnInput) -> Result<AssistantTurn, OrchestratorError> {
234        let publisher = TurnPublisher::null();
235        self.run(input, &publisher).await
236    }
237
238    /// Handles one turn, publishing its events through `sink` under the §18.5 gate:
239    /// understanding's steps as they happen, and nothing that states an outcome
240    /// before the commit.
241    ///
242    /// # Errors
243    ///
244    /// See [`Self::handle_turn`].
245    pub async fn handle_turn_streaming(
246        &self,
247        input: TurnInput,
248        sink: Arc<dyn crate::stream::TurnSink>,
249    ) -> Result<AssistantTurn, OrchestratorError> {
250        let publisher = TurnPublisher::new(sink);
251        self.run(input, &publisher).await
252    }
253
254    /// Handles one turn on a background task and streams its events (spec §18.5).
255    #[must_use]
256    pub fn stream_turn(self: Arc<Self>, input: TurnInput) -> TurnStream {
257        let (stream, sink) = TurnStream::channel();
258        let sink: Arc<dyn crate::stream::TurnSink> = Arc::new(sink);
259        tokio::spawn(async move {
260            let publisher = TurnPublisher::new(sink);
261            match self.run(input, &publisher).await {
262                Ok(turn) => publisher.completed(&turn),
263                Err(error) => publisher.failed(error_code(&error)),
264            }
265        });
266        stream
267    }
268
269    /// Decides what an unfinished turn needs, without acting on it (spec §23.1).
270    ///
271    /// # Errors
272    ///
273    /// [`OrchestratorError::Store`] when the phase marker or the journal could not be
274    /// read.
275    pub async fn plan_recovery(
276        &self,
277        account: &AccountId,
278        turn_id: &TurnId,
279    ) -> Result<RecoveryAction, OrchestratorError> {
280        self.recovery.decide(account, turn_id).await
281    }
282
283    /// Acts on that decision, and finishes the turn where it can (spec §23.1): pending
284    /// commands resume by idempotency key (I14), a committed turn regenerates only its
285    /// answer, an unknown external outcome is handed back (§16.5).
286    ///
287    /// # Errors
288    ///
289    /// [`OrchestratorError::Store`] when the journal, the conversation or the commit
290    /// could not be reached, and whatever composition returns.
291    pub async fn resume_turn(
292        &self,
293        account: &AccountId,
294        turn_id: &TurnId,
295    ) -> Result<ResumeOutcome, OrchestratorError> {
296        match self.recovery.decide(account, turn_id).await? {
297            RecoveryAction::Nothing { phase } => Ok(ResumeOutcome::Nothing { phase }),
298            RecoveryAction::RestartInterpretation => Ok(ResumeOutcome::Restartable),
299            RecoveryAction::ReconcileExternal { attempts, .. } => {
300                Ok(ResumeOutcome::AwaitingReconciliation { attempts })
301            }
302            RecoveryAction::RegenerateResponse {
303                events,
304                answer_tasks,
305            } => {
306                let turn = self
307                    .regenerate(account, turn_id, &events, &answer_tasks)
308                    .await?;
309                Ok(match turn {
310                    Some(turn) => ResumeOutcome::Regenerated {
311                        turn: Box::new(turn),
312                    },
313                    None => ResumeOutcome::Nothing {
314                        phase: TurnPhase::Delivered,
315                    },
316                })
317            }
318            RecoveryAction::ResumeCommands { entries } => {
319                let stored = self.recovery.stored_turn(account, turn_id).await?;
320                let batches = resume_batches(&stored.user.input.actor, *turn_id, &entries);
321                let execution = self
322                    .executor
323                    .execute(account, &batches, self.clock.now())
324                    .await?;
325                let bundle = execution
326                    .bundle()
327                    .with_turn_phase(*turn_id, TurnPhase::Committed);
328                self.executor.commit(account, bundle).await?;
329                let events =
330                    self.recovery.decide(account, turn_id).await.ok().and_then(
331                        |action| match action {
332                            RecoveryAction::RegenerateResponse {
333                                events,
334                                answer_tasks,
335                            } => Some((events, answer_tasks)),
336                            _ => None,
337                        },
338                    );
339                let turn = match events {
340                    Some((events, answer_tasks)) => {
341                        self.regenerate(account, turn_id, &events, &answer_tasks)
342                            .await?
343                    }
344                    None => None,
345                };
346                Ok(ResumeOutcome::Resumed {
347                    outcomes: execution.outcomes.clone(),
348                    turn: turn.map(Box::new),
349                })
350            }
351        }
352    }
353
354    /// Rebuilds the answer of a committed turn from the ledger, and persists it.
355    /// `None` when the turn already has one.
356    async fn regenerate(
357        &self,
358        account: &AccountId,
359        turn_id: &TurnId,
360        events: &[turnframe_store::events::StoredEvent],
361        answer_tasks: &[turnframe_core::reduce::AnswerTask],
362    ) -> Result<Option<AssistantTurn>, OrchestratorError> {
363        let stored = self.recovery.stored_turn(account, turn_id).await?;
364        if stored.assistant.is_some() {
365            return Ok(None);
366        }
367        // Grouped by the store, which carries a redaction through as one.
368        let groups = turnframe_store::events::group_for_receipts(events);
369        let interactions = self
370            .interactions
371            .open_for_conversation(account, &stored.user.input.conversation_id)
372            .await
373            .unwrap_or_default();
374        let input = CompositionInput::new(&stored.user.input)
375            .with_answer_tasks(answer_tasks)
376            .with_ledger(&groups)
377            .with_interactions(&interactions);
378        let composition = self.composer.compose(input).await?;
379        self.stores
380            .conversations()
381            .append_assistant_turn(account, composition.turn.clone())
382            .await
383            .map_err(OrchestratorError::Store)?;
384        self.stores
385            .conversations()
386            .set_turn_phase(account, turn_id, TurnPhase::Delivered)
387            .await
388            .map_err(OrchestratorError::Store)?;
389        Ok(Some(composition.turn))
390    }
391
392    async fn run(
393        &self,
394        input: TurnInput,
395        publisher: &TurnPublisher,
396    ) -> Result<AssistantTurn, OrchestratorError> {
397        let labels =
398            SignalLabels::none().with_effort(input.effort.unwrap_or(self.config.effort.default));
399        self.observer
400            .observe_labeled(&Signal::TurnReceived, &labels);
401        let stage = crate::signals::Stage::enter();
402        // A turn's state is large, and every caller awaits it: it lives on the heap.
403        let outcome = Box::pin(session::Session::new(self, input, publisher).run()).await;
404        stage.observe(self.observer.as_ref(), Signal::TurnDuration, &labels);
405        match outcome {
406            Ok(turn) => {
407                self.observer
408                    .observe_labeled(&Signal::TurnCompleted, &labels);
409                Ok(turn)
410            }
411            Err(error) => {
412                self.observer.observe_labeled(
413                    &Signal::TurnFailed,
414                    &labels.with_error_code(error_code(&error)),
415                );
416                Err(error)
417            }
418        }
419    }
420}
421
422/// What acting on a recovery decision achieved (spec §23.1).
423#[derive(Debug, Clone)]
424#[non_exhaustive]
425pub enum ResumeOutcome {
426    /// The turn was already finished.
427    Nothing {
428        /// The phase it is in.
429        phase: TurnPhase,
430    },
431    /// Nothing was ever admitted, so the turn may simply be submitted again.
432    Restartable,
433    /// Pending commands were resumed by idempotency key.
434    Resumed {
435        /// What each of them ended up doing.
436        outcomes: Vec<turnframe_core::replay::CommandOutcomeRecord>,
437        /// The answer, when the resumed turn could be finished as well.
438        turn: Option<Box<AssistantTurn>>,
439    },
440    /// The effects were already committed; only the answer was rebuilt.
441    Regenerated {
442        /// The answer, now persisted.
443        turn: Box<AssistantTurn>,
444    },
445    /// An external effect may or may not have happened; only the application can
446    /// settle it (§16.5, I15).
447    AwaitingReconciliation {
448        /// The attempts to settle.
449        attempts: Vec<turnframe_core::ids::AttemptId>,
450    },
451}
452
453/// Rebuilds the batches of a set of journal entries, one per case, so the executor
454/// sees the very commands that were admitted (spec §23.1).
455fn resume_batches(
456    actor: &ActorContext,
457    turn_id: TurnId,
458    entries: &[turnframe_store::journal::CommandJournalEntry],
459) -> Vec<CommandBatch<serde_json::Value>> {
460    let mut batches: IndexMap<CaseKey, CommandBatch<serde_json::Value>> = IndexMap::new();
461    for entry in entries {
462        let key = entry.case_ref.key();
463        let scope = turnframe_core::command::AtomicityScope::PerCase;
464        let batch_id = turnframe_core::ids::BatchId::derive(&turn_id, &key, &scope);
465        batches
466            .entry(key)
467            .or_insert_with(|| CommandBatch {
468                batch_id,
469                scope,
470                envelopes: Vec::new(),
471            })
472            .envelopes
473            .push(turnframe_core::command::CommandEnvelope {
474                command_id: entry.command_id,
475                turn_id,
476                actor: actor.clone(),
477                case_ref: entry.case_ref.clone(),
478                idempotency_key: entry.idempotency_key.clone(),
479                origin: entry.origin.clone(),
480                command: entry.command_payload.clone(),
481            });
482    }
483    batches.into_values().collect()
484}
485
486/// A stable code for an orchestrator failure, safe to put on the wire.
487#[must_use]
488pub fn error_code(error: &OrchestratorError) -> String {
489    use turnframe_core::error::ErrorClassification;
490    error.user_message_key().to_owned()
491}
492
493/// The model-authored text of a stored assistant turn, for the transcript.
494pub(crate) fn assistant_text(turn: &AssistantTurn) -> String {
495    reply_text(&turn.blocks)
496}
497
498/// The reply as the user read it: the transition, which carries the turn's answers, or the
499/// answers alone when no transition was written.
500fn reply_text(blocks: &[ResponseBlock]) -> String {
501    let said = |transitions: bool| {
502        blocks
503            .iter()
504            .filter_map(|block| match block {
505                ResponseBlock::Transition(transition) if transitions => {
506                    Some(transition.text.as_str())
507                }
508                ResponseBlock::Answer(answer) if !transitions => Some(answer.text.as_str()),
509                _ => None,
510            })
511            .collect::<Vec<_>>()
512            .join(" ")
513    };
514    let reply = said(true);
515    if reply.is_empty() { said(false) } else { reply }
516}
517
518#[cfg(test)]
519mod tests {
520    use super::*;
521    use turnframe_core::ids::BlockId;
522    use turnframe_core::plan::AnswerBasis;
523    use turnframe_core::response::{AnswerStatus, GeneratedAnswer, GeneratedTransition};
524
525    fn answer(text: &str) -> ResponseBlock {
526        ResponseBlock::Answer(GeneratedAnswer {
527            block_id: BlockId::from("answer:0"),
528            question_id: None,
529            text: text.to_owned(),
530            basis: AnswerBasis::CurrentCommittedState,
531            status: AnswerStatus::Answered,
532            facts_used: Vec::new(),
533            citations: Vec::new(),
534            enumerations: Vec::new(),
535        })
536    }
537
538    fn transition(text: &str) -> ResponseBlock {
539        ResponseBlock::Transition(GeneratedTransition {
540            block_id: BlockId::from("transition:0"),
541            text: text.to_owned(),
542            facts_used: Vec::new(),
543        })
544    }
545
546    #[test]
547    fn the_reply_is_read_back_once() {
548        let said = "The airline pays for it.";
549        assert_eq!(reply_text(&[answer(said), transition(said)]), said);
550        assert_eq!(reply_text(&[answer(said)]), said);
551    }
552}