Skip to main content

meerkat_mobkit/unified_runtime/
http.rs

1//! HTTP server and route assembly for the unified runtime.
2
3use 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
24/// Default cap for the stock reference app router.
25///
26/// SSE routes keep HTTP requests open for the lifetime of the subscription, so
27/// a demo-scale cap such as 20 can make the console look frozen while streams
28/// occupy all slots. Hosts that wrap MobKit in their own axum service may still
29/// choose a different outer limit, but the reference router itself defaults to
30/// a ceiling high enough for real console usage.
31pub 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        // Every SSE route shares the same `RuntimeDecisionState` the
83        // console RPC route uses: when `require_app_auth` is on, requests
84        // must carry a valid bearer / auth_token. Pre-fix, only the
85        // structural-events route gated; tier-2 (`/agents/{id}/events`)
86        // and tier-3 (`/mob/events`) shipped unauthenticated.
87        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}