use crate::backend::{OutcomeFeed, OutcomeSink};
use crate::content::ContentWriter;
use onlyne_acp::Agent;
use parking_lot::{Condvar, Mutex};
use std::collections::BTreeMap;
use std::path::PathBuf;
use std::sync::Arc;
use std::sync::atomic::{AtomicU64, AtomicUsize};
use std::time::Duration;
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct AcpOptions {
pub mode: String,
pub model: String,
pub reasoning_effort: String,
pub allow_permissions: bool,
}
impl AcpOptions {
pub(super) fn policy(&self) -> &'static str {
if self.allow_permissions {
"allow"
} else {
"deny"
}
}
}
#[derive(Clone)]
pub struct AcpBackend {
pub(super) options: AcpOptions,
pub(super) state: Arc<State>,
}
pub(super) struct State {
pub(super) agents: Mutex<BTreeMap<String, Arc<AgentSlot>>>,
pub(super) sessions: Mutex<BTreeMap<(String, u64, String), Arc<SessionEntry>>>,
pub(super) process: AtomicU64,
pub(super) sink: OutcomeSink,
pub(super) feed: OutcomeFeed,
pub(super) content: ContentWriter,
}
pub(super) struct AgentSlot {
pub(super) agent: Arc<Agent>,
pub(super) process: u64,
pub(super) live: AtomicUsize,
}
pub(super) struct SessionEntry {
pub(super) task_id: String,
pub(super) id: String,
pub(super) agent_key: String,
pub(super) process: u64,
pub(super) workdir: PathBuf,
pub(super) agent: Arc<Agent>,
pub(super) turn: Turn,
pub(super) refusals: Mutex<Vec<String>>,
}
impl SessionEntry {
pub(super) fn current_task(&self) -> String {
self.task_id.clone()
}
pub(super) fn key(&self) -> (String, u64, String) {
(self.agent_key.clone(), self.process, self.id.clone())
}
}
pub(super) struct Turn {
pub(super) phase: Mutex<TurnPhase>,
ended: Condvar,
}
pub(super) struct TurnPhase {
pub(super) live: bool,
generation: u64,
}
impl Turn {
pub(super) fn new() -> Self {
Turn {
phase: Mutex::new(TurnPhase {
live: false,
generation: 0,
}),
ended: Condvar::new(),
}
}
pub(super) fn begin(&self) -> Option<u64> {
let mut phase = self.phase.lock();
if phase.live {
return None;
}
phase.live = true;
phase.generation += 1;
Some(phase.generation)
}
pub(super) fn finish(&self) {
self.phase.lock().live = false;
self.ended.notify_all();
}
pub(super) fn generation(&self) -> u64 {
self.phase.lock().generation
}
pub(super) fn waited_out(&self, generation: u64, budget: Duration) -> bool {
let mut phase = self.phase.lock();
if !Turn::running(&phase, generation) {
return true;
}
self.ended.wait_while_for(
&mut phase,
|phase| phase.live && phase.generation == generation,
budget,
);
!Turn::running(&phase, generation)
}
fn running(phase: &TurnPhase, generation: u64) -> bool {
phase.live && phase.generation == generation
}
}