meerkat_mobkit/unified_runtime/
http.rs1use std::sync::Arc;
4
5use axum::Router;
6use axum::routing::get;
7
8use crate::console_aggregator::{
9 ConsoleVisibilityPolicy, HideImplicitDelegateMembersConsoleVisibilityPolicy,
10};
11use crate::http_console::{
12 console_frontend_router, console_json_router_with_runtime_events_and_policy,
13};
14use crate::http_flow_editor::protected_flow_editor_router_with_runtime_catalog;
15use crate::http_sse::{
16 agent_events_sse_router_with_access_and_priming, mob_events_sse_router_with_access_and_priming,
17 mob_structural_events_sse_router_with_access_and_priming,
18};
19use crate::runtime::RuntimeDecisionState;
20use tower::limit::ConcurrencyLimitLayer;
21
22use super::UnifiedRuntime;
23
24pub const DEFAULT_REFERENCE_APP_MAX_CONCURRENT_REQUESTS: usize = 1024;
32
33impl UnifiedRuntime {
34 pub fn build_console_json_router(&self, decisions: RuntimeDecisionState) -> Router {
35 self.build_console_json_router_with_policy(
36 decisions,
37 Arc::new(HideImplicitDelegateMembersConsoleVisibilityPolicy),
38 )
39 }
40
41 pub fn build_console_json_router_with_policy(
42 &self,
43 decisions: RuntimeDecisionState,
44 visibility_policy: Arc<dyn ConsoleVisibilityPolicy>,
45 ) -> Router {
46 console_json_router_with_runtime_events_and_policy(
47 decisions,
48 self.mob_runtime.clone(),
49 Some(self.module_runtime_handle()),
50 self.contact_directory.clone(),
51 self.event_log_store(),
52 self.gateway_peer_keys().cloned(),
53 Some(self.console_events()),
54 Some(self.console_log_store()),
55 Some(self.mob_events_store()),
56 Some(Arc::clone(self.metadata_table())),
57 self.identity_runtime().cloned(),
58 visibility_policy,
59 self.access_controller().cloned(),
60 self.memory_panel_store(),
61 )
62 }
63
64 pub fn build_console_frontend_router(&self) -> Router {
65 console_frontend_router()
66 }
67
68 pub fn build_reference_app_router(&self, decisions: RuntimeDecisionState) -> Router {
69 self.build_reference_app_router_with_console_visibility_policy(
70 decisions,
71 Arc::new(HideImplicitDelegateMembersConsoleVisibilityPolicy),
72 )
73 }
74
75 pub fn build_reference_app_router_with_console_visibility_policy(
76 &self,
77 decisions: RuntimeDecisionState,
78 visibility_policy: Arc<dyn ConsoleVisibilityPolicy>,
79 ) -> Router {
80 let agent_runtime = self.mob_runtime.clone();
81 let mob_runtime = self.mob_runtime.clone();
82 let flow_editor_decisions = decisions.clone();
88 let flow_editor_runtime_catalog = self.mobpack_runtime_catalog_state_snapshot();
89 let sse_decisions_a = decisions.clone();
90 let sse_decisions_b = decisions.clone();
91 let sse_decisions_c = decisions.clone();
92 let access = self.access_controller().cloned();
93 Router::new()
94 .route("/healthz", get(|| async { "ok" }))
95 .merge(self.build_console_frontend_router())
96 .merge(protected_flow_editor_router_with_runtime_catalog(
97 flow_editor_decisions,
98 flow_editor_runtime_catalog,
99 access.clone(),
100 ))
101 .merge(self.build_console_json_router_with_policy(decisions, visibility_policy))
102 .merge(agent_events_sse_router_with_access_and_priming(
103 Arc::new(move |agent_id| {
104 let runtime = agent_runtime.clone();
105 Box::pin(async move {
106 runtime
107 .handle()
108 .subscribe_agent_events(&crate::member_comms_id::mob_member_id(
109 &agent_id,
110 ))
111 .await
112 .map_err(Into::into)
113 })
114 }),
115 Some(sse_decisions_a),
116 access.clone(),
117 access.as_ref().map(|_| self.mob_runtime.clone()),
118 ))
119 .merge(mob_events_sse_router_with_access_and_priming(
120 Arc::new(move || {
121 let mob_runtime = mob_runtime.clone();
122 Box::pin(async move { mob_runtime.handle().subscribe_mob_events().await })
123 }),
124 Some(sse_decisions_b),
125 access.clone(),
126 access.as_ref().map(|_| self.mob_runtime.clone()),
127 ))
128 .merge(mob_structural_events_sse_router_with_access_and_priming(
129 self.mob_runtime.handle(),
130 self.mob_events_store(),
131 Some(sse_decisions_c),
132 access.clone(),
133 access.as_ref().map(|_| self.mob_runtime.clone()),
134 ))
135 .layer(ConcurrencyLimitLayer::new(
136 DEFAULT_REFERENCE_APP_MAX_CONCURRENT_REQUESTS,
137 ))
138 }
139}
140
141#[cfg(test)]
142mod tests {
143 use super::DEFAULT_REFERENCE_APP_MAX_CONCURRENT_REQUESTS;
144
145 #[test]
146 fn reference_router_default_concurrency_allows_sse_fanout() {
147 assert_eq!(DEFAULT_REFERENCE_APP_MAX_CONCURRENT_REQUESTS, 1024);
148 }
149}