Skip to main content

onlyne_client/backend/acp/
process.rs

1//! The agent process: one per distinct rendered command, shared by every
2//! session of that command, plus the permission responder that lives as long as
3//! it does.
4//!
5//! Process discipline lives here — the handshake that decides whether a command
6//! is a runtime at all, the reservation count that decides when a process is
7//! surplus, and the thread that answers `session/request_permission` on a policy
8//! rather than a fallback grant.
9
10use crate::backend::*;
11use crate::content::ContentWriter;
12use onlyne_acp::{
13    Agent, AgentOptions, ClientCapabilities, ClientInfo, Event, PermissionOption,
14    PermissionOutcome, PermissionRequest,
15};
16use parking_lot::Mutex;
17use std::path::Path;
18use std::sync::Weak;
19use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering};
20use std::thread;
21
22use super::journal::or_dash;
23use super::state::{AcpBackend, AcpOptions, AgentSlot, State};
24
25/// The client name reported in `initialize`.
26const CLIENT_NAME: &str = "onlyne-client";
27
28impl AcpBackend {
29    pub fn new(options: AcpOptions) -> Self {
30        let (sink, feed) = OutcomeFeed::channel();
31        AcpBackend {
32            options,
33            state: Arc::new(State {
34                agents: Mutex::new(BTreeMap::new()),
35                sessions: Mutex::new(BTreeMap::new()),
36                process: AtomicU64::new(0),
37                sink,
38                feed,
39                content: ContentWriter::default(),
40            }),
41        }
42    }
43
44    /// The argv this session's agent runs as: the role's rendered
45    /// `[client.runtime] command`, never a string handed to a shell.
46    pub(super) fn command_of(spec: &SpawnSpec) -> Result<Vec<String>> {
47        match spec.command.first() {
48            Some(program) if !program.trim().is_empty() => Ok(spec.command.clone()),
49            _ => Err(anyhow::anyhow!(
50                "acp: the role's `[client.runtime] command` is empty; there is nothing to run \
51                 for task {}",
52                spec.task_id
53            )),
54        }
55    }
56
57    /// The process serving this command, with one session already reserved on it
58    /// so a shared agent cannot be torn down between two sessions handing over.
59    ///
60    /// The handshake runs under the agent map: it is the one slow step that must
61    /// not be duplicated for a second session of the same command, and the map is
62    /// what makes "first one here" decidable.
63    pub(super) fn agent_for(
64        &self,
65        key: &str,
66        command: &[String],
67        cwd: &Path,
68        env: &BTreeMap<String, String>,
69    ) -> Result<Arc<AgentSlot>> {
70        let mut agents = self.state.agents.lock();
71        if let Some(slot) = agents.get(key) {
72            if !slot.agent.is_gone() {
73                slot.live.fetch_add(1, Ordering::SeqCst);
74                return Ok(Arc::clone(slot));
75            }
76            // A crashed process is not a runtime to share: drop it and start the
77            // replacement the next session deserves.
78            tracing::warn!(agent = %key, "acp: replacing an exited agent process");
79            agents.remove(key);
80        }
81        let agent = Arc::new(Agent::start(AgentOptions {
82            command: command.to_vec(),
83            cwd: Some(cwd.to_path_buf()),
84            env: env.clone(),
85        })?);
86        // The handshake is part of starting the process: an agent that refuses it
87        // is not a runtime. The handle goes to a thread rather than being dropped
88        // here because `AgentInner`'s teardown waits for the process it is ending,
89        // and this call runs under the client's dispatch lock.
90        if let Err(error) =
91            agent.initialize(ClientInfo::new(CLIENT_NAME), ClientCapabilities::default())
92        {
93            let _ = thread::Builder::new()
94                .name(format!("acp-drop {key}"))
95                .spawn(move || drop(agent));
96            return Err(anyhow::anyhow!("acp: {key} refused the handshake: {error}"));
97        }
98        let slot = Arc::new(AgentSlot {
99            live: AtomicUsize::new(1),
100            process: self.state.process.fetch_add(1, Ordering::Relaxed) + 1,
101            agent: Arc::clone(&agent),
102        });
103        spawn_responder(&slot, Arc::clone(&self.state), self.options.clone(), key);
104        agents.insert(key.to_string(), Arc::clone(&slot));
105        Ok(slot)
106    }
107}
108
109impl State {
110    /// One session of the process named by `process` went away, so one
111    /// reservation of that process is released. Past the last one the process
112    /// goes with it, on a thread that can afford to wait for it: see
113    /// [`State::reap`].
114    ///
115    /// The name is what makes the pairing exact. A session of a process that left
116    /// on its own arrives here after [`State::note_gone`] dropped that process, and
117    /// the key may already serve a replacement whose sessions never reserved it:
118    /// releasing one of those would tear down a live agent under its own sessions.
119    pub(super) fn retire(&self, key: &str, process: u64) {
120        let slot = {
121            let mut agents = self.agents.lock();
122            match agents.get(key) {
123                Some(slot) if slot.process == process => {
124                    if slot.live.fetch_sub(1, Ordering::SeqCst) > 1 {
125                        return;
126                    }
127                    agents.remove(key)
128                }
129                // Either no process answers to this command, or the one that
130                // reserved this session's slot left and a replacement took the
131                // key. The reservation went with the process it was made on; the
132                // one here is not this session's to release.
133                _ => return,
134            }
135        };
136        if let Some(slot) = slot {
137            self.reap(key, slot);
138        }
139    }
140
141    /// End a process nobody's session needs any more.
142    ///
143    /// [`onlyne_acp::Agent::shutdown`] is a bounded wait — end of stdin, then
144    /// `SIGTERM`, then `SIGKILL`, then the confirmation that the child was reaped
145    /// — and every caller that releases a session holds the client's dispatch
146    /// lock. The wait therefore belongs to its own thread. A handle another thread
147    /// still shares is simply dropped here, and `AgentInner`'s teardown runs
148    /// wherever the last share of it goes away.
149    fn reap(&self, key: &str, slot: Arc<AgentSlot>) {
150        let owned = key.to_string();
151        match thread::Builder::new()
152            .name(format!("acp-reap {owned}"))
153            .spawn(move || reap_slot(owned, slot))
154        {
155            Ok(_) => {}
156            // With no thread to wait on, the handle goes with this one, and its
157            // own teardown runs inline. That is the same work on the slow path,
158            // reachable only when the process cannot start a thread at all.
159            Err(error) => tracing::warn!(
160                error = %error,
161                agent = %key,
162                "acp: no reaper thread; the agent handle is dropped where it was released"
163            ),
164        }
165    }
166
167    /// The process left on its own. Every parked request of its sessions fails by
168    /// themselves, so this only keeps a new session from opening on a corpse.
169    fn note_gone(&self, key: &str) {
170        let slot = self.agents.lock().remove(key);
171        if let Some(slot) = slot {
172            tracing::warn!(
173                agent = %key,
174                pid = slot.agent.pid(),
175                process = slot.process,
176                sessions = slot.live.load(Ordering::SeqCst),
177                "acp: agent process exited"
178            );
179        }
180    }
181}
182
183/// Wait for the process to leave, on the thread that volunteered to.
184fn reap_slot(key: String, slot: Arc<AgentSlot>) {
185    match Arc::try_unwrap(slot) {
186        Ok(slot) => match Arc::try_unwrap(slot.agent) {
187            Ok(agent) => {
188                if let Err(error) = agent.shutdown() {
189                    tracing::warn!(agent = %key, error = %error, "acp: agent teardown reported");
190                }
191            }
192            Err(_) => tracing::debug!(
193                agent = %key,
194                "acp: agent handle still held by a live turn; its teardown runs with that share"
195            ),
196        },
197        Err(_) => tracing::debug!(
198            agent = %key,
199            "acp: agent slot still shared; its teardown follows the last handle"
200        ),
201    }
202}
203
204/// The permission responder, one per agent process.
205///
206/// It holds only a weak handle: the last strong one going away is what drops the
207/// `Agent`, and a strong handle here would keep a closed agent from ever being
208/// reaped, because the responder outlives the events it waits for.
209fn spawn_responder(slot: &AgentSlot, state: Arc<State>, options: AcpOptions, key: &str) {
210    let agent = Arc::downgrade(&slot.agent);
211    let (key, process) = (key.to_string(), slot.process);
212    if let Err(error) = thread::Builder::new()
213        .name(format!("acp-permissions {key}"))
214        .spawn(move || answer_permissions(agent, state, options, key, process))
215    {
216        tracing::warn!(
217            error = %error,
218            agent = %slot.agent.pid(),
219            "acp: no permission responder; an agent's request waits for its own timeout"
220        );
221    }
222}
223
224fn answer_permissions(
225    agent: Weak<Agent>,
226    state: Arc<State>,
227    options: AcpOptions,
228    key: String,
229    process: u64,
230) {
231    let Some(handle) = agent.upgrade() else {
232        return;
233    };
234    let events = handle.subscribe();
235    drop(handle);
236    while let Ok(event) = events.recv() {
237        match event {
238            Event::Permission(request) => {
239                let Some(handle) = agent.upgrade() else {
240                    return;
241                };
242                let (outcome, refusal) = decide(&options, &request);
243                if let Err(error) = handle.answer(request.request_id.clone(), outcome) {
244                    tracing::warn!(
245                        error = %error,
246                        acp_session = %request.session_id,
247                        "acp: the permission answer did not reach the agent"
248                    );
249                }
250                drop(handle);
251                if let Some(line) = refusal {
252                    record_refusal(&state, &key, process, &request.session_id, &options, line);
253                }
254            }
255            Event::Exited { detail } => {
256                tracing::warn!(
257                    agent = %key,
258                    detail = %detail,
259                    "acp: agent exited; its permission responder stops"
260                );
261                state.note_gone(&key);
262                return;
263            }
264            // Updates belong to the turn threads, each of which holds its own
265            // receiver; answering nothing is what keeps a turn's journal clean.
266            Event::Update { .. } => {}
267        }
268    }
269}
270
271/// Pick the answer one permission request gets, and the refusal line to record
272/// when that answer was not a grant.
273fn decide(
274    options: &AcpOptions,
275    request: &PermissionRequest,
276) -> (PermissionOutcome, Option<String>) {
277    let chosen = if options.allow_permissions {
278        request.option(PermissionOption::ALLOW_ONCE)
279    } else {
280        request
281            .option(PermissionOption::REJECT_ONCE)
282            .or_else(|| request.rejection_option())
283    };
284    match chosen {
285        Some(option) => (
286            PermissionOutcome::Selected {
287                option_id: option.option_id.clone(),
288            },
289            (!options.allow_permissions).then(|| refusal_line("refused", request, Some(option))),
290        ),
291        // Nothing single-use to answer with. `allow_always` is never picked here, so
292        // an agent that offers only a blanket grant gets no decision from this
293        // client, which is its own refusal.
294        None => (
295            PermissionOutcome::Cancelled,
296            Some(refusal_line("declined", request, None)),
297        ),
298    }
299}
300
301fn refusal_line(
302    verb: &str,
303    request: &PermissionRequest,
304    option: Option<&PermissionOption>,
305) -> String {
306    let field = |key: &str| {
307        request
308            .tool_call
309            .get(key)
310            .and_then(Value::as_str)
311            .unwrap_or_default()
312            .to_string()
313    };
314    let title = field("title");
315    let subject = if title.is_empty() {
316        format!("session {}", request.session_id)
317    } else {
318        title
319    };
320    format!(
321        "{verb} {subject} [kind={}, option={}]",
322        or_dash(field("kind")),
323        option.map_or("-".to_string(), |option| or_dash(option.kind.clone())),
324    )
325}
326
327/// File one permission refusal under the session that was asked. The agent command
328/// and the process it started scope the id, because the same command can be served
329/// by more than one process in turn and each of those is free to choose `sess-1`.
330fn record_refusal(
331    state: &State,
332    agent_key: &str,
333    process: u64,
334    session_id: &str,
335    options: &AcpOptions,
336    line: String,
337) {
338    let entry = state
339        .sessions
340        .lock()
341        .get(&(agent_key.to_string(), process, session_id.to_string()))
342        .cloned();
343    match entry {
344        Some(entry) => entry.refusals.lock().push(line),
345        None => tracing::debug!(
346            acp_session = %session_id,
347            policy = options.policy(),
348            "acp: refused a permission ask for a session this client already released"
349        ),
350    }
351}