Skip to main content

onlyne_client/backend/acp/
session.rs

1//! The `SessionBackend` surface: bring a session up on a shared agent process,
2//! report what it is, hand it one delivery, and let it go without parking the
3//! caller.
4//!
5//! Every method here answers its caller and nothing else; the facts a session
6//! produces arrive on [`crate::OutcomeSink`] from the turn and closer threads.
7
8use 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
21/// How long the detached closer waits for a running turn before it lets the agent
22/// decide the turn's end on its own.
23const TURN_CLOSE_BUDGET: Duration = Duration::from_secs(60);
24
25/// The instruction file a role's prose is written into.
26///
27/// The `agents.md` convention's own name, and the file this repository's own
28/// world uses. A file is the only vehicle: `onlyne-acp` has no instruction
29/// field on `session/new`, so the prose has to ride in a file the runtime
30/// reads. Claude Code supports `AGENTS.md`, so the prose lands.
31///
32/// The residual gap is narrow: an ACP agent that reads neither `AGENTS.md` nor
33/// `CLAUDE.md` gets no role prose at all.
34const PROSE_FILE: &str = "AGENTS.md";
35
36/// The marker pair one client-owned prose block sits between.
37///
38/// The block is this client's: the file belongs to the operator, and the two
39/// markers are the only bytes of it this code may find and replace.
40const PROSE_BEGIN: &str = "<!-- onlyne:role-prose:begin -->";
41const PROSE_END: &str = "<!-- onlyne:role-prose:end -->";
42
43impl AcpBackend {
44    /// Open the ACP session and apply the configured mode and model, so a session
45    /// is usable the moment this client reports it ready. A refusal of either
46    /// setting fails the spawn: an operator asked for that mode, and a session
47    /// running in a different one is not a partial success.
48    fn open_session(
49        &self,
50        slot: &AgentSlot,
51        spec: &SpawnSpec,
52        key: &str,
53    ) -> Result<Arc<SessionEntry>> {
54        // The role's prose goes into the workspace before the session opens: an
55        // agent that starts working the moment `session/new` returns still reads
56        // the instructions this client put there for it.
57        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    /// The live session a stored reference names: by the id the agent chose for it
113    /// under its own command and process first, by task id for a reference that
114    /// carries neither. ACP session ids are only unique within one agent process,
115    /// so two processes may hand out the same `sess-1` — a replacement of a
116    /// crashed agent included — and the key carries the command and this client's
117    /// name for the process for exactly that reason. A reference from a client run
118    /// that no longer holds the process names nothing here.
119    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    /// The live session one delivery or nudge is addressed to.
147    ///
148    /// One session serves one task, and the id it serves was bound when it
149    /// opened. A turn naming another task would journal itself under a task the
150    /// agent was never assigned, so the two ids can never be reconciled
151    /// afterwards; refuse it here, where the caller can still hear about it.
152    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    /// Start one turn on its own thread, with `prompt` verbatim as the text the
175    /// agent reads.
176    ///
177    /// The turn is claimed before the thread starts, so a second turn offered to
178    /// a session whose agent is still thinking is refused rather than queued
179    /// behind a turn it was never meant to join. `record` is the journal record
180    /// this turn leaves — `dispatch` for a delivery, `nudge` for the sentence
181    /// §3c sends — so an operator reading the file can tell which one it was.
182    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            // ACP v1 has no title path: naming a session is the terminal
226            // backends' trick, and the agent's own session id is what this one
227            // reports.
228            rename: false,
229        }
230    }
231
232    fn available(&self) -> Result<bool> {
233        // Nothing to discover: the command is the role's own config, and a program
234        // that is not installed fails its spawn naming the program.
235        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                // The reservation this spawn took on the chosen process goes back
258                // before the error does, so an agent that refuses its config is
259                // not left running for a session that never existed.
260                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            // A reference this client no longer holds. An ACP agent cannot be
303            // adopted across a client restart, so the process that answered for it
304            // is gone with the pipe it spoke on.
305            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    /// Hand one delivery to its session as a turn.
324    ///
325    /// The prompt is the rendered delivery text and nothing else: this client
326    /// appends no directive of its own, and the agent's own tools mount carries
327    /// whatever it owes back. The whole text lands in the journal's `dispatch`
328    /// record, which is the operator's proof of what was asked.
329    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    /// Tell a session whose turn ended without a completion, in the one words
336    /// this client owns (`docs/v2-CONTRACT.md` §3c).
337    ///
338    /// An ACP agent has no injection channel besides `session/prompt`, so the
339    /// sentence arrives as the next prompt, verbatim. It is not a delivery: no
340    /// payload, prose, or attachment is re-sent, nothing of the session's turn
341    /// bookkeeping is reset, and the journal names the record `nudge` so the two
342    /// are told apart in the operator's file.
343    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            // Already closed, or closed by the client run that held the process
352            // before this one. Either way there is nothing left to end.
353            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            // A notification, so it cannot block: the agent ends the turn and
362            // answers the parked prompt with `cancelled`, which is what lets the
363            // session close on its own terms instead of being cut off mid-frame.
364            if let Err(error) = entry.agent.cancel(&id) {
365                tracing::warn!(error = %error, acp_session = %id, "acp: cancel was not sent");
366            }
367        }
368        // Everything past this point waits on the agent, and every caller of
369        // `close` in the client holds its one dispatch lock — a settled task
370        // reaches here through `release_locked`, and so does the shutdown sweep,
371        // whose whole budget exists because a slow backend must not park a
372        // SIGTERM'd client. So the protocol close, the wait for the running turn,
373        // and the process teardown all go to a thread that holds nothing.
374        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
406/// End the ACP session on the agent, leaving the process for its other sessions.
407///
408/// This runs on the closer thread, so nothing it reports can reach a caller: an
409/// agent that will not end its own session is logged, not failed back, because the
410/// client's bookkeeping has already let the session go.
411fn end_session(entry: &SessionEntry) {
412    if entry.agent.is_gone() {
413        // The session went with the process.
414        return;
415    }
416    if !entry
417        .agent
418        .negotiated()
419        .is_some_and(|caps| caps.supports_close())
420    {
421        // The agent never advertised `session/close`. Our handle going away is all
422        // this client can do; the agent decides its own session's fate.
423        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        // A refusal that specific is an agent that already forgot the session,
430        // which is the state this call was asked to reach.
431        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
440/// Whether an error means the agent has nothing left for this call to close.
441fn 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
450/// A task id is long, and a thread name is read by an operator.
451fn short(id: &str) -> &str {
452    let from = id.len().saturating_sub(8);
453    id.get(from..).unwrap_or(id)
454}
455
456/// Write the role's prose into the workspace's instruction file.
457///
458/// The block between the markers is this client's and the rest of the file is
459/// the operator's: a block that is there is replaced where it stands, a file
460/// without one gets it appended as its own paragraph, and not a byte outside
461/// the markers is touched either way. Prose the role no longer has takes its
462/// block out, because a block that stays is prose this client would be standing
463/// behind after it stopped.
464///
465/// A read that fails for any reason but absence is an error: this runs while a
466/// session opens, and a spawn that cannot write the instructions it was asked
467/// to write must not report a session that will never read them.
468fn 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
501/// The byte span one client-owned block occupies, its markers and their own
502/// line endings included.
503///
504/// The span starts at the beginning of the line the begin marker sits on, so a
505/// replacement cannot leave the marker's own indentation behind. A begin marker
506/// with no end marker after it owns the rest of the file: everything from this
507/// client's own marker on is text this code wrote, and a half-written block is
508/// replaced rather than doubled.
509fn 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
525/// The file with the block added as its own paragraph.
526fn 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}