use crate::backend::*;
use crate::content::ContentWriter;
use onlyne_acp::{
Agent, AgentOptions, ClientCapabilities, ClientInfo, Event, PermissionOption,
PermissionOutcome, PermissionRequest,
};
use parking_lot::Mutex;
use std::path::Path;
use std::sync::Weak;
use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering};
use std::thread;
use super::journal::or_dash;
use super::state::{AcpBackend, AcpOptions, AgentSlot, State};
const CLIENT_NAME: &str = "onlyne-client";
impl AcpBackend {
pub fn new(options: AcpOptions) -> Self {
let (sink, feed) = OutcomeFeed::channel();
AcpBackend {
options,
state: Arc::new(State {
agents: Mutex::new(BTreeMap::new()),
sessions: Mutex::new(BTreeMap::new()),
process: AtomicU64::new(0),
sink,
feed,
content: ContentWriter::default(),
}),
}
}
pub(super) fn command_of(spec: &SpawnSpec) -> Result<Vec<String>> {
match spec.command.first() {
Some(program) if !program.trim().is_empty() => Ok(spec.command.clone()),
_ => Err(anyhow::anyhow!(
"acp: the role's `[client.runtime] command` is empty; there is nothing to run \
for task {}",
spec.task_id
)),
}
}
pub(super) fn agent_for(
&self,
key: &str,
command: &[String],
cwd: &Path,
env: &BTreeMap<String, String>,
) -> Result<Arc<AgentSlot>> {
let mut agents = self.state.agents.lock();
if let Some(slot) = agents.get(key) {
if !slot.agent.is_gone() {
slot.live.fetch_add(1, Ordering::SeqCst);
return Ok(Arc::clone(slot));
}
tracing::warn!(agent = %key, "acp: replacing an exited agent process");
agents.remove(key);
}
let agent = Arc::new(Agent::start(AgentOptions {
command: command.to_vec(),
cwd: Some(cwd.to_path_buf()),
env: env.clone(),
})?);
if let Err(error) =
agent.initialize(ClientInfo::new(CLIENT_NAME), ClientCapabilities::default())
{
let _ = thread::Builder::new()
.name(format!("acp-drop {key}"))
.spawn(move || drop(agent));
return Err(anyhow::anyhow!("acp: {key} refused the handshake: {error}"));
}
let slot = Arc::new(AgentSlot {
live: AtomicUsize::new(1),
process: self.state.process.fetch_add(1, Ordering::Relaxed) + 1,
agent: Arc::clone(&agent),
});
spawn_responder(&slot, Arc::clone(&self.state), self.options.clone(), key);
agents.insert(key.to_string(), Arc::clone(&slot));
Ok(slot)
}
}
impl State {
pub(super) fn retire(&self, key: &str, process: u64) {
let slot = {
let mut agents = self.agents.lock();
match agents.get(key) {
Some(slot) if slot.process == process => {
if slot.live.fetch_sub(1, Ordering::SeqCst) > 1 {
return;
}
agents.remove(key)
}
_ => return,
}
};
if let Some(slot) = slot {
self.reap(key, slot);
}
}
fn reap(&self, key: &str, slot: Arc<AgentSlot>) {
let owned = key.to_string();
match thread::Builder::new()
.name(format!("acp-reap {owned}"))
.spawn(move || reap_slot(owned, slot))
{
Ok(_) => {}
Err(error) => tracing::warn!(
error = %error,
agent = %key,
"acp: no reaper thread; the agent handle is dropped where it was released"
),
}
}
fn note_gone(&self, key: &str) {
let slot = self.agents.lock().remove(key);
if let Some(slot) = slot {
tracing::warn!(
agent = %key,
pid = slot.agent.pid(),
process = slot.process,
sessions = slot.live.load(Ordering::SeqCst),
"acp: agent process exited"
);
}
}
}
fn reap_slot(key: String, slot: Arc<AgentSlot>) {
match Arc::try_unwrap(slot) {
Ok(slot) => match Arc::try_unwrap(slot.agent) {
Ok(agent) => {
if let Err(error) = agent.shutdown() {
tracing::warn!(agent = %key, error = %error, "acp: agent teardown reported");
}
}
Err(_) => tracing::debug!(
agent = %key,
"acp: agent handle still held by a live turn; its teardown runs with that share"
),
},
Err(_) => tracing::debug!(
agent = %key,
"acp: agent slot still shared; its teardown follows the last handle"
),
}
}
fn spawn_responder(slot: &AgentSlot, state: Arc<State>, options: AcpOptions, key: &str) {
let agent = Arc::downgrade(&slot.agent);
let (key, process) = (key.to_string(), slot.process);
if let Err(error) = thread::Builder::new()
.name(format!("acp-permissions {key}"))
.spawn(move || answer_permissions(agent, state, options, key, process))
{
tracing::warn!(
error = %error,
agent = %slot.agent.pid(),
"acp: no permission responder; an agent's request waits for its own timeout"
);
}
}
fn answer_permissions(
agent: Weak<Agent>,
state: Arc<State>,
options: AcpOptions,
key: String,
process: u64,
) {
let Some(handle) = agent.upgrade() else {
return;
};
let events = handle.subscribe();
drop(handle);
while let Ok(event) = events.recv() {
match event {
Event::Permission(request) => {
let Some(handle) = agent.upgrade() else {
return;
};
let (outcome, refusal) = decide(&options, &request);
if let Err(error) = handle.answer(request.request_id.clone(), outcome) {
tracing::warn!(
error = %error,
acp_session = %request.session_id,
"acp: the permission answer did not reach the agent"
);
}
drop(handle);
if let Some(line) = refusal {
record_refusal(&state, &key, process, &request.session_id, &options, line);
}
}
Event::Exited { detail } => {
tracing::warn!(
agent = %key,
detail = %detail,
"acp: agent exited; its permission responder stops"
);
state.note_gone(&key);
return;
}
Event::Update { .. } => {}
}
}
}
fn decide(
options: &AcpOptions,
request: &PermissionRequest,
) -> (PermissionOutcome, Option<String>) {
let chosen = if options.allow_permissions {
request.option(PermissionOption::ALLOW_ONCE)
} else {
request
.option(PermissionOption::REJECT_ONCE)
.or_else(|| request.rejection_option())
};
match chosen {
Some(option) => (
PermissionOutcome::Selected {
option_id: option.option_id.clone(),
},
(!options.allow_permissions).then(|| refusal_line("refused", request, Some(option))),
),
None => (
PermissionOutcome::Cancelled,
Some(refusal_line("declined", request, None)),
),
}
}
fn refusal_line(
verb: &str,
request: &PermissionRequest,
option: Option<&PermissionOption>,
) -> String {
let field = |key: &str| {
request
.tool_call
.get(key)
.and_then(Value::as_str)
.unwrap_or_default()
.to_string()
};
let title = field("title");
let subject = if title.is_empty() {
format!("session {}", request.session_id)
} else {
title
};
format!(
"{verb} {subject} [kind={}, option={}]",
or_dash(field("kind")),
option.map_or("-".to_string(), |option| or_dash(option.kind.clone())),
)
}
fn record_refusal(
state: &State,
agent_key: &str,
process: u64,
session_id: &str,
options: &AcpOptions,
line: String,
) {
let entry = state
.sessions
.lock()
.get(&(agent_key.to_string(), process, session_id.to_string()))
.cloned();
match entry {
Some(entry) => entry.refusals.lock().push(line),
None => tracing::debug!(
acp_session = %session_id,
policy = options.policy(),
"acp: refused a permission ask for a session this client already released"
),
}
}