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