1use std::convert::Infallible;
4use std::future::Future;
5use std::pin::Pin;
6use std::sync::Arc;
7use std::time::Duration;
8
9use async_stream::stream;
10use axum::extract::{Path, Query, State};
11use axum::http::{HeaderMap, StatusCode, Uri, header};
12use axum::response::IntoResponse;
13use axum::response::sse::{Event, KeepAlive, Sse};
14use axum::routing::get;
15use axum::{Json, Router};
16use futures::StreamExt;
17use meerkat_core::AgentEvent;
18use meerkat_core::comms::EventStream;
19use meerkat_core::event::agent_event_type;
20use meerkat_mob::{MobEventRouterHandle, MobHandle};
21use serde::Deserialize;
22use serde_json::{Value, json};
23
24use crate::access::{ACTION_AGENT_VIEW, ACTION_MOB_OBSERVE, AccessController, AccessView};
25use crate::runtime::{RuntimeDecisionState, extract_bearer_token_from_header};
26use crate::unified_runtime::EventQuery;
27use crate::unified_runtime::mob_events::{MOB_EVENTS_STREAM_PATH, MobEventsStore};
28
29use crate::mob_handle_runtime::{MobRuntime, MobRuntimeError};
30use meerkat_core::comms::SendError;
31use meerkat_core::service::SessionError;
32use meerkat_mob::MobError;
33
34pub(crate) const DEFAULT_KEEP_ALIVE_INTERVAL: Duration = Duration::from_secs(15);
35pub(crate) const KEEP_ALIVE_TEXT: &str = "keep-alive";
36
37pub(crate) use crate::mob_handle_runtime::console_agent_event_payload;
38
39pub fn agent_event_sse(interaction_id: &str, seq: u64, event: &AgentEvent) -> Event {
40 let event_name = agent_event_name(event);
41 let payload = serde_json::to_string(&console_agent_event_payload(event))
42 .unwrap_or_else(|_| "{}".to_string());
43 Event::default()
44 .id(format!("{interaction_id}:{seq}"))
45 .event(event_name)
46 .data(payload)
47}
48
49fn agent_event_name(event: &AgentEvent) -> String {
50 serde_json::to_value(event)
51 .ok()
52 .and_then(|value| {
53 value
54 .as_object()
55 .and_then(|object| object.get("type"))
56 .and_then(Value::as_str)
57 .map(ToString::to_string)
58 })
59 .unwrap_or_else(|| "agent_event".to_string())
60}
61
62fn http_error(status: StatusCode, message: &str) -> (StatusCode, Json<Value>) {
63 (
64 status,
65 Json(json!({
66 "error": message
67 })),
68 )
69}
70
71fn map_runtime_error(error: MobRuntimeError) -> (StatusCode, Json<Value>) {
72 match error {
73 MobRuntimeError::InvalidInput(message) => http_error(StatusCode::BAD_REQUEST, message),
74 MobRuntimeError::Mob(
75 MobError::MemberNotFound(_)
76 | MobError::SessionError(SessionError::NotFound { .. })
77 | MobError::CommsError(SendError::PeerNotFound(_)),
78 ) => http_error(StatusCode::NOT_FOUND, "member_not_found"),
79 MobRuntimeError::Mob(MobError::SessionError(SessionError::Unsupported(_))) => {
80 http_error(StatusCode::UNPROCESSABLE_ENTITY, "unsupported")
81 }
82 _ => http_error(StatusCode::INTERNAL_SERVER_ERROR, "internal_server_error"),
83 }
84}
85
86pub type AgentEventSubscribeFuture =
91 Pin<Box<dyn Future<Output = Result<EventStream, MobRuntimeError>> + Send>>;
92
93pub type AgentEventSubscribeFn = Arc<dyn Fn(String) -> AgentEventSubscribeFuture + Send + Sync>;
94
95#[derive(Clone)]
96struct AgentSseState {
97 subscribe_fn: AgentEventSubscribeFn,
98 decisions: Option<RuntimeDecisionState>,
99 access: Option<AccessController>,
100 prime_runtime: Option<MobRuntime>,
105}
106
107pub fn agent_events_sse_router(
108 subscribe_fn: AgentEventSubscribeFn,
109 decisions: Option<RuntimeDecisionState>,
110) -> Router {
111 agent_events_sse_router_with_access(subscribe_fn, decisions, None)
112}
113
114pub fn agent_events_sse_router_with_access(
115 subscribe_fn: AgentEventSubscribeFn,
116 decisions: Option<RuntimeDecisionState>,
117 access: Option<AccessController>,
118) -> Router {
119 agent_events_sse_router_with_access_and_priming(subscribe_fn, decisions, access, None)
120}
121
122pub(crate) fn agent_events_sse_router_with_access_and_priming(
123 subscribe_fn: AgentEventSubscribeFn,
124 decisions: Option<RuntimeDecisionState>,
125 access: Option<AccessController>,
126 prime_runtime: Option<MobRuntime>,
127) -> Router {
128 Router::new()
129 .route("/agents/{agent_id}/events", get(agent_events_sse_handler))
130 .with_state(AgentSseState {
131 subscribe_fn,
132 decisions,
133 access,
134 prime_runtime,
135 })
136}
137
138async fn agent_events_sse_handler(
139 State(state): State<AgentSseState>,
140 headers: HeaderMap,
141 uri: Uri,
142 Path(agent_id): Path<String>,
143) -> Result<impl IntoResponse, (StatusCode, Json<Value>)> {
144 let access_view = sse_access_context(
145 state.decisions.as_ref(),
146 state.access.as_ref(),
147 &headers,
148 &uri,
149 )
150 .map_err(|()| sse_unauthorized("agent events stream requires a valid auth token"))?;
151 prime_sse_access_cache(state.prime_runtime.as_ref(), state.access.as_ref()).await;
152 if access_view
153 .as_ref()
154 .is_some_and(|view| view.enforced() && !view.allows_agent(ACTION_AGENT_VIEW, &agent_id))
155 {
156 return Err(sse_access_denied(ACTION_AGENT_VIEW));
157 }
158 let agent_id = agent_id.trim().to_string();
159 if agent_id.is_empty() {
160 return Err(http_error(
161 StatusCode::BAD_REQUEST,
162 "agent_id must not be empty",
163 ));
164 }
165
166 let event_stream = (state.subscribe_fn)(agent_id.clone())
167 .await
168 .map_err(map_runtime_error)?;
169
170 let stream = stream! {
171 let mut seq = 0_u64;
172 tokio::pin!(event_stream);
173 while let Some(envelope) = event_stream.next().await {
174 let event_name = agent_event_type(&envelope.payload).to_string();
175 let payload = serde_json::to_string(&console_agent_event_payload(&envelope.payload))
176 .unwrap_or_else(|_| "{}".to_string());
177 yield Ok::<Event, Infallible>(
178 Event::default()
179 .id(format!("{agent_id}:{seq}"))
180 .event(event_name)
181 .data(payload),
182 );
183 seq += 1;
184 }
185 };
186
187 Ok(Sse::new(stream).keep_alive(
188 KeepAlive::new()
189 .interval(DEFAULT_KEEP_ALIVE_INTERVAL)
190 .text(KEEP_ALIVE_TEXT),
191 ))
192}
193
194pub type MobEventSubscribeFuture =
201 Pin<Box<dyn Future<Output = Result<MobEventRouterHandle, meerkat_mob::MobError>> + Send>>;
202
203pub type MobEventSubscribeFn = Arc<dyn Fn() -> MobEventSubscribeFuture + Send + Sync>;
204
205#[derive(Clone)]
206struct MobSseState {
207 subscribe_fn: MobEventSubscribeFn,
208 decisions: Option<RuntimeDecisionState>,
209 access: Option<AccessController>,
210 prime_runtime: Option<MobRuntime>,
212}
213
214pub fn mob_events_sse_router(
215 subscribe_fn: MobEventSubscribeFn,
216 decisions: Option<RuntimeDecisionState>,
217) -> Router {
218 mob_events_sse_router_with_access(subscribe_fn, decisions, None)
219}
220
221pub fn mob_events_sse_router_with_access(
222 subscribe_fn: MobEventSubscribeFn,
223 decisions: Option<RuntimeDecisionState>,
224 access: Option<AccessController>,
225) -> Router {
226 mob_events_sse_router_with_access_and_priming(subscribe_fn, decisions, access, None)
227}
228
229pub(crate) fn mob_events_sse_router_with_access_and_priming(
230 subscribe_fn: MobEventSubscribeFn,
231 decisions: Option<RuntimeDecisionState>,
232 access: Option<AccessController>,
233 prime_runtime: Option<MobRuntime>,
234) -> Router {
235 Router::new()
236 .route("/mob/events", get(mob_events_sse_handler))
237 .with_state(MobSseState {
238 subscribe_fn,
239 decisions,
240 access,
241 prime_runtime,
242 })
243}
244
245async fn prime_sse_access_cache(
250 prime_runtime: Option<&MobRuntime>,
251 access: Option<&AccessController>,
252) {
253 if let (Some(runtime), Some(controller)) = (
254 prime_runtime,
255 access.filter(|controller| controller.enabled()),
256 ) {
257 crate::http_console::prime_access_cache_from_runtime(runtime, controller).await;
258 }
259}
260
261async fn mob_events_sse_handler(
262 State(state): State<MobSseState>,
263 headers: HeaderMap,
264 uri: Uri,
265) -> Result<impl IntoResponse, (StatusCode, Json<Value>)> {
266 let access_view = sse_access_context(
267 state.decisions.as_ref(),
268 state.access.as_ref(),
269 &headers,
270 &uri,
271 )
272 .map_err(|()| sse_unauthorized("mob events stream requires a valid auth token"))?;
273 prime_sse_access_cache(state.prime_runtime.as_ref(), state.access.as_ref()).await;
274 if access_view
281 .as_ref()
282 .is_some_and(|view| view.enforced() && !view.allows(ACTION_MOB_OBSERVE))
283 {
284 return Err(sse_access_denied(ACTION_MOB_OBSERVE));
285 }
286 let stream_view = access_view.filter(AccessView::enforced);
287 let reprime_runtime = state.prime_runtime.clone();
291 let reprime_access = state.access.clone();
292 let mut router_handle = (state.subscribe_fn)().await.map_err(|err| {
293 (
294 StatusCode::INTERNAL_SERVER_ERROR,
295 Json(json!({"error": format!("mob event subscription failed: {err}")})),
296 )
297 })?;
298
299 let stream = stream! {
300 let mut seq = 0_u64;
301 let mut reprimed: std::collections::HashSet<String> = std::collections::HashSet::new();
305 while let Some(attributed) = router_handle.event_rx.recv().await {
306 let source = crate::member_comms_id::runtime_event_alias(&attributed.source);
311 if let Some(view) = stream_view.as_ref()
317 && !view.knows_agent(&source)
318 && reprimed.insert(source.clone())
319 {
320 prime_sse_access_cache(reprime_runtime.as_ref(), reprime_access.as_ref()).await;
321 }
322 if stream_view
323 .as_ref()
324 .is_some_and(|view| !view.can_view_agent(&source))
325 {
326 continue;
327 }
328 let event_name = agent_event_type(&attributed.envelope.payload).to_string();
329 let data = json!({
330 "member_id": &source,
331 "source": &source,
332 "payload": console_agent_event_payload(&attributed.envelope.payload),
333 });
334 yield Ok::<Event, Infallible>(
335 Event::default()
336 .id(format!("mob:{seq}"))
337 .event(event_name)
338 .data(data.to_string()),
339 );
340 seq += 1;
341 }
342 };
343
344 Ok(Sse::new(stream).keep_alive(
345 KeepAlive::new()
346 .interval(DEFAULT_KEEP_ALIVE_INTERVAL)
347 .text(KEEP_ALIVE_TEXT),
348 ))
349}
350
351#[derive(Debug, Default, Deserialize)]
358pub struct MobStructuralStreamQuery {
359 #[serde(default)]
360 pub after_seq: Option<u64>,
361 #[serde(default)]
362 pub mob_id: Option<String>,
363 #[serde(default)]
364 pub run_id: Option<String>,
365 #[serde(default)]
366 pub step_id: Option<String>,
367 #[serde(default)]
368 pub identity: Option<String>,
369 #[serde(default)]
370 pub member_id: Option<String>,
371 #[serde(default)]
374 pub event_types: Option<String>,
375 #[serde(default)]
376 pub since_ms: Option<u64>,
377 #[serde(default)]
378 pub until_ms: Option<u64>,
379}
380
381impl MobStructuralStreamQuery {
382 fn into_event_query(self) -> EventQuery {
383 EventQuery {
384 since_ms: self.since_ms,
385 until_ms: self.until_ms,
386 member_id: self.member_id,
387 identity: self.identity,
388 mob_id: self.mob_id,
389 run_id: self.run_id,
390 step_id: self.step_id,
391 event_types: self
392 .event_types
393 .map(|raw| {
394 raw.split(',')
395 .map(str::trim)
396 .filter(|s| !s.is_empty())
397 .map(ToString::to_string)
398 .collect()
399 })
400 .unwrap_or_default(),
401 limit: None,
402 after_seq: self.after_seq,
403 }
404 }
405}
406
407#[derive(Clone)]
408struct MobStructuralSseState {
409 handle: MobHandle,
410 store: MobEventsStore,
411 prime_runtime: Option<MobRuntime>,
414 decisions: Option<RuntimeDecisionState>,
419 access: Option<AccessController>,
420}
421
422pub fn mob_structural_events_sse_router(
437 handle: MobHandle,
438 store: MobEventsStore,
439 decisions: Option<RuntimeDecisionState>,
440) -> Router {
441 mob_structural_events_sse_router_with_access(handle, store, decisions, None)
442}
443
444pub fn mob_structural_events_sse_router_with_access(
445 handle: MobHandle,
446 store: MobEventsStore,
447 decisions: Option<RuntimeDecisionState>,
448 access: Option<AccessController>,
449) -> Router {
450 mob_structural_events_sse_router_with_access_and_priming(handle, store, decisions, access, None)
451}
452
453pub(crate) fn mob_structural_events_sse_router_with_access_and_priming(
454 handle: MobHandle,
455 store: MobEventsStore,
456 decisions: Option<RuntimeDecisionState>,
457 access: Option<AccessController>,
458 prime_runtime: Option<MobRuntime>,
459) -> Router {
460 Router::new()
461 .route(
462 MOB_EVENTS_STREAM_PATH,
463 get(mob_structural_events_sse_handler),
464 )
465 .with_state(MobStructuralSseState {
466 handle,
467 store,
468 prime_runtime,
469 decisions,
470 access,
471 })
472}
473
474pub(crate) fn sse_access_context(
485 decisions: Option<&RuntimeDecisionState>,
486 access: Option<&AccessController>,
487 headers: &HeaderMap,
488 uri: &Uri,
489) -> Result<Option<AccessView>, ()> {
490 let Some(decisions) = decisions else {
491 return Ok(access.map(|controller| controller.view_for_subject(None)));
492 };
493 let bearer_token = headers
494 .get(header::AUTHORIZATION)
495 .and_then(|v| v.to_str().ok())
496 .and_then(extract_bearer_token_from_header)
497 .map(String::from);
498 let query_token = uri.query().and_then(|q| {
503 form_urlencoded::parse(q.as_bytes())
504 .find(|(key, _)| key == "auth_token")
505 .map(|(_, value)| value.into_owned())
506 });
507 let token = bearer_token.or(query_token);
508 if !decisions.console.require_app_auth {
509 let subject = token.as_deref().and_then(|token| {
512 crate::runtime::resolve_authorized_console_auth_from_token(decisions, token)
513 .map(|auth| auth.email)
514 });
515 return Ok(access.map(|controller| controller.view_for_subject(subject.as_deref())));
516 }
517 let token = token.ok_or(())?;
518 let auth =
519 crate::runtime::resolve_authorized_console_auth_from_token(decisions, &token).ok_or(())?;
520 Ok(access.map(|controller| controller.view_for_subject(Some(auth.email.as_str()))))
521}
522
523fn sse_unauthorized(reason: &str) -> (StatusCode, Json<Value>) {
524 (
525 StatusCode::UNAUTHORIZED,
526 Json(json!({
527 "error": "unauthorized",
528 "reason": reason,
529 })),
530 )
531}
532
533fn sse_access_denied(action: &str) -> (StatusCode, Json<Value>) {
534 (
535 StatusCode::FORBIDDEN,
536 Json(json!({
537 "error": "access_denied",
538 "action": action,
539 })),
540 )
541}
542
543async fn mob_structural_events_sse_handler(
544 State(state): State<MobStructuralSseState>,
545 headers: HeaderMap,
546 uri: Uri,
547 Query(params): Query<MobStructuralStreamQuery>,
548) -> Result<Sse<impl futures::Stream<Item = Result<Event, Infallible>>>, (StatusCode, Json<Value>)>
549{
550 let access_view = sse_access_context(
551 state.decisions.as_ref(),
552 state.access.as_ref(),
553 &headers,
554 &uri,
555 )
556 .map_err(|()| sse_unauthorized("mob_events stream requires a valid auth token"))?;
557 if state.prime_runtime.is_some() {
558 prime_sse_access_cache(state.prime_runtime.as_ref(), state.access.as_ref()).await;
559 } else if let Some(controller) = state
560 .access
561 .as_ref()
562 .filter(|controller| controller.enabled())
563 {
564 crate::http_console::prime_access_cache_from_handle(&state.handle, controller).await;
565 }
566 if access_view
573 .as_ref()
574 .is_some_and(|view| view.enforced() && !view.allows(ACTION_MOB_OBSERVE))
575 {
576 return Err(sse_access_denied(ACTION_MOB_OBSERVE));
577 }
578 let stream_view = access_view.filter(AccessView::enforced);
579 let query = params.into_event_query();
580 let events_view = state.handle.events();
581
582 let latest = events_view.latest_cursor().await.map_err(|err| {
583 (
584 StatusCode::INTERNAL_SERVER_ERROR,
585 Json(json!({
586 "error": "events_view_unavailable",
587 "detail": err.to_string(),
588 })),
589 )
590 })?;
591
592 if let Some(after_seq) = query.after_seq
593 && after_seq > latest
594 {
595 return Err((
596 StatusCode::GONE,
597 Json(json!({
598 "error": "event_query_stale",
599 "after_cursor": after_seq,
600 "latest_cursor": latest,
601 })),
602 ));
603 }
604
605 let after_cursor = query.after_seq.unwrap_or(latest);
606 let mut subscription = events_view
607 .subscribe_after(after_cursor)
608 .await
609 .map_err(|err| {
610 (
611 StatusCode::INTERNAL_SERVER_ERROR,
612 Json(json!({
613 "error": "subscribe_failed",
614 "detail": err.to_string(),
615 })),
616 )
617 })?;
618
619 let store = state.store;
620 let reprime_runtime = state.prime_runtime.clone();
623 let reprime_handle = state.handle.clone();
624 let reprime_access = state.access.clone();
625 let stream = stream! {
626 let mut reprimed: std::collections::HashSet<String> = std::collections::HashSet::new();
629 while let Some(event) = subscription.event_rx.recv().await {
630 let envelope = store.project_event_for_query(&event).await;
631 if !crate::unified_runtime::mob_events::envelope_matches(&envelope, &query) {
632 continue;
633 }
634 if let Some(identity) = envelope.agent_identity.as_deref() {
638 if let Some(view) = stream_view.as_ref()
644 && !view.knows_agent(identity)
645 && reprimed.insert(identity.to_string())
646 {
647 if reprime_runtime.is_some() {
648 prime_sse_access_cache(reprime_runtime.as_ref(), reprime_access.as_ref())
649 .await;
650 } else if let Some(controller) = reprime_access
651 .as_ref()
652 .filter(|controller| controller.enabled())
653 {
654 crate::http_console::prime_access_cache_from_handle(
655 &reprime_handle,
656 controller,
657 )
658 .await;
659 }
660 }
661 if stream_view
662 .as_ref()
663 .is_some_and(|view| !view.can_view_agent(identity))
664 {
665 continue;
666 }
667 }
668 let payload = serde_json::to_string(&envelope).unwrap_or_else(|_| "{}".to_string());
669 yield Ok::<Event, Infallible>(
670 Event::default()
671 .id(format!("mob-evt-{}", envelope.cursor))
672 .event(envelope.kind.clone())
673 .data(payload),
674 );
675 }
676 };
677
678 Ok(Sse::new(stream).keep_alive(
679 KeepAlive::new()
680 .interval(DEFAULT_KEEP_ALIVE_INTERVAL)
681 .text(KEEP_ALIVE_TEXT),
682 ))
683}