1mod 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
51pub trait TurnClock: Send + Sync + fmt::Debug {
53 fn now(&self) -> DateTime<Utc>;
55}
56
57#[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#[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#[async_trait]
86pub trait TurnConsequences: Send + Sync {
87 async fn following(
94 &self,
95 actor: &ActorContext,
96 batches: &[CommandBatch<serde_json::Value>],
97 ) -> Result<Vec<CommandBatch<serde_json::Value>>, StoreError>;
98}
99
100pub 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 #[must_use]
135 pub fn builder() -> OrchestratorBuilder {
136 OrchestratorBuilder::new()
137 }
138
139 #[must_use]
141 pub const fn config(&self) -> &OrchestratorConfig {
142 &self.config
143 }
144
145 #[must_use]
147 pub const fn stores(&self) -> &Stores {
148 &self.stores
149 }
150
151 #[must_use]
155 pub const fn workflows(&self) -> &Arc<WorkflowRegistry> {
156 &self.workflows
157 }
158
159 #[must_use]
161 pub const fn recovery(&self) -> &Recovery {
162 &self.recovery
163 }
164
165 #[must_use]
167 pub const fn interactions(&self) -> &InteractionEngine {
168 &self.interactions
169 }
170
171 #[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 #[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 pub async fn plan_turn(&self, input: TurnInput) -> Result<PlannedTurn, OrchestratorError> {
209 self.planner().plan(input).await
210 }
211
212 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 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 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 #[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 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 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 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 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 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#[derive(Debug, Clone)]
423#[non_exhaustive]
424pub enum ResumeOutcome {
425 Nothing {
427 phase: TurnPhase,
429 },
430 Restartable,
432 Resumed {
434 outcomes: Vec<turnframe_core::replay::CommandOutcomeRecord>,
436 turn: Option<Box<AssistantTurn>>,
438 },
439 Regenerated {
441 turn: Box<AssistantTurn>,
443 },
444 AwaitingReconciliation {
447 attempts: Vec<turnframe_core::ids::AttemptId>,
449 },
450}
451
452fn 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#[must_use]
487pub fn error_code(error: &OrchestratorError) -> String {
488 use turnframe_core::error::ErrorClassification;
489 error.user_message_key().to_owned()
490}
491
492pub(crate) fn assistant_text(turn: &AssistantTurn) -> String {
494 reply_text(&turn.blocks)
495}
496
497fn 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}