diff --git a/bin/pudoxd/src/main.rs b/bin/pudoxd/src/main.rs index 4c5be51..db87564 100644 --- a/bin/pudoxd/src/main.rs +++ b/bin/pudoxd/src/main.rs @@ -77,10 +77,13 @@ async fn webhook_handler( let session_id = payload.agent_session.id.clone(); let (tx, rx) = mpsc::channel::(32); - tokio::spawn(run_session(rx, session_id.clone())); + tokio::spawn(run_session(app.sessions.clone(), rx, session_id.clone())); let _ = tx.send(SessionMessage::AgentEvent(payload)).await; - app.sessions.write().await.insert(session_id, Session { tx }); + app.sessions + .write() + .await + .insert(session_id, Session { tx }); } "Webhook processed" diff --git a/bin/pudoxd/src/session.rs b/bin/pudoxd/src/session.rs index 8cbfa83..f40510a 100644 --- a/bin/pudoxd/src/session.rs +++ b/bin/pudoxd/src/session.rs @@ -1,4 +1,6 @@ -use tokio::sync::mpsc; +use std::collections::HashMap; +use std::sync::Arc; +use tokio::sync::{mpsc, RwLock}; use crate::linear::WebhookPayload; use crate::sandbox::create_sandbox; @@ -32,10 +34,7 @@ impl EmbryoState { let sandbox = create_sandbox(); sandbox.prepare(session_id).await; - tracing::info!( - "Transitioning to Running for session {}", - session_id - ); + tracing::info!("Transitioning to Running for session {}", session_id); SessionStage::Running(RunningState) } } @@ -50,7 +49,11 @@ impl RunningState { } } -pub async fn run_session(mut rx: mpsc::Receiver, session_id: String) { +pub async fn run_session( + sessions: Arc>>, + mut rx: mpsc::Receiver, + session_id: String, +) { let mut stage = SessionStage::Embryo(EmbryoState); while let Some(msg) = rx.recv().await { @@ -60,5 +63,6 @@ pub async fn run_session(mut rx: mpsc::Receiver, session_id: Str }; } - tracing::info!("Session {} channel closed", session_id); + sessions.write().await.remove(&session_id); + tracing::info!("Session {} removed from state", session_id); }