use crate::backend::*;
use parking_lot::Mutex;
use serde_json::json;
use std::fs;
use std::path::Path;
use std::thread;
use std::time::Duration;
use super::journal::Journal;
use super::mcp;
use super::state::{AcpBackend, AgentSlot, SessionEntry, Turn};
use super::turn::run_turn;
const TURN_CLOSE_BUDGET: Duration = Duration::from_secs(60);
const PROSE_FILE: &str = "AGENTS.md";
const PROSE_BEGIN: &str = "<!-- onlyne:role-prose:begin -->";
const PROSE_END: &str = "<!-- onlyne:role-prose:end -->";
impl AcpBackend {
fn open_session(
&self,
slot: &AgentSlot,
spec: &SpawnSpec,
key: &str,
) -> Result<Arc<SessionEntry>> {
write_role_prose(&spec.cwd, &spec.prose)?;
let start = slot.agent.new_session(&spec.cwd, vec![mcp::mount(spec)?])?;
if !self.options.mode.is_empty() {
slot.agent
.set_mode(&start.session_id, &self.options.mode)
.map_err(|error| {
anyhow::anyhow!(
"acp: {} rejected mode {:?}: {error}",
start.session_id,
self.options.mode
)
})?;
}
for (config, value) in [
("model", self.options.model.as_str()),
("reasoning_effort", self.options.reasoning_effort.as_str()),
] {
if value.is_empty() {
continue;
}
slot.agent
.set_config_option(&start.session_id, config, value)
.map_err(|error| {
anyhow::anyhow!(
"acp: {} rejected {config} {:?}: {error}",
start.session_id,
value
)
})?;
}
let entry = Arc::new(SessionEntry {
task_id: spec.task_id.clone(),
id: start.session_id,
agent_key: key.to_string(),
process: slot.process,
workdir: spec.cwd.clone(),
agent: Arc::clone(&slot.agent),
turn: Turn::new(),
refusals: Mutex::new(Vec::new()),
});
self.state
.sessions
.lock()
.insert(entry.key(), Arc::clone(&entry));
tracing::info!(
task = %spec.task_id,
acp_session = %entry.id,
pid = entry.agent.pid(),
process = entry.process,
agent = %key,
"acp session opened"
);
Ok(entry)
}
fn entry_of(&self, session: &SessionRef) -> Option<Arc<SessionEntry>> {
let sessions = self.state.sessions.lock();
let agent = session
.backend_ref
.get("agent")
.and_then(Value::as_str)
.unwrap_or_default();
let process = session.backend_ref.get("process").and_then(Value::as_u64);
if let (Some(id), Some(process)) = (
session.backend_ref.get("id").and_then(Value::as_str),
process,
) && let Some(found) = sessions.get(&(agent.to_string(), process, id.to_string()))
{
return Some(Arc::clone(found));
}
sessions
.values()
.find(|entry| entry.current_task() == session.task_id)
.cloned()
}
fn take_entry(&self, session: &SessionRef) -> Option<Arc<SessionEntry>> {
let found = self.entry_of(session)?;
let key = found.key();
self.state.sessions.lock().remove(&key)
}
fn live_entry(&self, session: &SessionRef, task_id: &str) -> Result<Arc<SessionEntry>> {
let entry = self.entry_of(session).ok_or_else(|| {
anyhow::anyhow!(
"acp: no live session for task {task_id}; its agent is not running here"
)
})?;
if entry.current_task() != task_id {
return Err(anyhow::anyhow!(
"acp: session {} serves task {}; task {task_id} needs a session of its own",
entry.id,
entry.current_task(),
));
}
if entry.agent.is_gone() {
return Err(anyhow::anyhow!(
"acp: agent {} exited before task {task_id} was delivered",
entry.agent_key
));
}
Ok(entry)
}
fn start_turn(
&self,
entry: &Arc<SessionEntry>,
task_id: &str,
prompt: String,
record: &'static str,
) -> Result<()> {
if entry.turn.begin().is_none() {
return Err(anyhow::anyhow!(
"acp: session {} is still running a turn for task {}",
entry.id,
entry.current_task()
));
}
let thread_entry = Arc::clone(entry);
let sink = self.state.sink.clone();
let content = self.state.content.clone();
let policy = self.options.policy();
if let Err(error) = thread::Builder::new()
.name(format!("acp-turn {}", short(task_id)))
.spawn(move || run_turn(thread_entry, sink, content, prompt, record, policy))
{
entry.turn.finish();
return Err(anyhow::anyhow!(
"acp: task could not start its turn thread: {error}"
));
}
Ok(())
}
}
impl SessionBackend for AcpBackend {
fn name(&self) -> &'static str {
"acp"
}
fn capabilities(&self) -> Capabilities {
Capabilities {
spawn: true,
attach: true,
probe: true,
close: true,
focus: false,
rename: false,
}
}
fn available(&self) -> Result<bool> {
Ok(true)
}
fn self_driven(&self) -> bool {
true
}
fn set_content_sink(&self, sink: Arc<dyn crate::content::ContentSink>) {
self.state.content.set_sink(sink);
}
fn outcomes(&self) -> Option<OutcomeFeed> {
Some(self.state.feed.clone())
}
fn spawn(&self, spec: SpawnSpec) -> Result<SessionRef> {
let command = Self::command_of(&spec)?;
let key = command.join(" ");
let slot = self.agent_for(&key, &command, &spec.cwd, &spec.env)?;
let entry = match self.open_session(&slot, &spec, &key) {
Ok(entry) => entry,
Err(error) => {
self.state.retire(&key, slot.process);
return Err(error);
}
};
let journal = Journal::new(
&spec.cwd,
&spec.task_id,
&entry.id,
self.state.content.clone(),
);
Ok(SessionRef {
task_id: spec.task_id.clone(),
backend: self.name().into(),
backend_ref: json!({
"id": entry.id,
"pid": entry.agent.pid(),
"process": entry.process,
"agent": key,
"log": journal.log.to_string_lossy(),
"events": journal.events.to_string_lossy(),
}),
generation: 1,
})
}
fn attach(&self, session: &SessionRef) -> Result<SessionRef> {
match self.entry_of(session) {
Some(entry) if !entry.agent.is_gone() => Ok(session.clone()),
Some(entry) => Err(anyhow::anyhow!(
"acp session {} is gone (agent {} exited)",
entry.id,
entry.agent_key
)),
None => Err(anyhow::anyhow!(
"acp session {} is not held by this client",
session.task_id
)),
}
}
fn probe(&self, session: &SessionRef) -> Result<ResourceProbe> {
let Some(entry) = self.entry_of(session) else {
return Ok(ResourceProbe {
alive: false,
attached: false,
detail: Some(json!({"reason": "no agent handle in this client"})),
});
};
let gone = entry.agent.is_gone();
Ok(ResourceProbe {
alive: !gone,
attached: !gone,
detail: Some(json!({
"pid": entry.agent.pid(),
"acp_session": entry.id,
"turn": entry.turn.phase.lock().live,
})),
})
}
fn deliver(&self, session: &SessionRef, task_id: &str, prompt: &str) -> Result<()> {
let entry = self.live_entry(session, task_id)?;
let task_id = entry.current_task();
self.start_turn(&entry, &task_id, prompt.to_string(), "dispatch")
}
fn nudge(&self, session: &SessionRef, task_id: &str, text: &str) -> Result<()> {
let entry = self.live_entry(session, task_id)?;
let task_id = entry.current_task();
self.start_turn(&entry, &task_id, text.to_string(), "nudge")
}
fn close(&self, session: &SessionRef, reason: CloseReason, _force: bool) -> Result<()> {
let Some(entry) = self.take_entry(session) else {
tracing::debug!(task = %session.task_id, ?reason, "acp session already gone");
return Ok(());
};
let key = entry.agent_key.clone();
let process = entry.process;
let id = entry.id.clone();
let generation = entry.turn.generation();
if entry.turn.phase.lock().live {
if let Err(error) = entry.agent.cancel(&id) {
tracing::warn!(error = %error, acp_session = %id, "acp: cancel was not sent");
}
}
let state = Arc::clone(&self.state);
let (closing, reap_key) = (id.clone(), key.clone());
if let Err(error) = thread::Builder::new()
.name(format!("acp-close {closing}"))
.spawn(move || {
if !entry.turn.waited_out(generation, TURN_CLOSE_BUDGET) {
tracing::warn!(
acp_session = %closing,
"acp: the turn outlasted its close budget; the agent decides its end"
);
}
end_session(&entry);
state.retire(&reap_key, process);
})
{
tracing::warn!(
error = %error,
acp_session = %id,
"acp: no closer thread; the agent decides its own session's end"
);
self.state.retire(&key, process);
}
tracing::info!(
task = %session.task_id,
acp_session = %id,
?reason,
"acp session closed"
);
Ok(())
}
}
fn end_session(entry: &SessionEntry) {
if entry.agent.is_gone() {
return;
}
if !entry
.agent
.negotiated()
.is_some_and(|caps| caps.supports_close())
{
tracing::debug!(acp_session = %entry.id, "acp: agent offers no session/close");
return;
}
match entry.agent.close_session(&entry.id) {
Ok(()) => {}
Err(_) if entry.agent.is_gone() => {}
Err(error) if tolerated(&error) => {
tracing::debug!(error = %error, acp_session = %entry.id, "acp: close refused");
}
Err(error) => {
tracing::warn!(error = %error, acp_session = %entry.id, "acp: close reported");
}
}
}
fn tolerated(error: &anyhow::Error) -> bool {
error
.downcast_ref::<onlyne_acp::RpcError>()
.is_some_and(|error| {
error.code == onlyne_acp::RpcError::METHOD_NOT_FOUND
|| error.code == onlyne_acp::RpcError::INVALID_PARAMS
})
}
fn short(id: &str) -> &str {
let from = id.len().saturating_sub(8);
id.get(from..).unwrap_or(id)
}
fn write_role_prose(workdir: &Path, prose: &str) -> Result<()> {
let path = workdir.join(PROSE_FILE);
let existing = match fs::read_to_string(&path) {
Ok(text) => Some(text),
Err(error) if error.kind() == std::io::ErrorKind::NotFound => None,
Err(error) => {
return Err(anyhow::anyhow!(
"acp: {} could not be read: {error}",
path.display()
));
}
};
let block = (!prose.trim().is_empty()).then(|| format!("{PROSE_BEGIN}\n{prose}\n{PROSE_END}"));
let next = match existing {
None => match block {
Some(block) => format!("{block}\n"),
None => return Ok(()),
},
Some(text) => match block_span(&text) {
Some((start, stop)) => match block {
Some(block) => format!("{}{block}{}", &text[..start], &text[stop..]),
None => format!("{}{}", &text[..start], &text[stop..]),
},
None => match block {
Some(block) => append_block(&text, &block),
None => return Ok(()),
},
},
};
fs::write(&path, next)
.map_err(|error| anyhow::anyhow!("acp: {} could not be written: {error}", path.display()))
}
fn block_span(text: &str) -> Option<(usize, usize)> {
let begin = text.find(PROSE_BEGIN)?;
let start = text[..begin].rfind('\n').map(|at| at + 1).unwrap_or(0);
let stop = match text[begin..].find(PROSE_END) {
Some(at) => {
let end = begin + at + PROSE_END.len();
text[end..]
.find('\n')
.map(|at| end + at + 1)
.unwrap_or(text.len())
}
None => text.len(),
};
Some((start, stop))
}
fn append_block(text: &str, block: &str) -> String {
let gap = if text.is_empty() || text.ends_with("\n\n") {
""
} else if text.ends_with('\n') {
"\n"
} else {
"\n\n"
};
format!("{text}{gap}{block}\n")
}