use std::collections::{HashMap, VecDeque};
use std::path::PathBuf;
use std::sync::Arc;
use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use tokio::sync::RwLock;
use ractor::{Actor, ActorRef};
use crate::config::OuijaConfig;
use crate::persistence::OuijaSettings;
use crate::project_index::ProjectInfo;
use crate::scheduler::{ScheduledTask, TaskRun};
use crate::transport::Transport;
pub fn sanitize_session_id(name: &str) -> String {
name.to_lowercase()
.chars()
.map(|c| {
if c.is_ascii_alphanumeric() || c == '-' {
c
} else {
'-'
}
})
.collect::<String>()
.trim_matches('-')
.to_string()
}
pub fn expand_tilde(path: &str) -> String {
if let Some(rest) = path.strip_prefix("~/") {
let home = std::env::var("HOME").unwrap_or_else(|_| "/tmp".into());
format!("{home}/{rest}")
} else {
path.to_string()
}
}
pub fn resolve_project_root(path: &str) -> &str {
if let Some(idx) = path.find("/.claude/worktrees/") {
&path[..idx]
} else if let Some(idx) = path.find("/.ouija/worktrees/") {
&path[..idx]
} else {
path
}
}
type TransportMap = HashMap<String, Arc<dyn Transport>>;
#[derive(Debug)]
pub struct DuplicateNode(pub String);
impl std::fmt::Display for DuplicateNode {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "{}", self.0)
}
}
impl std::error::Error for DuplicateNode {}
#[derive(Clone, Debug, PartialEq, Eq, serde::Serialize)]
pub struct SessionSnapshot {
pub id: String,
pub origin: String,
pub role: Option<String>,
pub bulletin: Option<String>,
}
pub type SharedState = Arc<AppState>;
pub struct AppState {
pub config: OuijaConfig,
pub protocol: RwLock<crate::daemon_protocol::DaemonState>,
pub nodes: RwLock<HashMap<String, NodeInfo>>,
pub message_log: RwLock<VecDeque<LogEntry>>,
pub log_file: PathBuf,
transports: RwLock<TransportMap>,
pub settings: RwLock<OuijaSettings>,
pub scheduled_tasks: RwLock<HashMap<String, ScheduledTask>>,
pub task_runs: RwLock<VecDeque<TaskRun>>,
pane_queues: std::sync::Mutex<
HashMap<String, tokio::sync::mpsc::UnboundedSender<crate::tmux::InjectRequest>>,
>,
log_file_lock: std::sync::Mutex<()>,
task_run_log_lock: std::sync::Mutex<()>,
connected_npubs: std::sync::Mutex<HashMap<String, String>>,
last_reciprocated: std::sync::Mutex<HashMap<String, std::time::Instant>>,
session_agents: RwLock<HashMap<String, ActorRef<crate::session_agent::SessionMsg>>>,
pub project_index: RwLock<HashMap<String, ProjectInfo>>,
pending_commands: std::sync::Mutex<Vec<(String, tokio::sync::oneshot::Sender<String>)>>,
cached_assistant_panes: RwLock<Vec<crate::tmux::TmuxPane>>,
pub perfire_worktree_panes: RwLock<HashMap<String, String>>,
pub backends: crate::backend::BackendRegistry,
pub http_client: reqwest::Client,
pub pending_prompts: std::sync::Mutex<std::collections::HashMap<String, (String, String)>>,
pub session_diff_baselines: std::sync::Mutex<HashMap<String, Vec<SessionSnapshot>>>,
}
impl std::fmt::Debug for AppState {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("AppState")
.field("config", &self.config)
.finish_non_exhaustive()
}
}
#[derive(Clone, Debug, Serialize, Deserialize)]
pub struct SessionMetadata {
#[serde(default)]
pub vim_mode: bool,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub project_dir: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub role: Option<String>,
#[serde(default = "default_true")]
pub networked: bool,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub last_metadata_update: Option<DateTime<Utc>>,
#[serde(
default,
skip_serializing_if = "Option::is_none",
alias = "claude_session_id"
)]
pub backend_session_id: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub backend: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub project_description: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub bulletin: Option<String>,
#[serde(default, skip_serializing_if = "std::ops::Not::not")]
pub worktree: bool,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub model: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub reminder: Option<String>,
#[serde(
default,
skip_serializing_if = "Option::is_none",
alias = "original_prompt"
)]
pub prompt: Option<String>,
#[serde(default, alias = "loop_iteration")]
pub iteration: u64,
#[serde(default, skip_serializing_if = "Vec::is_empty", alias = "loop_log")]
pub iteration_log: Vec<crate::daemon_protocol::IterationLogEntry>,
#[serde(
default,
skip_serializing_if = "Option::is_none",
alias = "last_loop_next"
)]
pub last_iteration_at: Option<i64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub on_fire: Option<crate::scheduler::OnFire>,
}
fn default_true() -> bool {
true
}
impl Default for SessionMetadata {
fn default() -> Self {
Self {
vim_mode: false,
project_dir: None,
role: None,
networked: true,
last_metadata_update: None,
backend_session_id: None,
backend: None,
project_description: None,
bulletin: None,
worktree: false,
model: None,
reminder: None,
prompt: None,
iteration: 0,
iteration_log: Vec::new(),
last_iteration_at: None,
on_fire: None,
}
}
}
#[derive(Clone, Debug, Serialize)]
pub struct Session {
pub id: String,
pub pane: Option<String>,
pub origin: SessionOrigin,
pub registered_at: DateTime<Utc>,
pub last_activity_at: DateTime<Utc>,
pub metadata: SessionMetadata,
}
#[derive(Clone, Debug, Serialize, Deserialize)]
pub enum SessionOrigin {
Local,
Remote(String),
Human(String),
}
#[derive(Clone, Debug, Serialize)]
pub struct NodeInfo {
pub name: String,
pub daemon_id: String,
pub connected_at: DateTime<Utc>,
}
#[derive(Clone, Debug, Serialize)]
pub struct LogEntry {
pub timestamp: DateTime<Utc>,
pub from: String,
pub to: String,
pub message: String,
pub delivered: bool,
}
const MAX_LOG: usize = 100;
const MAX_TASK_RUNS: usize = 200;
const MAX_NAME_SUFFIX: u32 = 100;
const RECIPROCATE_DEBOUNCE_SECS: u64 = 30;
impl AppState {
#[cfg(test)]
pub fn new_for_test() -> Arc<Self> {
Arc::new(Self {
config: crate::config::OuijaConfig {
name: "test".into(),
npub: "npub1test".into(),
port: 0,
data_dir: std::path::PathBuf::from("/tmp/ouija-test-agent"),
config_dir: std::path::PathBuf::from("/tmp/ouija-test-agent"),
},
protocol: RwLock::new(crate::daemon_protocol::DaemonState::new(
"npub1test".into(),
"test".into(),
)),
nodes: RwLock::new(HashMap::new()),
message_log: RwLock::new(VecDeque::with_capacity(MAX_LOG)),
log_file: std::path::PathBuf::from("/tmp/ouija-test-agent/messages.jsonl"),
transports: RwLock::new(HashMap::new()),
settings: RwLock::new(Default::default()),
scheduled_tasks: RwLock::new(HashMap::new()),
task_runs: RwLock::new(VecDeque::with_capacity(MAX_TASK_RUNS)),
pane_queues: std::sync::Mutex::new(HashMap::new()),
log_file_lock: std::sync::Mutex::new(()),
task_run_log_lock: std::sync::Mutex::new(()),
connected_npubs: std::sync::Mutex::new(HashMap::new()),
last_reciprocated: std::sync::Mutex::new(HashMap::new()),
session_agents: RwLock::new(HashMap::new()),
project_index: RwLock::new(HashMap::new()),
pending_commands: std::sync::Mutex::new(Vec::new()),
cached_assistant_panes: RwLock::new(Vec::new()),
perfire_worktree_panes: RwLock::new(HashMap::new()),
backends: crate::backend::BackendRegistry::default_registry(),
http_client: reqwest::Client::new(),
pending_prompts: std::sync::Mutex::new(std::collections::HashMap::new()),
session_diff_baselines: std::sync::Mutex::new(HashMap::new()),
})
}
pub fn new(config: OuijaConfig) -> SharedState {
let log_file = config.data_dir.join("messages.jsonl");
let settings = crate::persistence::load_settings(&config.config_dir).unwrap_or_default();
let scheduled_tasks = crate::persistence::load_tasks(&config.data_dir).unwrap_or_default();
let protocol =
crate::daemon_protocol::DaemonState::new(config.npub.clone(), config.name.clone());
Arc::new(Self {
config,
protocol: RwLock::new(protocol),
nodes: RwLock::new(HashMap::new()),
message_log: RwLock::new(VecDeque::with_capacity(MAX_LOG)),
log_file,
transports: RwLock::new(HashMap::new()),
settings: RwLock::new(settings),
scheduled_tasks: RwLock::new(scheduled_tasks),
task_runs: RwLock::new(VecDeque::with_capacity(MAX_TASK_RUNS)),
pane_queues: std::sync::Mutex::new(HashMap::new()),
log_file_lock: std::sync::Mutex::new(()),
task_run_log_lock: std::sync::Mutex::new(()),
connected_npubs: std::sync::Mutex::new(HashMap::new()),
last_reciprocated: std::sync::Mutex::new(HashMap::new()),
session_agents: RwLock::new(HashMap::new()),
project_index: RwLock::new(HashMap::new()),
pending_commands: std::sync::Mutex::new(Vec::new()),
cached_assistant_panes: RwLock::new(Vec::new()),
perfire_worktree_panes: RwLock::new(HashMap::new()),
backends: crate::backend::BackendRegistry::default_registry(),
http_client: reqwest::Client::new(),
pending_prompts: std::sync::Mutex::new(std::collections::HashMap::new()),
session_diff_baselines: std::sync::Mutex::new(HashMap::new()),
})
}
pub async fn backend_for_session(
&self,
session_id: &str,
) -> std::sync::Arc<dyn crate::backend::CodingAssistant> {
let backend_name = self
.protocol
.read()
.await
.sessions
.get(session_id)
.and_then(|s| s.metadata.backend.as_deref())
.map(String::from);
match backend_name {
Some(name) => self
.backends
.get(&name)
.unwrap_or_else(|| self.backends.default()),
None => self.backends.default(),
}
}
pub async fn detect_backend_in_pane(&self, pane: &str) -> Option<String> {
let mut backend_process_names: Vec<(String, Vec<String>)> = Vec::new();
for name in self.backends.available() {
if let Some(b) = self.backends.get(name) {
let pnames: Vec<String> = b.process_names().iter().map(|s| s.to_string()).collect();
backend_process_names.push((name.to_string(), pnames));
}
}
let pane = pane.to_string();
tokio::task::spawn_blocking(move || {
use std::process::Command;
let output = Command::new("tmux")
.args(["display-message", "-t", &pane, "-p", "#{pane_pid}"])
.output()
.ok()?;
if !output.status.success() {
return None;
}
let pane_pid: u32 = String::from_utf8_lossy(&output.stdout)
.trim()
.parse()
.ok()?;
let output = Command::new("ps")
.args(["-eo", "pid,ppid,comm"])
.output()
.ok()?;
let stdout = String::from_utf8_lossy(&output.stdout);
let mut children: std::collections::HashMap<u32, Vec<u32>> =
std::collections::HashMap::new();
let mut names: std::collections::HashMap<u32, String> =
std::collections::HashMap::new();
for line in stdout.lines().skip(1) {
let mut parts = line.split_whitespace();
let (Some(pid_s), Some(ppid_s), Some(comm)) =
(parts.next(), parts.next(), parts.next())
else {
continue;
};
let (Ok(pid), Ok(ppid)) = (pid_s.parse::<u32>(), ppid_s.parse::<u32>()) else {
continue;
};
children.entry(ppid).or_default().push(pid);
names.insert(pid, comm.to_string());
}
let mut stack = vec![pane_pid];
while let Some(pid) = stack.pop() {
if let Some(comm) = names.get(&pid) {
for (backend_name, pnames) in &backend_process_names {
for pn in pnames {
if comm == pn || comm.strip_prefix('.') == Some(pn.as_str()) {
return Some(backend_name.clone());
}
}
}
}
if let Some(kids) = children.get(&pid) {
stack.extend(kids);
}
}
None
})
.await
.ok()
.flatten()
}
pub async fn find_session_by_pane(&self, pane: &str) -> Option<String> {
let proto = self.protocol.read().await;
proto
.sessions
.values()
.find(|s| s.pane.as_deref() == Some(pane))
.map(|s| s.id.clone())
}
pub async fn find_session_by_pane_or_backend_sid(
&self,
pane: Option<&str>,
backend_sid: Option<&str>,
) -> Option<String> {
let proto = self.protocol.read().await;
proto
.sessions
.values()
.find(|s| {
pane.is_some_and(|p| s.pane.as_deref() == Some(p))
|| backend_sid
.is_some_and(|b| s.metadata.backend_session_id.as_deref() == Some(b))
})
.map(|s| s.id.clone())
}
pub fn apply_and_execute(
self: &Arc<Self>,
event: crate::daemon_protocol::Event,
) -> std::pin::Pin<
Box<dyn std::future::Future<Output = Vec<crate::daemon_protocol::Effect>> + Send + '_>,
> {
Box::pin(self._apply_and_execute(event))
}
async fn _apply_and_execute(
self: &Arc<Self>,
event: crate::daemon_protocol::Event,
) -> Vec<crate::daemon_protocol::Effect> {
use crate::daemon_protocol::{Effect, LogLevel};
let effects = {
let mut state = self.protocol.write().await;
state.apply(event)
};
for effect in &effects {
match effect {
Effect::Broadcast(msg) => {
crate::transport::broadcast(self, msg).await;
}
Effect::BroadcastSessionList => {
crate::transport::broadcast_local_sessions(self).await;
}
Effect::InjectMessage {
session_id,
pane,
message,
vim_mode,
} => {
let _ = crate::tmux::locked_inject(self, session_id, pane, message, *vim_mode)
.await;
}
Effect::SetTmuxVar { pane, name, value } => {
let p = pane.clone();
let n = name.clone();
let v = value.clone();
tokio::task::spawn_blocking(move || crate::tmux_var::set(&p, &n, &v));
}
Effect::ClearTmuxVar { pane, name } => {
let p = pane.clone();
let n = name.clone();
tokio::task::spawn_blocking(move || crate::tmux_var::clear(&p, &n));
}
Effect::RenameWindow { pane, name } => {
let p = pane.clone();
let n = name.clone();
tokio::task::spawn_blocking(move || crate::tmux::rename_window(&p, &n));
}
Effect::EnableAutoRename { pane } => {
let p = pane.clone();
tokio::task::spawn_blocking(move || crate::tmux::enable_automatic_rename(&p));
}
Effect::SpawnAgent { session_id, pane } => {
self.spawn_session_agent(session_id, pane).await;
}
Effect::StopAgent { session_id } => {
if let Some(agent) = self
.session_agents
.write()
.await
.remove(session_id.as_str())
{
agent.stop(None);
}
}
Effect::RenameAgent { old_id, new_id } => {
let mut agents = self.session_agents.write().await;
if let Some(agent) = agents.remove(old_id.as_str()) {
let _ = agent.cast(crate::session_agent::SessionMsg::Renamed {
new_id: new_id.clone(),
});
agents.insert(new_id.clone(), agent);
}
}
Effect::ClearPendingReplies { removed_ids } => {
self.clear_orphaned_pending_replies(removed_ids).await;
}
Effect::Persist => {
let proto = self.protocol.read().await;
self.persist_protocol_state(&proto);
}
Effect::CleanupWorktree { project_dir } => {
let dir = project_dir.clone();
tokio::task::spawn(async move {
Self::cleanup_worktree_dir(&dir).await;
});
}
Effect::SendToHuman { npub, message } => {
let _ = crate::nostr_transport::send_plain_dm(self, npub, message).await;
}
Effect::ExecuteCommand { command, daemon_id } => {
tracing::info!("received command from {daemon_id}: {command}");
let state = Arc::clone(self);
let cmd = command.clone();
tokio::spawn(async move {
let result =
crate::nostr_transport::handle_human_command(&state, &cmd).await;
let reply = crate::protocol::WireMessage::CommandResult {
command: cmd,
result,
daemon_id: state.config.npub.clone(),
};
crate::transport::broadcast(&state, &reply).await;
});
}
Effect::ExecuteSessionStart {
name,
worktree,
project_dir,
prompt,
reminder,
from,
expects_reply,
daemon_id: sender_id,
} => {
tracing::info!("received session_start from {sender_id}: {name}");
let state = Arc::clone(self);
let name = name.clone();
let worktree = *worktree;
let project_dir = project_dir.clone();
let prompt = prompt.clone();
let reminder = reminder.clone();
let from = from.clone();
let expects_reply = *expects_reply;
tokio::spawn(async move {
let (result, _prompt_msg_id) = crate::nostr_transport::start_session(
&state,
&name,
worktree,
project_dir.as_deref(),
prompt.as_deref(),
from.as_deref(),
expects_reply,
None,
None,
reminder.as_deref(),
None, None, )
.await;
let reply = crate::protocol::WireMessage::CommandResult {
command: format!("/start {name}"),
result,
daemon_id: state.config.npub.clone(),
};
crate::transport::broadcast(&state, &reply).await;
});
}
Effect::ExecuteSessionRestart {
name,
fresh,
prompt,
reminder,
from,
expects_reply,
daemon_id: sender_id,
} => {
tracing::info!("received session_restart from {sender_id}: {name}");
let state = Arc::clone(self);
let name = name.clone();
let fresh = fresh.unwrap_or(false);
let prompt = prompt.clone();
let reminder = reminder.clone();
let from = from.clone();
let expects_reply = *expects_reply;
tokio::spawn(async move {
let (result, _prompt_msg_id) = crate::nostr_transport::restart_session(
&state,
&name,
fresh,
prompt.as_deref(),
from.as_deref(),
expects_reply,
None,
None,
reminder.as_deref(),
)
.await;
let reply = crate::protocol::WireMessage::CommandResult {
command: format!("/restart {name}"),
result,
daemon_id: state.config.npub.clone(),
};
crate::transport::broadcast(&state, &reply).await;
});
}
Effect::DeliverCommandResult {
daemon_id,
command,
result,
} => {
tracing::info!("command result from {daemon_id}: {command} -> {result}");
self.deliver_command_result(daemon_id, command, result)
.await;
}
Effect::RecordNode {
daemon_id,
daemon_name,
} => {
self.nodes.write().await.insert(
daemon_id.clone(),
NodeInfo {
name: daemon_name.clone(),
daemon_id: daemon_id.clone(),
connected_at: Utc::now(),
},
);
}
Effect::Reciprocate { daemon_id } => {
if self.should_reciprocate(daemon_id) {
tracing::info!("reciprocating session list to {daemon_id}");
crate::transport::broadcast_local_sessions(self).await;
}
}
Effect::LogMessage {
from,
to,
message,
delivered,
transport,
} => {
self.log_message(
from.clone(),
to.clone(),
message.clone(),
*delivered,
transport,
)
.await;
}
Effect::Log { level, message } => match level {
LogLevel::Info => tracing::info!("{message}"),
LogLevel::Warn => tracing::warn!("{message}"),
LogLevel::Debug => tracing::debug!("{message}"),
},
Effect::RegisterOk { .. }
| Effect::SendDelivered { .. }
| Effect::SendFailed { .. }
| Effect::RenameOk { .. }
| Effect::RenameFailed { .. }
| Effect::RemoveOk { .. }
| Effect::RemoveFailed { .. } => {}
}
}
effects
}
pub(crate) fn persist_protocol_state(&self, proto: &crate::daemon_protocol::DaemonState) {
let sessions: HashMap<String, Session> = proto
.sessions
.iter()
.map(|(k, entry)| {
let session = Session {
id: entry.id.clone(),
pane: entry.pane.clone(),
origin: match &entry.origin {
crate::daemon_protocol::Origin::Local => SessionOrigin::Local,
crate::daemon_protocol::Origin::Remote(d) => {
SessionOrigin::Remote(d.clone())
}
crate::daemon_protocol::Origin::Human(n) => SessionOrigin::Human(n.clone()),
},
registered_at: Utc::now(),
last_activity_at: Utc::now(),
metadata: SessionMetadata {
vim_mode: entry.metadata.vim_mode,
project_dir: entry.metadata.project_dir.clone(),
role: entry.metadata.role.clone(),
networked: entry.metadata.networked,
bulletin: entry.metadata.bulletin.clone(),
worktree: entry.metadata.worktree,
reminder: entry.metadata.reminder.clone(),
prompt: entry.metadata.prompt.clone(),
iteration: entry.metadata.iteration,
iteration_log: entry.metadata.iteration_log.clone(),
..Default::default()
},
};
(k.clone(), session)
})
.collect();
self.persist_sessions_from(&sessions);
}
pub(crate) async fn cleanup_worktree_dir(dir: &str) {
let dir_owned = dir.to_string();
let dir_clone = dir.to_string();
let repo = match tokio::task::spawn_blocking(move || {
std::process::Command::new("git")
.args(["-C", &dir_clone, "rev-parse", "--show-toplevel"])
.output()
.ok()
.filter(|o| o.status.success())
.map(|o| String::from_utf8_lossy(&o.stdout).trim().to_string())
})
.await
{
Ok(Some(r)) if !r.is_empty() => r,
_ => {
tracing::info!("worktree {dir_owned} not inside a git repo, skipping cleanup");
return;
}
};
let dir_clone = dir_owned.clone();
let has_changes = tokio::task::spawn_blocking(move || {
std::process::Command::new("git")
.args(["-C", &dir_clone, "status", "--porcelain"])
.output()
.map(|o| !o.stdout.is_empty())
.unwrap_or(true)
})
.await
.unwrap_or(true);
if has_changes {
tracing::info!("worktree {dir_owned} has uncommitted changes, keeping it");
return;
}
tracing::info!("cleaning up worktree: {dir_owned}");
let _ = tokio::task::spawn_blocking(move || {
let _ = std::process::Command::new("git")
.args(["-C", &repo, "worktree", "remove", &dir_owned, "--force"])
.status();
})
.await;
}
pub fn try_add_node(&self, npub: &str, name: &str) -> Result<(), DuplicateNode> {
let mut connected = self
.connected_npubs
.lock()
.expect("connected_npubs poisoned");
if let Some(existing) = connected.get(npub) {
return Err(DuplicateNode(existing.clone()));
}
connected.insert(npub.to_string(), name.to_string());
Ok(())
}
pub async fn disconnect_node(&self, daemon_id: &str) -> usize {
self.connected_npubs
.lock()
.expect("connected_npubs poisoned")
.remove(daemon_id);
for t in self.transports().await.values() {
t.deauthorize_peer(daemon_id).await;
}
self.nodes.write().await.remove(daemon_id);
let mut proto = self.protocol.write().await;
let to_remove: Vec<String> = proto.sessions
.iter()
.filter(|(_, s)| matches!(&s.origin, crate::daemon_protocol::Origin::Remote(d) if d == daemon_id))
.map(|(key, _)| key.clone())
.collect();
let count = to_remove.len();
for key in &to_remove {
proto.sessions.remove(key);
}
drop(proto);
if let Ok(mut conns) = crate::persistence::load_connections(&self.config.data_dir) {
conns.retain(|c| c.daemon_npub.as_deref() != Some(daemon_id));
let data = serde_json::to_string(&conns).unwrap_or_default();
let _ = std::fs::write(
self.config.data_dir.join("connections.json"),
data.as_bytes(),
);
}
count
}
pub fn enqueue_inject(&self, req: crate::tmux::InjectRequest) {
let pane_key = req.pane.clone();
let mut queues = self.pane_queues.lock().expect("pane_queues poisoned");
let req = if let Some(tx) = queues.get(&pane_key) {
match tx.send(req) {
Ok(()) => return,
Err(e) => {
queues.remove(&pane_key);
e.0
}
}
} else {
req
};
let (tx, rx) = tokio::sync::mpsc::unbounded_channel();
tx.send(req).expect("fresh channel cannot be closed");
tokio::spawn(crate::tmux::pane_inject_loop(rx));
queues.insert(pane_key, tx);
}
pub async fn transports(&self) -> TransportMap {
self.transports.read().await.clone()
}
pub async fn transport_by_name(&self, name: &str) -> Option<Arc<dyn Transport>> {
self.transports.read().await.get(name).cloned()
}
pub async fn add_transport(&self, t: Arc<dyn Transport>) {
self.transports
.write()
.await
.insert(t.transport_name().to_string(), t);
}
pub async fn spawn_session_agent(self: &Arc<Self>, id: &str, pane: &str) {
if let Some(old) = self.session_agents.write().await.remove(id) {
old.stop(None);
}
let agent = crate::session_agent::SessionAgent {
app_state: Arc::clone(self),
};
let args = crate::session_agent::SessionAgentArgs {
session_id: id.to_string(),
pane: pane.to_string(),
};
match Actor::spawn(None, agent, args).await {
Ok((actor_ref, _handle)) => {
self.session_agents
.write()
.await
.insert(id.to_string(), actor_ref);
tracing::info!("spawned session agent for {id}");
}
Err(e) => {
tracing::error!("failed to spawn session agent for {id}: {e}");
}
}
}
pub async fn notify_agent(&self, session_id: &str, msg: crate::session_agent::SessionMsg) {
let agent = {
let agents = self.session_agents.read().await;
agents.get(session_id).cloned()
};
if let Some(agent) = agent {
let _ = agent.cast(msg);
}
}
pub async fn query_agent_pending_replies(
&self,
session_id: &str,
) -> Vec<crate::daemon_protocol::PendingReplyEntry> {
let agents = self.session_agents.read().await;
if let Some(agent) = agents.get(session_id) {
ractor::call!(agent, crate::session_agent::SessionMsg::GetPendingReplies)
.unwrap_or_default()
} else {
Vec::new()
}
}
pub async fn drain_agent_compact_continuation(&self, session_id: &str) -> Option<String> {
let agents = self.session_agents.read().await;
if let Some(agent) = agents.get(session_id) {
ractor::call!(
agent,
crate::session_agent::SessionMsg::DrainPendingCompactContinuation
)
.unwrap_or(None)
} else {
None
}
}
pub(crate) async fn clear_orphaned_pending_replies(&self, removed_ids: &[String]) {
let mut proto = self.protocol.write().await;
proto.clear_orphaned_replies(removed_ids);
}
pub async fn collect_excess_idle_sessions(&self) -> Vec<String> {
let max = self.settings.read().await.max_local_sessions as usize;
if max == 0 {
return vec![];
}
let proto = self.protocol.read().await;
let local: Vec<_> = proto
.sessions
.values()
.filter(|s| matches!(s.origin, crate::daemon_protocol::Origin::Local))
.collect();
if local.len() <= max {
return vec![];
}
let excess = local.len() - max;
let mut stale: Vec<_> = local
.into_iter()
.filter(|s| s.metadata.is_stale())
.collect();
stale.sort_by_key(|s| s.metadata.last_metadata_update.unwrap_or(s.registered_at));
stale.iter().take(excess).map(|s| s.id.clone()).collect()
}
pub fn persist_sessions_from(&self, sessions: &HashMap<String, Session>) {
let persisted: Vec<_> = sessions
.values()
.filter_map(crate::persistence::PersistedSession::from_session)
.collect();
if let Err(e) = crate::persistence::save_sessions(&self.config.data_dir, &persisted) {
tracing::warn!("failed to persist sessions: {e}");
}
}
pub async fn cached_assistant_panes(&self) -> Vec<crate::tmux::TmuxPane> {
self.cached_assistant_panes.read().await.clone()
}
pub async fn scan_and_autoregister_panes(self: &Arc<Self>) {
let names: Vec<String> = self.backends.all_process_names();
let panes = match tokio::task::spawn_blocking(move || {
let name_refs: Vec<&str> = names.iter().map(|s| s.as_str()).collect();
crate::tmux::find_assistant_panes(&name_refs)
})
.await
.unwrap_or_else(|e| Err(anyhow::anyhow!("spawn_blocking join error: {e}")))
{
Ok(p) => p,
Err(e) => {
tracing::warn!("tmux scan failed: {e}");
return;
}
};
*self.cached_assistant_panes.write().await = panes.clone();
let auto_register = self.settings.read().await.auto_register;
if !auto_register {
return;
}
let (mut registered_panes, mut id_to_pane) = {
let proto = self.protocol.read().await;
let registered: std::collections::HashSet<String> = proto
.sessions
.values()
.filter(|s| matches!(s.origin, crate::daemon_protocol::Origin::Local))
.filter_map(|s| s.pane.clone())
.collect();
let id_to_pane: std::collections::HashMap<String, Option<String>> = proto
.sessions
.iter()
.map(|(id, s)| (id.clone(), s.pane.clone()))
.collect();
(registered, id_to_pane)
};
for pane in &panes {
if registered_panes.contains(&pane.pane_id) {
continue;
}
let pane_id_check = pane.pane_id.clone();
let has_ouija_id = tokio::task::spawn_blocking(move || {
std::process::Command::new("tmux")
.args(["show-options", "-pv", "-t", &pane_id_check, "@ouija_id"])
.output()
.map(|o| o.status.success() && !o.stdout.is_empty())
.unwrap_or(false)
})
.await
.unwrap_or(false);
if has_ouija_id {
continue;
}
let Some(ref path) = pane.pane_current_path else {
continue;
};
let project_root = resolve_project_root(path);
let basename = std::path::Path::new(project_root)
.file_name()
.and_then(|n| n.to_str())
.unwrap_or("unknown");
let base_id = sanitize_session_id(basename);
if base_id.is_empty() {
continue;
}
let mut id = base_id.clone();
let mut suffix = 2u32;
while let Some(existing_pane) = id_to_pane.get(&id) {
if existing_pane.as_deref() == Some(pane.pane_id.as_str()) {
break; }
id = format!("{base_id}-{suffix}");
suffix += 1;
if suffix > MAX_NAME_SUFFIX {
tracing::warn!("could not find available name for pane {}", pane.pane_id);
break;
}
}
let proto_meta = crate::daemon_protocol::SessionMeta {
project_dir: Some(project_root.to_string()),
role: Some(format!("working on {basename}")),
..Default::default()
};
tracing::info!("auto-registering pane {} as '{id}'", pane.pane_id);
self.apply_and_execute(crate::daemon_protocol::Event::Register {
id: id.clone(),
pane: Some(pane.pane_id.clone()),
metadata: proto_meta,
})
.await;
id_to_pane.insert(id.clone(), Some(pane.pane_id.clone()));
registered_panes.insert(pane.pane_id.clone());
}
}
pub fn should_reciprocate(&self, daemon_id: &str) -> bool {
let mut map = self
.last_reciprocated
.lock()
.expect("last_reciprocated poisoned");
let now = std::time::Instant::now();
if let Some(last) = map.get(daemon_id) {
if now.duration_since(*last) < std::time::Duration::from_secs(RECIPROCATE_DEBOUNCE_SECS)
{
return false;
}
}
map.insert(daemon_id.to_string(), now);
true
}
#[allow(dead_code)]
pub fn register_pending_command(
&self,
command: String,
) -> tokio::sync::oneshot::Receiver<String> {
let (tx, rx) = tokio::sync::oneshot::channel();
self.pending_commands
.lock()
.expect("pending_commands poisoned")
.push((command, tx));
rx
}
pub async fn deliver_command_result(&self, _daemon_id: &str, command: &str, result: &str) {
let tx = {
let mut pending = self
.pending_commands
.lock()
.expect("pending_commands poisoned");
pending
.iter()
.position(|(cmd, _)| cmd == command)
.map(|idx| pending.remove(idx).1)
};
if let Some(tx) = tx {
let _ = tx.send(result.to_string());
}
}
pub async fn local_session_hash(&self) -> u64 {
use std::hash::{Hash, Hasher};
let proto = self.protocol.read().await;
let mut entries: Vec<(&str, bool, Option<&str>, Option<&str>)> = proto
.sessions
.values()
.filter(|s| matches!(s.origin, crate::daemon_protocol::Origin::Local))
.map(|s| {
(
s.id.as_str(),
s.metadata.networked,
s.metadata.role.as_deref(),
s.metadata.bulletin.as_deref(),
)
})
.collect();
entries.sort_by_key(|(id, _, _, _)| *id);
let mut hasher = std::collections::hash_map::DefaultHasher::new();
entries.hash(&mut hasher);
hasher.finish()
}
pub async fn add_task(&self, task: ScheduledTask) {
let mut tasks = self.scheduled_tasks.write().await;
tasks.insert(task.id.clone(), task);
self.persist_tasks_from(&tasks);
}
pub async fn remove_task(&self, id: &str) -> Option<ScheduledTask> {
let mut tasks = self.scheduled_tasks.write().await;
let removed = tasks.remove(id);
if removed.is_some() {
self.persist_tasks_from(&tasks);
}
removed
}
pub async fn update_task(&self, id: &str, f: impl FnOnce(&mut ScheduledTask)) {
let mut tasks = self.scheduled_tasks.write().await;
if let Some(task) = tasks.get_mut(id) {
f(task);
self.persist_tasks_from(&tasks);
}
}
pub async fn log_task_run(&self, run: TaskRun) {
{
let _guard = self
.task_run_log_lock
.lock()
.expect("task_run_log_lock poisoned");
if let Err(e) = crate::persistence::append_task_run(&self.config.data_dir, &run) {
tracing::warn!("failed to append task run: {e}");
}
}
let mut runs = self.task_runs.write().await;
if runs.len() >= MAX_TASK_RUNS {
runs.pop_front();
}
runs.push_back(run);
}
pub fn persist_tasks_from(&self, tasks: &HashMap<String, ScheduledTask>) {
if let Err(e) = crate::persistence::save_tasks(&self.config.data_dir, tasks) {
tracing::warn!("failed to persist tasks: {e}");
}
}
pub async fn log_message(
&self,
from: String,
to: String,
message: String,
delivered: bool,
method: &str,
) {
let ts = Utc::now().to_rfc3339_opts(chrono::SecondsFormat::Secs, true);
let line = serde_json::json!({
"ts": ts,
"from": from,
"to": to,
"method": method,
"delivered": delivered,
});
{
let _guard = self.log_file_lock.lock().expect("log_file_lock poisoned");
if let Ok(mut f) = std::fs::OpenOptions::new()
.create(true)
.append(true)
.open(&self.log_file)
{
use std::io::Write;
let _ = writeln!(f, "{}", line);
}
}
let entry = LogEntry {
timestamp: Utc::now(),
from,
to,
message,
delivered,
};
let mut log = self.message_log.write().await;
if log.len() >= MAX_LOG {
log.pop_front();
}
log.push_back(entry);
}
pub fn opencode_serve_port(&self) -> u16 {
self.config.port + 320
}
}
#[cfg(test)]
pub(crate) mod tests {
use super::*;
#[test]
fn resolve_project_root_normal_path() {
assert_eq!(
resolve_project_root("/Users/dan/code/myproject"),
"/Users/dan/code/myproject"
);
}
#[test]
fn resolve_project_root_worktree_path() {
assert_eq!(
resolve_project_root("/Users/dan/code/chess-reader/.claude/worktrees/feature-branch"),
"/Users/dan/code/chess-reader"
);
}
#[test]
fn resolve_project_root_linux_worktree() {
assert_eq!(
resolve_project_root("/home/daniel/code/ouija/.claude/worktrees/auto-register"),
"/home/daniel/code/ouija"
);
}
#[test]
fn resolve_project_root_ouija_worktree() {
assert_eq!(
resolve_project_root("/home/daniel/code/ouija/.ouija/worktrees/feature-x"),
"/home/daniel/code/ouija"
);
}
pub(crate) fn test_config() -> OuijaConfig {
let dir = tempfile::tempdir().unwrap();
let path = dir.keep();
OuijaConfig {
name: "test".into(),
data_dir: path.clone(),
config_dir: path,
port: 0,
npub: "npub1test".into(),
}
}
async fn proto_register(state: &Arc<AppState>, id: &str, pane: Option<&str>) {
state
.apply_and_execute(crate::daemon_protocol::Event::Register {
id: id.into(),
pane: pane.map(Into::into),
metadata: crate::daemon_protocol::SessionMeta::default(),
})
.await;
}
#[tokio::test]
async fn register_session_basic() {
let state = AppState::new(test_config());
proto_register(&state, "s1", Some("%1")).await;
let proto = state.protocol.read().await;
let sessions = &proto.sessions;
assert_eq!(sessions.len(), 1);
assert!(sessions.contains_key("s1"));
}
#[tokio::test]
async fn register_session_dedup_by_pane() {
let state = AppState::new(test_config());
proto_register(&state, "old", Some("%1")).await;
proto_register(&state, "new", Some("%1")).await;
let proto = state.protocol.read().await;
let sessions = &proto.sessions;
assert_eq!(sessions.len(), 1);
assert!(sessions.contains_key("new"));
assert!(!sessions.contains_key("old"));
}
#[tokio::test]
async fn register_session_same_id_different_pane_updates() {
let state = AppState::new(test_config());
proto_register(&state, "s1", Some("%1")).await;
let effects = state
.apply_and_execute(crate::daemon_protocol::Event::Register {
id: "s1".into(),
pane: Some("%2".into()),
metadata: crate::daemon_protocol::SessionMeta::default(),
})
.await;
assert!(
effects
.iter()
.any(|e| matches!(e, crate::daemon_protocol::Effect::RegisterOk { .. }))
);
let proto = state.protocol.read().await;
let sessions = &proto.sessions;
assert_eq!(sessions.len(), 1);
assert_eq!(sessions.get("s1").unwrap().pane.as_deref(), Some("%2"));
}
#[tokio::test]
async fn register_session_same_id_same_pane_updates() {
let state = AppState::new(test_config());
proto_register(&state, "s1", Some("%1")).await;
state
.apply_and_execute(crate::daemon_protocol::Event::Register {
id: "s1".into(),
pane: Some("%1".into()),
metadata: crate::daemon_protocol::SessionMeta {
vim_mode: true,
..Default::default()
},
})
.await;
let proto = state.protocol.read().await;
let sessions = &proto.sessions;
assert!(sessions.get("s1").unwrap().metadata.vim_mode);
}
#[tokio::test]
async fn rename_session_basic() {
let state = AppState::new(test_config());
proto_register(&state, "old", Some("%1")).await;
state
.apply_and_execute(crate::daemon_protocol::Event::Rename {
old_id: "old".into(),
new_id: "new".into(),
})
.await;
let proto = state.protocol.read().await;
let sessions = &proto.sessions;
assert!(!sessions.contains_key("old"));
assert!(sessions.contains_key("new"));
}
#[tokio::test]
async fn rename_session_rejects_slash() {
let state = AppState::new(test_config());
proto_register(&state, "s1", Some("%1")).await;
let effects = state
.apply_and_execute(crate::daemon_protocol::Event::Rename {
old_id: "s1".into(),
new_id: "has/slash".into(),
})
.await;
assert!(
effects
.iter()
.any(|e| matches!(e, crate::daemon_protocol::Effect::RenameFailed { .. }))
);
assert!(state.protocol.read().await.sessions.contains_key("s1"));
}
#[tokio::test]
async fn rename_nonexistent_returns_none() {
let state = AppState::new(test_config());
let effects = state
.apply_and_execute(crate::daemon_protocol::Event::Rename {
old_id: "nope".into(),
new_id: "new".into(),
})
.await;
assert!(
effects
.iter()
.any(|e| matches!(e, crate::daemon_protocol::Effect::RenameFailed { .. }))
);
}
#[tokio::test]
async fn remove_session_basic() {
let state = AppState::new(test_config());
proto_register(&state, "s1", Some("%1")).await;
state
.apply_and_execute(crate::daemon_protocol::Event::Remove {
id: "s1".into(),
keep_worktree: false,
})
.await;
assert!(state.protocol.read().await.sessions.is_empty());
}
#[tokio::test]
async fn remove_nonexistent_is_noop() {
let state = AppState::new(test_config());
let effects = state
.apply_and_execute(crate::daemon_protocol::Event::Remove {
id: "nope".into(),
keep_worktree: false,
})
.await;
assert!(
effects
.iter()
.any(|e| matches!(e, crate::daemon_protocol::Effect::RemoveFailed { .. }))
);
}
#[tokio::test]
async fn remove_remote_session_fails() {
let state = AppState::new(test_config());
{
let mut proto = state.protocol.write().await;
proto.sessions.insert(
"remote/s1".into(),
crate::daemon_protocol::SessionEntry {
id: "remote/s1".into(),
origin: crate::daemon_protocol::Origin::Remote("remote".into()),
..Default::default()
},
);
}
let effects = state
.apply_and_execute(crate::daemon_protocol::Event::Remove {
id: "remote/s1".into(),
keep_worktree: false,
})
.await;
assert!(
effects
.iter()
.any(|e| matches!(e, crate::daemon_protocol::Effect::RemoveFailed { .. }))
);
assert_eq!(state.protocol.read().await.sessions.len(), 1);
}
fn test_entry(
id: &str,
pane: Option<&str>,
origin: crate::daemon_protocol::Origin,
metadata: crate::daemon_protocol::SessionMeta,
) -> crate::daemon_protocol::SessionEntry {
crate::daemon_protocol::SessionEntry {
id: id.into(),
pane: pane.map(Into::into),
origin,
metadata,
..Default::default()
}
}
#[tokio::test]
async fn log_message_caps_at_max() {
let state = AppState::new(test_config());
for i in 0..150 {
state
.log_message("from".into(), "to".into(), format!("msg {i}"), true, "test")
.await;
}
let log = state.message_log.read().await;
assert_eq!(log.len(), MAX_LOG);
}
#[tokio::test]
async fn local_session_hash_changes_on_networked_toggle() {
let state = AppState::new(test_config());
proto_register(&state, "s1", Some("%1")).await;
let hash_networked = state.local_session_hash().await;
{
let mut proto = state.protocol.write().await;
proto.sessions.get_mut("s1").unwrap().metadata.networked = false;
}
let hash_not_networked = state.local_session_hash().await;
assert_ne!(hash_networked, hash_not_networked);
}
#[tokio::test]
async fn disconnect_node_removes_sessions() {
let state = AppState::new(test_config());
{
let mut proto = state.protocol.write().await;
proto.sessions.insert(
"remote/s1".into(),
test_entry(
"remote/s1",
None,
crate::daemon_protocol::Origin::Remote("npub1remote".into()),
crate::daemon_protocol::SessionMeta::default(),
),
);
proto.sessions.insert(
"remote/s2".into(),
test_entry(
"remote/s2",
None,
crate::daemon_protocol::Origin::Remote("npub1remote".into()),
crate::daemon_protocol::SessionMeta::default(),
),
);
}
state.nodes.write().await.insert(
"npub1remote".into(),
NodeInfo {
name: "remote".into(),
daemon_id: "npub1remote".into(),
connected_at: Utc::now(),
},
);
state.try_add_node("npub1remote", "remote").unwrap();
let removed = state.disconnect_node("npub1remote").await;
assert_eq!(removed, 2);
assert!(state.protocol.read().await.sessions.is_empty());
assert!(state.nodes.read().await.is_empty());
}
#[test]
fn session_metadata_networked_defaults_true() {
let meta = SessionMetadata::default();
assert!(meta.networked);
}
#[test]
fn session_metadata_networked_serde_default() {
let json = r#"{"vim_mode": false}"#;
let meta: SessionMetadata = serde_json::from_str(json).unwrap();
assert!(meta.networked);
}
#[test]
fn session_origin_human_round_trip() {
let origin = SessionOrigin::Human("npub1abc".into());
let json = serde_json::to_string(&origin).unwrap();
let parsed: SessionOrigin = serde_json::from_str(&json).unwrap();
assert!(matches!(parsed, SessionOrigin::Human(npub) if npub == "npub1abc"));
}
#[test]
fn session_origin_human_deserializes() {
let json = r#"{"Human":"npub1xyz"}"#;
let origin: SessionOrigin = serde_json::from_str(json).unwrap();
assert!(matches!(origin, SessionOrigin::Human(npub) if npub == "npub1xyz"));
}
#[tokio::test]
async fn update_session_metadata_sets_role() {
let state = AppState::new(test_config());
proto_register(&state, "s1", Some("%1")).await;
state
.apply_and_execute(crate::daemon_protocol::Event::UpdateMetadata {
id: "s1".into(),
role: Some("debugging auth".into()),
bulletin: None,
project_dir: None,
networked: None,
})
.await;
let proto = state.protocol.read().await;
assert_eq!(
proto.sessions["s1"].metadata.role.as_deref(),
Some("debugging auth")
);
}
#[tokio::test]
async fn local_session_hash_changes_on_role_update() {
let state = AppState::new(test_config());
proto_register(&state, "s1", Some("%1")).await;
let hash_before = state.local_session_hash().await;
state
.apply_and_execute(crate::daemon_protocol::Event::UpdateMetadata {
id: "s1".into(),
role: Some("new role".into()),
bulletin: None,
project_dir: None,
networked: None,
})
.await;
let hash_after = state.local_session_hash().await;
assert_ne!(hash_before, hash_after);
}
#[tokio::test]
async fn update_metadata_sets_bulletin() {
let state = AppState::new(test_config());
proto_register(&state, "s1", Some("%1")).await;
state
.apply_and_execute(crate::daemon_protocol::Event::UpdateMetadata {
id: "s1".into(),
role: None,
bulletin: Some("offering review".into()),
project_dir: None,
networked: None,
})
.await;
let proto = state.protocol.read().await;
assert_eq!(
proto.sessions["s1"].metadata.bulletin.as_deref(),
Some("offering review")
);
}
#[tokio::test]
async fn excess_idle_disabled_when_zero() {
let state = AppState::new(test_config());
proto_register(&state, "s1", Some("%1")).await;
assert!(state.collect_excess_idle_sessions().await.is_empty());
}
#[tokio::test]
async fn excess_idle_no_eviction_at_limit() {
let state = AppState::new(test_config());
state.settings.write().await.max_local_sessions = 2;
proto_register(&state, "s1", Some("%1")).await;
proto_register(&state, "s2", Some("%2")).await;
assert!(state.collect_excess_idle_sessions().await.is_empty());
}
#[tokio::test]
async fn excess_idle_evicts_when_over_limit() {
use crate::daemon_protocol::{Origin, SessionMeta};
let state = AppState::new(test_config());
state.settings.write().await.max_local_sessions = 2;
{
let mut proto = state.protocol.write().await;
for name in &["a", "b", "c"] {
proto.sessions.insert(
name.to_string(),
test_entry(name, Some("%1"), Origin::Local, SessionMeta::default()),
);
}
}
let evicted = state.collect_excess_idle_sessions().await;
assert_eq!(evicted.len(), 1);
}
#[tokio::test]
async fn excess_idle_ignores_remote_and_human() {
use crate::daemon_protocol::{Origin, SessionMeta};
let state = AppState::new(test_config());
state.settings.write().await.max_local_sessions = 1;
{
let mut proto = state.protocol.write().await;
proto.sessions.insert(
"local".into(),
test_entry("local", Some("%1"), Origin::Local, SessionMeta::default()),
);
proto.sessions.insert(
"remote/r1".into(),
test_entry(
"remote/r1",
None,
Origin::Remote("npub1x".into()),
SessionMeta::default(),
),
);
proto.sessions.insert(
"human".into(),
test_entry(
"human",
None,
Origin::Human("npub1h".into()),
SessionMeta::default(),
),
);
}
assert!(state.collect_excess_idle_sessions().await.is_empty());
}
}