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 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}