1use std::fmt;
4use std::sync::Arc;
5
6use turnframe_core::flow::WorkflowRegistry;
7use turnframe_core::knowledge::KnowledgeProvider;
8use turnframe_core::locale::Locale;
9use turnframe_core::observe::{NoopObserver, Observer};
10use turnframe_core::policy::PolicySnapshot;
11use turnframe_core::prompt::{PromptSelector, PromptSource};
12use turnframe_core::turn::AttachmentSource;
13use turnframe_provider::router::{PolicyRouter, ProviderPool, ProviderRouter};
14use turnframe_store::stores::Stores;
15use turnframe_tasks::{RecordPolicy, TaskEngine};
16use turnframe_understand::{TurnUnderstander, Understander};
17
18use super::{
19 CaseDirectory, Orchestrator, StaticCaseDirectory, SystemTurnClock, TurnClock, TurnConsequences,
20};
21use crate::attachments::AttachmentCopy;
22use crate::compose::Composer;
23use crate::config::{OrchestrationMode, OrchestratorConfig};
24use crate::execute::CommandExecutor;
25use crate::interactions::InteractionEngine;
26use crate::policy::{ConfirmationCopy, PolicyEngine};
27use crate::recover::Recovery;
28use crate::reduce::NoticeCopy;
29use crate::resolve::{CaseIdFactory, DerivedCaseIdFactory};
30use crate::trace::TurnTrace;
31
32#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
34#[non_exhaustive]
35pub enum BuildError {
36 #[error("no {part} was supplied to the orchestrator builder")]
38 Missing {
39 part: &'static str,
41 },
42 #[error(transparent)]
44 Config(crate::config::ConfigError),
45 #[error(
49 "both `mode` and `config` were given to the orchestrator builder, and they write the \
50 same field: set the mode on the configuration, or pass no configuration"
51 )]
52 ConflictingMode,
53 #[error("the {copy} has no {locale} text for: {}", sentences.join(", "))]
56 CopyMissing {
57 locale: String,
59 copy: &'static str,
61 sentences: Vec<&'static str>,
63 },
64}
65
66impl From<crate::config::ConfigError> for BuildError {
67 fn from(value: crate::config::ConfigError) -> Self {
68 Self::Config(value)
69 }
70}
71
72pub struct OrchestratorBuilder {
74 workflows: Option<Arc<WorkflowRegistry>>,
75 providers: Option<Arc<ProviderPool>>,
76 stores: Option<Stores>,
77 directory: Option<Arc<dyn CaseDirectory>>,
78 consequences: Option<Arc<dyn TurnConsequences>>,
79 knowledge: Option<Arc<dyn KnowledgeProvider>>,
80 attachment_source: Option<Arc<dyn AttachmentSource>>,
81 policy: PolicySnapshot,
82 observer: Arc<dyn Observer>,
83 config: OrchestratorConfig,
84 mode_set: bool,
85 config_set: bool,
86 clock: Arc<dyn TurnClock>,
87 case_ids: Arc<dyn CaseIdFactory>,
88 understander: Option<Arc<dyn TurnUnderstander>>,
89 composer: Option<Composer>,
90 prompt_source: Option<Arc<dyn PromptSource>>,
91 prompt_selector: PromptSelector,
92 notice_copy: NoticeCopy,
93 attachment_copy: AttachmentCopy,
94 confirmation_copy: ConfirmationCopy,
95 locales: Vec<Locale>,
96 trace: Option<Arc<dyn TurnTrace>>,
97}
98
99impl fmt::Debug for OrchestratorBuilder {
100 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
101 f.debug_struct("OrchestratorBuilder")
102 .field("workflows", &self.workflows.is_some())
103 .field("providers", &self.providers.is_some())
104 .field("stores", &self.stores.is_some())
105 .field("mode", &self.config.mode)
106 .finish_non_exhaustive()
107 }
108}
109
110impl Default for OrchestratorBuilder {
111 fn default() -> Self {
112 Self {
113 workflows: None,
114 providers: None,
115 stores: None,
116 directory: None,
117 consequences: None,
118 knowledge: None,
119 attachment_source: None,
120 policy: PolicySnapshot::conservative(),
121 observer: Arc::new(NoopObserver),
122 config: OrchestratorConfig::conservative(),
123 mode_set: false,
124 config_set: false,
125 clock: Arc::new(SystemTurnClock),
126 case_ids: Arc::new(DerivedCaseIdFactory),
127 understander: None,
128 composer: None,
129 prompt_source: None,
130 prompt_selector: PromptSelector::Latest,
131 notice_copy: NoticeCopy::standard(),
132 attachment_copy: AttachmentCopy::standard(),
133 confirmation_copy: ConfirmationCopy::standard(),
134 locales: Vec::new(),
135 trace: None,
136 }
137 }
138}
139
140impl OrchestratorBuilder {
141 #[must_use]
143 pub fn new() -> Self {
144 Self::default()
145 }
146
147 #[must_use]
149 pub fn workflows(mut self, workflows: Arc<WorkflowRegistry>) -> Self {
150 self.workflows = Some(workflows);
151 self
152 }
153
154 #[must_use]
156 pub fn providers(mut self, providers: Arc<ProviderPool>) -> Self {
157 self.providers = Some(providers);
158 self
159 }
160
161 #[must_use]
163 pub fn stores(mut self, stores: Stores) -> Self {
164 self.stores = Some(stores);
165 self
166 }
167
168 #[must_use]
170 pub fn consequences(mut self, consequences: Arc<dyn TurnConsequences>) -> Self {
171 self.consequences = Some(consequences);
172 self
173 }
174
175 #[must_use]
177 pub fn case_directory(mut self, directory: Arc<dyn CaseDirectory>) -> Self {
178 self.directory = Some(directory);
179 self
180 }
181
182 #[must_use]
184 pub fn knowledge(mut self, knowledge: Arc<dyn KnowledgeProvider>) -> Self {
185 self.knowledge = Some(knowledge);
186 self
187 }
188
189 #[must_use]
192 pub fn attachments(mut self, source: Arc<dyn AttachmentSource>) -> Self {
193 self.attachment_source = Some(source);
194 self
195 }
196
197 #[must_use]
199 pub fn policy(mut self, policy: PolicySnapshot) -> Self {
200 self.policy = policy;
201 self
202 }
203
204 #[must_use]
206 pub fn observer(mut self, observer: Arc<dyn Observer>) -> Self {
207 self.observer = observer;
208 self
209 }
210
211 #[must_use]
214 pub fn mode(mut self, mode: OrchestrationMode) -> Self {
215 self.config.mode = mode;
216 self.mode_set = true;
217 self
218 }
219
220 #[must_use]
222 pub fn config(mut self, config: OrchestratorConfig) -> Self {
223 self.config = config;
224 self.config_set = true;
225 self
226 }
227
228 #[must_use]
230 pub fn clock(mut self, clock: Arc<dyn TurnClock>) -> Self {
231 self.clock = clock;
232 self
233 }
234
235 #[must_use]
237 pub fn case_id_factory(mut self, factory: Arc<dyn CaseIdFactory>) -> Self {
238 self.case_ids = factory;
239 self
240 }
241
242 #[must_use]
246 pub fn understander(mut self, understander: Arc<dyn TurnUnderstander>) -> Self {
247 self.understander = Some(understander);
248 self
249 }
250
251 #[must_use]
257 pub fn prompt_source(mut self, source: Arc<dyn PromptSource>) -> Self {
258 self.prompt_source = Some(source);
259 self
260 }
261
262 #[must_use]
265 pub fn prompt_selector(mut self, selector: PromptSelector) -> Self {
266 self.prompt_selector = selector;
267 self
268 }
269
270 #[must_use]
272 pub fn composer(mut self, composer: Composer) -> Self {
273 self.composer = Some(composer);
274 self
275 }
276
277 #[must_use]
281 pub fn notice_copy(mut self, copy: NoticeCopy) -> Self {
282 self.notice_copy = copy;
283 self
284 }
285
286 #[must_use]
288 pub fn attachment_copy(mut self, copy: AttachmentCopy) -> Self {
289 self.attachment_copy = copy;
290 self
291 }
292
293 #[must_use]
296 pub fn locales<L: Into<Locale>>(mut self, locales: impl IntoIterator<Item = L>) -> Self {
297 self.locales = locales.into_iter().map(Into::into).collect();
298 self
299 }
300
301 #[must_use]
304 pub fn confirmation_copy(mut self, copy: ConfirmationCopy) -> Self {
305 self.confirmation_copy = copy;
306 self
307 }
308
309 #[must_use]
314 pub fn trace(mut self, trace: Arc<dyn TurnTrace>) -> Self {
315 self.trace = Some(trace);
316 self
317 }
318
319 pub fn build(self) -> Result<Orchestrator, BuildError> {
327 if self.mode_set && self.config_set {
328 return Err(BuildError::ConflictingMode);
329 }
330 self.config.validate()?;
331 let workflows = self
332 .workflows
333 .ok_or(BuildError::Missing { part: "workflows" })?;
334 let providers = self
335 .providers
336 .ok_or(BuildError::Missing { part: "providers" })?;
337 let stores = self.stores.ok_or(BuildError::Missing { part: "stores" })?;
338 let directory = self
339 .directory
340 .unwrap_or_else(|| Arc::new(StaticCaseDirectory::new()));
341 let router: Arc<dyn ProviderRouter> = Arc::new(PolicyRouter::new(providers));
342 let engine = task_engine(
343 &router,
344 &self.config,
345 self.prompt_source
346 .map(|source| (source, self.prompt_selector)),
347 &self.observer,
348 );
349 let understander = self.understander.unwrap_or_else(|| {
350 Arc::new(
351 Understander::new(engine.clone()).with_settings(self.config.understanding.settings),
352 )
353 });
354 let composer = self
355 .composer
356 .unwrap_or_else(|| {
357 let composer = Composer::new(
358 Arc::clone(&workflows),
359 Arc::clone(&router),
360 self.config.narration,
361 );
362 match self.knowledge.clone() {
363 Some(knowledge) => composer.with_knowledge(knowledge),
364 None => composer,
365 }
366 })
367 .with_tasks(engine);
368 let mut copies: Vec<&dyn crate::copy::ServerCopy> = vec![
369 &self.notice_copy,
370 &self.attachment_copy,
371 &self.confirmation_copy,
372 ];
373 copies.extend(composer.server_copy());
374 for locale in &self.locales {
375 for copy in &copies {
376 let sentences = crate::copy::missing(*copy, locale);
377 if !sentences.is_empty() {
378 return Err(BuildError::CopyMissing {
379 locale: locale.to_string(),
380 copy: copy.name(),
381 sentences,
382 });
383 }
384 }
385 }
386 let executor = CommandExecutor::new(
387 Arc::clone(&workflows),
388 Arc::clone(stores.journal()),
389 Arc::clone(stores.commit()),
390 Arc::clone(stores.outbox()),
391 self.config.execution,
392 );
393 let interactions =
394 InteractionEngine::new(Arc::clone(stores.interactions()), self.config.interaction)
395 .with_observer(Arc::clone(&self.observer));
396 let recovery = Recovery::new(
397 Arc::clone(stores.conversations()),
398 Arc::clone(stores.journal()),
399 Arc::clone(stores.events()),
400 Arc::clone(stores.replay()),
401 );
402 let policy_engine = PolicyEngine::new(&self.config).with_copy(self.confirmation_copy);
403 Ok(Orchestrator {
404 consequences: self.consequences,
405 workflows,
406 stores,
407 directory,
408 understander,
409 composer,
410 executor,
411 interactions,
412 recovery,
413 policy_engine,
414 policy: self.config.policy_snapshot(self.policy),
415 observer: self.observer,
416 attachment_source: self.attachment_source,
417 clock: self.clock,
418 case_ids: self.case_ids,
419 config: self.config,
420 notice_copy: self.notice_copy,
421 attachment_copy: self.attachment_copy,
422 trace: self.trace,
423 })
424 }
425}
426
427fn task_engine(
429 router: &Arc<dyn ProviderRouter>,
430 config: &OrchestratorConfig,
431 prompts: Option<(Arc<dyn PromptSource>, PromptSelector)>,
432 observer: &Arc<dyn Observer>,
433) -> TaskEngine {
434 let mut records = RecordPolicy::default();
435 records.keep_prompts = config.privacy.store_model_prompts;
436 records.keep_raw_output = config.privacy.store_model_prompts;
437 let mut engine = TaskEngine::builder(Arc::clone(router))
438 .profiles(config.understanding.tasks.clone())
439 .records(records)
440 .observer(Arc::clone(observer));
441 if let Some((source, selector)) = prompts {
442 engine = engine.prompts(source, selector);
443 }
444 engine.build()
445}
446
447#[cfg(test)]
448mod tests {
449 use super::*;
450
451 #[test]
452 fn a_builder_given_both_mode_and_config_refuses_to_build() {
453 let built = Orchestrator::builder()
454 .mode(OrchestrationMode::Deterministic)
455 .config(OrchestratorConfig::conservative())
456 .build();
457 assert!(matches!(built, Err(BuildError::ConflictingMode)));
458 }
459
460 #[test]
461 fn either_setter_alone_fails_only_for_the_missing_parts() {
462 for built in [
463 Orchestrator::builder()
464 .mode(OrchestrationMode::Deterministic)
465 .build(),
466 Orchestrator::builder()
467 .config(OrchestratorConfig::conservative())
468 .build(),
469 ] {
470 assert!(
471 matches!(built, Err(BuildError::Missing { .. })),
472 "{built:?}"
473 );
474 }
475 }
476}