Skip to main content

stasis/application/runtime/
stasis_runtime_builder.rs

1use std::sync::Arc;
2
3use async_trait::async_trait;
4
5use crate::application::orchestration::tool_registry::{InMemoryToolRegistry, StasisTool};
6use crate::application::runtime::agent_session_job_handler::AgentSessionJobHandler;
7use crate::application::runtime::agent_turn_job_handler::AgentTurnJobHandler;
8use crate::application::runtime::chat_client_middleware::ChatClientMiddleware;
9use crate::application::runtime::concurrent_pattern_job_handler::ConcurrentPatternJobHandler;
10use crate::application::runtime::coordinator_failover_job_handler::CoordinatorFailoverJobHandler;
11use crate::application::runtime::default_chat_middlewares::{
12    CacheChatMiddleware, LoggingChatMiddleware, TelemetryChatMiddleware,
13    ToolCallInterceptionChatMiddleware,
14};
15use crate::application::runtime::grapheme_echo_job_handler::GraphemeEchoJobHandler;
16use crate::application::runtime::grapheme_healthcheck_job_handler::GraphemeHealthcheckJobHandler;
17use crate::application::runtime::grapheme_job_handler::GraphemeJobHandler;
18use crate::application::runtime::grapheme_textops_job_handler::GraphemeTextOpsJobHandler;
19use crate::application::runtime::handoff_pattern_job_handler::HandoffPatternJobHandler;
20use crate::application::runtime::in_memory_runtime::{JobExecutionOutcome, JobHandler};
21use crate::application::runtime::memory_aggregate_job_handler::MemoryAggregateJobHandler;
22use crate::application::runtime::memory_find_job_handler::MemoryFindJobHandler;
23use crate::application::runtime::memory_recall_job_handler::MemoryRecallJobHandler;
24use crate::application::runtime::memory_rollup_job_handler::MemoryRollupJobHandler;
25use crate::application::runtime::memory_schema_job_handler::MemorySchemaJobHandler;
26use crate::application::runtime::memory_transform_job_handler::MemoryTransformJobHandler;
27use crate::application::runtime::orchestrator_pattern_job_handler::OrchestratorPatternJobHandler;
28use crate::application::runtime::prompt_chat_job_handler::PromptChatJobHandler;
29use crate::application::runtime::queue_ownership_rebalance_job_handler::QueueOwnershipRebalanceJobHandler;
30use crate::application::runtime::runtime_factory::{
31    RuntimeBackend, RuntimeComposition, RuntimeFactory,
32};
33use crate::application::runtime::sequential_pattern_job_handler::SequentialPatternJobHandler;
34use crate::application::runtime::tool_loop_job_handler::ToolLoopJobHandler;
35use crate::application::telemetry::operation::OperationTelemetry;
36use crate::domain::errors::Result;
37use crate::ports::outbound::ai_chat_client::AiChatClient;
38use crate::ports::outbound::ai_chat_response_cache::AiChatResponseCache;
39use crate::ports::outbound::ai_chat_tool_interceptor::AiChatToolInterceptor;
40use crate::ports::outbound::memory::memory_context_reader::MemoryContextReader;
41use crate::ports::outbound::memory::memory_context_writer::MemoryContextWriter;
42use crate::ports::outbound::memory::identity_memory_store::IdentityMemoryStore;
43use crate::ports::outbound::memory::memory_operations::MemoryOperations;
44use crate::ports::outbound::runtime::cluster_node_store::ClusterNodeStore;
45use crate::ports::outbound::runtime::delivery_endpoint_store::DeliveryEndpointStore;
46use crate::ports::outbound::runtime::endpoint_delivery_status_store::EndpointDeliveryStatusStore;
47use crate::ports::outbound::runtime::endpoint_routing_policy::EndpointRoutingPolicy;
48use crate::ports::outbound::runtime::endpoint_transport_publisher::EndpointTransportPublisher;
49use crate::ports::outbound::runtime::runtime_metrics::RuntimeMetrics;
50use crate::ports::outbound::runtime::runtime_telemetry::RuntimeTelemetry;
51use crate::ports::outbound::runtime::runtime_tracing::RuntimeTracing;
52use crate::ports::outbound::runtime::thread_store::ThreadStore;
53
54#[derive(Clone)]
55struct DelegatingJobHandler {
56    inner: Arc<dyn JobHandler>,
57}
58
59#[async_trait]
60impl JobHandler for DelegatingJobHandler {
61    fn job_type(&self) -> &'static str {
62        self.inner.job_type()
63    }
64
65    async fn execute(&self, job: &crate::domain::runtime::job::Job) -> Result<JobExecutionOutcome> {
66        self.inner.execute(job).await
67    }
68}
69
70#[derive(Clone)]
71pub struct StasisRuntimeBuilder {
72    backend: RuntimeBackend,
73    chat_client: Option<Arc<dyn AiChatClient>>,
74    chat_middlewares: Vec<Arc<dyn ChatClientMiddleware>>,
75    memory_context_reader: Option<Arc<dyn MemoryContextReader>>,
76    memory_context_writer: Option<Arc<dyn MemoryContextWriter>>,
77    identity_memory_store: Option<Arc<dyn IdentityMemoryStore>>,
78    memory_operations: Option<Arc<dyn MemoryOperations>>,
79    thread_store: Option<Arc<dyn ThreadStore>>,
80    cluster_node_store: Option<Arc<dyn ClusterNodeStore>>,
81    delivery_endpoint_store: Option<Arc<dyn DeliveryEndpointStore>>,
82    endpoint_delivery_status_store: Option<Arc<dyn EndpointDeliveryStatusStore>>,
83    endpoint_transport_publishers: Vec<Arc<dyn EndpointTransportPublisher>>,
84    endpoint_routing_policy: Option<Arc<dyn EndpointRoutingPolicy>>,
85    enable_endpoint_routing_delivery: bool,
86    enable_locus_memory: bool,
87    tool_registry: InMemoryToolRegistry,
88    include_grapheme_handlers: bool,
89    include_prompt_handler: bool,
90    include_tool_loop_handler: bool,
91    include_agent_handlers: bool,
92    include_memory_operation_handlers: bool,
93    include_orchestration_pattern_handlers: bool,
94    include_cluster_control_handlers: bool,
95    extra_handlers: Vec<Arc<dyn JobHandler>>,
96    runtime_telemetry_metrics: Option<Arc<dyn RuntimeMetrics>>,
97    runtime_telemetry_tracing: Option<Arc<dyn RuntimeTracing>>,
98    explicit_telemetry_chat_middleware: bool,
99}
100
101macro_rules! define_arc_option_setter {
102    ($fn_name:ident, $field:ident, $ty:ty) => {
103        pub fn $fn_name(mut self, value: Arc<$ty>) -> Self {
104            self.$field = Some(value);
105            self
106        }
107    };
108}
109
110macro_rules! define_enable_flag_setter {
111    ($fn_name:ident, $field:ident) => {
112        pub fn $fn_name(mut self) -> Self {
113            self.$field = true;
114            self
115        }
116    };
117}
118
119macro_rules! define_disable_flag_setter {
120    ($fn_name:ident, $field:ident) => {
121        pub fn $fn_name(mut self) -> Self {
122            self.$field = false;
123            self
124        }
125    };
126}
127
128impl StasisRuntimeBuilder {
129    pub fn new(backend: RuntimeBackend) -> Self {
130        Self {
131            backend,
132            chat_client: None,
133            chat_middlewares: Vec::new(),
134            memory_context_reader: None,
135            memory_context_writer: None,
136            identity_memory_store: None,
137            memory_operations: None,
138            thread_store: None,
139            cluster_node_store: None,
140            delivery_endpoint_store: None,
141            endpoint_delivery_status_store: None,
142            endpoint_transport_publishers: Vec::new(),
143            endpoint_routing_policy: None,
144            enable_endpoint_routing_delivery: false,
145            enable_locus_memory: false,
146            tool_registry: InMemoryToolRegistry::default(),
147            include_grapheme_handlers: true,
148            include_prompt_handler: true,
149            include_tool_loop_handler: true,
150            include_agent_handlers: true,
151            include_memory_operation_handlers: true,
152            include_orchestration_pattern_handlers: true,
153            include_cluster_control_handlers: true,
154            extra_handlers: Vec::new(),
155            runtime_telemetry_metrics: None,
156            runtime_telemetry_tracing: None,
157            explicit_telemetry_chat_middleware: false,
158        }
159    }
160
161    define_arc_option_setter!(with_chat_client, chat_client, dyn AiChatClient);
162
163    pub fn with_chat_middleware<M: ChatClientMiddleware + 'static>(
164        mut self,
165        middleware: M,
166    ) -> Self {
167        self.chat_middlewares.push(Arc::new(middleware));
168        self
169    }
170
171    pub fn with_chat_middleware_arc(mut self, middleware: Arc<dyn ChatClientMiddleware>) -> Self {
172        self.chat_middlewares.push(middleware);
173        self
174    }
175
176    pub fn with_logging_chat_middleware(self) -> Self {
177        self.with_chat_middleware(LoggingChatMiddleware)
178    }
179
180    pub fn with_telemetry_chat_middleware(mut self, metrics: Arc<dyn RuntimeMetrics>) -> Self {
181        self.explicit_telemetry_chat_middleware = true;
182        self.with_chat_middleware(TelemetryChatMiddleware::new(metrics))
183    }
184
185    pub fn with_runtime_telemetry<T: RuntimeTelemetry + 'static>(
186        mut self,
187        telemetry: Arc<T>,
188    ) -> Self {
189        self.runtime_telemetry_metrics = Some(telemetry.clone());
190        self.runtime_telemetry_tracing = Some(telemetry);
191        self
192    }
193
194    #[cfg(feature = "otel")]
195    pub fn with_otel_from_env(self) -> Result<Self> {
196        let telemetry = crate::infrastructure::telemetry::OpenTelemetryTelemetry::from_env()?;
197        Ok(self.with_runtime_telemetry(telemetry))
198    }
199
200    #[cfg(not(feature = "otel"))]
201    pub fn with_otel_from_env(self) -> Result<Self> {
202        Err(crate::domain::errors::StasisError::PortFailure(
203            "OpenTelemetry support requires the `otel` Cargo feature".to_string(),
204        ))
205    }
206
207    pub fn with_cache_chat_middleware(self, cache: Arc<dyn AiChatResponseCache>) -> Self {
208        self.with_chat_middleware(CacheChatMiddleware::new(cache))
209    }
210
211    pub fn with_tool_call_interception_chat_middleware(
212        self,
213        interceptor: Arc<dyn AiChatToolInterceptor>,
214    ) -> Self {
215        self.with_chat_middleware(ToolCallInterceptionChatMiddleware::new(interceptor))
216    }
217
218    define_arc_option_setter!(
219        with_memory_context_reader,
220        memory_context_reader,
221        dyn MemoryContextReader
222    );
223    define_arc_option_setter!(
224        with_memory_context_writer,
225        memory_context_writer,
226        dyn MemoryContextWriter
227    );
228    define_enable_flag_setter!(with_locus_memory, enable_locus_memory);
229    define_arc_option_setter!(
230        with_identity_memory_store,
231        identity_memory_store,
232        dyn IdentityMemoryStore
233    );
234    define_arc_option_setter!(with_memory_operations, memory_operations, dyn MemoryOperations);
235    define_arc_option_setter!(with_thread_store, thread_store, dyn ThreadStore);
236    define_arc_option_setter!(with_cluster_node_store, cluster_node_store, dyn ClusterNodeStore);
237    define_arc_option_setter!(
238        with_delivery_endpoint_store,
239        delivery_endpoint_store,
240        dyn DeliveryEndpointStore
241    );
242    define_arc_option_setter!(
243        with_endpoint_delivery_status_store,
244        endpoint_delivery_status_store,
245        dyn EndpointDeliveryStatusStore
246    );
247
248    pub fn with_endpoint_transport_publisher<P: EndpointTransportPublisher + 'static>(
249        mut self,
250        transport: P,
251    ) -> Self {
252        self.endpoint_transport_publishers.push(Arc::new(transport));
253        self
254    }
255
256    pub fn with_endpoint_transport_publisher_arc(
257        mut self,
258        transport: Arc<dyn EndpointTransportPublisher>,
259    ) -> Self {
260        self.endpoint_transport_publishers.push(transport);
261        self
262    }
263
264    define_enable_flag_setter!(with_endpoint_routing_delivery, enable_endpoint_routing_delivery);
265
266    pub fn with_endpoint_routing_policy<P: EndpointRoutingPolicy + 'static>(
267        mut self,
268        policy: P,
269    ) -> Self {
270        self.endpoint_routing_policy = Some(Arc::new(policy));
271        self
272    }
273
274    define_arc_option_setter!(
275        with_endpoint_routing_policy_arc,
276        endpoint_routing_policy,
277        dyn EndpointRoutingPolicy
278    );
279
280    pub fn with_tool<T: StasisTool + 'static>(self, tool: T) -> Result<Self> {
281        self.tool_registry.register_tool(tool)?;
282        Ok(self)
283    }
284
285    pub fn with_extra_handler<H: JobHandler + 'static>(mut self, handler: H) -> Self {
286        self.extra_handlers.push(Arc::new(handler));
287        self
288    }
289
290    define_disable_flag_setter!(without_grapheme_handlers, include_grapheme_handlers);
291    define_disable_flag_setter!(without_prompt_handler, include_prompt_handler);
292    define_disable_flag_setter!(without_tool_loop_handler, include_tool_loop_handler);
293    define_disable_flag_setter!(without_agent_handlers, include_agent_handlers);
294    define_disable_flag_setter!(
295        without_memory_operation_handlers,
296        include_memory_operation_handlers
297    );
298    define_disable_flag_setter!(
299        without_orchestration_pattern_handlers,
300        include_orchestration_pattern_handlers
301    );
302    define_disable_flag_setter!(without_cluster_control_handlers, include_cluster_control_handlers);
303
304    pub async fn build(self) -> Result<RuntimeComposition> {
305        let mut runtime = RuntimeFactory::build(self.backend).await?;
306        let mut chat_middlewares = self.chat_middlewares;
307
308        if let (Some(metrics), Some(tracing)) = (
309            self.runtime_telemetry_metrics.clone(),
310            self.runtime_telemetry_tracing.clone(),
311        ) {
312            runtime.replace_telemetry(metrics.clone(), tracing.clone());
313            if !self.explicit_telemetry_chat_middleware {
314                chat_middlewares.push(Arc::new(
315                    TelemetryChatMiddleware::new(metrics.clone()).with_tracing(tracing),
316                ));
317            }
318        }
319
320        let workflow_engine = RuntimeFactory::default_workflow_engine();
321        let chat_client = self
322            .chat_client
323            .unwrap_or_else(RuntimeFactory::default_chat_client);
324        let chat_client = Self::compose_chat_client(chat_client, &chat_middlewares);
325        let (memory_context_reader, memory_context_writer, memory_operations) =
326            RuntimeFactory::ensure_locus_memory_adapters(
327                self.enable_locus_memory,
328                self.memory_context_reader,
329                self.memory_context_writer,
330                self.memory_operations,
331            )
332            .await?;
333        let identity_memory_store = self.identity_memory_store;
334        let default_thread_store = self.thread_store.clone();
335        let configured_cluster_store = self.cluster_node_store.clone();
336        let configured_endpoint_store = self.delivery_endpoint_store.clone();
337        let configured_endpoint_status_store = self.endpoint_delivery_status_store.clone();
338        let configured_endpoint_transports = self.endpoint_transport_publishers.clone();
339        let configured_endpoint_routing_policy = self.endpoint_routing_policy.clone();
340
341        let tool_registry = Arc::new(self.tool_registry);
342        let operation_telemetry = self
343            .runtime_telemetry_metrics
344            .as_ref()
345            .and_then(|metrics| {
346                self.runtime_telemetry_tracing.as_ref().map(|tracing| {
347                    OperationTelemetry::new(metrics.clone(), tracing.clone())
348                })
349            });
350
351        match &runtime {
352            RuntimeComposition::InMemory(rt) => {
353                let thread_store =
354                    RuntimeFactory::resolve_thread_store(&runtime, default_thread_store.clone());
355                let cluster_store = RuntimeFactory::resolve_cluster_node_store(
356                    &runtime,
357                    configured_cluster_store.clone(),
358                );
359
360                if self.enable_endpoint_routing_delivery {
361                    let endpoint_store = RuntimeFactory::resolve_delivery_endpoint_store(
362                        &runtime,
363                        configured_endpoint_store.clone(),
364                    );
365                    let status_store = RuntimeFactory::resolve_endpoint_delivery_status_store(
366                        &runtime,
367                        configured_endpoint_status_store.clone(),
368                    );
369
370                    let routing_publisher = RuntimeFactory::build_endpoint_routing_publisher(
371                        endpoint_store,
372                        status_store,
373                        &configured_endpoint_transports,
374                        configured_endpoint_routing_policy.clone(),
375                    );
376
377                    rt.register_event_publisher(routing_publisher)?;
378                }
379
380                if self.include_grapheme_handlers {
381                    rt.register_handler(
382                        GraphemeJobHandler::new(workflow_engine.clone())
383                            .with_operation_telemetry(operation_telemetry.clone()),
384                    )?;
385                    rt.register_handler(GraphemeHealthcheckJobHandler::new(
386                        workflow_engine.clone(),
387                    ))?;
388                    rt.register_handler(GraphemeEchoJobHandler::new(workflow_engine.clone()))?;
389                    rt.register_handler(GraphemeTextOpsJobHandler::new(workflow_engine.clone()))?;
390                }
391
392                if self.include_prompt_handler {
393                    rt.register_handler(PromptChatJobHandler::new_with_memory_and_identity(
394                        chat_client.clone(),
395                        memory_context_reader.clone(),
396                        memory_context_writer.clone(),
397                        identity_memory_store.clone(),
398                    ))?;
399                }
400
401                if self.include_tool_loop_handler {
402                    rt.register_handler(ToolLoopJobHandler::new_with_memory_and_identity(
403                        chat_client.clone(),
404                        tool_registry.clone(),
405                        memory_context_reader.clone(),
406                        memory_context_writer.clone(),
407                        identity_memory_store.clone(),
408                    ))?;
409                }
410
411                if self.include_agent_handlers {
412                    rt.register_handler(AgentTurnJobHandler::new_with_memory_and_identity(
413                        chat_client.clone(),
414                        tool_registry.clone(),
415                        memory_context_reader.clone(),
416                        memory_context_writer.clone(),
417                        identity_memory_store.clone(),
418                    ))?;
419                    rt.register_handler(AgentSessionJobHandler::new_with_memory_and_identity(
420                        chat_client.clone(),
421                        tool_registry.clone(),
422                        memory_context_reader.clone(),
423                        memory_context_writer.clone(),
424                        identity_memory_store.clone(),
425                    ))?;
426                }
427
428                if self.include_memory_operation_handlers {
429                    if let Some(reader) = memory_context_reader.clone() {
430                        rt.register_handler(
431                            MemoryRecallJobHandler::new(reader.clone())
432                                .with_operation_telemetry(operation_telemetry.clone()),
433                        )?;
434                        rt.register_handler(MemoryFindJobHandler::new(reader))?;
435                    }
436                    if let Some(operations) = memory_operations.clone() {
437                        rt.register_handler(MemoryAggregateJobHandler::new(operations.clone()))?;
438                        rt.register_handler(MemoryTransformJobHandler::new(operations.clone()))?;
439                        rt.register_handler(MemoryRollupJobHandler::new(operations.clone()))?;
440                        rt.register_handler(MemorySchemaJobHandler::new(operations))?;
441                    }
442                }
443
444                if self.include_orchestration_pattern_handlers {
445                    rt.register_handler(ConcurrentPatternJobHandler::new_with_thread_store(
446                        chat_client.clone(),
447                        Some(thread_store.clone()),
448                    ))?;
449                    rt.register_handler(HandoffPatternJobHandler::new_with_thread_store(
450                        chat_client.clone(),
451                        Some(thread_store.clone()),
452                    ))?;
453                    rt.register_handler(OrchestratorPatternJobHandler::new_with_thread_store(
454                        chat_client.clone(),
455                        Some(thread_store.clone()),
456                    ))?;
457                    rt.register_handler(SequentialPatternJobHandler::new_with_thread_store(
458                        chat_client.clone(),
459                        Some(thread_store),
460                    ))?;
461                }
462
463                if self.include_cluster_control_handlers {
464                    rt.register_handler(CoordinatorFailoverJobHandler::new(cluster_store.clone()))?;
465                    rt.register_handler(QueueOwnershipRebalanceJobHandler::new(cluster_store))?;
466                }
467
468                for handler in &self.extra_handlers {
469                    rt.register_handler(DelegatingJobHandler {
470                        inner: handler.clone(),
471                    })?;
472                }
473            }
474            RuntimeComposition::Surreal(rt) => {
475                let thread_store =
476                    RuntimeFactory::resolve_thread_store(&runtime, default_thread_store.clone());
477                let cluster_store = RuntimeFactory::resolve_cluster_node_store(
478                    &runtime,
479                    configured_cluster_store.clone(),
480                );
481
482                if self.enable_endpoint_routing_delivery {
483                    let endpoint_store = RuntimeFactory::resolve_delivery_endpoint_store(
484                        &runtime,
485                        configured_endpoint_store.clone(),
486                    );
487                    let status_store = RuntimeFactory::resolve_endpoint_delivery_status_store(
488                        &runtime,
489                        configured_endpoint_status_store.clone(),
490                    );
491
492                    let routing_publisher = RuntimeFactory::build_endpoint_routing_publisher(
493                        endpoint_store,
494                        status_store,
495                        &configured_endpoint_transports,
496                        configured_endpoint_routing_policy.clone(),
497                    );
498
499                    rt.register_event_publisher(routing_publisher)?;
500                }
501
502                if self.include_grapheme_handlers {
503                    rt.register_handler(
504                        GraphemeJobHandler::new(workflow_engine.clone())
505                            .with_operation_telemetry(operation_telemetry.clone()),
506                    )?;
507                    rt.register_handler(GraphemeHealthcheckJobHandler::new(
508                        workflow_engine.clone(),
509                    ))?;
510                    rt.register_handler(GraphemeEchoJobHandler::new(workflow_engine.clone()))?;
511                    rt.register_handler(GraphemeTextOpsJobHandler::new(workflow_engine.clone()))?;
512                }
513
514                if self.include_prompt_handler {
515                    rt.register_handler(PromptChatJobHandler::new_with_memory_and_identity(
516                        chat_client.clone(),
517                        memory_context_reader.clone(),
518                        memory_context_writer.clone(),
519                        identity_memory_store.clone(),
520                    ))?;
521                }
522
523                if self.include_tool_loop_handler {
524                    rt.register_handler(ToolLoopJobHandler::new_with_memory_and_identity(
525                        chat_client.clone(),
526                        tool_registry.clone(),
527                        memory_context_reader.clone(),
528                        memory_context_writer.clone(),
529                        identity_memory_store.clone(),
530                    ))?;
531                }
532
533                if self.include_agent_handlers {
534                    rt.register_handler(AgentTurnJobHandler::new_with_memory_and_identity(
535                        chat_client.clone(),
536                        tool_registry.clone(),
537                        memory_context_reader.clone(),
538                        memory_context_writer.clone(),
539                        identity_memory_store.clone(),
540                    ))?;
541                    rt.register_handler(AgentSessionJobHandler::new_with_memory_and_identity(
542                        chat_client.clone(),
543                        tool_registry.clone(),
544                        memory_context_reader.clone(),
545                        memory_context_writer.clone(),
546                        identity_memory_store.clone(),
547                    ))?;
548                }
549
550                if self.include_memory_operation_handlers {
551                    if let Some(reader) = memory_context_reader.clone() {
552                        rt.register_handler(
553                            MemoryRecallJobHandler::new(reader.clone())
554                                .with_operation_telemetry(operation_telemetry.clone()),
555                        )?;
556                        rt.register_handler(MemoryFindJobHandler::new(reader))?;
557                    }
558                    if let Some(operations) = memory_operations.clone() {
559                        rt.register_handler(MemoryAggregateJobHandler::new(operations.clone()))?;
560                        rt.register_handler(MemoryTransformJobHandler::new(operations.clone()))?;
561                        rt.register_handler(MemoryRollupJobHandler::new(operations.clone()))?;
562                        rt.register_handler(MemorySchemaJobHandler::new(operations))?;
563                    }
564                }
565
566                if self.include_orchestration_pattern_handlers {
567                    rt.register_handler(ConcurrentPatternJobHandler::new_with_thread_store(
568                        chat_client.clone(),
569                        Some(thread_store.clone()),
570                    ))?;
571                    rt.register_handler(HandoffPatternJobHandler::new_with_thread_store(
572                        chat_client.clone(),
573                        Some(thread_store.clone()),
574                    ))?;
575                    rt.register_handler(OrchestratorPatternJobHandler::new_with_thread_store(
576                        chat_client.clone(),
577                        Some(thread_store.clone()),
578                    ))?;
579                    rt.register_handler(SequentialPatternJobHandler::new_with_thread_store(
580                        chat_client.clone(),
581                        Some(thread_store),
582                    ))?;
583                }
584
585                if self.include_cluster_control_handlers {
586                    rt.register_handler(CoordinatorFailoverJobHandler::new(cluster_store.clone()))?;
587                    rt.register_handler(QueueOwnershipRebalanceJobHandler::new(cluster_store))?;
588                }
589
590                for handler in &self.extra_handlers {
591                    rt.register_handler(DelegatingJobHandler {
592                        inner: handler.clone(),
593                    })?;
594                }
595            }
596        }
597
598        Ok(runtime)
599    }
600
601    /// Builds and returns the primary runtime facade.
602    pub async fn build_stasis_runtime(self) -> Result<crate::sdk::runtime_sdk::StasisRuntime> {
603        let composition = self.build().await?;
604        Ok(crate::sdk::runtime_sdk::RuntimeSdk::new(composition))
605    }
606
607    fn compose_chat_client(
608        chat_client: Arc<dyn AiChatClient>,
609        middlewares: &[Arc<dyn ChatClientMiddleware>],
610    ) -> Arc<dyn AiChatClient> {
611        let mut wrapped = chat_client;
612        for middleware in middlewares.iter().rev() {
613            wrapped = middleware.wrap(wrapped);
614        }
615        wrapped
616    }
617}
618
619#[cfg(test)]
620mod tests {
621    use std::sync::{Arc, RwLock};
622
623    use async_trait::async_trait;
624    use chrono::Utc;
625
626    use crate::application::runtime::in_memory_runtime::{JobExecutionOutcome, JobHandler};
627    use crate::application::runtime::runtime_factory::{RuntimeBackend, RuntimeComposition};
628    use crate::domain::errors::{Result, StasisError};
629    use crate::domain::runtime::delivery_endpoint::{
630        DeliveryEndpoint, DeliveryProtocol, NewDeliveryEndpoint,
631    };
632    use crate::domain::runtime::job::{BackoffPolicy, NewJob};
633    use crate::domain::runtime::outbox::OutboxEvent;
634    use crate::infrastructure::runtime::in_memory_delivery_endpoint_store::InMemoryDeliveryEndpointStore;
635    use crate::ports::outbound::runtime::delivery_endpoint_store::DeliveryEndpointStore;
636    use crate::ports::outbound::runtime::endpoint_transport_publisher::EndpointTransportPublisher;
637
638    use super::StasisRuntimeBuilder;
639
640    #[derive(Clone)]
641    struct SuccessHandler;
642
643    #[async_trait]
644    impl JobHandler for SuccessHandler {
645        fn job_type(&self) -> &'static str {
646            "test.success"
647        }
648
649        async fn execute(
650            &self,
651            _job: &crate::domain::runtime::job::Job,
652        ) -> Result<JobExecutionOutcome> {
653            Ok(JobExecutionOutcome::Success {
654                sttp_output_node_id: "sttp:out:test".to_string(),
655                execution_id: Some("exec:test".to_string()),
656                diagnostics: None,
657            })
658        }
659    }
660
661    #[derive(Clone)]
662    struct RecordingTransport {
663        calls: Arc<RwLock<Vec<String>>>,
664    }
665
666    #[async_trait]
667    impl EndpointTransportPublisher for RecordingTransport {
668        fn supports(&self, protocol: &DeliveryProtocol) -> bool {
669            matches!(protocol, DeliveryProtocol::HttpWebhook)
670        }
671
672        async fn publish_to_endpoint(
673            &self,
674            endpoint: &DeliveryEndpoint,
675            _event: &OutboxEvent,
676        ) -> Result<()> {
677            let mut calls = self
678                .calls
679                .write()
680                .map_err(|_| StasisError::PortFailure("calls lock poisoned".to_string()))?;
681            calls.push(endpoint.endpoint_id.clone());
682            Ok(())
683        }
684    }
685
686    #[tokio::test]
687    async fn builder_wires_endpoint_routing_delivery_for_in_memory_runtime() {
688        let endpoint_store = InMemoryDeliveryEndpointStore::default();
689        endpoint_store
690            .insert(NewDeliveryEndpoint {
691                endpoint_id: "endpoint.webhook.builder".to_string(),
692                name: "Builder Webhook".to_string(),
693                protocol: DeliveryProtocol::HttpWebhook,
694                target: "https://example.com/hook".to_string(),
695                metadata: None,
696                created_at: Utc::now(),
697            })
698            .await
699            .expect("endpoint should insert");
700
701        let calls = Arc::new(RwLock::new(Vec::new()));
702        let runtime = StasisRuntimeBuilder::new(RuntimeBackend::InMemory)
703            .with_delivery_endpoint_store(Arc::new(endpoint_store))
704            .with_endpoint_transport_publisher(RecordingTransport {
705                calls: Arc::clone(&calls),
706            })
707            .with_endpoint_routing_delivery()
708            .with_extra_handler(SuccessHandler)
709            .without_grapheme_handlers()
710            .without_prompt_handler()
711            .without_tool_loop_handler()
712            .without_agent_handlers()
713            .without_memory_operation_handlers()
714            .without_orchestration_pattern_handlers()
715            .build()
716            .await
717            .expect("runtime should build");
718
719        let RuntimeComposition::InMemory(rt) = runtime else {
720            panic!("expected in-memory runtime composition");
721        };
722
723        let now = Utc::now();
724        rt.enqueue(NewJob {
725            id: "job-builder-routing".to_string(),
726            queue: "default".to_string(),
727            job_type: "test.success".to_string(),
728            payload_ref: "sttp:in:test".to_string(),
729            priority: 100,
730            max_attempts: 1,
731            idempotency_key: "idem-builder-routing".to_string(),
732            correlation_id: "corr-builder-routing".to_string(),
733            causation_id: "cause-builder-routing".to_string(),
734            trace_id: "trace-builder-routing".to_string(),
735            sttp_input_node_id: "sttp:in:test".to_string(),
736            scheduled_at: now,
737            backoff_policy: BackoffPolicy::default(),
738        })
739        .await
740        .expect("job should enqueue");
741
742        rt.process_once("default", "worker-builder", now)
743            .await
744            .expect("process should succeed");
745
746        let published = rt
747            .publish_pending_events(10, now)
748            .await
749            .expect("publish should succeed");
750        assert_eq!(published, 1);
751
752        let calls = calls.read().expect("calls read lock should succeed");
753        assert_eq!(calls.len(), 1);
754        assert_eq!(calls[0], "endpoint.webhook.builder");
755    }
756}