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}