diff --git a/server/src/api/events.rs b/server/src/api/events.rs new file mode 100644 index 0000000..eea514d --- /dev/null +++ b/server/src/api/events.rs @@ -0,0 +1,44 @@ +use crate::AppState; +use axum::extract::{Query, State}; +use axum::response::IntoResponse; +use std::collections::HashMap; +use std::sync::Arc; +use crate::state::GenericEvent; + +pub async fn wait_for_event_handler( + State(state): State>, + Query(params): Query>, +) -> impl IntoResponse { + let topic = params.get("topic").cloned(); + let session_id = params.get("session_id").cloned(); + + let mut rx = state.handler.state.event_bus_tx.subscribe(); + + loop { + match rx.recv().await { + Ok(event) => { + let topic_matches = topic.as_ref().map_or(true, |t| t == &event.topic); + let session_matches = session_id.as_ref().map_or(true, |s| Some(s) == event.session_id.as_ref()); + + if topic_matches && session_matches { + return axum::Json(event); + } + } + Err(_) => { + return axum::Json(GenericEvent { + topic: "error".to_string(), + session_id: None, + payload: serde_json::json!({"error": "Event bus lagged or closed"}), + }); + } + } + } +} + +pub async fn post_event_handler( + State(state): State>, + axum::Json(event): axum::Json, +) -> impl IntoResponse { + let _ = state.handler.state.event_bus_tx.send(event); + axum::Json(serde_json::json!({"status": "ok"})) +} diff --git a/server/src/api/mod.rs b/server/src/api/mod.rs index 6d09d65..7423866 100644 --- a/server/src/api/mod.rs +++ b/server/src/api/mod.rs @@ -2,3 +2,4 @@ pub mod rest; pub mod setup; pub mod telemetry; pub mod ws; +pub mod events; diff --git a/server/src/api/setup.rs b/server/src/api/setup.rs index 9dbbbf5..6d5449c 100644 --- a/server/src/api/setup.rs +++ b/server/src/api/setup.rs @@ -24,6 +24,8 @@ pub fn create_router(app_state: Arc) -> Router { .route("/ws", get(ws_handler)) .route("/health", get(health_handler)) .route("/nvim/telemetry", post(nvim_telemetry_handler)) + .route("/events/wait", get(crate::api::events::wait_for_event_handler)) + .route("/events", post(crate::api::events::post_event_handler)) .route("/gate/verify", get(gate_verify_handler)) .route("/gate/set", post(gate_set_handler)) .route( diff --git a/server/src/api/telemetry.rs b/server/src/api/telemetry.rs index dc88e7a..627ae8b 100644 --- a/server/src/api/telemetry.rs +++ b/server/src/api/telemetry.rs @@ -45,5 +45,18 @@ pub async fn nvim_telemetry_handler( let _ = tx.try_send(msg_str.clone()); } + if payload.event == "BufWritePost" { + if let Some(ref file_path) = payload.file { + let normalized_file = file_path.replace("\\", "/"); + let topic = format!("nvim:save:{}", normalized_file); + let event = crate::state::GenericEvent { + topic, + session_id: Some(payload.session_id.clone()), + payload: serde_json::json!(&payload), + }; + let _ = state.handler.state.event_bus_tx.send(event); + } + } + axum::Json(serde_json::json!({"status": "ok"})) } diff --git a/server/src/state.rs b/server/src/state.rs index bd530dc..c80fc5c 100644 --- a/server/src/state.rs +++ b/server/src/state.rs @@ -5,6 +5,13 @@ use std::collections::HashMap; use std::path::PathBuf; use std::sync::{Arc, RwLock}; +#[derive(Clone, Debug, serde::Serialize, serde::Deserialize)] +pub struct GenericEvent { + pub topic: String, + pub session_id: Option, + pub payload: serde_json::Value, +} + pub struct MemoryState { pub base_dir: PathBuf, pub graph: Store, @@ -29,6 +36,7 @@ pub struct MemoryState { pub context_workspaces: Store>, pub recent_activities: Store>, pub activity_tx: tokio::sync::broadcast::Sender, + pub event_bus_tx: tokio::sync::broadcast::Sender, } impl MemoryState { @@ -72,6 +80,7 @@ impl MemoryState { context_workspaces: Store::new("context_workspaces", db.clone()), recent_activities: Store::new("recent_activities", db.clone()), activity_tx: tokio::sync::broadcast::channel(100).0, + event_bus_tx: tokio::sync::broadcast::channel(1000).0, } }