Name update
This commit is contained in:
@@ -0,0 +1,66 @@
|
||||
//! Server-Sent Events stream for the Terminal tab's live log pane.
|
||||
//!
|
||||
//! On connection we emit the recent ring buffer (so the UI doesn't open
|
||||
//! to a blank pane), then forward every new line from the broadcast
|
||||
//! channel. Slow clients that fall behind get a `lagged` event and
|
||||
//! resume — better than dropping the connection mid-tail.
|
||||
|
||||
use crate::state::AppState;
|
||||
use axum::{
|
||||
extract::State,
|
||||
response::sse::{Event, KeepAlive, Sse},
|
||||
Json,
|
||||
};
|
||||
use futures::stream::{Stream, StreamExt};
|
||||
use pxeforge_core::LogLine;
|
||||
use serde_json::json;
|
||||
use std::convert::Infallible;
|
||||
use std::time::Duration;
|
||||
use tokio_stream::wrappers::BroadcastStream;
|
||||
|
||||
/// SSE handler. Each `data:` payload is a JSON object matching `LogLine`.
|
||||
pub async fn stream(
|
||||
State(state): State<AppState>,
|
||||
) -> Sse<impl Stream<Item = Result<Event, Infallible>>> {
|
||||
// 1. Snapshot the recent buffer first so a fresh UI sees context.
|
||||
let recent = state.log_bus.recent();
|
||||
let recent_stream = futures::stream::iter(
|
||||
recent
|
||||
.into_iter()
|
||||
.map(|l| Ok(Event::default().data(line_json(&l)))),
|
||||
);
|
||||
|
||||
// 2. Then live updates. BroadcastStream yields Result<T, Lagged>; on
|
||||
// a lagged client we send a synthetic event so the UI can flag it
|
||||
// rather than silently dropping data.
|
||||
let rx = state.log_bus.subscribe();
|
||||
let live = BroadcastStream::new(rx).map(|res| match res {
|
||||
Ok(line) => Ok(Event::default().data(line_json(&line))),
|
||||
Err(tokio_stream::wrappers::errors::BroadcastStreamRecvError::Lagged(n)) => Ok(Event::default()
|
||||
.event("lagged")
|
||||
.data(json!({ "skipped": n }).to_string())),
|
||||
});
|
||||
|
||||
Sse::new(recent_stream.chain(live))
|
||||
.keep_alive(KeepAlive::new().interval(Duration::from_secs(15)))
|
||||
}
|
||||
|
||||
/// Plain JSON snapshot of the recent buffer, for clients that prefer a
|
||||
/// pull-based fetch over an SSE subscription.
|
||||
pub async fn recent(State(state): State<AppState>) -> Json<serde_json::Value> {
|
||||
Json(json!({ "lines": state.log_bus.recent() }))
|
||||
}
|
||||
|
||||
/// Drop the in-memory ring buffer. Live subscribers are unaffected (they
|
||||
/// keep streaming new lines as they arrive).
|
||||
pub async fn clear(State(state): State<AppState>) -> Json<serde_json::Value> {
|
||||
state.log_bus.clear();
|
||||
state
|
||||
.log_bus
|
||||
.push("info", "pxeforge::terminal", "log buffer cleared by operator");
|
||||
Json(json!({ "ok": true }))
|
||||
}
|
||||
|
||||
fn line_json(l: &LogLine) -> String {
|
||||
serde_json::to_string(l).unwrap_or_else(|_| "{}".to_string())
|
||||
}
|
||||
Reference in New Issue
Block a user