Skip to main content

onlyne_client/backend/acp/
state.rs

1//! What an ACP backend holds: the client's options, the handle that owns them,
2//! and the per-process and per-session bookkeeping underneath it. No I/O.
3//!
4//! A session is named by the triple (agent command, this client's name for the
5//! process, agent-chosen id): ACP ids are unique within one agent process, and
6//! one command can be served by several processes in turn, so the command alone
7//! does not own an id.
8
9use crate::backend::{OutcomeFeed, OutcomeSink};
10use crate::content::ContentWriter;
11use onlyne_acp::Agent;
12use parking_lot::{Condvar, Mutex};
13use std::collections::BTreeMap;
14use std::path::PathBuf;
15use std::sync::Arc;
16use std::sync::atomic::{AtomicU64, AtomicUsize};
17use std::time::Duration;
18
19/// The `[client.acp]` table, in the shape this crate can hold without a config
20/// dependency: every field already defaulted, and the permission mode reduced to
21/// the one decision a backend makes with it.
22#[derive(Debug, Clone, Default, PartialEq, Eq)]
23pub struct AcpOptions {
24    /// `session/set_mode` id. Empty leaves the agent's own default.
25    pub mode: String,
26    /// `model` config option value. Empty leaves the agent's default.
27    pub model: String,
28    /// `reasoning_effort` config option value. Empty leaves the default.
29    pub reasoning_effort: String,
30    /// Whether this client answers an agent's permission request with a grant.
31    /// Off by default: the refusal is recorded and a supervisor decides.
32    pub allow_permissions: bool,
33}
34
35impl AcpOptions {
36    /// The word a fault record names as the policy behind a refusal.
37    pub(super) fn policy(&self) -> &'static str {
38        if self.allow_permissions {
39            "allow"
40        } else {
41            "deny"
42        }
43    }
44}
45
46/// The backend: an agent per command key, an ACP session per task, one outcome
47/// stream for the whole role.
48#[derive(Clone)]
49pub struct AcpBackend {
50    pub(super) options: AcpOptions,
51    pub(super) state: Arc<State>,
52}
53
54pub(super) struct State {
55    /// Live agent processes, keyed by the rendered command that started them.
56    pub(super) agents: Mutex<BTreeMap<String, Arc<AgentSlot>>>,
57    /// Every session this client holds, keyed by the agent command, the process
58    /// of that command the id came from, and the id that process chose for the
59    /// session. ACP ids are unique within a process, not across processes, so
60    /// the pair of command and id is not enough: when a crashed agent is
61    /// replaced under the same command, the replacement is free to hand out
62    /// `sess-1` again, and the session that id names is a different one. A task
63    /// id names none of them stably, which is why it is not in the key.
64    pub(super) sessions: Mutex<BTreeMap<(String, u64, String), Arc<SessionEntry>>>,
65    /// The names handed to agent processes as this client starts them. A process
66    /// keeps its name for as long as this client means the term to cover it, so
67    /// neither its sessions nor its reservations can be mistaken for another
68    /// process's. The pid is not this: the operating system recycles those.
69    pub(super) process: AtomicU64,
70    pub(super) sink: OutcomeSink,
71    pub(super) feed: OutcomeFeed,
72    /// Serializes journal cursors across every task served by this role.
73    pub(super) content: ContentWriter,
74}
75
76pub(super) struct AgentSlot {
77    pub(super) agent: Arc<Agent>,
78    /// This client's name for the process behind this slot; see [`State::process`].
79    pub(super) process: u64,
80    /// Sessions of this process still held by this client, so the last one to
81    /// leave can take the process with it. One reservation is taken for every
82    /// slot [`AcpBackend::agent_for`] hands out and released by exactly one
83    /// [`State::retire`] of the session that took it — including a session whose
84    /// open failed after the process was chosen — and never by a session of the
85    /// process this one replaced.
86    pub(super) live: AtomicUsize,
87}
88
89/// One ACP session, plus the turn state this client keeps for it.
90pub(super) struct SessionEntry {
91    /// The task this session serves. One session runs one task, and that task
92    /// owns its own journal. Bound when the session opens and never rewritten:
93    /// a second task takes a second session.
94    pub(super) task_id: String,
95    /// The id the agent gave this session; every later request is keyed by it.
96    pub(super) id: String,
97    /// The command key of the process serving this session.
98    pub(super) agent_key: String,
99    /// The name of that process, of the three that name this session.
100    pub(super) process: u64,
101    /// The directory the agent runs in, which is where its journal lives.
102    pub(super) workdir: PathBuf,
103    pub(super) agent: Arc<Agent>,
104    pub(super) turn: Turn,
105    /// Permission asks refused since the turn began, written by the responder
106    /// thread and taken by the turn thread when the turn ends. Per turn, because
107    /// the asks interleave with the updates of the one parked prompt they belong
108    /// to, and a refusal has to name the turn that produced it.
109    pub(super) refusals: Mutex<Vec<String>>,
110}
111
112impl SessionEntry {
113    pub(super) fn current_task(&self) -> String {
114        self.task_id.clone()
115    }
116
117    /// The key this session holds in [`State::sessions`].
118    pub(super) fn key(&self) -> (String, u64, String) {
119        (self.agent_key.clone(), self.process, self.id.clone())
120    }
121}
122
123/// The turn bookkeeping of one session: whether a turn is in flight, and which
124/// turn a waiter is waiting out.
125pub(super) struct Turn {
126    pub(super) phase: Mutex<TurnPhase>,
127    ended: Condvar,
128}
129
130pub(super) struct TurnPhase {
131    pub(super) live: bool,
132    generation: u64,
133}
134
135impl Turn {
136    pub(super) fn new() -> Self {
137        Turn {
138            phase: Mutex::new(TurnPhase {
139                live: false,
140                generation: 0,
141            }),
142            ended: Condvar::new(),
143        }
144    }
145
146    /// Claim the session for a turn, refusing one already in flight. Answers the
147    /// generation a later waiter names.
148    pub(super) fn begin(&self) -> Option<u64> {
149        let mut phase = self.phase.lock();
150        if phase.live {
151            return None;
152        }
153        phase.live = true;
154        phase.generation += 1;
155        Some(phase.generation)
156    }
157
158    /// Mark the turn over and wake anyone waiting out the close of a session.
159    pub(super) fn finish(&self) {
160        self.phase.lock().live = false;
161        self.ended.notify_all();
162    }
163
164    pub(super) fn generation(&self) -> u64 {
165        self.phase.lock().generation
166    }
167
168    /// Whether the turn named by `generation` is over: true when it ended, false
169    /// when the wait ran out. A newer generation counts as ended too, because a
170    /// turn cannot start before the one before it stopped.
171    pub(super) fn waited_out(&self, generation: u64, budget: Duration) -> bool {
172        let mut phase = self.phase.lock();
173        if !Turn::running(&phase, generation) {
174            return true;
175        }
176        self.ended.wait_while_for(
177            &mut phase,
178            |phase| phase.live && phase.generation == generation,
179            budget,
180        );
181        !Turn::running(&phase, generation)
182    }
183
184    fn running(phase: &TurnPhase, generation: u64) -> bool {
185        phase.live && phase.generation == generation
186    }
187}