use std::path::PathBuf;
use std::sync::{Arc, Mutex as StdMutex};
use codewhale_protocol::event_msg::EventMsg;
use codewhale_protocol::ids::{SessionId, ThreadId};
use codewhale_protocol::op::{Op, OpEnvelope};
use codewhale_state::StateStore;
use tokio::sync::mpsc;
use crate::ids::ThreadId as CoreThreadId;
use crate::journal::Journal;
use crate::session::{Session, Thread};
pub mod thread;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum CancelReason {
User,
External,
Preempted,
Internal,
}
#[derive(Clone)]
pub struct EngineHandle {
pub tx_op: mpsc::Sender<OpEnvelope>,
pub rx_event: Arc<tokio::sync::RwLock<mpsc::Receiver<EventMsg>>>,
cancel_token: Arc<StdMutex<tokio_util::sync::CancellationToken>>,
}
impl EngineHandle {
pub async fn send(&self, op: OpEnvelope) -> anyhow::Result<()> {
self.tx_op
.send(op)
.await
.map_err(|e| anyhow::anyhow!("{e}"))?;
Ok(())
}
pub fn cancel(&self) {
self.cancel_with_reason(CancelReason::User);
}
pub fn cancel_with_reason(&self, _reason: CancelReason) {
if let Ok(token) = self.cancel_token.lock() {
token.cancel();
}
}
pub async fn steer(
&self,
thread_id: ThreadId,
content: impl Into<String>,
) -> anyhow::Result<()> {
let env = OpEnvelope {
op_id: format!("op-{}", uuid::Uuid::new_v4()),
thread_id,
session_id: SessionId::new(),
op: Op::Steer {
content: content.into(),
},
};
self.tx_op
.send(env)
.await
.map_err(|e| anyhow::anyhow!("{e}"))?;
Ok(())
}
}
#[derive(Debug, Clone)]
pub struct EngineConfig {
pub workspace: PathBuf,
pub model: String,
pub model_provider: String,
pub thread_id: ThreadId,
pub session_id: SessionId,
pub max_steps: u32,
}
impl Default for EngineConfig {
fn default() -> Self {
Self {
workspace: PathBuf::from("."),
model: "deepseek-v4-flash".to_string(),
model_provider: "deepseek".to_string(),
thread_id: ThreadId::new(),
session_id: SessionId::new(),
max_steps: 32,
}
}
}
pub struct Engine {
rx_op: mpsc::Receiver<OpEnvelope>,
tx_event: mpsc::Sender<EventMsg>,
journal: Journal,
session: Session,
thread: Thread,
}
const ENGINE_OP_CHANNEL_CAPACITY: usize = 32;
const ENGINE_EVENT_CHANNEL_CAPACITY: usize = 128;
impl Engine {
#[must_use]
pub fn new(config: EngineConfig, _state: StateStore) -> (Self, EngineHandle) {
let (tx_op, rx_op) = mpsc::channel(ENGINE_OP_CHANNEL_CAPACITY);
let (tx_event, rx_event) = mpsc::channel(ENGINE_EVENT_CHANNEL_CAPACITY);
let thread = Thread::new(
CoreThreadId::from_string(config.thread_id.as_str().to_string()),
config.workspace.clone(),
config.model.clone(),
);
let session = Session::new(
CoreThreadId::from_string(config.thread_id.as_str().to_string()),
config.workspace.clone(),
config.model.clone(),
);
let handle = EngineHandle {
tx_op,
rx_event: Arc::new(tokio::sync::RwLock::new(rx_event)),
cancel_token: Arc::new(StdMutex::new(tokio_util::sync::CancellationToken::new())),
};
let engine = Self {
rx_op,
tx_event,
journal: Journal::new(),
session,
thread,
};
(engine, handle)
}
pub async fn run(mut self) {
while let Some(env) = self.rx_op.recv().await {
let _ = self
.tx_event
.send(EventMsg::TurnStarted {
thread_id: env.thread_id.clone(),
session_id: env.session_id.clone(),
turn_id: format!("turn-{}", uuid::Uuid::new_v4()),
})
.await;
match env.op {
Op::SendMessage { content, .. } => {
self.journal.append("user", serde_json::json!(content));
self.thread.leaf_id = self.journal.leaf_id.clone();
self.session.bump_revision();
let turn_id = format!("turn-{}", uuid::Uuid::new_v4());
let _ = self
.tx_event
.send(EventMsg::TurnComplete {
thread_id: env.thread_id.clone(),
session_id: env.session_id.clone(),
turn_id,
status: "completed".to_string(),
error: None,
})
.await;
}
Op::Steer { content } => {
self.journal.append("user", serde_json::json!(content));
self.thread.leaf_id = self.journal.leaf_id.clone();
}
Op::Shutdown | Op::Cancel => break,
_ => {}
}
}
}
}
pub fn spawn_engine(config: EngineConfig, state: StateStore) -> EngineHandle {
let (engine, handle) = Engine::new(config, state);
let handle_clone = handle.clone();
tokio::spawn(async move {
engine.run().await;
});
handle_clone
}
pub fn spawn_supervised(config: EngineConfig, state: StateStore) -> EngineHandle {
spawn_engine(config, state)
}
pub fn spawn_headless_thread(
workspace: PathBuf,
model: impl Into<String>,
state: StateStore,
) -> (EngineHandle, ThreadId, SessionId) {
let thread_id = ThreadId::new();
let session_id = SessionId::new();
let config = EngineConfig {
workspace,
model: model.into(),
model_provider: "deepseek".to_string(),
thread_id: thread_id.clone(),
session_id: session_id.clone(),
max_steps: 32,
};
let handle = spawn_engine(config, state);
(handle, thread_id, session_id)
}
#[cfg(test)]
mod tests {
use super::*;
use codewhale_state::StateStore;
#[tokio::test]
async fn headless_session_can_be_started_with_no_tui() {
let dir = tempfile::tempdir().unwrap();
let state = StateStore::open(Some(dir.path().join("state.db"))).unwrap();
let (handle, thread_id, _session_id) =
spawn_headless_thread(dir.path().to_path_buf(), "deepseek-v4-flash", state);
let env = OpEnvelope {
op_id: "op-1".into(),
thread_id: thread_id.clone(),
session_id: SessionId::new(),
op: Op::SendMessage {
content: "hello".into(),
mode: "agent".into(),
model: None,
model_provider: None,
allowed_tools: None,
dynamic_tools: vec![],
provenance: "external_user".into(),
},
};
handle.send(env).await.unwrap();
drop(handle);
}
}