feat(adr): implement ADR-0109, ADR-0110, and ADR-0111
- ADR-0109: Incremental Background AST Indexing & Differential Graph Updates via Tree-Sitter - ADR-0110: Real-Time SSE Live Activity & Log Stream (/api/events/stream) - ADR-0111: Automated Stale Symbol Pruning & Graph Tombstoning (sweep_graph_health auto_prune_stale_files)
This commit is contained in:
1 parent
d792b50343
commit
40cab6142b
6 files changed
+146
-7
No files matched your search
@@ -2,8 +2,11 @@ use crate::AppState;
|
||||
use crate::state::GenericEvent;
|
||||
use axum::extract::{Query, State};
|
||||
use axum::response::IntoResponse;
|
||||
use axum::response::sse::{Event, KeepAlive, Sse};
|
||||
use std::collections::HashMap;
|
||||
use std::sync::Arc;
|
||||
use tokio_stream::StreamExt;
|
||||
use tokio_stream::wrappers::BroadcastStream;
|
||||
|
||||
pub async fn wait_for_event_handler(
|
||||
State(state): State<Arc<AppState>>,
|
||||
@@ -49,6 +52,36 @@ pub async fn post_event_handler(
|
||||
axum::Json(serde_json::json!({"status": "ok"}))
|
||||
}
|
||||
|
||||
/// Real-time Server-Sent Events (SSE) stream for agent execution events and observability (ADR-0110)
|
||||
pub async fn sse_events_handler(
|
||||
State(state): State<Arc<AppState>>,
|
||||
Query(params): Query<HashMap<String, String>>,
|
||||
) -> impl IntoResponse {
|
||||
let topic_filter = params.get("topic").cloned();
|
||||
let session_filter = params.get("session_id").cloned();
|
||||
let rx = state.handler.state.event_bus_tx.subscribe();
|
||||
|
||||
let stream = BroadcastStream::new(rx).filter_map(move |msg| match msg {
|
||||
Ok(event) => {
|
||||
let topic_matches = topic_filter.as_ref().is_none_or(|t| t == &event.topic);
|
||||
let session_matches = session_filter
|
||||
.as_ref()
|
||||
.is_none_or(|s| Some(s) == event.session_id.as_ref());
|
||||
if topic_matches && session_matches {
|
||||
let json_data = serde_json::to_string(&event).unwrap_or_default();
|
||||
Some(Ok::<_, std::convert::Infallible>(
|
||||
Event::default().event(event.topic).data(json_data),
|
||||
))
|
||||
} else {
|
||||
None
|
||||
}
|
||||
}
|
||||
Err(_) => None,
|
||||
});
|
||||
|
||||
Sse::new(stream).keep_alive(KeepAlive::default())
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
Reference in new issue
Block a user