1use crate::backend::*;
9use parking_lot::Mutex;
10use serde_json::json;
11use std::fs;
12use std::path::Path;
13use std::thread;
14use std::time::Duration;
15
16use super::journal::Journal;
17use super::mcp;
18use super::state::{AcpBackend, AgentSlot, SessionEntry, Turn};
19use super::turn::run_turn;
20
21const TURN_CLOSE_BUDGET: Duration = Duration::from_secs(60);
24
25const PROSE_FILE: &str = "AGENTS.md";
35
36const PROSE_BEGIN: &str = "<!-- onlyne:role-prose:begin -->";
41const PROSE_END: &str = "<!-- onlyne:role-prose:end -->";
42
43impl AcpBackend {
44 fn open_session(
49 &self,
50 slot: &AgentSlot,
51 spec: &SpawnSpec,
52 key: &str,
53 ) -> Result<Arc<SessionEntry>> {
54 write_role_prose(&spec.cwd, &spec.prose)?;
58 let start = slot.agent.new_session(&spec.cwd, vec![mcp::mount(spec)?])?;
59 if !self.options.mode.is_empty() {
60 slot.agent
61 .set_mode(&start.session_id, &self.options.mode)
62 .map_err(|error| {
63 anyhow::anyhow!(
64 "acp: {} rejected mode {:?}: {error}",
65 start.session_id,
66 self.options.mode
67 )
68 })?;
69 }
70 for (config, value) in [
71 ("model", self.options.model.as_str()),
72 ("reasoning_effort", self.options.reasoning_effort.as_str()),
73 ] {
74 if value.is_empty() {
75 continue;
76 }
77 slot.agent
78 .set_config_option(&start.session_id, config, value)
79 .map_err(|error| {
80 anyhow::anyhow!(
81 "acp: {} rejected {config} {:?}: {error}",
82 start.session_id,
83 value
84 )
85 })?;
86 }
87 let entry = Arc::new(SessionEntry {
88 task_id: spec.task_id.clone(),
89 id: start.session_id,
90 agent_key: key.to_string(),
91 process: slot.process,
92 workdir: spec.cwd.clone(),
93 agent: Arc::clone(&slot.agent),
94 turn: Turn::new(),
95 refusals: Mutex::new(Vec::new()),
96 });
97 self.state
98 .sessions
99 .lock()
100 .insert(entry.key(), Arc::clone(&entry));
101 tracing::info!(
102 task = %spec.task_id,
103 acp_session = %entry.id,
104 pid = entry.agent.pid(),
105 process = entry.process,
106 agent = %key,
107 "acp session opened"
108 );
109 Ok(entry)
110 }
111
112 fn entry_of(&self, session: &SessionRef) -> Option<Arc<SessionEntry>> {
120 let sessions = self.state.sessions.lock();
121 let agent = session
122 .backend_ref
123 .get("agent")
124 .and_then(Value::as_str)
125 .unwrap_or_default();
126 let process = session.backend_ref.get("process").and_then(Value::as_u64);
127 if let (Some(id), Some(process)) = (
128 session.backend_ref.get("id").and_then(Value::as_str),
129 process,
130 ) && let Some(found) = sessions.get(&(agent.to_string(), process, id.to_string()))
131 {
132 return Some(Arc::clone(found));
133 }
134 sessions
135 .values()
136 .find(|entry| entry.current_task() == session.task_id)
137 .cloned()
138 }
139
140 fn take_entry(&self, session: &SessionRef) -> Option<Arc<SessionEntry>> {
141 let found = self.entry_of(session)?;
142 let key = found.key();
143 self.state.sessions.lock().remove(&key)
144 }
145
146 fn live_entry(&self, session: &SessionRef, task_id: &str) -> Result<Arc<SessionEntry>> {
153 let entry = self.entry_of(session).ok_or_else(|| {
154 anyhow::anyhow!(
155 "acp: no live session for task {task_id}; its agent is not running here"
156 )
157 })?;
158 if entry.current_task() != task_id {
159 return Err(anyhow::anyhow!(
160 "acp: session {} serves task {}; task {task_id} needs a session of its own",
161 entry.id,
162 entry.current_task(),
163 ));
164 }
165 if entry.agent.is_gone() {
166 return Err(anyhow::anyhow!(
167 "acp: agent {} exited before task {task_id} was delivered",
168 entry.agent_key
169 ));
170 }
171 Ok(entry)
172 }
173
174 fn start_turn(
183 &self,
184 entry: &Arc<SessionEntry>,
185 task_id: &str,
186 prompt: String,
187 record: &'static str,
188 ) -> Result<()> {
189 if entry.turn.begin().is_none() {
190 return Err(anyhow::anyhow!(
191 "acp: session {} is still running a turn for task {}",
192 entry.id,
193 entry.current_task()
194 ));
195 }
196 let thread_entry = Arc::clone(entry);
197 let sink = self.state.sink.clone();
198 let content = self.state.content.clone();
199 let policy = self.options.policy();
200 if let Err(error) = thread::Builder::new()
201 .name(format!("acp-turn {}", short(task_id)))
202 .spawn(move || run_turn(thread_entry, sink, content, prompt, record, policy))
203 {
204 entry.turn.finish();
205 return Err(anyhow::anyhow!(
206 "acp: task could not start its turn thread: {error}"
207 ));
208 }
209 Ok(())
210 }
211}
212
213impl SessionBackend for AcpBackend {
214 fn name(&self) -> &'static str {
215 "acp"
216 }
217
218 fn capabilities(&self) -> Capabilities {
219 Capabilities {
220 spawn: true,
221 attach: true,
222 probe: true,
223 close: true,
224 focus: false,
225 rename: false,
229 }
230 }
231
232 fn available(&self) -> Result<bool> {
233 Ok(true)
236 }
237
238 fn self_driven(&self) -> bool {
239 true
240 }
241
242 fn set_content_sink(&self, sink: Arc<dyn crate::content::ContentSink>) {
243 self.state.content.set_sink(sink);
244 }
245
246 fn outcomes(&self) -> Option<OutcomeFeed> {
247 Some(self.state.feed.clone())
248 }
249
250 fn spawn(&self, spec: SpawnSpec) -> Result<SessionRef> {
251 let command = Self::command_of(&spec)?;
252 let key = command.join(" ");
253 let slot = self.agent_for(&key, &command, &spec.cwd, &spec.env)?;
254 let entry = match self.open_session(&slot, &spec, &key) {
255 Ok(entry) => entry,
256 Err(error) => {
257 self.state.retire(&key, slot.process);
261 return Err(error);
262 }
263 };
264 let journal = Journal::new(
265 &spec.cwd,
266 &spec.task_id,
267 &entry.id,
268 self.state.content.clone(),
269 );
270 Ok(SessionRef {
271 task_id: spec.task_id.clone(),
272 backend: self.name().into(),
273 backend_ref: json!({
274 "id": entry.id,
275 "pid": entry.agent.pid(),
276 "process": entry.process,
277 "agent": key,
278 "log": journal.log.to_string_lossy(),
279 "events": journal.events.to_string_lossy(),
280 }),
281 generation: 1,
282 })
283 }
284
285 fn attach(&self, session: &SessionRef) -> Result<SessionRef> {
286 match self.entry_of(session) {
287 Some(entry) if !entry.agent.is_gone() => Ok(session.clone()),
288 Some(entry) => Err(anyhow::anyhow!(
289 "acp session {} is gone (agent {} exited)",
290 entry.id,
291 entry.agent_key
292 )),
293 None => Err(anyhow::anyhow!(
294 "acp session {} is not held by this client",
295 session.task_id
296 )),
297 }
298 }
299
300 fn probe(&self, session: &SessionRef) -> Result<ResourceProbe> {
301 let Some(entry) = self.entry_of(session) else {
302 return Ok(ResourceProbe {
306 alive: false,
307 attached: false,
308 detail: Some(json!({"reason": "no agent handle in this client"})),
309 });
310 };
311 let gone = entry.agent.is_gone();
312 Ok(ResourceProbe {
313 alive: !gone,
314 attached: !gone,
315 detail: Some(json!({
316 "pid": entry.agent.pid(),
317 "acp_session": entry.id,
318 "turn": entry.turn.phase.lock().live,
319 })),
320 })
321 }
322
323 fn deliver(&self, session: &SessionRef, task_id: &str, prompt: &str) -> Result<()> {
330 let entry = self.live_entry(session, task_id)?;
331 let task_id = entry.current_task();
332 self.start_turn(&entry, &task_id, prompt.to_string(), "dispatch")
333 }
334
335 fn nudge(&self, session: &SessionRef, task_id: &str, text: &str) -> Result<()> {
344 let entry = self.live_entry(session, task_id)?;
345 let task_id = entry.current_task();
346 self.start_turn(&entry, &task_id, text.to_string(), "nudge")
347 }
348
349 fn close(&self, session: &SessionRef, reason: CloseReason, _force: bool) -> Result<()> {
350 let Some(entry) = self.take_entry(session) else {
351 tracing::debug!(task = %session.task_id, ?reason, "acp session already gone");
354 return Ok(());
355 };
356 let key = entry.agent_key.clone();
357 let process = entry.process;
358 let id = entry.id.clone();
359 let generation = entry.turn.generation();
360 if entry.turn.phase.lock().live {
361 if let Err(error) = entry.agent.cancel(&id) {
365 tracing::warn!(error = %error, acp_session = %id, "acp: cancel was not sent");
366 }
367 }
368 let state = Arc::clone(&self.state);
375 let (closing, reap_key) = (id.clone(), key.clone());
376 if let Err(error) = thread::Builder::new()
377 .name(format!("acp-close {closing}"))
378 .spawn(move || {
379 if !entry.turn.waited_out(generation, TURN_CLOSE_BUDGET) {
380 tracing::warn!(
381 acp_session = %closing,
382 "acp: the turn outlasted its close budget; the agent decides its end"
383 );
384 }
385 end_session(&entry);
386 state.retire(&reap_key, process);
387 })
388 {
389 tracing::warn!(
390 error = %error,
391 acp_session = %id,
392 "acp: no closer thread; the agent decides its own session's end"
393 );
394 self.state.retire(&key, process);
395 }
396 tracing::info!(
397 task = %session.task_id,
398 acp_session = %id,
399 ?reason,
400 "acp session closed"
401 );
402 Ok(())
403 }
404}
405
406fn end_session(entry: &SessionEntry) {
412 if entry.agent.is_gone() {
413 return;
415 }
416 if !entry
417 .agent
418 .negotiated()
419 .is_some_and(|caps| caps.supports_close())
420 {
421 tracing::debug!(acp_session = %entry.id, "acp: agent offers no session/close");
424 return;
425 }
426 match entry.agent.close_session(&entry.id) {
427 Ok(()) => {}
428 Err(_) if entry.agent.is_gone() => {}
429 Err(error) if tolerated(&error) => {
432 tracing::debug!(error = %error, acp_session = %entry.id, "acp: close refused");
433 }
434 Err(error) => {
435 tracing::warn!(error = %error, acp_session = %entry.id, "acp: close reported");
436 }
437 }
438}
439
440fn tolerated(error: &anyhow::Error) -> bool {
442 error
443 .downcast_ref::<onlyne_acp::RpcError>()
444 .is_some_and(|error| {
445 error.code == onlyne_acp::RpcError::METHOD_NOT_FOUND
446 || error.code == onlyne_acp::RpcError::INVALID_PARAMS
447 })
448}
449
450fn short(id: &str) -> &str {
452 let from = id.len().saturating_sub(8);
453 id.get(from..).unwrap_or(id)
454}
455
456fn write_role_prose(workdir: &Path, prose: &str) -> Result<()> {
469 let path = workdir.join(PROSE_FILE);
470 let existing = match fs::read_to_string(&path) {
471 Ok(text) => Some(text),
472 Err(error) if error.kind() == std::io::ErrorKind::NotFound => None,
473 Err(error) => {
474 return Err(anyhow::anyhow!(
475 "acp: {} could not be read: {error}",
476 path.display()
477 ));
478 }
479 };
480 let block = (!prose.trim().is_empty()).then(|| format!("{PROSE_BEGIN}\n{prose}\n{PROSE_END}"));
481 let next = match existing {
482 None => match block {
483 Some(block) => format!("{block}\n"),
484 None => return Ok(()),
485 },
486 Some(text) => match block_span(&text) {
487 Some((start, stop)) => match block {
488 Some(block) => format!("{}{block}{}", &text[..start], &text[stop..]),
489 None => format!("{}{}", &text[..start], &text[stop..]),
490 },
491 None => match block {
492 Some(block) => append_block(&text, &block),
493 None => return Ok(()),
494 },
495 },
496 };
497 fs::write(&path, next)
498 .map_err(|error| anyhow::anyhow!("acp: {} could not be written: {error}", path.display()))
499}
500
501fn block_span(text: &str) -> Option<(usize, usize)> {
510 let begin = text.find(PROSE_BEGIN)?;
511 let start = text[..begin].rfind('\n').map(|at| at + 1).unwrap_or(0);
512 let stop = match text[begin..].find(PROSE_END) {
513 Some(at) => {
514 let end = begin + at + PROSE_END.len();
515 text[end..]
516 .find('\n')
517 .map(|at| end + at + 1)
518 .unwrap_or(text.len())
519 }
520 None => text.len(),
521 };
522 Some((start, stop))
523}
524
525fn append_block(text: &str, block: &str) -> String {
527 let gap = if text.is_empty() || text.ends_with("\n\n") {
528 ""
529 } else if text.ends_with('\n') {
530 "\n"
531 } else {
532 "\n\n"
533 };
534 format!("{text}{gap}{block}\n")
535}