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 knowledge: self.composer.has_knowledge(),
200 }
201 }
202
203 pub async fn plan_turn(&self, input: TurnInput) -> Result<PlannedTurn, OrchestratorError> {
210 self.planner().plan(input).await
211 }
212
213 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 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 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 #[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 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 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 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 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 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#[derive(Debug, Clone)]
424#[non_exhaustive]
425pub enum ResumeOutcome {
426 Nothing {
428 phase: TurnPhase,
430 },
431 Restartable,
433 Resumed {
435 outcomes: Vec<turnframe_core::replay::CommandOutcomeRecord>,
437 turn: Option<Box<AssistantTurn>>,
439 },
440 Regenerated {
442 turn: Box<AssistantTurn>,
444 },
445 AwaitingReconciliation {
448 attempts: Vec<turnframe_core::ids::AttemptId>,
450 },
451}
452
453fn 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#[must_use]
488pub fn error_code(error: &OrchestratorError) -> String {
489 use turnframe_core::error::ErrorClassification;
490 error.user_message_key().to_owned()
491}
492
493pub(crate) fn assistant_text(turn: &AssistantTurn) -> String {
495 reply_text(&turn.blocks)
496}
497
498fn 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}