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