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