diff --git a/bin/pudoxd/src/main.rs b/bin/pudoxd/src/main.rs index 4a61a98..4c5be51 100644 --- a/bin/pudoxd/src/main.rs +++ b/bin/pudoxd/src/main.rs @@ -5,15 +5,11 @@ use std::sync::Arc; use tokio::sync::{mpsc, RwLock}; mod linear; -use crate::linear::{Linear, WebhookPayload}; - -enum SessionMessage { - AgentEvent(WebhookPayload), -} +mod sandbox; +mod session; -struct Session { - tx: mpsc::Sender, -} +use crate::linear::{Linear, WebhookPayload}; +use crate::session::{run_session, Session, SessionMessage}; #[derive(Clone)] struct App { @@ -79,22 +75,9 @@ async fn webhook_handler( } let session_id = payload.agent_session.id.clone(); - let (tx, mut rx) = mpsc::channel::(32); - - let task_session_id = session_id.clone(); - tokio::spawn(async move { - while let Some(msg) = rx.recv().await { - match msg { - SessionMessage::AgentEvent(p) => { - tracing::info!( - "Session {}: forwarding event to agent runner (stub)", - p.agent_session.id - ); - } - } - } - tracing::info!("Session {} channel closed", task_session_id); - }); + let (tx, rx) = mpsc::channel::(32); + + tokio::spawn(run_session(rx, session_id.clone())); let _ = tx.send(SessionMessage::AgentEvent(payload)).await; app.sessions.write().await.insert(session_id, Session { tx }); diff --git a/bin/pudoxd/src/sandbox.rs b/bin/pudoxd/src/sandbox.rs new file mode 100644 index 0000000..13d2385 --- /dev/null +++ b/bin/pudoxd/src/sandbox.rs @@ -0,0 +1,32 @@ +pub enum Sandbox { + Local, + Microvm, +} + +impl Sandbox { + pub async fn prepare(&self, session_id: &str) { + match self { + Sandbox::Local => { + tracing::info!("LocalSandbox: preparing for session {} (stub)", session_id); + } + Sandbox::Microvm => { + unimplemented!("MicrovmSandbox not yet implemented"); + } + } + } + + pub async fn spawn_agent(&self, session_id: &str) { + match self { + Sandbox::Local => { + tracing::info!("LocalSandbox: spawning agent for session {} (stub)", session_id); + } + Sandbox::Microvm => { + unimplemented!("MicrovmSandbox not yet implemented"); + } + } + } +} + +pub fn create_sandbox() -> Sandbox { + Sandbox::Local +} diff --git a/bin/pudoxd/src/session.rs b/bin/pudoxd/src/session.rs new file mode 100644 index 0000000..8cbfa83 --- /dev/null +++ b/bin/pudoxd/src/session.rs @@ -0,0 +1,64 @@ +use tokio::sync::mpsc; + +use crate::linear::WebhookPayload; +use crate::sandbox::create_sandbox; + +pub enum SessionMessage { + AgentEvent(WebhookPayload), +} + +pub struct Session { + pub tx: mpsc::Sender, +} + +enum SessionStage { + Embryo(EmbryoState), + Running(RunningState), +} + +struct EmbryoState; + +impl EmbryoState { + async fn handle(self, _msg: SessionMessage, session_id: &str) -> SessionStage { + tracing::info!("Processing event for session {}", session_id); + + let _repository = "user/repo"; + tracing::info!( + "Resolved repository {} for session {}", + _repository, + session_id + ); + + let sandbox = create_sandbox(); + sandbox.prepare(session_id).await; + + tracing::info!( + "Transitioning to Running for session {}", + session_id + ); + SessionStage::Running(RunningState) + } +} + +struct RunningState; + +impl RunningState { + async fn handle(self, _msg: SessionMessage, session_id: &str) -> SessionStage { + tracing::info!("Agent active for session {}", session_id); + std::future::pending::<()>().await; + SessionStage::Running(self) + } +} + +pub async fn run_session(mut rx: mpsc::Receiver, session_id: String) { + let mut stage = SessionStage::Embryo(EmbryoState); + + while let Some(msg) = rx.recv().await { + stage = match stage { + SessionStage::Embryo(s) => s.handle(msg, &session_id).await, + SessionStage::Running(s) => s.handle(msg, &session_id).await, + }; + } + + tracing::info!("Session {} channel closed", session_id); +}