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 resolve_unique_session_id(
id_to_pane: &HashMap<String, Option<String>>,
base_id: &str,
target_pane: Option<&str>,
) -> String {
let mut id = base_id.to_string();
let mut suffix = 2u32;
while let Some(existing_pane) = id_to_pane.get(&id) {
if target_pane.is_some() && existing_pane.as_deref() == target_pane {
return id;
}
id = format!("{base_id}-{suffix}");
if suffix > MAX_NAME_SUFFIX {
tracing::warn!(
"resolve_unique_session_id: exhausted suffixes 2..={MAX_NAME_SUFFIX} for base '{base_id}', returning '{id}'"
);
return id;
}
suffix += 1;
}
id
}
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 {}
pub type SharedState = Arc<AppState>;
#[derive(Clone, Debug)]
pub(crate) struct EffectDeliveryFailure {
reason: String,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub(crate) enum DeliveryOutcome {
Accepted,
Rejected(String),
Ambiguous(String),
}
fn prompt_async_failure_reason(
decision: crate::nostr_transport::PromptAsyncFallbackDecision,
) -> String {
format!("prompt_async request failed: {decision:?}")
}
fn http_delivery_attempt_failure(
decision: crate::nostr_transport::PromptAsyncFallbackDecision,
) -> DeliveryOutcome {
let reason = prompt_async_failure_reason(decision);
match decision {
crate::nostr_transport::PromptAsyncFallbackDecision::DefiniteNonAcceptance => {
DeliveryOutcome::Rejected(reason)
}
crate::nostr_transport::PromptAsyncFallbackDecision::Ambiguous => {
DeliveryOutcome::Ambiguous(reason)
}
}
}
async fn session_owns_pane(state: &AppState, session_id: &str, pane: &str) -> Result<(), String> {
let registered_pane = {
let proto = state.protocol.read().await;
proto
.sessions
.get(session_id)
.and_then(|session| session.pane.clone())
};
if registered_pane.as_deref() == Some(pane) {
Ok(())
} else {
Err(format!("pane {pane} is not owned by session {session_id}"))
}
}
async fn deliver_raw_tmux_for_session(
state: &AppState,
request: &InjectDeliveryRequest<'_>,
inject_config: Option<crate::backend::InjectConfig>,
tui_pattern: Option<String>,
) -> DeliveryOutcome {
if let Err(reason) = session_owns_pane(state, request.session_id, request.pane).await {
return DeliveryOutcome::Rejected(reason);
}
let result = match inject_config {
Some(inject_config) => {
crate::tmux::locked_inject_raw_tmux_with_config(
state,
request.pane,
request.message,
request.vim_mode,
inject_config,
tui_pattern,
)
.await
}
None => {
crate::tmux::locked_inject_raw_tmux(
state,
request.session_id,
request.pane,
request.message,
request.vim_mode,
)
.await
}
};
result
.map(|()| DeliveryOutcome::Accepted)
.unwrap_or_else(|error| DeliveryOutcome::Rejected(error.to_string()))
}
async fn deliver_by_current_session_plan(
state: &AppState,
request: &InjectDeliveryRequest<'_>,
) -> DeliveryOutcome {
match crate::tmux::session_delivery_plan(state, request.session_id, request.pane).await {
crate::tmux::SessionDeliveryPlan::Http(delivery) => {
if let Err(reason) = session_owns_pane(state, request.session_id, request.pane).await {
return DeliveryOutcome::Rejected(reason);
}
crate::tmux::deliver_via_http(
state,
&delivery.backend_session_id,
delivery.project_dir.as_deref(),
request.message,
delivery.model.as_deref(),
delivery.effort.as_deref(),
)
.await
.map(|()| DeliveryOutcome::Accepted)
.unwrap_or_else(http_delivery_attempt_failure)
}
crate::tmux::SessionDeliveryPlan::RawTmux {
inject_config,
tui_pattern,
} => deliver_raw_tmux_for_session(state, request, Some(inject_config), tui_pattern).await,
crate::tmux::SessionDeliveryPlan::Unavailable(reason) => DeliveryOutcome::Rejected(reason),
}
}
pub(crate) struct InjectDeliveryRequest<'a> {
pub session_id: &'a str,
pub pane: &'a str,
pub message: &'a str,
pub vim_mode: bool,
pub delivery_method: Option<&'a str>,
pub recorded_method: Option<&'a str>,
}
pub(crate) async fn deliver_inject_message_effect(
state: &Arc<AppState>,
request: InjectDeliveryRequest<'_>,
) -> DeliveryOutcome {
let method = request.delivery_method.or(request.recorded_method);
match method {
Some("http") => deliver_by_current_session_plan(state, &request).await,
Some("tmux") => deliver_raw_tmux_for_session(state, &request, None, None).await,
_ => deliver_by_current_session_plan(state, &request).await,
}
}
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>)>>,
pub(crate) cached_assistant_panes: RwLock<Vec<crate::tmux::TmuxPane>>,
pub perfire_worktree_panes: RwLock<HashMap<String, String>>,
sweep_in_progress: std::sync::atomic::AtomicBool,
sweep_backoff_until: std::sync::Mutex<Option<std::time::Instant>>,
pub backends: crate::backend::BackendRegistry,
pub http_client: reqwest::Client,
pub pending_prompts: std::sync::Mutex<std::collections::HashMap<String, PendingPrompt>>,
compact_in_progress: std::sync::Mutex<std::collections::HashSet<String>>,
soft_restart_in_progress: std::sync::Mutex<std::collections::HashSet<String>>,
}
pub(crate) struct CompactInProgressGuard<'a> {
state: &'a AppState,
key: String,
}
impl Drop for CompactInProgressGuard<'_> {
fn drop(&mut self) {
self.state
.compact_in_progress
.lock()
.expect("compact_in_progress mutex poisoned")
.remove(&self.key);
}
}
pub(crate) struct SoftRestartInProgressGuard<'a> {
state: &'a AppState,
session_id: String,
}
impl Drop for SoftRestartInProgressGuard<'_> {
fn drop(&mut self) {
self.state
.soft_restart_in_progress
.lock()
.expect("soft_restart_in_progress mutex poisoned")
.remove(&self.session_id);
}
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct PendingPrompt {
pub pane_id: String,
pub prompt: String,
pub backend_session_id: Option<String>,
}
impl PendingPrompt {
pub fn new(pane_id: String, prompt: String, backend_session_id: Option<String>) -> Self {
Self {
pane_id,
prompt,
backend_session_id,
}
}
}
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 opencode_binding: Option<crate::daemon_protocol::OpenCodeBinding>,
#[serde(default)]
pub restart_generation: u64,
#[serde(default)]
pub session_incarnation: i64,
#[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 effort: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub codex_home: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub reminder: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub parent_session: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub idle_policy: Option<crate::daemon_protocol::IdlePolicy>,
#[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>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub worktree_present: Option<bool>,
}
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,
opencode_binding: None,
restart_generation: 0,
session_incarnation: 0,
project_description: None,
bulletin: None,
worktree: false,
model: None,
effort: None,
codex_home: None,
reminder: None,
parent_session: None,
idle_policy: None,
prompt: None,
iteration: 0,
iteration_log: Vec::new(),
last_iteration_at: None,
on_fire: None,
worktree_present: 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()),
sweep_in_progress: std::sync::atomic::AtomicBool::new(false),
sweep_backoff_until: std::sync::Mutex::new(None),
backends: crate::backend::BackendRegistry::default_registry(),
http_client: reqwest::Client::new(),
pending_prompts: std::sync::Mutex::new(std::collections::HashMap::new()),
compact_in_progress: std::sync::Mutex::new(std::collections::HashSet::new()),
soft_restart_in_progress: std::sync::Mutex::new(std::collections::HashSet::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()),
sweep_in_progress: std::sync::atomic::AtomicBool::new(false),
sweep_backoff_until: std::sync::Mutex::new(None),
backends: crate::backend::BackendRegistry::default_registry(),
http_client: reqwest::Client::new(),
pending_prompts: std::sync::Mutex::new(std::collections::HashMap::new()),
compact_in_progress: std::sync::Mutex::new(std::collections::HashSet::new()),
soft_restart_in_progress: std::sync::Mutex::new(std::collections::HashSet::new()),
})
}
pub(crate) fn try_acquire_compact_in_progress(
&self,
key: &str,
) -> Option<CompactInProgressGuard<'_>> {
let mut compact_in_progress = self
.compact_in_progress
.lock()
.expect("compact_in_progress mutex poisoned");
if !compact_in_progress.insert(key.to_string()) {
return None;
}
Some(CompactInProgressGuard {
state: self,
key: key.to_string(),
})
}
pub(crate) fn try_acquire_soft_restart_in_progress(
&self,
session_id: &str,
) -> Option<SoftRestartInProgressGuard<'_>> {
let mut soft_restart_in_progress = self
.soft_restart_in_progress
.lock()
.expect("soft_restart_in_progress mutex poisoned");
if !soft_restart_in_progress.insert(session_id.to_string()) {
return None;
}
Some(SoftRestartInProgressGuard {
state: self,
session_id: session_id.to_string(),
})
}
pub(crate) fn is_soft_restart_in_progress(&self, session_id: &str) -> bool {
self.soft_restart_in_progress
.lock()
.expect("soft_restart_in_progress mutex poisoned")
.contains(session_id)
}
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 backend_process_names: Vec<(String, Vec<String>)> =
self.backends.all_backend_process_names();
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> {
let (effects, rollback) = {
let mut state = self.protocol.write().await;
let mut rollback = FailedEffectSendRollback::capture_for_event(&state, &event);
let effects = state.apply(event);
if let Some(rollback) = &mut rollback {
rollback.capture_after_send(&state);
rollback.reserve_sender_state_after_send(&mut state);
}
(effects, rollback)
};
let delivery_failure = self.execute_effects(&effects).await;
if let Some(failure) = delivery_failure {
self.clear_pending_reply_for_failed_effect_delivery(&effects)
.await;
self.rollback_sender_state_for_failed_effect_delivery(rollback)
.await;
return rewrite_send_delivery_failure(effects, &failure.reason);
}
if effects
.iter()
.any(|effect| matches!(effect, crate::daemon_protocol::Effect::SendFailed { .. }))
{
self.rollback_sender_state_for_failed_effect_delivery(rollback)
.await;
return effects;
}
self.finalize_successful_effect_delivery(rollback).await;
effects
}
pub(crate) async fn execute_effects(
self: &Arc<Self>,
effects: &[crate::daemon_protocol::Effect],
) -> Option<EffectDeliveryFailure> {
use crate::daemon_protocol::{Effect, LogLevel};
let recorded_method = effects.iter().find_map(|effect| match effect {
Effect::SendDelivered { method, .. } => Some(method.as_str()),
_ => None,
});
let recorded_http_delivery = effects.iter().find_map(|effect| match effect {
Effect::SendDelivered { http_delivery, .. } => http_delivery.as_ref(),
_ => None,
});
let mut delivery_failure = None;
for effect in effects {
match effect {
Effect::Broadcast(msg) => {
if delivery_failure.is_some() {
if let crate::protocol::WireMessage::SessionSendAck {
from,
to,
delivered: true,
daemon_id,
} = msg
{
let failed_ack = crate::protocol::WireMessage::SessionSendAck {
from: from.clone(),
to: to.clone(),
delivered: false,
daemon_id: daemon_id.clone(),
};
crate::transport::broadcast(self, &failed_ack).await;
continue;
}
}
crate::transport::broadcast(self, msg).await;
}
Effect::BroadcastSessionList => {
crate::transport::broadcast_local_sessions(self).await;
}
Effect::InjectMessage {
session_id,
pane,
message,
vim_mode,
delivery_method,
..
} => {
let outcome = deliver_inject_message_effect(
self,
InjectDeliveryRequest {
session_id,
pane,
message,
vim_mode: *vim_mode,
delivery_method: delivery_method.as_deref(),
recorded_method,
},
)
.await;
match outcome {
DeliveryOutcome::Accepted => {}
DeliveryOutcome::Rejected(reason) => {
tracing::warn!(session = %session_id, "message delivery failed: {reason}");
delivery_failure.get_or_insert(EffectDeliveryFailure { reason });
}
DeliveryOutcome::Ambiguous(reason) => {
tracing::warn!(session = %session_id, "message delivery outcome ambiguous; preserving delivered state: {reason}");
}
}
}
Effect::DeliverHttpMessage {
session_id,
message,
http_delivery,
..
} => match Some(http_delivery).or(recorded_http_delivery) {
Some(delivery) => {
if let Err(decision) = crate::tmux::deliver_via_http(
self,
&delivery.backend_session_id,
delivery.project_dir.as_deref(),
message,
delivery.model.as_deref(),
delivery.effort.as_deref(),
)
.await
{
match http_delivery_attempt_failure(decision) {
DeliveryOutcome::Accepted => {}
DeliveryOutcome::Rejected(reason) => {
tracing::warn!(session = %session_id, "http delivery failed: {reason}");
delivery_failure
.get_or_insert(EffectDeliveryFailure { reason });
}
DeliveryOutcome::Ambiguous(reason) => {
tracing::warn!(session = %session_id, "http delivery outcome ambiguous; preserving delivered state: {reason}");
}
}
}
}
None => {
let error = anyhow::anyhow!(
"http delivery skipped: no recorded backend_session_id on send"
);
tracing::warn!(session = %session_id, "{error}");
delivery_failure.get_or_insert_with(|| EffectDeliveryFailure {
reason: error.to_string(),
});
}
},
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, None, reminder.as_deref(),
None, None, None, None, false, )
.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, None, reminder.as_deref(),
crate::nostr_transport::ParentSessionOverride::PreservePrevious,
None, )
.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,
} => {
let delivered = if delivery_failure.is_some() {
false
} else {
*delivered
};
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::RegisterFailed { .. }
| Effect::SendDelivered { .. }
| Effect::SendFailed { .. }
| Effect::RenameOk { .. }
| Effect::RenameFailed { .. }
| Effect::RemoveOk { .. }
| Effect::RemoveFailed { .. } => {}
}
}
delivery_failure
}
async fn clear_pending_reply_for_failed_effect_delivery(
&self,
effects: &[crate::daemon_protocol::Effect],
) {
let Some((to, msg_id, from)) = effects.iter().find_map(|effect| match effect {
crate::daemon_protocol::Effect::SendDelivered {
from, to, msg_id, ..
} => Some((to.clone(), *msg_id, Some(from.clone()))),
crate::daemon_protocol::Effect::InjectMessage {
session_id,
pending_reply_msg_id,
pending_reply_from,
..
} => pending_reply_msg_id
.map(|msg_id| (session_id.clone(), msg_id, pending_reply_from.clone())),
crate::daemon_protocol::Effect::DeliverHttpMessage {
session_id,
pending_reply_msg_id,
pending_reply_from,
..
} => pending_reply_msg_id
.map(|msg_id| (session_id.clone(), msg_id, pending_reply_from.clone())),
_ => None,
}) else {
return;
};
let Some(from) = from else {
return;
};
let mut proto = self.protocol.write().await;
let Some(pending) = proto.pending_replies.get_mut(&to) else {
return;
};
pending.retain(|entry| entry.msg_id != msg_id || entry.from != from);
if pending.is_empty() {
proto.pending_replies.remove(&to);
}
}
async fn rollback_sender_state_for_failed_effect_delivery(
&self,
rollback: Option<FailedEffectSendRollback>,
) {
let Some(rollback) = rollback else {
return;
};
let mut proto = self.protocol.write().await;
if rollback.sender_state_reserved() {
return;
}
if let Some(entry) = rollback.pending_reply_before_send {
let current_entry =
proto
.pending_replies
.get(&rollback.sender_id)
.and_then(|pending| {
pending
.iter()
.find(|pending| pending.msg_id == entry.msg_id)
.cloned()
});
if rollback.pending_reply_after_send.as_ref() == Some(¤t_entry) {
let pending = proto
.pending_replies
.entry(rollback.sender_id.clone())
.or_default();
if let Some(existing) = pending
.iter_mut()
.find(|pending| pending.msg_id == entry.msg_id)
{
*existing = entry;
} else {
pending.push(entry);
}
}
}
if rollback.done {
let current_reminder = proto
.sessions
.get(&rollback.sender_id)
.and_then(|session| session.metadata.reminder.clone());
if rollback.sender_reminder_after_send.as_ref() == Some(¤t_reminder)
&& let Some(session) = proto.sessions.get_mut(&rollback.sender_id)
{
session.metadata.reminder = rollback.sender_reminder.flatten();
}
}
}
async fn finalize_successful_effect_delivery(
&self,
rollback: Option<FailedEffectSendRollback>,
) {
let Some(rollback) = rollback else {
return;
};
if !rollback.done {
return;
}
let mut proto = self.protocol.write().await;
if let Some(entry) = rollback.pending_reply_before_send {
if let Some(pending) = proto.pending_replies.get_mut(&rollback.sender_id) {
pending.retain(|pending| pending.msg_id != entry.msg_id);
if pending.is_empty() {
proto.pending_replies.remove(&rollback.sender_id);
}
}
}
if rollback.sender_reminder.is_some()
&& let Some(session) = proto.sessions.get_mut(&rollback.sender_id)
{
session.metadata.reminder = None;
}
}
pub(crate) fn persist_protocol_state(&self, proto: &crate::daemon_protocol::DaemonState) {
let sessions: HashMap<String, Session> = proto
.sessions
.iter()
.map(|(k, entry)| {
let m = &entry.metadata;
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: m.vim_mode,
project_dir: m.project_dir.clone(),
role: m.role.clone(),
networked: m.networked,
last_metadata_update: m
.last_metadata_update
.and_then(|ts| chrono::DateTime::from_timestamp(ts, 0)),
backend_session_id: m.backend_session_id.clone(),
backend: m.backend.clone(),
opencode_binding: m.opencode_binding.clone(),
restart_generation: m.restart_generation,
session_incarnation: m.session_incarnation,
project_description: m.project_description.clone(),
bulletin: m.bulletin.clone(),
worktree: m.worktree,
model: m.model.clone(),
effort: m.effort.clone(),
codex_home: m.codex_home.clone(),
reminder: m.reminder.clone(),
parent_session: m.parent_session.clone(),
idle_policy: m.idle_policy.clone(),
prompt: m.prompt.clone(),
iteration: m.iteration,
iteration_log: m.iteration_log.clone(),
last_iteration_at: m.last_iteration_at,
on_fire: m.on_fire.clone(),
worktree_present: m.worktree_present,
},
};
(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 async fn try_set_pending_compact_continuation(
&self,
session_id: &str,
text: String,
) -> bool {
let agents = self.session_agents.read().await;
if let Some(agent) = agents.get(session_id) {
ractor::call!(
agent,
crate::session_agent::SessionMsg::TrySetPendingCompactContinuation,
text
)
.unwrap_or(false)
} else {
false
}
}
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 async fn sweep_worktree_presence(self: &Arc<Self>) {
{
let mut backoff = self.sweep_backoff_until.lock().unwrap();
if let Some(until) = *backoff {
if std::time::Instant::now() < until {
tracing::debug!(
"worktree sweep in backoff window after recent timeout, skipping"
);
return;
}
*backoff = None;
self.sweep_in_progress
.store(false, std::sync::atomic::Ordering::Relaxed);
}
}
let sessions_with_dirs: Vec<(String, String)> = {
let proto = self.protocol.read().await;
proto
.sessions
.values()
.filter(|s| {
matches!(s.origin, crate::daemon_protocol::Origin::Local)
&& s.metadata.project_dir.is_some()
})
.filter_map(|s| Some((s.id.clone(), s.metadata.project_dir.clone()?)))
.collect()
};
if sessions_with_dirs.is_empty() {
return;
}
if self
.sweep_in_progress
.swap(true, std::sync::atomic::Ordering::Relaxed)
{
tracing::debug!("worktree sweep already in progress, skipping");
return;
}
let unique_dirs: Vec<String> = {
let mut dirs: Vec<String> = sessions_with_dirs.iter().map(|(_, d)| d.clone()).collect();
dirs.sort();
dirs.dedup();
dirs
};
const SWEEP_TIMEOUT_SECS: u64 = 30;
const SWEEP_BACKOFF_SECS: u64 = 300;
let unique_dirs = unique_dirs.clone();
let presence_map: std::collections::HashMap<String, bool> = match tokio::time::timeout(
std::time::Duration::from_secs(SWEEP_TIMEOUT_SECS),
tokio::task::spawn_blocking(move || {
let mut map = std::collections::HashMap::new();
for dir in unique_dirs {
let presence = match std::fs::metadata(&dir) {
Ok(m) if m.is_dir() => Some(true),
Ok(_) => {
tracing::debug!("worktree path exists but is not a directory: {}", dir);
None }
Err(e) if e.kind() == std::io::ErrorKind::NotFound => Some(false),
Err(e) => {
tracing::debug!("worktree stat failed for {}: {}", dir, e);
None }
};
if let Some(p) = presence {
map.insert(dir, p);
}
}
map
}),
)
.await
{
Ok(Ok(m)) => m,
Ok(Err(e)) => {
tracing::warn!("worktree sweep spawn_blocking failed: {e}");
self.sweep_in_progress
.store(false, std::sync::atomic::Ordering::Relaxed);
return;
}
Err(_) => {
tracing::warn!(
"worktree sweep timed out after {SWEEP_TIMEOUT_SECS}s - possible hung mount; \
backing off for {SWEEP_BACKOFF_SECS}s"
);
*self.sweep_backoff_until.lock().unwrap() = Some(
std::time::Instant::now() + std::time::Duration::from_secs(SWEEP_BACKOFF_SECS),
);
return;
}
};
let updates: Vec<(String, String, bool)> = sessions_with_dirs
.into_iter()
.filter_map(|(id, dir)| presence_map.get(&dir).map(|p| (id, dir.clone(), *p)))
.collect();
if !updates.is_empty() {
let _ = self
.apply_and_execute(crate::daemon_protocol::Event::MarkWorktreePresence { updates })
.await;
}
self.sweep_in_progress
.store(false, std::sync::atomic::Ordering::Relaxed);
}
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 list_assistant_panes(&self) -> Vec<crate::tmux::TmuxPane> {
if cfg!(test) {
return self.cached_assistant_panes().await;
}
let names: Vec<String> = self.backends.all_process_names();
tokio::task::spawn_blocking(move || {
let refs: Vec<&str> = names.iter().map(|s| s.as_str()).collect();
crate::tmux::find_assistant_panes(&refs).unwrap_or_default()
})
.await
.unwrap_or_default()
}
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 id = resolve_unique_session_id(&id_to_pane, &base_id, Some(pane.pane_id.as_str()));
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
}
}
fn rewrite_send_delivery_failure(
effects: Vec<crate::daemon_protocol::Effect>,
reason: &str,
) -> Vec<crate::daemon_protocol::Effect> {
effects
.into_iter()
.map(|effect| match effect {
crate::daemon_protocol::Effect::SendDelivered { from, to, .. } => {
crate::daemon_protocol::Effect::SendFailed {
from,
to,
reason: reason.to_string(),
renamed_to: None,
}
}
crate::daemon_protocol::Effect::LogMessage {
from,
to,
message,
delivered: true,
transport,
} => crate::daemon_protocol::Effect::LogMessage {
from,
to,
message,
delivered: false,
transport,
},
other => other,
})
.collect()
}
struct FailedEffectSendRollback {
sender_id: String,
pending_reply_before_send: Option<crate::daemon_protocol::PendingReplyEntry>,
pending_reply_after_send: Option<Option<crate::daemon_protocol::PendingReplyEntry>>,
sender_reminder: Option<Option<String>>,
sender_reminder_after_send: Option<Option<String>>,
sender_state_reserved: bool,
done: bool,
}
impl FailedEffectSendRollback {
fn capture_for_event(
proto: &crate::daemon_protocol::DaemonState,
event: &crate::daemon_protocol::Event,
) -> Option<Self> {
let crate::daemon_protocol::Event::Send {
from,
responds_to,
done,
..
} = event
else {
return None;
};
let pending_reply_before_send = responds_to.and_then(|msg_id| {
proto
.pending_replies
.get(from)
.and_then(|pending| pending.iter().find(|entry| entry.msg_id == msg_id).cloned())
});
Some(Self {
sender_id: from.clone(),
pending_reply_before_send,
pending_reply_after_send: None,
sender_reminder: done.then(|| {
proto
.sessions
.get(from)
.and_then(|session| session.metadata.reminder.clone())
}),
sender_reminder_after_send: None,
sender_state_reserved: false,
done: *done,
})
}
fn capture_after_send(&mut self, proto: &crate::daemon_protocol::DaemonState) {
if let Some(before) = &self.pending_reply_before_send {
self.pending_reply_after_send = Some(
proto
.pending_replies
.get(&self.sender_id)
.and_then(|pending| {
pending
.iter()
.find(|entry| entry.msg_id == before.msg_id)
.cloned()
}),
);
}
if self.done {
self.sender_reminder_after_send = Some(
proto
.sessions
.get(&self.sender_id)
.and_then(|session| session.metadata.reminder.clone()),
);
}
}
fn reserve_sender_state_after_send(&mut self, proto: &mut crate::daemon_protocol::DaemonState) {
if !self.done {
return;
}
if let Some(entry) = self.pending_reply_before_send.clone()
&& self.pending_reply_after_send == Some(None)
{
proto
.pending_replies
.entry(self.sender_id.clone())
.or_default()
.push(entry);
self.sender_state_reserved = true;
}
if self.sender_reminder.is_some()
&& self.sender_reminder_after_send == Some(None)
&& let Some(session) = proto.sessions.get_mut(&self.sender_id)
{
session.metadata.reminder = self.sender_reminder.clone().flatten();
self.sender_state_reserved = true;
}
}
fn sender_state_reserved(&self) -> bool {
self.sender_state_reserved
}
}
#[cfg(test)]
pub(crate) mod tests {
use super::*;
use crate::daemon_protocol::Origin;
#[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"
);
}
#[test]
fn resolve_unique_session_id_no_conflicts_returns_base() {
let map: HashMap<String, Option<String>> = HashMap::new();
assert_eq!(
resolve_unique_session_id(&map, "ouija", Some("%17")),
"ouija"
);
}
#[test]
fn resolve_unique_session_id_same_pane_returns_base_idempotent() {
let mut map = HashMap::new();
map.insert("ouija".into(), Some("%17".into()));
assert_eq!(
resolve_unique_session_id(&map, "ouija", Some("%17")),
"ouija"
);
}
#[test]
fn resolve_unique_session_id_distinct_pane_bumps_suffix() {
let mut map = HashMap::new();
map.insert("ouija".into(), Some("%17".into()));
assert_eq!(
resolve_unique_session_id(&map, "ouija", Some("%18")),
"ouija-2"
);
}
#[test]
fn resolve_unique_session_id_walks_through_taken_suffixes() {
let mut map = HashMap::new();
map.insert("ouija".into(), Some("%17".into()));
map.insert("ouija-2".into(), Some("%18".into()));
assert_eq!(
resolve_unique_session_id(&map, "ouija", Some("%19")),
"ouija-3"
);
}
#[test]
fn resolve_unique_session_id_no_target_pane_treats_existing_as_conflict() {
let mut map = HashMap::new();
map.insert("ouija".into(), None);
assert_eq!(resolve_unique_session_id(&map, "ouija", None), "ouija-2");
}
#[test]
fn resolve_unique_session_id_overflow_returns_last_attempted_id() {
let mut map = HashMap::new();
map.insert("ouija".into(), Some("%1".into()));
for n in 2..=MAX_NAME_SUFFIX {
map.insert(format!("ouija-{n}"), Some(format!("%{n}")));
}
let resolved = resolve_unique_session_id(&map, "ouija", Some("%9999"));
assert!(
resolved.starts_with("ouija"),
"expected resolved id to start with the base, got: {resolved}"
);
}
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(),
}
}
fn dead_opencode_serve_config() -> OuijaConfig {
test_config()
}
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 execute_effects_uses_recorded_tmux_method_for_send_inject() {
use axum::Router;
use axum::extract::State as AxumState;
use axum::http::StatusCode;
use axum::routing::post;
use std::sync::Arc as StdArc;
use std::sync::atomic::{AtomicUsize, Ordering};
use tokio::net::TcpListener;
async fn prompt_async(AxumState(calls): AxumState<StdArc<AtomicUsize>>) -> StatusCode {
calls.fetch_add(1, Ordering::SeqCst);
StatusCode::NO_CONTENT
}
let calls = StdArc::new(AtomicUsize::new(0));
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let port = listener.local_addr().unwrap().port();
let app = Router::new()
.route("/session/{session_id}/prompt_async", post(prompt_async))
.with_state(calls.clone());
let server = tokio::spawn(async move {
axum::serve(listener, app).await.unwrap();
});
let mut config = test_config();
config.port = port - 320;
let state = AppState::new(config);
{
let mut proto = state.protocol.write().await;
proto.sessions.insert(
"oc".into(),
crate::daemon_protocol::SessionEntry {
id: "oc".into(),
pane: Some("%1".into()),
origin: Origin::Local,
metadata: crate::daemon_protocol::SessionMeta {
backend: Some("opencode".into()),
backend_session_id: Some("ses_live".into()),
opencode_binding: Some(
crate::daemon_protocol::OpenCodeBinding::StrongManaged,
),
..Default::default()
},
registered_at: 0,
},
);
}
let effects = vec![
crate::daemon_protocol::Effect::InjectMessage {
session_id: "oc".into(),
pane: "%1".into(),
message: "hello".into(),
vim_mode: false,
delivery_method: None,
http_delivery: None,
pending_reply_msg_id: None,
pending_reply_from: None,
},
crate::daemon_protocol::Effect::SendDelivered {
from: "sender".into(),
to: "oc".into(),
method: "tmux".into(),
msg_id: 7,
http_delivery: None,
},
];
state.execute_effects(&effects).await;
assert_eq!(calls.load(Ordering::SeqCst), 0);
server.abort();
}
#[tokio::test]
async fn normal_tmux_inject_rejects_pane_not_owned_by_session() {
let state = AppState::new_for_test();
proto_register(&state, "target", Some("%1")).await;
let outcome = deliver_inject_message_effect(
&state,
InjectDeliveryRequest {
session_id: "target",
pane: "%2",
message: "hello",
vim_mode: false,
delivery_method: Some("tmux"),
recorded_method: None,
},
)
.await;
assert!(
matches!(outcome, DeliveryOutcome::Rejected(ref reason) if reason.contains("pane %2 is not owned by session target")),
"expected stale pane rejection, got {outcome:?}"
);
}
#[tokio::test]
async fn methodless_inject_rejects_pane_not_owned_by_session() {
let state = AppState::new_for_test();
proto_register(&state, "target", Some("%1")).await;
let outcome = deliver_inject_message_effect(
&state,
InjectDeliveryRequest {
session_id: "target",
pane: "%2",
message: "hello",
vim_mode: false,
delivery_method: None,
recorded_method: None,
},
)
.await;
assert!(
matches!(outcome, DeliveryOutcome::Rejected(ref reason) if reason.contains("pane %2 is not owned by session target")),
"expected stale pane rejection, got {outcome:?}"
);
}
#[tokio::test]
async fn methodless_inject_stale_pane_marks_delivery_failed() {
let state = AppState::new_for_test();
proto_register(&state, "sender", Some("%9")).await;
proto_register(&state, "target", Some("%1")).await;
let effects = vec![
crate::daemon_protocol::Effect::InjectMessage {
session_id: "target".into(),
pane: "%1".into(),
message: "hello".into(),
vim_mode: false,
delivery_method: None,
http_delivery: None,
pending_reply_msg_id: None,
pending_reply_from: None,
},
crate::daemon_protocol::Effect::LogMessage {
from: "sender".into(),
to: "target".into(),
message: "hello".into(),
delivered: true,
transport: "nostr".into(),
},
];
{
let mut proto = state.protocol.write().await;
proto.sessions.get_mut("target").unwrap().pane = Some("%2".into());
}
let failure = state.execute_effects(&effects).await;
assert!(
failure.as_ref().is_some_and(|failure| failure
.reason
.contains("pane %1 is not owned by session target")),
"expected stale pane failure, got {failure:?}"
);
let log = state.message_log.read().await;
assert_eq!(log.len(), 1);
assert!(!log[0].delivered);
}
#[tokio::test]
async fn execute_effects_delivers_http_from_recorded_snapshot_without_live_session() {
use axum::Router;
use axum::extract::State as AxumState;
use axum::http::StatusCode;
use axum::routing::post;
use std::sync::Arc as StdArc;
use std::sync::atomic::{AtomicUsize, Ordering};
use tokio::net::TcpListener;
async fn prompt_async(AxumState(calls): AxumState<StdArc<AtomicUsize>>) -> StatusCode {
calls.fetch_add(1, Ordering::SeqCst);
StatusCode::NO_CONTENT
}
let calls = StdArc::new(AtomicUsize::new(0));
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let port = listener.local_addr().unwrap().port();
let app = Router::new()
.route("/session/{session_id}/prompt_async", post(prompt_async))
.with_state(calls.clone());
let server = tokio::spawn(async move {
axum::serve(listener, app).await.unwrap();
});
let mut config = test_config();
config.port = port - 320;
let state = AppState::new(config);
let effects = vec![
crate::daemon_protocol::Effect::DeliverHttpMessage {
session_id: "oc".into(),
message: "hello".into(),
http_delivery: crate::daemon_protocol::HttpDeliverySnapshot {
backend_session_id: "ses_live".into(),
project_dir: None,
model: None,
effort: None,
},
pending_reply_msg_id: None,
pending_reply_from: None,
},
crate::daemon_protocol::Effect::SendDelivered {
from: "sender".into(),
to: "oc".into(),
method: "http".into(),
msg_id: 8,
http_delivery: Some(crate::daemon_protocol::HttpDeliverySnapshot {
backend_session_id: "ses_recorded".into(),
project_dir: None,
model: None,
effort: None,
}),
},
];
state.execute_effects(&effects).await;
assert_eq!(calls.load(Ordering::SeqCst), 1);
server.abort();
}
#[tokio::test]
async fn execute_effects_reports_strong_opencode_inject_failure_without_recorded_method() {
let state = AppState::new(dead_opencode_serve_config());
{
let mut proto = state.protocol.write().await;
proto.sessions.insert(
"oc".into(),
crate::daemon_protocol::SessionEntry {
id: "oc".into(),
pane: Some("%1".into()),
origin: Origin::Local,
metadata: crate::daemon_protocol::SessionMeta {
backend: Some("opencode".into()),
backend_session_id: Some("ses_live".into()),
opencode_binding: Some(
crate::daemon_protocol::OpenCodeBinding::StrongManaged,
),
..Default::default()
},
registered_at: 0,
},
);
}
let effects = vec![crate::daemon_protocol::Effect::InjectMessage {
session_id: "oc".into(),
pane: "%1".into(),
message: "hello".into(),
vim_mode: false,
delivery_method: Some("http".into()),
http_delivery: Some(crate::daemon_protocol::HttpDeliverySnapshot {
backend_session_id: "ses_live".into(),
project_dir: None,
model: None,
effort: None,
}),
pending_reply_msg_id: None,
pending_reply_from: None,
}];
let failure = state.execute_effects(&effects).await;
assert!(
failure
.as_ref()
.is_some_and(|failure| failure.reason.contains("prompt_async request failed")),
"expected observable HTTP delivery failure, got {failure:?}"
);
}
#[tokio::test]
async fn execute_effects_revalidates_http_inject_against_current_opencode_binding() {
use axum::Router;
use axum::extract::State as AxumState;
use axum::http::StatusCode;
use axum::routing::post;
use std::sync::Arc as StdArc;
use std::sync::atomic::{AtomicUsize, Ordering};
use tokio::net::TcpListener;
async fn prompt_async(AxumState(calls): AxumState<StdArc<AtomicUsize>>) -> StatusCode {
calls.fetch_add(1, Ordering::SeqCst);
StatusCode::NO_CONTENT
}
let calls = StdArc::new(AtomicUsize::new(0));
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let port = listener.local_addr().unwrap().port();
let app = Router::new()
.route("/session/{session_id}/prompt_async", post(prompt_async))
.with_state(calls.clone());
let server = tokio::spawn(async move {
axum::serve(listener, app).await.unwrap();
});
let mut config = test_config();
config.port = port - 320;
let state = AppState::new(config);
{
let mut proto = state.protocol.write().await;
proto.sessions.insert(
"oc".into(),
crate::daemon_protocol::SessionEntry {
id: "oc".into(),
pane: Some("%1".into()),
origin: Origin::Local,
metadata: crate::daemon_protocol::SessionMeta {
backend: Some("opencode".into()),
backend_session_id: Some("ses_live".into()),
opencode_binding: Some(
crate::daemon_protocol::OpenCodeBinding::WeakAdopted,
),
..Default::default()
},
registered_at: 0,
},
);
}
let effects = vec![crate::daemon_protocol::Effect::InjectMessage {
session_id: "oc".into(),
pane: "%1".into(),
message: "hello".into(),
vim_mode: false,
delivery_method: Some("http".into()),
http_delivery: Some(crate::daemon_protocol::HttpDeliverySnapshot {
backend_session_id: "ses_live".into(),
project_dir: None,
model: None,
effort: None,
}),
pending_reply_msg_id: None,
pending_reply_from: None,
}];
let failure = state.execute_effects(&effects).await;
assert!(failure.is_none());
assert_eq!(
calls.load(Ordering::SeqCst),
0,
"stale/forged HTTP inject metadata must not bypass the shared OpenCode delivery gate"
);
server.abort();
}
#[tokio::test]
async fn execute_effects_rejects_strong_opencode_http_inject_after_session_moves_panes() {
use axum::Router;
use axum::extract::State as AxumState;
use axum::http::StatusCode;
use axum::routing::post;
use std::sync::Arc as StdArc;
use std::sync::atomic::{AtomicUsize, Ordering};
use tokio::net::TcpListener;
async fn prompt_async(AxumState(calls): AxumState<StdArc<AtomicUsize>>) -> StatusCode {
calls.fetch_add(1, Ordering::SeqCst);
StatusCode::NO_CONTENT
}
let calls = StdArc::new(AtomicUsize::new(0));
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let port = listener.local_addr().unwrap().port();
let app = Router::new()
.route("/session/{session_id}/prompt_async", post(prompt_async))
.with_state(calls.clone());
let server = tokio::spawn(async move {
axum::serve(listener, app).await.unwrap();
});
let mut config = test_config();
config.port = port - 320;
let state = AppState::new(config);
{
let mut proto = state.protocol.write().await;
proto.sessions.insert(
"oc".into(),
crate::daemon_protocol::SessionEntry {
id: "oc".into(),
pane: Some("%1".into()),
origin: Origin::Local,
metadata: crate::daemon_protocol::SessionMeta {
backend: Some("opencode".into()),
backend_session_id: Some("ses_live".into()),
opencode_binding: Some(
crate::daemon_protocol::OpenCodeBinding::StrongManaged,
),
..Default::default()
},
registered_at: 0,
},
);
}
let effects = vec![crate::daemon_protocol::Effect::InjectMessage {
session_id: "oc".into(),
pane: "%1".into(),
message: "hello".into(),
vim_mode: false,
delivery_method: Some("http".into()),
http_delivery: Some(crate::daemon_protocol::HttpDeliverySnapshot {
backend_session_id: "ses_live".into(),
project_dir: None,
model: None,
effort: None,
}),
pending_reply_msg_id: None,
pending_reply_from: None,
}];
{
let mut proto = state.protocol.write().await;
proto.sessions.get_mut("oc").unwrap().pane = Some("%2".into());
}
let failure = state.execute_effects(&effects).await;
assert!(
failure.as_ref().is_some_and(|failure| failure
.reason
.contains("pane %1 is not owned by session oc")),
"expected stale pane rejection, got {failure:?}"
);
assert_eq!(
calls.load(Ordering::SeqCst),
0,
"stale HTTP inject must not call prompt_async"
);
server.abort();
}
#[tokio::test]
async fn incoming_weak_opencode_inject_uses_apply_time_delivery_method() {
let state = AppState::new(dead_opencode_serve_config());
let effects = {
let mut proto = state.protocol.write().await;
proto.sessions.insert(
"oc".into(),
crate::daemon_protocol::SessionEntry {
id: "oc".into(),
pane: Some("%17".into()),
origin: Origin::Local,
metadata: crate::daemon_protocol::SessionMeta {
backend: Some("opencode".into()),
backend_session_id: Some("ses_old".into()),
opencode_binding: Some(
crate::daemon_protocol::OpenCodeBinding::WeakAdopted,
),
..Default::default()
},
registered_at: 0,
},
);
proto.apply(crate::daemon_protocol::Event::IncomingWire {
msg: crate::protocol::WireMessage::SessionSend {
from: "remote".into(),
to: "oc".into(),
message: "hello".into(),
expects_reply: false,
msg_id: 42,
responds_to: None,
done: false,
},
sender_npub: Some("npub1remote".into()),
})
};
{
let mut proto = state.protocol.write().await;
let session = proto.sessions.get_mut("oc").unwrap();
session.metadata.backend_session_id = Some("ses_new".into());
session.metadata.opencode_binding =
Some(crate::daemon_protocol::OpenCodeBinding::StrongManaged);
}
let failure = state.execute_effects(&effects).await;
assert!(failure.is_none());
}
#[tokio::test]
async fn execute_effects_broadcasts_failure_ack_after_inject_failure() {
use std::sync::Arc as StdArc;
use std::sync::atomic::{AtomicUsize, Ordering};
struct CountingTransport {
broadcasts: StdArc<AtomicUsize>,
failure_acks: StdArc<AtomicUsize>,
}
#[async_trait::async_trait]
impl crate::transport::Transport for CountingTransport {
fn as_any(&self) -> &dyn std::any::Any {
self
}
async fn broadcast(&self, msg: &crate::protocol::WireMessage) -> bool {
self.broadcasts.fetch_add(1, Ordering::SeqCst);
if matches!(
msg,
crate::protocol::WireMessage::SessionSendAck {
delivered: false,
..
}
) {
self.failure_acks.fetch_add(1, Ordering::SeqCst);
}
true
}
async fn connect(
&self,
_ticket: &str,
_state: Arc<AppState>,
_wait: bool,
) -> anyhow::Result<()> {
Ok(())
}
async fn ticket_string(&self) -> Option<String> {
None
}
async fn regenerate(
&self,
_config_dir: &std::path::Path,
_data_dir: &std::path::Path,
) -> anyhow::Result<String> {
Ok("ticket".into())
}
fn endpoint_id(&self) -> Option<String> {
None
}
fn is_ready(&self) -> bool {
true
}
fn transport_name(&self) -> &'static str {
"counting"
}
}
let state = AppState::new(dead_opencode_serve_config());
let broadcasts = StdArc::new(AtomicUsize::new(0));
let failure_acks = StdArc::new(AtomicUsize::new(0));
state
.add_transport(StdArc::new(CountingTransport {
broadcasts: broadcasts.clone(),
failure_acks: failure_acks.clone(),
}))
.await;
{
let mut proto = state.protocol.write().await;
proto.sessions.insert(
"oc".into(),
crate::daemon_protocol::SessionEntry {
id: "oc".into(),
pane: Some("%1".into()),
origin: Origin::Local,
metadata: crate::daemon_protocol::SessionMeta {
backend: Some("opencode".into()),
backend_session_id: Some("ses_live".into()),
opencode_binding: Some(
crate::daemon_protocol::OpenCodeBinding::StrongManaged,
),
..Default::default()
},
registered_at: 0,
},
);
}
let effects = vec![
crate::daemon_protocol::Effect::InjectMessage {
session_id: "oc".into(),
pane: "%1".into(),
message: "hello".into(),
vim_mode: false,
delivery_method: Some("http".into()),
http_delivery: Some(crate::daemon_protocol::HttpDeliverySnapshot {
backend_session_id: "ses_live".into(),
project_dir: None,
model: None,
effort: None,
}),
pending_reply_msg_id: None,
pending_reply_from: None,
},
crate::daemon_protocol::Effect::Broadcast(
crate::protocol::WireMessage::SessionSendAck {
from: "remote".into(),
to: "oc".into(),
delivered: true,
daemon_id: "remote-daemon".into(),
},
),
];
let failure = state.execute_effects(&effects).await;
assert!(failure.is_some());
assert_eq!(broadcasts.load(Ordering::SeqCst), 1);
assert_eq!(failure_acks.load(Ordering::SeqCst), 1);
}
#[tokio::test]
async fn execute_effects_does_not_rewrite_ack_after_ambiguous_http_inject_failure() {
use axum::Router;
use axum::http::StatusCode;
use axum::routing::post;
use std::sync::Arc as StdArc;
use std::sync::atomic::{AtomicUsize, Ordering};
struct CountingTransport {
success_acks: StdArc<AtomicUsize>,
failure_acks: StdArc<AtomicUsize>,
}
#[async_trait::async_trait]
impl crate::transport::Transport for CountingTransport {
fn as_any(&self) -> &dyn std::any::Any {
self
}
async fn broadcast(&self, msg: &crate::protocol::WireMessage) -> bool {
match msg {
crate::protocol::WireMessage::SessionSendAck {
delivered: true, ..
} => {
self.success_acks.fetch_add(1, Ordering::SeqCst);
}
crate::protocol::WireMessage::SessionSendAck {
delivered: false, ..
} => {
self.failure_acks.fetch_add(1, Ordering::SeqCst);
}
_ => {}
}
true
}
async fn connect(
&self,
_ticket: &str,
_state: Arc<AppState>,
_wait: bool,
) -> anyhow::Result<()> {
Ok(())
}
async fn ticket_string(&self) -> Option<String> {
None
}
async fn regenerate(
&self,
_config_dir: &std::path::Path,
_data_dir: &std::path::Path,
) -> anyhow::Result<String> {
Ok("ticket".into())
}
fn endpoint_id(&self) -> Option<String> {
None
}
fn is_ready(&self) -> bool {
true
}
fn transport_name(&self) -> &'static str {
"counting"
}
}
async fn prompt_async() -> StatusCode {
StatusCode::INTERNAL_SERVER_ERROR
}
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
let port = listener.local_addr().unwrap().port();
let app = Router::new().route("/session/{session_id}/prompt_async", post(prompt_async));
let server = tokio::spawn(async move {
axum::serve(listener, app).await.unwrap();
});
let mut config = test_config();
config.port = port.checked_sub(320).unwrap();
let state = AppState::new(config);
{
let mut proto = state.protocol.write().await;
proto.sessions.insert(
"oc".into(),
crate::daemon_protocol::SessionEntry {
id: "oc".into(),
pane: Some("%1".into()),
origin: Origin::Local,
metadata: crate::daemon_protocol::SessionMeta {
backend: Some("opencode".into()),
backend_session_id: Some("ses_live".into()),
opencode_binding: Some(
crate::daemon_protocol::OpenCodeBinding::StrongManaged,
),
..Default::default()
},
registered_at: 0,
},
);
}
let success_acks = StdArc::new(AtomicUsize::new(0));
let failure_acks = StdArc::new(AtomicUsize::new(0));
state
.add_transport(StdArc::new(CountingTransport {
success_acks: success_acks.clone(),
failure_acks: failure_acks.clone(),
}))
.await;
let effects = vec![
crate::daemon_protocol::Effect::InjectMessage {
session_id: "oc".into(),
pane: "%1".into(),
message: "hello".into(),
vim_mode: false,
delivery_method: Some("http".into()),
http_delivery: Some(crate::daemon_protocol::HttpDeliverySnapshot {
backend_session_id: "ses_live".into(),
project_dir: None,
model: None,
effort: None,
}),
pending_reply_msg_id: None,
pending_reply_from: None,
},
crate::daemon_protocol::Effect::Broadcast(
crate::protocol::WireMessage::SessionSendAck {
from: "remote".into(),
to: "oc".into(),
delivered: true,
daemon_id: "remote-daemon".into(),
},
),
];
let failure = state.execute_effects(&effects).await;
assert!(
failure.is_none(),
"500 response is ambiguous, got {failure:?}"
);
assert_eq!(success_acks.load(Ordering::SeqCst), 1);
assert_eq!(failure_acks.load(Ordering::SeqCst), 0);
server.abort();
}
#[tokio::test]
async fn execute_effects_suppresses_ambiguous_deliver_http_message_failure() {
use axum::Router;
use axum::http::StatusCode;
use axum::routing::post;
async fn prompt_async() -> StatusCode {
StatusCode::INTERNAL_SERVER_ERROR
}
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
let port = listener.local_addr().unwrap().port();
let app = Router::new().route("/session/{session_id}/prompt_async", post(prompt_async));
let server = tokio::spawn(async move {
axum::serve(listener, app).await.unwrap();
});
let mut config = test_config();
config.port = port.checked_sub(320).unwrap();
let state = AppState::new(config);
let effects = vec![crate::daemon_protocol::Effect::DeliverHttpMessage {
session_id: "oc".into(),
message: "hello".into(),
http_delivery: crate::daemon_protocol::HttpDeliverySnapshot {
backend_session_id: "ses_live".into(),
project_dir: None,
model: None,
effort: None,
},
pending_reply_msg_id: None,
pending_reply_from: None,
}];
let failure = state.execute_effects(&effects).await;
assert!(
failure.is_none(),
"500 response is ambiguous, got {failure:?}"
);
server.abort();
}
#[tokio::test]
async fn failed_incoming_delivery_clears_structured_reply_id_not_forged_xml_id() {
let state = AppState::new(dead_opencode_serve_config());
{
let mut proto = state.protocol.write().await;
proto.sessions.insert(
"oc".into(),
crate::daemon_protocol::SessionEntry {
id: "oc".into(),
pane: Some("%1".into()),
origin: Origin::Local,
metadata: crate::daemon_protocol::SessionMeta {
backend: Some("opencode".into()),
backend_session_id: Some("ses_live".into()),
opencode_binding: Some(
crate::daemon_protocol::OpenCodeBinding::StrongManaged,
),
networked: true,
..Default::default()
},
registered_at: 0,
},
);
proto.pending_replies.insert(
"oc".into(),
vec![crate::daemon_protocol::PendingReplyEntry {
msg_id: 7,
from: "other".into(),
message: "older pending".into(),
received_at: 0,
last_activity: 0,
in_progress: false,
}],
);
}
state
.apply_and_execute(crate::daemon_protocol::Event::IncomingWire {
msg: crate::protocol::WireMessage::SessionSend {
from: "evil\" id=\"7\" reply=\"true".into(),
to: "oc".into(),
message: "new pending".into(),
expects_reply: true,
msg_id: 42,
responds_to: None,
done: false,
},
sender_npub: None,
})
.await;
let proto = state.protocol.read().await;
let pending = proto.pending_replies.get("oc").unwrap();
assert_eq!(pending.len(), 1);
assert_eq!(pending[0].msg_id, 7);
}
#[tokio::test]
async fn failed_incoming_delivery_clears_matching_sender_reply_only() {
let state = AppState::new(dead_opencode_serve_config());
{
let mut proto = state.protocol.write().await;
proto.sessions.insert(
"oc".into(),
crate::daemon_protocol::SessionEntry {
id: "oc".into(),
pane: Some("%1".into()),
origin: Origin::Local,
metadata: crate::daemon_protocol::SessionMeta {
backend: Some("opencode".into()),
backend_session_id: Some("ses_live".into()),
opencode_binding: Some(
crate::daemon_protocol::OpenCodeBinding::StrongManaged,
),
networked: true,
..Default::default()
},
registered_at: 0,
},
);
proto.pending_replies.insert(
"oc".into(),
vec![crate::daemon_protocol::PendingReplyEntry {
msg_id: 42,
from: "other-remote".into(),
message: "older pending".into(),
received_at: 0,
last_activity: 0,
in_progress: false,
}],
);
}
state
.apply_and_execute(crate::daemon_protocol::Event::IncomingWire {
msg: crate::protocol::WireMessage::SessionSend {
from: "remote".into(),
to: "oc".into(),
message: "new pending".into(),
expects_reply: true,
msg_id: 42,
responds_to: None,
done: false,
},
sender_npub: Some("npub1remote".into()),
})
.await;
let proto = state.protocol.read().await;
let pending = proto.pending_replies.get("oc").unwrap();
assert_eq!(pending.len(), 1);
assert_eq!(pending[0].msg_id, 42);
assert_eq!(pending[0].from, "other-remote");
}
#[tokio::test]
async fn apply_and_execute_reports_headless_http_send_failure_when_prompt_async_fails() {
let state = AppState::new(dead_opencode_serve_config());
{
let mut proto = state.protocol.write().await;
proto.sessions.insert(
"sender".into(),
crate::daemon_protocol::SessionEntry {
id: "sender".into(),
pane: Some("%1".into()),
origin: Origin::Local,
metadata: crate::daemon_protocol::SessionMeta::default(),
registered_at: 0,
},
);
proto.sessions.insert(
"oc".into(),
crate::daemon_protocol::SessionEntry {
id: "oc".into(),
pane: None,
origin: Origin::Local,
metadata: crate::daemon_protocol::SessionMeta {
backend: Some("opencode".into()),
backend_session_id: Some("ses_headless".into()),
opencode_binding: Some(
crate::daemon_protocol::OpenCodeBinding::StrongManaged,
),
..Default::default()
},
registered_at: 0,
},
);
}
let effects = state
.apply_and_execute(crate::daemon_protocol::Event::Send {
from: "sender".into(),
to: "oc".into(),
message: "hello".into(),
expects_reply: true,
responds_to: None,
done: false,
})
.await;
assert!(effects.iter().any(|effect| {
matches!(
effect,
crate::daemon_protocol::Effect::SendFailed { reason, .. }
if reason.contains("prompt_async request failed")
)
}));
assert!(
!effects.iter().any(|effect| matches!(
effect,
crate::daemon_protocol::Effect::SendDelivered { .. }
))
);
let log = state.message_log.read().await;
assert_eq!(log.len(), 1);
assert!(!log[0].delivered);
drop(log);
let proto = state.protocol.read().await;
assert!(!proto.pending_replies.contains_key("oc"));
}
#[tokio::test]
async fn apply_and_execute_clears_incoming_pending_reply_after_inject_failure() {
let state = AppState::new(dead_opencode_serve_config());
{
let mut proto = state.protocol.write().await;
proto.sessions.insert(
"oc".into(),
crate::daemon_protocol::SessionEntry {
id: "oc".into(),
pane: Some("%17".into()),
origin: Origin::Local,
metadata: crate::daemon_protocol::SessionMeta {
backend: Some("opencode".into()),
backend_session_id: Some("ses_incoming".into()),
opencode_binding: Some(
crate::daemon_protocol::OpenCodeBinding::StrongManaged,
),
..Default::default()
},
registered_at: 0,
},
);
}
let effects = state
.apply_and_execute(crate::daemon_protocol::Event::IncomingWire {
msg: crate::protocol::WireMessage::SessionSend {
from: "remote".into(),
to: "oc".into(),
message: "hello".into(),
expects_reply: true,
msg_id: 42,
responds_to: None,
done: false,
},
sender_npub: Some("npub1remote".into()),
})
.await;
assert!(effects.iter().any(|effect| matches!(
effect,
crate::daemon_protocol::Effect::LogMessage {
delivered: false,
..
}
)));
let proto = state.protocol.read().await;
assert!(!proto.pending_replies.contains_key("oc"));
}
#[tokio::test]
async fn apply_and_execute_clears_incoming_pending_reply_after_headless_http_failure() {
let state = AppState::new(dead_opencode_serve_config());
{
let mut proto = state.protocol.write().await;
proto.sessions.insert(
"oc".into(),
crate::daemon_protocol::SessionEntry {
id: "oc".into(),
pane: None,
origin: Origin::Local,
metadata: crate::daemon_protocol::SessionMeta {
backend: Some("opencode".into()),
backend_session_id: Some("ses_headless".into()),
opencode_binding: Some(
crate::daemon_protocol::OpenCodeBinding::StrongManaged,
),
networked: true,
..Default::default()
},
registered_at: 0,
},
);
}
let effects = state
.apply_and_execute(crate::daemon_protocol::Event::IncomingWire {
msg: crate::protocol::WireMessage::SessionSend {
from: "remote".into(),
to: "oc".into(),
message: "hello".into(),
expects_reply: true,
msg_id: 42,
responds_to: None,
done: false,
},
sender_npub: Some("npub1remote".into()),
})
.await;
assert!(effects.iter().any(|effect| matches!(
effect,
crate::daemon_protocol::Effect::LogMessage {
delivered: false,
..
}
)));
let proto = state.protocol.read().await;
assert!(!proto.pending_replies.contains_key("oc"));
}
#[tokio::test]
async fn apply_and_execute_restores_sender_reply_state_after_delivery_failure() {
let state = AppState::new(dead_opencode_serve_config());
{
let mut proto = state.protocol.write().await;
proto.sessions.insert(
"sender".into(),
crate::daemon_protocol::SessionEntry {
id: "sender".into(),
pane: Some("%1".into()),
origin: Origin::Local,
metadata: crate::daemon_protocol::SessionMeta {
reminder: Some("keep working".into()),
..Default::default()
},
registered_at: 0,
},
);
proto.sessions.insert(
"oc".into(),
crate::daemon_protocol::SessionEntry {
id: "oc".into(),
pane: None,
origin: Origin::Local,
metadata: crate::daemon_protocol::SessionMeta {
backend: Some("opencode".into()),
backend_session_id: Some("ses_headless".into()),
opencode_binding: Some(
crate::daemon_protocol::OpenCodeBinding::StrongManaged,
),
..Default::default()
},
registered_at: 0,
},
);
proto.pending_replies.insert(
"sender".into(),
vec![crate::daemon_protocol::PendingReplyEntry {
msg_id: 7,
from: "requester".into(),
message: "please respond".into(),
received_at: 100,
last_activity: 100,
in_progress: false,
}],
);
}
let effects = state
.apply_and_execute(crate::daemon_protocol::Event::Send {
from: "sender".into(),
to: "oc".into(),
message: "done, but unreachable".into(),
expects_reply: false,
responds_to: Some(7),
done: true,
})
.await;
assert!(effects.iter().any(|effect| {
matches!(
effect,
crate::daemon_protocol::Effect::SendFailed { reason, .. }
if reason.contains("prompt_async request failed")
)
}));
let proto = state.protocol.read().await;
let pending = proto.pending_replies.get("sender").unwrap();
assert_eq!(pending.len(), 1);
assert_eq!(pending[0].msg_id, 7);
assert_eq!(
proto.sessions["sender"].metadata.reminder.as_deref(),
Some("keep working")
);
}
#[tokio::test]
async fn apply_and_execute_restores_sender_state_after_send_failed_before_delivery() {
let state = AppState::new_for_test();
{
let mut proto = state.protocol.write().await;
proto.sessions.insert(
"sender".into(),
crate::daemon_protocol::SessionEntry {
id: "sender".into(),
pane: Some("%1".into()),
origin: Origin::Local,
metadata: crate::daemon_protocol::SessionMeta {
reminder: Some("keep working".into()),
..Default::default()
},
registered_at: 0,
},
);
proto.pending_replies.insert(
"sender".into(),
vec![crate::daemon_protocol::PendingReplyEntry {
msg_id: 7,
from: "requester".into(),
message: "please respond".into(),
received_at: 100,
last_activity: 100,
in_progress: false,
}],
);
}
let effects = state
.apply_and_execute(crate::daemon_protocol::Event::Send {
from: "sender".into(),
to: "missing".into(),
message: "done, but missing".into(),
expects_reply: false,
responds_to: Some(7),
done: true,
})
.await;
assert!(effects.iter().any(|effect| matches!(
effect,
crate::daemon_protocol::Effect::SendFailed { to, .. } if to == "missing"
)));
let proto = state.protocol.read().await;
let pending = proto.pending_replies.get("sender").unwrap();
assert_eq!(pending.len(), 1);
assert_eq!(pending[0].msg_id, 7);
assert_eq!(
proto.sessions["sender"].metadata.reminder.as_deref(),
Some("keep working")
);
}
#[tokio::test]
async fn apply_and_execute_does_not_restore_concurrently_cleared_sender_reply_state() {
use axum::Router;
use axum::extract::State as AxumState;
use axum::http::StatusCode;
use axum::routing::post;
use std::sync::Arc as StdArc;
use tokio::sync::Notify;
#[derive(Clone)]
struct Gate {
started: StdArc<Notify>,
release: StdArc<Notify>,
}
async fn prompt_async(AxumState(gate): AxumState<Gate>) -> StatusCode {
gate.started.notify_one();
gate.release.notified().await;
StatusCode::NOT_FOUND
}
let gate = Gate {
started: StdArc::new(Notify::new()),
release: StdArc::new(Notify::new()),
};
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
let port = listener.local_addr().unwrap().port();
let app = Router::new()
.route("/session/{session_id}/prompt_async", post(prompt_async))
.with_state(gate.clone());
let server = tokio::spawn(async move {
axum::serve(listener, app).await.unwrap();
});
let mut config = test_config();
config.port = port.checked_sub(320).unwrap();
let state = AppState::new(config);
{
let mut proto = state.protocol.write().await;
proto.sessions.insert(
"sender".into(),
crate::daemon_protocol::SessionEntry {
id: "sender".into(),
pane: Some("%1".into()),
origin: Origin::Local,
metadata: crate::daemon_protocol::SessionMeta {
reminder: Some("keep working".into()),
..Default::default()
},
registered_at: 0,
},
);
proto.sessions.insert(
"oc".into(),
crate::daemon_protocol::SessionEntry {
id: "oc".into(),
pane: None,
origin: Origin::Local,
metadata: crate::daemon_protocol::SessionMeta {
backend: Some("opencode".into()),
backend_session_id: Some("ses_headless".into()),
opencode_binding: Some(
crate::daemon_protocol::OpenCodeBinding::StrongManaged,
),
..Default::default()
},
registered_at: 0,
},
);
proto.pending_replies.insert(
"sender".into(),
vec![crate::daemon_protocol::PendingReplyEntry {
msg_id: 7,
from: "requester".into(),
message: "please respond".into(),
received_at: 100,
last_activity: 100,
in_progress: false,
}],
);
}
let delivery = tokio::spawn({
let state = state.clone();
async move {
state
.apply_and_execute(crate::daemon_protocol::Event::Send {
from: "sender".into(),
to: "oc".into(),
message: "done, but unreachable".into(),
expects_reply: false,
responds_to: Some(7),
done: true,
})
.await
}
});
gate.started.notified().await;
{
let mut proto = state.protocol.write().await;
proto.pending_replies.remove("sender");
proto.sessions.get_mut("sender").unwrap().metadata.reminder = None;
}
gate.release.notify_one();
let effects = delivery.await.unwrap();
assert!(effects.iter().any(|effect| {
matches!(effect, crate::daemon_protocol::Effect::SendFailed { reason, .. } if reason.contains("prompt_async"))
}));
let proto = state.protocol.read().await;
assert!(!proto.pending_replies.contains_key("sender"));
assert_eq!(proto.sessions["sender"].metadata.reminder, None);
server.abort();
}
#[tokio::test]
async fn successful_delivery_clears_sender_state_by_msg_id_after_mutations() {
let state = AppState::new_for_test();
let original_entry = crate::daemon_protocol::PendingReplyEntry {
msg_id: 7,
from: "requester".into(),
message: "please respond".into(),
received_at: 100,
last_activity: 100,
in_progress: false,
};
{
let mut proto = state.protocol.write().await;
proto.sessions.insert(
"sender".into(),
crate::daemon_protocol::SessionEntry {
id: "sender".into(),
pane: Some("%1".into()),
origin: Origin::Local,
metadata: crate::daemon_protocol::SessionMeta {
reminder: Some("keep working (activity tick)".into()),
..Default::default()
},
registered_at: 0,
},
);
proto.pending_replies.insert(
"sender".into(),
vec![crate::daemon_protocol::PendingReplyEntry {
last_activity: 200,
in_progress: true,
..original_entry.clone()
}],
);
}
state
.finalize_successful_effect_delivery(Some(FailedEffectSendRollback {
sender_id: "sender".into(),
pending_reply_before_send: Some(original_entry),
pending_reply_after_send: None,
sender_reminder: Some(Some("keep working".into())),
sender_reminder_after_send: None,
sender_state_reserved: false,
done: true,
}))
.await;
let proto = state.protocol.read().await;
assert!(!proto.pending_replies.contains_key("sender"));
assert_eq!(proto.sessions["sender"].metadata.reminder, None);
}
#[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 persist_protocol_state_round_trips_all_metadata_fields() {
let config = test_config();
let state = AppState::new(config.clone());
let meta = crate::daemon_protocol::SessionMeta {
project_dir: Some("/tmp/proj".into()),
role: Some("worker".into()),
networked: false,
bulletin: Some("available".into()),
last_metadata_update: Some(1_700_000_100),
backend_session_id: Some("oc_abc123".into()),
backend: Some("opencode".into()),
opencode_binding: Some(crate::daemon_protocol::OpenCodeBinding::StrongManaged),
restart_generation: 7,
session_incarnation: 11,
project_description: Some("test project".into()),
vim_mode: true,
worktree: true,
model: Some("openrouter/sonnet".into()),
effort: Some("max".into()),
codex_home: None,
reminder: Some("remember to...".into()),
parent_session: Some("parent".into()),
idle_policy: Some(crate::daemon_protocol::IdlePolicy::AskParentWhenDone),
prompt: Some("do the thing".into()),
iteration: 3,
iteration_log: vec![],
last_iteration_at: Some(1_700_000_000),
on_fire: Some(crate::scheduler::OnFire::NewSession),
worktree_present: Some(false),
};
state
.apply_and_execute(crate::daemon_protocol::Event::Register {
id: "s1".into(),
pane: Some("%1".into()),
metadata: meta,
})
.await;
{
let proto = state.protocol.read().await;
state.persist_protocol_state(&proto);
}
let loaded = crate::persistence::load_sessions(&config.data_dir)
.expect("load_sessions after persist");
let s = loaded
.iter()
.find(|p| p.id == "s1")
.expect("session s1 not persisted");
assert_eq!(
s.metadata.model.as_deref(),
Some("openrouter/sonnet"),
"model dropped by persist"
);
assert_eq!(
s.metadata.effort.as_deref(),
Some("max"),
"effort dropped by persist"
);
assert_eq!(
s.metadata.backend.as_deref(),
Some("opencode"),
"backend dropped by persist"
);
assert_eq!(
s.metadata.backend_session_id.as_deref(),
Some("oc_abc123"),
"backend_session_id dropped by persist"
);
assert_eq!(
s.metadata.opencode_binding,
Some(crate::daemon_protocol::OpenCodeBinding::StrongManaged),
"opencode_binding dropped by persist"
);
assert_eq!(
s.metadata.restart_generation, 7,
"restart_generation dropped by persist"
);
assert_eq!(
s.metadata.project_description.as_deref(),
Some("test project"),
"project_description dropped by persist"
);
assert!(
s.metadata.last_metadata_update.is_some(),
"last_metadata_update dropped by persist"
);
assert_eq!(
s.metadata.last_iteration_at,
Some(1_700_000_000),
"last_iteration_at dropped by persist"
);
assert!(s.metadata.on_fire.is_some(), "on_fire dropped by persist");
assert_eq!(
s.metadata.role.as_deref(),
Some("worker"),
"role dropped by persist"
);
assert_eq!(
s.metadata.bulletin.as_deref(),
Some("available"),
"bulletin dropped by persist"
);
assert_eq!(
s.metadata.reminder.as_deref(),
Some("remember to..."),
"reminder preserved"
);
assert_eq!(
s.metadata.prompt.as_deref(),
Some("do the thing"),
"prompt preserved"
);
assert!(s.metadata.vim_mode, "vim_mode preserved");
assert!(s.metadata.worktree, "worktree preserved");
assert!(!s.metadata.networked, "networked=false preserved");
assert_eq!(s.metadata.iteration, 3, "iteration preserved");
assert_eq!(
s.metadata.worktree_present,
Some(false),
"worktree_present dropped by persist (issue #661)"
);
let hydrated = crate::daemon_protocol::metadata_to_session_meta_for_test(&s.metadata);
assert_eq!(hydrated.model.as_deref(), Some("openrouter/sonnet"));
assert_eq!(hydrated.effort.as_deref(), Some("max"));
assert_eq!(hydrated.backend.as_deref(), Some("opencode"));
assert_eq!(hydrated.backend_session_id.as_deref(), Some("oc_abc123"));
assert_eq!(
hydrated.opencode_binding,
Some(crate::daemon_protocol::OpenCodeBinding::StrongManaged)
);
assert_eq!(hydrated.restart_generation, 7);
assert!(hydrated.on_fire.is_some());
assert_eq!(hydrated.last_iteration_at, Some(1_700_000_000));
assert_eq!(hydrated.last_metadata_update, Some(1_700_000_100));
assert_eq!(hydrated.worktree_present, Some(false));
}
#[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());
}
#[tokio::test]
async fn sweep_worktree_presence_sets_true_for_existing_dir() {
let state = AppState::new_for_test();
let tempdir = tempfile::tempdir().unwrap();
let project_dir = tempdir.path().to_str().unwrap().to_string();
{
let mut proto = state.protocol.write().await;
proto.sessions.insert(
"local/s1".into(),
crate::daemon_protocol::SessionEntry {
id: "local/s1".into(),
pane: Some("%1".into()),
origin: Origin::Local,
metadata: crate::daemon_protocol::SessionMeta {
project_dir: Some(project_dir.clone()),
..Default::default()
},
registered_at: 0,
},
);
}
state.sweep_worktree_presence().await;
{
let proto = state.protocol.read().await;
let session = proto.sessions.get("local/s1").unwrap();
assert_eq!(
session.metadata.worktree_present,
Some(true),
"existing dir should show as present"
);
}
}
#[tokio::test]
async fn sweep_worktree_presence_sets_false_for_missing_dir() {
let state = AppState::new_for_test();
let missing_dir = "/tmp/ouija-test-nonexistent-dir-12345".to_string();
{
let mut proto = state.protocol.write().await;
proto.sessions.insert(
"local/s1".into(),
crate::daemon_protocol::SessionEntry {
id: "local/s1".into(),
pane: Some("%1".into()),
origin: Origin::Local,
metadata: crate::daemon_protocol::SessionMeta {
project_dir: Some(missing_dir.clone()),
..Default::default()
},
registered_at: 0,
},
);
}
state.sweep_worktree_presence().await;
{
let proto = state.protocol.read().await;
let session = proto.sessions.get("local/s1").unwrap();
assert_eq!(
session.metadata.worktree_present,
Some(false),
"missing dir should show as absent"
);
}
}
#[tokio::test]
async fn sweep_worktree_presence_skips_non_local() {
let state = AppState::new_for_test();
let tempdir = tempfile::tempdir().unwrap();
let project_dir = tempdir.path().to_str().unwrap().to_string();
{
let mut proto = state.protocol.write().await;
proto.sessions.insert(
"remote/s1".into(),
crate::daemon_protocol::SessionEntry {
id: "remote/s1".into(),
pane: None,
origin: Origin::Remote("npub1x".into()),
metadata: crate::daemon_protocol::SessionMeta {
project_dir: Some(project_dir.clone()),
..Default::default()
},
registered_at: 0,
},
);
proto.sessions.insert(
"local/s1".into(),
crate::daemon_protocol::SessionEntry {
id: "local/s1".into(),
pane: Some("%1".into()),
origin: Origin::Local,
metadata: crate::daemon_protocol::SessionMeta {
project_dir: Some(project_dir),
..Default::default()
},
registered_at: 0,
},
);
}
state.sweep_worktree_presence().await;
{
let proto = state.protocol.read().await;
let local = proto.sessions.get("local/s1").unwrap();
assert_eq!(local.metadata.worktree_present, Some(true));
let remote = proto.sessions.get("remote/s1").unwrap();
assert_eq!(remote.metadata.worktree_present, None);
}
}
#[tokio::test]
async fn sweep_worktree_presence_respects_backoff_after_timeout() {
let state = AppState::new_for_test();
state
.sweep_in_progress
.store(true, std::sync::atomic::Ordering::Relaxed);
*state.sweep_backoff_until.lock().unwrap() =
Some(std::time::Instant::now() + std::time::Duration::from_secs(60));
let tempdir = tempfile::tempdir().unwrap();
let project_dir = tempdir.path().to_str().unwrap().to_string();
{
let mut proto = state.protocol.write().await;
proto.sessions.insert(
"local/s1".into(),
crate::daemon_protocol::SessionEntry {
id: "local/s1".into(),
pane: Some("%1".into()),
origin: Origin::Local,
metadata: crate::daemon_protocol::SessionMeta {
project_dir: Some(project_dir),
..Default::default()
},
registered_at: 0,
},
);
}
state.sweep_worktree_presence().await;
{
let proto = state.protocol.read().await;
let session = proto.sessions.get("local/s1").unwrap();
assert_eq!(
session.metadata.worktree_present, None,
"sweep should be skipped during backoff window"
);
}
assert!(
state
.sweep_in_progress
.load(std::sync::atomic::Ordering::Relaxed),
"sweep_in_progress flag must remain set during backoff (orphan thread still holds it)"
);
}
#[tokio::test]
async fn sweep_worktree_presence_clears_expired_backoff_and_runs() {
let state = AppState::new_for_test();
state
.sweep_in_progress
.store(true, std::sync::atomic::Ordering::Relaxed);
*state.sweep_backoff_until.lock().unwrap() =
Some(std::time::Instant::now() - std::time::Duration::from_secs(1));
let tempdir = tempfile::tempdir().unwrap();
let project_dir = tempdir.path().to_str().unwrap().to_string();
{
let mut proto = state.protocol.write().await;
proto.sessions.insert(
"local/s1".into(),
crate::daemon_protocol::SessionEntry {
id: "local/s1".into(),
pane: Some("%1".into()),
origin: Origin::Local,
metadata: crate::daemon_protocol::SessionMeta {
project_dir: Some(project_dir),
..Default::default()
},
registered_at: 0,
},
);
}
state.sweep_worktree_presence().await;
{
let proto = state.protocol.read().await;
let session = proto.sessions.get("local/s1").unwrap();
assert_eq!(
session.metadata.worktree_present,
Some(true),
"sweep should run after backoff window expired"
);
}
assert!(
state.sweep_backoff_until.lock().unwrap().is_none(),
"backoff_until must be cleared once the window expires"
);
}
#[tokio::test]
async fn sweep_worktree_presence_empty_snapshot_does_not_clear_dedup_flag() {
let state = AppState::new_for_test();
state
.sweep_in_progress
.store(true, std::sync::atomic::Ordering::Relaxed);
state.sweep_worktree_presence().await;
assert!(
state
.sweep_in_progress
.load(std::sync::atomic::Ordering::Relaxed),
"empty-snapshot early return must not clear sweep_in_progress flag it never owned"
);
}
#[tokio::test]
async fn sweep_worktree_presence_follows_symlinks() {
let state = AppState::new_for_test();
let real_dir = tempfile::tempdir().unwrap();
let real_path = real_dir.path();
let symlink_path = real_path.join("symlink_to_dir");
std::os::unix::fs::symlink(real_path, &symlink_path).unwrap();
let project_dir = symlink_path.to_str().unwrap().to_string();
{
let mut proto = state.protocol.write().await;
proto.sessions.insert(
"local/s1".into(),
crate::daemon_protocol::SessionEntry {
id: "local/s1".into(),
pane: Some("%1".into()),
origin: Origin::Local,
metadata: crate::daemon_protocol::SessionMeta {
project_dir: Some(project_dir),
..Default::default()
},
registered_at: 0,
},
);
}
state.sweep_worktree_presence().await;
{
let proto = state.protocol.read().await;
let session = proto.sessions.get("local/s1").unwrap();
assert_eq!(
session.metadata.worktree_present,
Some(true),
"symlink to existing dir should show as present"
);
}
}
}