onlyne_client/backend/acp/
process.rs1use 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
25const 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 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 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 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 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 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 _ => return,
134 }
135 };
136 if let Some(slot) = slot {
137 self.reap(key, slot);
138 }
139 }
140
141 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 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 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
183fn 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
204fn 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 Event::Update { .. } => {}
267 }
268 }
269}
270
271fn 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 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
327fn 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}