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