feat(events): add generic Event Bus and /events/wait endpoint for agent wakeups
This commit is contained in:
1 parent
b8d58b2e64
commit
efcab4a7b4
5 files changed
+69
No files matched your search
@@ -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<Arc<AppState>>,
|
||||||
|
Query(params): Query<HashMap<String, String>>,
|
||||||
|
) -> 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<Arc<AppState>>,
|
||||||
|
axum::Json(event): axum::Json<GenericEvent>,
|
||||||
|
) -> impl IntoResponse {
|
||||||
|
let _ = state.handler.state.event_bus_tx.send(event);
|
||||||
|
axum::Json(serde_json::json!({"status": "ok"}))
|
||||||
|
}
|
||||||
@@ -2,3 +2,4 @@ pub mod rest;
|
|||||||
pub mod setup;
|
pub mod setup;
|
||||||
pub mod telemetry;
|
pub mod telemetry;
|
||||||
pub mod ws;
|
pub mod ws;
|
||||||
|
pub mod events;
|
||||||
@@ -24,6 +24,8 @@ pub fn create_router(app_state: Arc<AppState>) -> Router {
|
|||||||
.route("/ws", get(ws_handler))
|
.route("/ws", get(ws_handler))
|
||||||
.route("/health", get(health_handler))
|
.route("/health", get(health_handler))
|
||||||
.route("/nvim/telemetry", post(nvim_telemetry_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/verify", get(gate_verify_handler))
|
||||||
.route("/gate/set", post(gate_set_handler))
|
.route("/gate/set", post(gate_set_handler))
|
||||||
.route(
|
.route(
|
||||||
|
|||||||
@@ -45,5 +45,18 @@ pub async fn nvim_telemetry_handler(
|
|||||||
let _ = tx.try_send(msg_str.clone());
|
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"}))
|
axum::Json(serde_json::json!({"status": "ok"}))
|
||||||
}
|
}
|
||||||
@@ -5,6 +5,13 @@ use std::collections::HashMap;
|
|||||||
use std::path::PathBuf;
|
use std::path::PathBuf;
|
||||||
use std::sync::{Arc, RwLock};
|
use std::sync::{Arc, RwLock};
|
||||||
|
|
||||||
|
#[derive(Clone, Debug, serde::Serialize, serde::Deserialize)]
|
||||||
|
pub struct GenericEvent {
|
||||||
|
pub topic: String,
|
||||||
|
pub session_id: Option<String>,
|
||||||
|
pub payload: serde_json::Value,
|
||||||
|
}
|
||||||
|
|
||||||
pub struct MemoryState {
|
pub struct MemoryState {
|
||||||
pub base_dir: PathBuf,
|
pub base_dir: PathBuf,
|
||||||
pub graph: Store<KnowledgeGraph>,
|
pub graph: Store<KnowledgeGraph>,
|
||||||
@@ -29,6 +36,7 @@ pub struct MemoryState {
|
|||||||
pub context_workspaces: Store<Vec<ContextWorkspace>>,
|
pub context_workspaces: Store<Vec<ContextWorkspace>>,
|
||||||
pub recent_activities: Store<std::collections::VecDeque<serde_json::Value>>,
|
pub recent_activities: Store<std::collections::VecDeque<serde_json::Value>>,
|
||||||
pub activity_tx: tokio::sync::broadcast::Sender<String>,
|
pub activity_tx: tokio::sync::broadcast::Sender<String>,
|
||||||
|
pub event_bus_tx: tokio::sync::broadcast::Sender<GenericEvent>,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl MemoryState {
|
impl MemoryState {
|
||||||
@@ -72,6 +80,7 @@ impl MemoryState {
|
|||||||
context_workspaces: Store::new("context_workspaces", db.clone()),
|
context_workspaces: Store::new("context_workspaces", db.clone()),
|
||||||
recent_activities: Store::new("recent_activities", db.clone()),
|
recent_activities: Store::new("recent_activities", db.clone()),
|
||||||
activity_tx: tokio::sync::broadcast::channel(100).0,
|
activity_tx: tokio::sync::broadcast::channel(100).0,
|
||||||
|
event_bus_tx: tokio::sync::broadcast::channel(1000).0,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
Reference in new issue
Block a user