lattice_ai/acp/handle.rs
1//! Editor-thread handle onto the AI supervisor task (AI-1b).
2//!
3//! `AiClientHandle` is the only thing the editor thread touches: it sends
4//! `AiCmd`s into the supervisor's channel and reads an `ArcSwap<AiState>`
5//! snapshot. Both are non-blocking -- no protocol I/O, no locks held across
6//! `.await`, ever runs on the editor thread. The supervisor itself (owning
7//! the provider child, `Connection`, and `SessionId`) lives in
8//! `supervisor.rs`.
9
10use std::sync::Arc;
11use std::sync::atomic::AtomicUsize;
12
13use arc_swap::ArcSwap;
14use tokio::sync::mpsc;
15
16use lattice_agent::SessionKey;
17
18use crate::acp::providers::ProviderConfig;
19
20/// Editor-visible snapshot of the active AI session, if any. Cheap to clone
21/// and compare; read via [`AiClientHandle::snapshot`].
22#[derive(Debug, Clone, Default, PartialEq, Eq)]
23pub struct AiState {
24 pub running: bool,
25 pub provider: Option<&'static str>,
26 /// Active session key (provider + per-provider index), if a session is
27 /// open.
28 pub session: Option<SessionKey>,
29 /// AU‑5: trust mode. `false` (the default) is **review** mode — file edits
30 /// are gated on a diff verdict and un-reviewable mutating ops are denied.
31 /// `true` auto-grants every permission request without the diff gate.
32 pub auto_accept: bool,
33 /// AUX‑4: number of prompts currently queued behind an in-flight one.
34 pub queue_len: usize,
35}
36
37/// Commands the handle sends into the supervisor's command loop.
38pub(crate) enum AiCmd {
39 Start(ProviderConfig),
40 Prompt(String),
41 /// AU‑3: interrupt the active turn without ending the session. The
42 /// supervisor forwards an ACP `session/cancel`; the session stays open.
43 Interrupt,
44 /// AU‑5: set trust mode. `true` auto-grants every permission request; `false`
45 /// restores review mode (diff-gated edits, denied un-reviewable ops).
46 SetAutoAccept(bool),
47 Stop,
48}
49
50/// Clone-able handle onto a running (or idle) AI supervisor task.
51///
52/// Fields are `pub(crate)` so a later task's `commands.rs` tests can build a
53/// handle directly from a channel + `ArcSwap` without going through
54/// `spawn`.
55#[derive(Clone)]
56pub struct AiClientHandle {
57 pub(crate) cmd_tx: mpsc::UnboundedSender<AiCmd>,
58 pub(crate) state: Arc<ArcSwap<AiState>>,
59 /// AUX‑4: shared with the supervisor loop; read by the headerline renderer.
60 pub queue_len: Arc<AtomicUsize>,
61}
62
63impl AiClientHandle {
64 /// AUX‑4: expose the live queue length for the headerline.
65 pub fn queue_len(&self) -> usize {
66 self.queue_len.load(std::sync::atomic::Ordering::Relaxed)
67 }
68
69 /// Ask the supervisor to start `provider`. Non-blocking; the result
70 /// surfaces later via [`AiClientHandle::snapshot`] and the provider's
71 /// `AiLogger` ring.
72 pub fn start(&self, provider: ProviderConfig) {
73 let _ = self.cmd_tx.send(AiCmd::Start(provider));
74 }
75
76 /// Ask the supervisor to send `text` as a prompt on the active session.
77 /// Non-blocking; if no session is open the supervisor drops the prompt
78 /// and logs a `Warn`-level "prompt dropped: no active session" record
79 /// to the subsystem-wide `AiLogger` ring instead of sending it.
80 pub fn prompt(&self, text: String) {
81 let _ = self.cmd_tx.send(AiCmd::Prompt(text));
82 }
83
84 /// AU‑3: interrupt the active turn without ending the session.
85 /// Non-blocking; the supervisor forwards an ACP `session/cancel`. If no
86 /// session is open it's a no-op. Distinct from [`AiClientHandle::stop`],
87 /// which tears the session (and provider child) down.
88 pub fn interrupt(&self) {
89 let _ = self.cmd_tx.send(AiCmd::Interrupt);
90 }
91
92 /// AU‑5: set trust mode. Non-blocking; the supervisor applies the flag and
93 /// republishes `AiState`.
94 pub fn set_auto_accept(&self, on: bool) {
95 let _ = self.cmd_tx.send(AiCmd::SetAutoAccept(on));
96 }
97
98 /// AU‑5: flip trust mode, returning the value it was flipped to (from the
99 /// current snapshot). Non-blocking; the supervisor applies + republishes.
100 pub fn toggle_auto_accept(&self) -> bool {
101 let next = !self.snapshot().auto_accept;
102 self.set_auto_accept(next);
103 next
104 }
105
106 /// Ask the supervisor to stop the active session. Non-blocking.
107 pub fn stop(&self) {
108 let _ = self.cmd_tx.send(AiCmd::Stop);
109 }
110
111 /// Read the current state snapshot. Never blocks.
112 pub fn snapshot(&self) -> AiState {
113 (**self.state.load()).clone()
114 }
115}
116
117#[cfg(test)]
118mod tests {
119 use super::*;
120
121 #[test]
122 fn idle_snapshot_and_nonblocking_sends() {
123 let (cmd_tx, cmd_rx) = mpsc::unbounded_channel();
124 let state = Arc::new(ArcSwap::from_pointee(AiState::default()));
125 let handle = AiClientHandle {
126 cmd_tx,
127 state,
128 queue_len: Arc::new(AtomicUsize::new(0)),
129 };
130
131 assert_eq!(handle.snapshot(), AiState::default());
132
133 // Drop the receiver -- sends must not panic even though nothing is
134 // listening.
135 drop(cmd_rx);
136 handle.prompt("hi".into());
137 handle.stop();
138 }
139
140 /// AU‑5: `toggle_auto_accept` flips against the current snapshot, returns
141 /// the new value, and sends a matching `SetAutoAccept`.
142 #[test]
143 fn toggle_auto_accept_flips_and_sends() {
144 let (cmd_tx, mut cmd_rx) = mpsc::unbounded_channel();
145 let state = Arc::new(ArcSwap::from_pointee(AiState::default()));
146 let handle = AiClientHandle {
147 cmd_tx,
148 state,
149 queue_len: Arc::new(AtomicUsize::new(0)),
150 };
151
152 // Default is review mode (false) → toggling turns trust on.
153 assert!(handle.toggle_auto_accept());
154 assert!(matches!(cmd_rx.try_recv(), Ok(AiCmd::SetAutoAccept(true))));
155
156 // Reflect the applied state, then toggle back off.
157 handle.state.store(Arc::new(AiState {
158 auto_accept: true,
159 ..AiState::default()
160 }));
161 assert!(!handle.toggle_auto_accept());
162 assert!(matches!(cmd_rx.try_recv(), Ok(AiCmd::SetAutoAccept(false))));
163 }
164
165 // ── AUX‑4: queue_len ──
166
167 #[test]
168 fn queue_len_defaults_zero() {
169 let handle = AiClientHandle {
170 cmd_tx: mpsc::unbounded_channel().0,
171 state: Arc::new(ArcSwap::from_pointee(AiState::default())),
172 queue_len: Arc::new(AtomicUsize::new(0)),
173 };
174 assert_eq!(handle.queue_len(), 0);
175 assert_eq!(handle.snapshot().queue_len, 0);
176 }
177
178 #[test]
179 fn queue_len_accessible_on_handle() {
180 let ql = Arc::new(AtomicUsize::new(3));
181 let handle = AiClientHandle {
182 cmd_tx: mpsc::unbounded_channel().0,
183 state: Arc::new(ArcSwap::from_pointee(AiState::default())),
184 queue_len: ql.clone(),
185 };
186 assert_eq!(handle.queue_len(), 3);
187 // Mutating the atomic is reflected in the handle's live reader.
188 ql.store(5, std::sync::atomic::Ordering::Relaxed);
189 assert_eq!(handle.queue_len(), 5);
190 }
191}