Skip to main content

meerkat_mobkit/
http_sse.rs

1//! Server-Sent Events (SSE) streaming endpoints for agent and mob observation.
2
3use 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
86// ---------------------------------------------------------------------------
87// Tier 2: Per-agent persistent SSE  (MK-005)
88// ---------------------------------------------------------------------------
89
90pub 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    /// Live mob runtime used to prime the access attribute cache at
101    /// connection time (roster plus spawn-registered console metadata), so
102    /// label/role/lineage rules resolve without a prior
103    /// `/console/experience` call. `None` keeps the route behaviour unchanged.
104    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
194// ---------------------------------------------------------------------------
195// Tier 3: Mob-merged SSE  (MK-006)
196// ---------------------------------------------------------------------------
197
198/// Meerkat 0.7: mob event-router subscription is fallible (machine command
199/// faults surface as `MobError` instead of panicking inside the router).
200pub 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    /// See [`AgentSseState::prime_runtime`].
211    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
245/// Refresh the access attribute cache from the live roster and the spawn
246/// registry before an SSE stream applies any per-agent filter, so
247/// label/role/lineage rules resolve without depending on a prior
248/// `/console/experience` call.
249async 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    // `mob.observe` gates access to the merged stream surface. The events
275    // flowing through it carry the same rich per-agent payload as
276    // `/agents/{id}/events`, so each one is still filtered by `agent.view`
277    // on its source — otherwise a `mob.observe` grant would silently defeat
278    // a per-agent view denial. "Observe everything" is expressed by also
279    // granting `agent.view` on `*`.
280    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    // Captured so the long-lived stream can re-prime the shared attribute cache
288    // for members spawned AFTER the one-time subscribe prime (the view reads
289    // attributes through the controller's shared cache).
290    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        // Agents we've already attempted a cache re-prime for, so an agent that
302        // genuinely has no roster attributes does not trigger a re-prime on
303        // every event.
304        let mut reprimed: std::collections::HashSet<String> = std::collections::HashSet::new();
305        while let Some(attributed) = router_handle.event_rx.recv().await {
306            // Decode the comms-safe roster member id back to the public
307            // alias space: SDK `EventStream` consumers filter by alias, and
308            // fail-closed per-agent ABAC view rules are written against
309            // aliases — an encoded id would silently drop both.
310            let source = crate::member_comms_id::runtime_event_alias(&attributed.source);
311            // Cold-cache fail-open guard: a member spawned after the one-time
312            // subscribe prime has no cached attributes, so a label/role-scoped
313            // `agent.view` deny would NOT match and the member's events would
314            // leak. Re-prime the shared cache once per newly-seen unknown agent
315            // before the decision so the deny resolves fail-closed.
316            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// ---------------------------------------------------------------------------
352// Structural mob events: per-client meerkat ledger subscription
353// ---------------------------------------------------------------------------
354
355/// Query parameters for `/mobkit/mob_events/stream`. Mirrors
356/// [`EventQuery`] for the field filters; cursor pagination is `after_seq`.
357#[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    /// Comma-separated list of event-kind labels to keep
372    /// (e.g. `flow_started,step_completed`). Empty / absent = all.
373    #[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    /// See [`AgentSseState::prime_runtime`]. `None` falls back to a
412    /// roster-only prime from `handle`.
413    prime_runtime: Option<MobRuntime>,
414    /// Optional auth context. When `Some`, requests are gated by the
415    /// same `require_app_auth` toggle the console RPC route uses; when
416    /// `None`, the route is unauthenticated (in-process or trusted
417    /// embedding).
418    decisions: Option<RuntimeDecisionState>,
419    access: Option<AccessController>,
420}
421
422/// Per-client SSE subscription to the meerkat structural-event ledger.
423///
424/// Each connection opens its own `MobEventsView::subscribe_after` so
425/// catch-up and live tail share the same ordered stream and there is
426/// no race window between snapshot and live subscription. Stale
427/// cursors are rejected with HTTP 410 Gone before any SSE handshake.
428/// Filtering matches the `mobkit/mob_events/query` predicate; the
429/// per-client `MobEventsSubscription` is dropped when the client
430/// disconnects, which cancels the upstream forwarder automatically.
431///
432/// When `decisions` is `Some` and `decisions.console.require_app_auth`
433/// is on, every request must carry a valid bearer token (Authorization
434/// header or `auth_token` query param) — same gate the console RPC
435/// route uses. `None` opts out of auth (e.g. trusted local embedding).
436pub 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
474/// Shared auth gate for every SSE route in mobkit. When `decisions` is
475/// `Some(_)` and `decisions.console.require_app_auth` is on, the
476/// request must carry a valid bearer / `auth_token` token; otherwise
477/// the route is open. Used by `mob_structural_events_sse_router`,
478/// `interaction_stream_router`, and the agent-/mob-event tier 2/3
479/// routers.
480///
481/// On success returns the caller's [`AccessView`] when an
482/// [`AccessController`] is wired (anonymous view on open routes), so the
483/// SSE handlers can apply per-agent ABAC checks. `Err(())` means 401.
484pub(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    // Parse the query string with `form_urlencoded` so percent-encoded
499    // tokens (e.g. base64 padding `=` re-encoded as `%3D`) decode
500    // correctly, and so substring-shadowing values like `xauth_token=`
501    // don't masquerade as `auth_token=`.
502    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        // Open route: identify callers that volunteered a valid token so
510        // per-user ABAC grants apply; everyone else is anonymous.
511        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    // Structural events span the whole mob: require the mob-wide
567    // observation grant, mirroring `mobkit/mob_events/query`. Envelopes
568    // attributed to a specific agent are additionally filtered by
569    // `agent.view` on that agent (below), so `mob.observe` cannot surface
570    // the lifecycle of an agent the caller is denied. Mob-level envelopes
571    // with no agent attribution flow under `mob.observe` alone.
572    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    // Captured so the long-lived stream can re-prime the shared attribute cache
621    // for members spawned AFTER the one-time subscribe prime.
622    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        // Agents we've already attempted a cache re-prime for, so an agent that
627        // genuinely has no roster attributes does not re-prime on every event.
628        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            // Agent-attributed structural events are gated by `agent.view`
635            // on their agent; mob-level events (no attribution) pass on the
636            // `mob.observe` grant alone.
637            if let Some(identity) = envelope.agent_identity.as_deref() {
638                // Cold-cache fail-open guard: a member spawned after the
639                // one-time subscribe prime has no cached attributes, so a
640                // label/role-scoped `agent.view` deny would NOT match and the
641                // member's lifecycle would leak. Re-prime the shared cache once
642                // per newly-seen unknown agent before deciding.
643                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}