Skip to main content

turnframe_runtime/orchestrator/
builder.rs

1//! Collecting the parts of an [`Orchestrator`].
2
3use 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/// Why an [`Orchestrator`] could not be built.
33#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
34#[non_exhaustive]
35pub enum BuildError {
36    /// A mandatory part was not supplied.
37    #[error("no {part} was supplied to the orchestrator builder")]
38    Missing {
39        /// Which part.
40        part: &'static str,
41    },
42    /// The configuration is one the library refuses to run.
43    #[error(transparent)]
44    Config(crate::config::ConfigError),
45    /// [`OrchestratorBuilder::mode`] and [`OrchestratorBuilder::config`] were both
46    /// given. They write the same field, so the builder refuses instead of letting
47    /// call order pick one.
48    #[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    /// A language declared with [`OrchestratorBuilder::locales`] has no text for some
54    /// of the server's own sentences.
55    #[error("the {copy} has no {locale} text for: {}", sentences.join(", "))]
56    CopyMissing {
57        /// The declared language.
58        locale: String,
59        /// The copy that lacks it.
60        copy: &'static str,
61        /// The sentences without it, by field.
62        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
72/// Collects everything one [`Orchestrator`] needs.
73pub 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    /// An empty builder with the conservative configuration.
142    #[must_use]
143    pub fn new() -> Self {
144        Self::default()
145    }
146
147    /// The workflows this runtime hosts.
148    #[must_use]
149    pub fn workflows(mut self, workflows: Arc<WorkflowRegistry>) -> Self {
150        self.workflows = Some(workflows);
151        self
152    }
153
154    /// The configured providers.
155    #[must_use]
156    pub fn providers(mut self, providers: Arc<ProviderPool>) -> Self {
157        self.providers = Some(providers);
158        self
159    }
160
161    /// The persistence layer.
162    #[must_use]
163    pub fn stores(mut self, stores: Stores) -> Self {
164        self.stores = Some(stores);
165        self
166    }
167
168    /// Declares what a turn's writes imply on other cases. See [`TurnConsequences`].
169    #[must_use]
170    pub fn consequences(mut self, consequences: Arc<dyn TurnConsequences>) -> Self {
171        self.consequences = Some(consequences);
172        self
173    }
174
175    /// The directory of addressable cases.
176    #[must_use]
177    pub fn case_directory(mut self, directory: Arc<dyn CaseDirectory>) -> Self {
178        self.directory = Some(directory);
179        self
180    }
181
182    /// The knowledge provider answers may rest on (spec §19.2).
183    #[must_use]
184    pub fn knowledge(mut self, knowledge: Arc<dyn KnowledgeProvider>) -> Self {
185        self.knowledge = Some(knowledge);
186        self
187    }
188
189    /// Where the bytes of the turn's files come from. Without one, no model is shown
190    /// a file.
191    #[must_use]
192    pub fn attachments(mut self, source: Arc<dyn AttachmentSource>) -> Self {
193        self.attachment_source = Some(source);
194        self
195    }
196
197    /// The policy snapshot commands are judged against (spec §14.3).
198    #[must_use]
199    pub fn policy(mut self, policy: PolicySnapshot) -> Self {
200        self.policy = policy;
201        self
202    }
203
204    /// Where metrics go (spec §26.2).
205    #[must_use]
206    pub fn observer(mut self, observer: Arc<dyn Observer>) -> Self {
207        self.observer = observer;
208        self
209    }
210
211    /// How much autonomy the model gets (spec §11.1). Refused beside
212    /// [`Self::config`]; see [`BuildError::ConflictingMode`].
213    #[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    /// The whole configuration, mode included. Refused beside [`Self::mode`].
221    #[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    /// Replaces the runtime clock.
229    #[must_use]
230    pub fn clock(mut self, clock: Arc<dyn TurnClock>) -> Self {
231        self.clock = clock;
232        self
233    }
234
235    /// Replaces the factory that mints identifiers for new cases.
236    #[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    /// Replaces what understands each turn. The default is an [`Understander`] over
243    /// the configured providers, with the profiles and settings of
244    /// [`UnderstandingConfig`](crate::config::UnderstandingConfig).
245    #[must_use]
246    pub fn understander(mut self, understander: Arc<dyn TurnUnderstander>) -> Self {
247        self.understander = Some(understander);
248        self
249    }
250
251    /// Lets a prompt source supply the instructions of every model task.
252    ///
253    /// Without one the runtime uses the text compiled into the library. Each task asks
254    /// for its own name, `understand.<task>` or `narrate.<task>`, and its record cites
255    /// the prompt it ran under. It applies to the stages this builder creates.
256    #[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    /// Which version of each prompt the source is asked for. Pin one, so an edit in a
263    /// registry cannot change a running system.
264    #[must_use]
265    pub fn prompt_selector(mut self, selector: PromptSelector) -> Self {
266        self.prompt_selector = selector;
267        self
268    }
269
270    /// Replaces the composition stage.
271    #[must_use]
272    pub fn composer(mut self, composer: Composer) -> Self {
273        self.composer = Some(composer);
274        self
275    }
276
277    /// Replaces the notices the reducer writes itself. English by default; a
278    /// deployment in another language wants this, the composer's copy and
279    /// [`Self::confirmation_copy`].
280    #[must_use]
281    pub fn notice_copy(mut self, copy: NoticeCopy) -> Self {
282        self.notice_copy = copy;
283        self
284    }
285
286    /// Replaces what a user is told about a file the model was not shown.
287    #[must_use]
288    pub fn attachment_copy(mut self, copy: AttachmentCopy) -> Self {
289        self.attachment_copy = copy;
290        self
291    }
292
293    /// The languages this deployment serves: building fails while any of the server's own
294    /// sentences has no text in one of them. The built-in copy speaks English and Italian.
295    #[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    /// Replaces the copy on the cards the policy engine raises, the two buttons of
302    /// every confirmation card included.
303    #[must_use]
304    pub fn confirmation_copy(mut self, copy: ConfirmationCopy) -> Self {
305        self.confirmation_copy = copy;
306        self
307    }
308
309    /// Reports every turn's events to `trace`: for local debugging, since a trace holds
310    /// the users' words. Wrap the providers in
311    /// [`TracedProvider`](turnframe_provider::trace::TracedProvider) to add the model
312    /// calls. See [`crate::trace`].
313    #[must_use]
314    pub fn trace(mut self, trace: Arc<dyn TurnTrace>) -> Self {
315        self.trace = Some(trace);
316        self
317    }
318
319    /// Builds the orchestrator.
320    ///
321    /// # Errors
322    ///
323    /// [`BuildError::Missing`] for the first mandatory part not supplied,
324    /// [`BuildError::Config`] for a configuration the library refuses, and
325    /// [`BuildError::ConflictingMode`].
326    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
427/// The engine every model task of a turn runs on, understanding's and narration's.
428fn 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}