Skip to main content

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}