use axum::extract::State;
use axum::http::StatusCode;
use axum::response::Json;
use serde::Deserialize;
use serde_json::{Value, json};
use crate::state::SharedState;
#[derive(Debug, Deserialize)]
pub struct PaneBody {
#[serde(default)]
pub pane: Option<String>,
#[serde(default)]
pub backend_session_id: Option<String>,
}
pub async fn session_end(
State(state): State<SharedState>,
Json(body): Json<PaneBody>,
) -> (StatusCode, Json<Value>) {
let result = session_end_inner(&state, body).await;
(StatusCode::OK, Json(result))
}
async fn session_end_inner(
state: &std::sync::Arc<crate::state::AppState>,
body: PaneBody,
) -> Value {
let session = {
let proto = state.protocol.read().await;
let found = proto
.sessions
.values()
.find(|s| {
body.pane
.as_deref()
.is_some_and(|p| s.pane.as_deref() == Some(p))
|| body
.backend_session_id
.as_deref()
.is_some_and(|b| s.metadata.backend_session_id.as_deref() == Some(b))
})
.cloned();
match found {
Some(s) => s,
None => return json!({ "skipped": "no session" }),
}
};
let age = chrono::Utc::now().timestamp() - session.registered_at;
if session.registered_at > 0 && age < 5 {
return json!({ "skipped": format!("recently registered ({}s ago)", age) });
}
let id = session.id.clone();
state
.apply_and_execute(crate::daemon_protocol::Event::Remove {
id: id.clone(),
keep_worktree: true,
})
.await;
let pane = session.pane.unwrap_or_default();
tokio::task::spawn_blocking(move || {
let _ = std::process::Command::new("tmux")
.args(["set-option", "-pu", "-t", &pane, "@ouija_id"])
.status();
});
json!({ "removed": id })
}
pub async fn hook_stop(
State(state): State<SharedState>,
Json(body): Json<PaneBody>,
) -> (StatusCode, Json<Value>) {
let result = hook_stop_inner(&state, body).await;
(StatusCode::OK, Json(result))
}
async fn hook_stop_inner(state: &std::sync::Arc<crate::state::AppState>, body: PaneBody) -> Value {
if let Some(id) = state
.find_session_by_pane_or_backend_sid(
body.pane.as_deref(),
body.backend_session_id.as_deref(),
)
.await
{
state
.notify_agent(&id, crate::session_agent::SessionMsg::Stopped)
.await;
}
json!({ "ok": true })
}
pub async fn prompt_submit(
State(state): State<SharedState>,
Json(body): Json<PaneBody>,
) -> (StatusCode, Json<Value>) {
let result = prompt_submit_inner(&state, body).await;
(StatusCode::OK, Json(result))
}
async fn prompt_submit_inner(
state: &std::sync::Arc<crate::state::AppState>,
body: PaneBody,
) -> Value {
if let Some(id) = state
.find_session_by_pane_or_backend_sid(
body.pane.as_deref(),
body.backend_session_id.as_deref(),
)
.await
{
state
.notify_agent(&id, crate::session_agent::SessionMsg::Active)
.await;
}
json!({ "output": "" })
}
#[derive(Debug, Deserialize)]
#[allow(dead_code)] pub struct PreToolUseBody {
#[serde(default)]
pub pane: Option<String>,
#[serde(default)]
pub backend_session_id: Option<String>,
#[serde(default)]
pub tool_name: Option<String>,
}
pub async fn pre_tool_use(
State(state): State<SharedState>,
Json(body): Json<PreToolUseBody>,
) -> (StatusCode, Json<Value>) {
let result = pre_tool_use_inner(&state, body).await;
(StatusCode::OK, Json(result))
}
async fn pre_tool_use_inner(
state: &std::sync::Arc<crate::state::AppState>,
body: PreToolUseBody,
) -> Value {
if let Some(id) = state
.find_session_by_pane_or_backend_sid(
body.pane.as_deref(),
body.backend_session_id.as_deref(),
)
.await
{
state
.notify_agent(&id, crate::session_agent::SessionMsg::Active)
.await;
}
json!({ "block": false })
}
pub async fn post_compact(
State(state): State<SharedState>,
Json(body): Json<PaneBody>,
) -> (StatusCode, Json<Value>) {
let result = post_compact_inner(&state, body).await;
(StatusCode::OK, Json(result))
}
async fn post_compact_inner(
state: &std::sync::Arc<crate::state::AppState>,
body: PaneBody,
) -> Value {
let session_id = match state
.find_session_by_pane_or_backend_sid(
body.pane.as_deref(),
body.backend_session_id.as_deref(),
)
.await
{
Some(id) => id,
None => return json!({ "ok": true, "continuation_injected": false }),
};
let continuation = state.drain_agent_compact_continuation(&session_id).await;
let Some(continuation) = continuation else {
return json!({ "ok": true, "continuation_injected": false });
};
let pane = {
let proto = state.protocol.read().await;
proto.sessions.get(&session_id).and_then(|s| s.pane.clone())
};
let Some(pane) = pane else {
return json!({ "ok": true, "continuation_injected": false, "error": "no pane" });
};
if let Err(e) =
crate::tmux::locked_inject(state, &session_id, &pane, &continuation, false).await
{
tracing::warn!(
session = %session_id,
"post-compact continuation injection failed: {e}"
);
return json!({ "ok": false, "error": e.to_string() });
}
json!({ "ok": true, "continuation_injected": true })
}
#[derive(Debug, Deserialize)]
pub struct SessionStartBody {
pub pane: String,
pub cwd: String,
#[serde(default)]
pub backend_session_id: Option<String>,
}
pub async fn session_start(
State(state): State<SharedState>,
Json(body): Json<SessionStartBody>,
) -> (StatusCode, Json<Value>) {
let result = session_start_inner(&state, body).await;
(StatusCode::OK, Json(result))
}
fn mesh_instructions_for_backend(backend: Option<&str>, public_id: &str) -> String {
if backend != Some("codex-cli") {
return String::new();
}
format!(
"You are on the Ouija mesh. Message other sessions with the `ouija` CLI \
(NOT your own messaging tools — they cannot reach the mesh).\n\
Your public Ouija id is `{public_id}`. Pass it as `--from {public_id}` on \
every command so the mesh knows who is sending.\n\n\
- `ouija ls` — list reachable sessions (targets for messages).\n\
- `ouija ask <target> \"question\" --from {public_id}` — send a question that \
expects a reply; the command returns after delivery.\n\
- `ouija tell <target> \"note\" --from {public_id}` — fire-and-forget message.\n\
- `ouija reply <target> <msg-id> \"answer\" --from {public_id}` — answer a \
`<msg ... reply=\"true\">` you received (the sender is blocked until you reply).\n\n\
For generated or multi-line message text, use `--stdin` instead of putting the \
message in shell quotes.\n\n\
Incoming messages arrive as `<msg from=\"...\" id=\"N\" reply=\"true\">text</msg>`; \
reply to those with `reply=\"true\"` using their `id`. Replies to your asks are pushed \
into this session later as `<msg ... re=\"N\">...</msg>`. If that reply is your only \
remaining blocker, end your turn and wait for the pushed message; do not poll the \
message log, status, or pane output unless you are debugging suspected delivery failure."
)
}
async fn session_start_inner(
state: &std::sync::Arc<crate::state::AppState>,
body: SessionStartBody,
) -> Value {
if !state.settings.read().await.auto_register {
return json!({ "skipped": "auto_register disabled", "output": "" });
}
if let Some(existing_id) = state.find_session_by_pane(&body.pane).await {
let existing_backend = {
let proto = state.protocol.read().await;
proto
.sessions
.get(&existing_id)
.and_then(|s| s.metadata.backend.clone())
};
if let Some(backend_session_id) =
normalize_backend_session_id(body.backend_session_id.as_deref())
{
let updated = {
let proto = state.protocol.read().await;
proto.sessions.get(&existing_id).and_then(|session| {
let needs_binding = session.metadata.backend_session_id.as_deref()
!= Some(backend_session_id.as_str());
needs_binding.then(|| {
let mut metadata = session.metadata.clone();
metadata.backend_session_id = Some(backend_session_id.clone());
metadata
})
})
};
if let Some(metadata) = updated {
state
.apply_and_execute(crate::daemon_protocol::Event::Register {
id: existing_id.clone(),
pane: Some(body.pane.clone()),
metadata,
})
.await;
}
}
let output = mesh_instructions_for_backend(existing_backend.as_deref(), &existing_id);
return json!({
"registered": existing_id,
"output": output,
});
}
let project_root = crate::state::resolve_project_root(&body.cwd);
let basename = std::path::Path::new(project_root)
.file_name()
.and_then(|n| n.to_str())
.unwrap_or("unnamed");
let base_id = crate::state::sanitize_session_id(basename);
if base_id.is_empty() {
return json!({ "error": "could not derive session name", "output": "" });
}
let id = {
let proto = state.protocol.read().await;
let id_to_pane: std::collections::HashMap<String, Option<String>> = proto
.sessions
.iter()
.map(|(id, s)| (id.clone(), s.pane.clone()))
.collect();
crate::state::resolve_unique_session_id(&id_to_pane, &base_id, Some(&body.pane))
};
let detected_backend = state.detect_backend_in_pane(&body.pane).await;
let backend_session_id = match normalize_backend_session_id(body.backend_session_id.as_deref())
{
Some(session_id) => Some(session_id),
None if detected_backend.as_deref() == Some("opencode") => {
resolve_opencode_session_id(state, project_root).await
}
None => None,
};
let output = mesh_instructions_for_backend(detected_backend.as_deref(), &id);
let role = format!("working on {basename}");
let proto_meta = crate::daemon_protocol::SessionMeta {
project_dir: Some(project_root.to_string()),
role: Some(role),
backend: detected_backend,
backend_session_id,
..Default::default()
};
state
.apply_and_execute(crate::daemon_protocol::Event::Register {
id: id.clone(),
pane: Some(body.pane.clone()),
metadata: proto_meta,
})
.await;
json!({
"registered": id,
"output": output,
})
}
fn normalize_backend_session_id(value: Option<&str>) -> Option<String> {
value
.map(str::trim)
.filter(|v| !v.is_empty())
.map(String::from)
}
async fn resolve_opencode_session_id(
state: &std::sync::Arc<crate::state::AppState>,
project_dir: &str,
) -> Option<String> {
let port = state.opencode_serve_port();
let url = format!("http://127.0.0.1:{port}/session");
let resp = state
.http_client
.get(&url)
.timeout(std::time::Duration::from_secs(3))
.send()
.await
.ok()?;
if !resp.status().is_success() {
return None;
}
let sessions: Vec<serde_json::Value> = resp.json().await.ok()?;
sessions
.iter()
.filter(|s| s["directory"].as_str() == Some(project_dir))
.max_by_key(|s| {
s["time"]["updated"]
.as_i64()
.or_else(|| s["time"]["created"].as_i64())
.unwrap_or(0)
})
.and_then(|s| s["id"].as_str().map(String::from))
}
#[cfg(test)]
mod tests {
use super::*;
#[tokio::test]
async fn session_end_removes_old_session() {
let state = crate::state::AppState::new_for_test();
state
.apply_and_execute(crate::daemon_protocol::Event::Register {
id: "test-session".into(),
pane: Some("%99".into()),
metadata: crate::daemon_protocol::SessionMeta::default(),
})
.await;
assert!(state.find_session_by_pane("%99").await.is_some());
{
let mut proto = state.protocol.write().await;
if let Some(s) = proto.sessions.get_mut("test-session") {
s.registered_at = chrono::Utc::now().timestamp() - 10;
}
}
let body = PaneBody {
pane: Some("%99".into()),
backend_session_id: None,
};
let result = session_end_inner(&state, body).await;
assert!(result.get("removed").is_some());
assert!(state.find_session_by_pane("%99").await.is_none());
}
#[tokio::test]
async fn session_end_rejects_recently_registered() {
let state = crate::state::AppState::new_for_test();
state
.apply_and_execute(crate::daemon_protocol::Event::Register {
id: "fresh".into(),
pane: Some("%99".into()),
metadata: crate::daemon_protocol::SessionMeta::default(),
})
.await;
let body = PaneBody {
pane: Some("%99".into()),
backend_session_id: None,
};
let result = session_end_inner(&state, body).await;
assert!(result.get("skipped").is_some());
assert!(state.find_session_by_pane("%99").await.is_some());
}
#[tokio::test]
async fn session_end_no_session() {
let state = crate::state::AppState::new_for_test();
let body = PaneBody {
pane: Some("%999".into()),
backend_session_id: None,
};
let result = session_end_inner(&state, body).await;
assert!(result.get("skipped").is_some());
}
#[tokio::test]
async fn hook_stop_no_session_returns_ok() {
let state = crate::state::AppState::new_for_test();
let body = PaneBody {
pane: Some("%999".into()),
backend_session_id: None,
};
let result = hook_stop_inner(&state, body).await;
assert_eq!(result, json!({ "ok": true }));
}
#[tokio::test]
async fn prompt_submit_returns_empty_for_unknown_pane() {
let state = crate::state::AppState::new_for_test();
let body = PaneBody {
pane: Some("%999".into()),
backend_session_id: None,
};
let result = prompt_submit_inner(&state, body).await;
assert_eq!(result["output"], "");
}
#[tokio::test]
async fn pre_tool_use_no_session_allows() {
let state = crate::state::AppState::new_for_test();
let body = PreToolUseBody {
pane: Some("%999".into()),
backend_session_id: None,
tool_name: Some("AskUserQuestion".into()),
};
let result = pre_tool_use_inner(&state, body).await;
assert_eq!(result["block"], false);
}
#[tokio::test]
async fn pre_tool_use_signals_activity_for_registered_session() {
let state = crate::state::AppState::new_for_test();
state
.apply_and_execute(crate::daemon_protocol::Event::Register {
id: "tool-activity".into(),
pane: Some("%42".into()),
metadata: crate::daemon_protocol::SessionMeta {
reminder: Some("keep working".into()),
..Default::default()
},
})
.await;
state.settings.write().await.idle_timeout_secs = 1;
state
.notify_agent("tool-activity", crate::session_agent::SessionMsg::Stopped)
.await;
tokio::time::sleep(std::time::Duration::from_millis(500)).await;
let body = PreToolUseBody {
pane: Some("%42".into()),
backend_session_id: None,
tool_name: Some("Bash".into()),
};
let result = pre_tool_use_inner(&state, body).await;
assert_eq!(result["block"], false);
tokio::time::sleep(std::time::Duration::from_millis(100)).await;
assert!(state.find_session_by_pane("%42").await.is_some());
}
#[tokio::test]
async fn pre_tool_use_accepts_backend_session_id() {
let state = crate::state::AppState::new_for_test();
state
.apply_and_execute(crate::daemon_protocol::Event::Register {
id: "oc-session".into(),
pane: None,
metadata: crate::daemon_protocol::SessionMeta {
backend_session_id: Some("oc-uuid-123".into()),
..Default::default()
},
})
.await;
let body = PreToolUseBody {
pane: None,
backend_session_id: Some("oc-uuid-123".into()),
tool_name: Some("bash".into()),
};
let result = pre_tool_use_inner(&state, body).await;
assert_eq!(result["block"], false);
}
#[test]
fn mesh_instructions_only_for_codex() {
let codex = mesh_instructions_for_backend(Some("codex-cli"), "feat/123-worker");
assert!(codex.contains("ouija ls"), "{codex}");
assert!(codex.contains("ouija ask"), "{codex}");
assert!(codex.contains("ouija tell"), "{codex}");
assert!(codex.contains("ouija reply"), "{codex}");
assert!(codex.contains("returns after delivery"), "{codex}");
assert!(codex.contains("do not poll"), "{codex}");
assert!(
codex.contains("--from feat/123-worker"),
"must teach the resolved public id as --from: {codex}"
);
assert_eq!(mesh_instructions_for_backend(Some("claude-code"), "x"), "");
assert_eq!(mesh_instructions_for_backend(Some("opencode"), "x"), "");
assert_eq!(mesh_instructions_for_backend(None, "x"), "");
}
#[tokio::test]
async fn session_start_onboards_already_registered_codex_session() {
let state = crate::state::AppState::new_for_test();
state
.apply_and_execute(crate::daemon_protocol::Event::Register {
id: "feat/worker".into(),
pane: Some("%70".into()),
metadata: crate::daemon_protocol::SessionMeta {
backend: Some("codex-cli".into()),
..Default::default()
},
})
.await;
let body = SessionStartBody {
pane: "%70".into(),
cwd: "/home/user/code/proj".into(),
backend_session_id: Some("codex-thread-1".into()),
};
let result = session_start_inner(&state, body).await;
assert_eq!(result["registered"], "feat/worker");
let output = result["output"].as_str().unwrap();
assert!(
output.contains("ouija ls"),
"codex must be onboarded: {output}"
);
assert!(
output.contains("--from feat/worker"),
"must use the authoritative registered id: {output}"
);
let proto = state.protocol.read().await;
let session = proto.sessions.get("feat/worker").unwrap();
assert_eq!(
session.metadata.backend_session_id.as_deref(),
Some("codex-thread-1")
);
}
#[tokio::test]
async fn session_start_binds_identity_for_any_registered_backend() {
let state = crate::state::AppState::new_for_test();
state
.apply_and_execute(crate::daemon_protocol::Event::Register {
id: "claude-worker".into(),
pane: Some("%71".into()),
metadata: crate::daemon_protocol::SessionMeta {
backend: Some("claude-code".into()),
..Default::default()
},
})
.await;
let body = SessionStartBody {
pane: "%71".into(),
cwd: "/home/user/code/proj".into(),
backend_session_id: Some("claude-session-1".into()),
};
let result = session_start_inner(&state, body).await;
assert_eq!(result["registered"], "claude-worker");
assert_eq!(result["output"], "");
let proto = state.protocol.read().await;
let session = proto.sessions.get("claude-worker").unwrap();
assert_eq!(
session.metadata.backend_session_id.as_deref(),
Some("claude-session-1")
);
}
#[tokio::test]
async fn session_start_registers_new_session() {
let state = crate::state::AppState::new_for_test();
let body = SessionStartBody {
pane: "%50".into(),
cwd: "/home/user/code/myproject".into(),
backend_session_id: None,
};
let result = session_start_inner(&state, body).await;
assert_eq!(result["registered"], "myproject");
assert_eq!(result["output"], "");
}
#[tokio::test]
async fn session_start_skips_already_registered() {
let state = crate::state::AppState::new_for_test();
state
.apply_and_execute(crate::daemon_protocol::Event::Register {
id: "existing".into(),
pane: Some("%50".into()),
metadata: crate::daemon_protocol::SessionMeta::default(),
})
.await;
let body = SessionStartBody {
pane: "%50".into(),
cwd: "/home/user/code/existing".into(),
backend_session_id: None,
};
let result = session_start_inner(&state, body).await;
assert_eq!(result["registered"], "existing");
let proto = state.protocol.read().await;
let count = proto.sessions.len();
assert_eq!(count, 1, "should still have exactly 1 session, got {count}");
}
#[tokio::test]
async fn session_start_resolves_worktree_path() {
let state = crate::state::AppState::new_for_test();
let body = SessionStartBody {
pane: "%50".into(),
cwd: "/home/user/code/ouija/.ouija/worktrees/feature-x".into(),
backend_session_id: None,
};
let result = session_start_inner(&state, body).await;
assert_eq!(result["registered"], "ouija");
}
#[tokio::test]
async fn post_compact_no_session_returns_ok() {
let state = crate::state::AppState::new_for_test();
let body = PaneBody {
pane: Some("%999".into()),
backend_session_id: None,
};
let result = post_compact_inner(&state, body).await;
assert_eq!(result["ok"], true);
assert_eq!(result["continuation_injected"], false);
}
#[tokio::test]
async fn post_compact_no_pending_continuation() {
let state = crate::state::AppState::new_for_test();
state
.apply_and_execute(crate::daemon_protocol::Event::Register {
id: "compact-test".into(),
pane: Some("%88".into()),
metadata: crate::daemon_protocol::SessionMeta::default(),
})
.await;
let body = PaneBody {
pane: Some("%88".into()),
backend_session_id: None,
};
let result = post_compact_inner(&state, body).await;
assert_eq!(result["ok"], true);
assert_eq!(result["continuation_injected"], false);
}
#[tokio::test]
async fn post_compact_drains_and_clears_continuation() {
let state = crate::state::AppState::new_for_test();
state
.apply_and_execute(crate::daemon_protocol::Event::Register {
id: "drain-test".into(),
pane: Some("%77".into()),
metadata: crate::daemon_protocol::SessionMeta::default(),
})
.await;
let acquired = state
.try_set_pending_compact_continuation("drain-test", "Continue working.".into())
.await;
assert!(acquired, "slot should be empty for a fresh session");
let continuation = state.drain_agent_compact_continuation("drain-test").await;
assert_eq!(continuation.as_deref(), Some("Continue working."));
let continuation = state.drain_agent_compact_continuation("drain-test").await;
assert_eq!(continuation, None);
}
}