lattice_agent/write_bus.rs
1//! IDE-protocol I3: the write-request payload + per-item handler.
2//!
3//! Writes mutate the editor, which must happen on the editor (actor) thread.
4//! The WS task hands an [`EditorWriteRequest`] to the editor thread over
5//! the generic inbound bus ([`lattice_mode::inbound`], via
6//! [`SubsystemBoot::inbound`](lattice_mode::SubsystemBoot::inbound)) and awaits
7//! a oneshot reply.
8//!
9//! BC.3b: the bespoke `ClaudeCodeInboundBus` + per-tick `make_drain` were
10//! replaced by the generic `InboundBus<EditorWriteRequest>` primitive,
11//! whose `send` wakes the actor off-keystroke (the wake is baked into the
12//! sender — structurally impossible to forget, paramount #4) and whose per-tick
13//! drain runs each request through [`make_handler`]. This module now owns only
14//! the write payload + mapping logic; the channel + drain + wake are
15//! the shared primitive.
16//!
17//! [`make_handler`] validates + maps each request to an EXISTING `Effect` and
18//! resolves its oneshot — optimistic-ack: `ok=true` on a valid map, `ok=false`
19//! on an unknown / non-active target (option C, design §2: per-buffer save/close
20//! targeting lands with the diff/tab work).
21
22use std::path::PathBuf;
23
24use lattice_grammar::Utf16Pos;
25use lattice_grammar::effect::Effect;
26use tokio::sync::oneshot;
27
28use crate::state_cache::EditorStateHandle;
29
30/// A write request's payload.
31#[derive(Debug)]
32pub enum InboundKind {
33 /// Open `path`. `column` (a UTF-16 cursor position from the agent's
34 /// `selection.start`, `None` when absent) is carried unconverted — the host
35 /// resolves it to a byte offset against the opened line. BC.8c follow-up:
36 /// maps to the HOST-APPLIED `OpenBufferAtColumn`, not the peer-applied
37 /// `OpenBufferAt` (which the inbound tick path discards, so openFile never
38 /// actually opened before this fix).
39 OpenFile {
40 path: PathBuf,
41 column: Option<Utf16Pos>,
42 },
43 /// Save the document for `path` (option C: only when it's the active buffer).
44 SaveDocument { path: PathBuf },
45 /// D-fix.6: close the tab from connection `origin_session`. Maps to the
46 /// host-applied [`Effect::CloseSessionDiffs`] — the host rejects that
47 /// connection's programmatic diff session(s) (presentation-agnostic; keyed
48 /// on `origin_session`, NOT `tab_name`), falling back to the legacy
49 /// active-buffer file-close only when `tab_name` matches the active path
50 /// and no diff was torn down.
51 CloseTab {
52 origin_session: u64,
53 tab_name: String,
54 },
55 /// D-fix.6: `closeAllDiffTabs` from connection `origin_session` — reject
56 /// every programmatic diff that connection opened. Maps to
57 /// [`Effect::CloseAllSessionDiffs`].
58 CloseAllDiffTabs { origin_session: u64 },
59}
60
61/// The drain's reply to the WS task. Optimistic-ack: `ok` reflects whether the
62/// request mapped to a valid Effect, not the eventual apply result.
63#[derive(Debug, Clone, PartialEq, Eq)]
64pub struct InboundReply {
65 pub ok: bool,
66 pub message: Option<String>,
67}
68
69impl InboundReply {
70 fn ok() -> Self {
71 Self {
72 ok: true,
73 message: None,
74 }
75 }
76 fn fail(msg: impl Into<String>) -> Self {
77 Self {
78 ok: false,
79 message: Some(msg.into()),
80 }
81 }
82}
83
84/// One write request: payload + the oneshot the drain resolves. Mirrors LSP's
85/// `InboundShowDocument`.
86#[derive(Debug)]
87pub struct EditorWriteRequest {
88 pub kind: InboundKind,
89 pub response: oneshot::Sender<InboundReply>,
90}
91
92/// Build the per-item handler for the generic inbound primitive
93/// ([`SubsystemBoot::inbound`](lattice_mode::SubsystemBoot::inbound)). Maps one
94/// request to its `Effect` (or none, on `ok=false`), resolves the oneshot, and
95/// returns the effect(s) for the host to apply. The generic bus owns the
96/// channel, the per-tick `try_recv` loop, and the off-keystroke wake.
97pub fn make_handler(
98 cache: EditorStateHandle,
99) -> impl FnMut(EditorWriteRequest) -> Vec<Effect> + Send + 'static {
100 move |req| {
101 let (effect, reply) = map_request(&req.kind, &cache);
102 // A dropped response receiver (agent gone) is fine — log-and-skip.
103 let _ = req.response.send(reply);
104 effect.into_iter().collect()
105 }
106}
107
108/// The active buffer's path, if any (option-C targeting).
109fn active_path(cache: &EditorStateHandle) -> Option<PathBuf> {
110 let g = cache.lock().unwrap_or_else(|e| e.into_inner());
111 let active = g.active.as_ref()?;
112 g.open_buffers
113 .get(&active.buffer)
114 .and_then(|b| b.path.clone())
115}
116
117/// Map a write request to an existing Effect + an optimistic reply.
118fn map_request(kind: &InboundKind, cache: &EditorStateHandle) -> (Option<Effect>, InboundReply) {
119 match kind {
120 InboundKind::OpenFile { path, column } => (
121 // BC.8c follow-up: host-applied open (works on the inbound tick
122 // path, where peer-applied `OpenBufferAt` is discarded). The host
123 // does do_edit + the UTF-16→byte cursor conversion against the
124 // opened line; `column = None` opens without forcing the cursor.
125 Some(Effect::OpenBufferAtColumn {
126 path: Some(path.clone()),
127 column: *column,
128 force: false,
129 }),
130 InboundReply::ok(),
131 ),
132 InboundKind::SaveDocument { path } => {
133 if active_path(cache).as_deref() == Some(path.as_path()) {
134 (Some(Effect::SaveBuffer { path: None }), InboundReply::ok())
135 } else {
136 (
137 None,
138 InboundReply::fail(
139 "saveDocument: target is not the active buffer (I3 limitation)",
140 ),
141 )
142 }
143 }
144 // D-fix.6: the diff-vs-buffer decision is HOST-side now (only the host
145 // knows the open programmatic diffs + their `origin_session`). Emit the
146 // host-applied effect carrying the connection id; the host rejects that
147 // connection's diff session(s), else falls back to the active-buffer
148 // file-close via `tab_name`. Optimistic-ack `ok` (the close is
149 // fire-and-forget; the host does the right thing regardless).
150 InboundKind::CloseTab {
151 origin_session,
152 tab_name,
153 } => (
154 Some(Effect::CloseSessionDiffs {
155 origin_session: *origin_session,
156 tab_name: tab_name.clone(),
157 }),
158 InboundReply::ok(),
159 ),
160 InboundKind::CloseAllDiffTabs { origin_session } => (
161 Some(Effect::CloseAllSessionDiffs {
162 origin_session: *origin_session,
163 }),
164 InboundReply::ok(),
165 ),
166 }
167}
168
169#[cfg(test)]
170mod tests {
171 #![allow(clippy::unwrap_used)]
172 use super::*;
173 use crate::state_cache::EditorStateCache;
174 use lattice_protocol::ids::DocumentId;
175 use lattice_protocol::{Event, SelectionSet};
176 use std::sync::{Arc, Mutex};
177 use tokio::sync::Notify;
178
179 fn cache_with_active(path: &str) -> EditorStateHandle {
180 let mut s = EditorStateCache::default();
181 s.apply_event(&Event::DocumentOpened {
182 id: DocumentId::new(1),
183 path: Some(PathBuf::from(path)),
184 version: 1,
185 text: String::new(),
186 });
187 s.apply_event(&Event::SelectionsChanged {
188 id: DocumentId::new(1),
189 version: 1,
190 selections: SelectionSet::default(),
191 });
192 Arc::new(Mutex::new(s))
193 }
194
195 fn empty_cache() -> EditorStateHandle {
196 Arc::new(Mutex::new(EditorStateCache::default()))
197 }
198
199 #[test]
200 fn open_file_maps_to_host_applied_open_ok() {
201 // No selection → column None (open only). Maps to the HOST-APPLIED
202 // OpenBufferAtColumn so it actually opens on the inbound tick path.
203 let (e, r) = map_request(
204 &InboundKind::OpenFile {
205 path: PathBuf::from("/a.rs"),
206 column: None,
207 },
208 &empty_cache(),
209 );
210 assert!(matches!(
211 e,
212 Some(Effect::OpenBufferAtColumn { column: None, .. })
213 ));
214 assert!(r.ok);
215 }
216
217 #[test]
218 fn open_file_with_selection_carries_utf16_column() {
219 let (e, _r) = map_request(
220 &InboundKind::OpenFile {
221 path: PathBuf::from("/a.rs"),
222 column: Some(Utf16Pos { line: 3, col: 7 }),
223 },
224 &empty_cache(),
225 );
226 assert!(matches!(
227 e,
228 Some(Effect::OpenBufferAtColumn {
229 column: Some(Utf16Pos { line: 3, col: 7 }),
230 ..
231 })
232 ));
233 }
234
235 #[test]
236 fn save_active_buffer_maps_to_save_ok() {
237 let (e, r) = map_request(
238 &InboundKind::SaveDocument {
239 path: PathBuf::from("/a.rs"),
240 },
241 &cache_with_active("/a.rs"),
242 );
243 assert!(matches!(e, Some(Effect::SaveBuffer { .. })));
244 assert!(r.ok);
245 }
246
247 #[test]
248 fn save_non_active_buffer_is_ok_false_no_effect() {
249 let (e, r) = map_request(
250 &InboundKind::SaveDocument {
251 path: PathBuf::from("/other.rs"),
252 },
253 &cache_with_active("/a.rs"),
254 );
255 assert!(e.is_none());
256 assert!(!r.ok);
257 }
258
259 #[test]
260 fn close_tab_maps_to_session_scoped_diff_teardown() {
261 // D-fix.6: close_tab now ALWAYS maps to the host-applied
262 // `CloseSessionDiffs` carrying the connection id (the host decides
263 // diff-vs-buffer with its diff state). Optimistic-ack ok.
264 let (e, r) = map_request(
265 &InboundKind::CloseTab {
266 origin_session: 7,
267 tab_name: "/a.rs".to_string(),
268 },
269 &cache_with_active("/a.rs"),
270 );
271 match e {
272 Some(Effect::CloseSessionDiffs {
273 origin_session,
274 tab_name,
275 }) => {
276 assert_eq!(origin_session, 7, "scoped to the originating connection");
277 assert_eq!(tab_name, "/a.rs");
278 }
279 other => panic!("expected CloseSessionDiffs, got {other:?}"),
280 }
281 assert!(r.ok);
282 }
283
284 #[test]
285 fn close_all_diff_tabs_maps_to_session_scoped_bulk_teardown() {
286 let (e, r) = map_request(
287 &InboundKind::CloseAllDiffTabs { origin_session: 9 },
288 &empty_cache(),
289 );
290 assert!(matches!(
291 e,
292 Some(Effect::CloseAllSessionDiffs { origin_session: 9 })
293 ));
294 assert!(r.ok);
295 }
296
297 #[test]
298 fn handler_maps_request_and_resolves_oneshot() {
299 // The generic inbound primitive (`make_inbound`) owns the channel + the
300 // per-tick drain + the wake (those are pinned by lattice-mode's inbound
301 // tests); this pins claude's per-item handler — it maps the request to
302 // an Effect and resolves the oneshot.
303 let (bus, mut drain) = lattice_mode::inbound::make_inbound(
304 Arc::new(Notify::new()),
305 make_handler(empty_cache()),
306 );
307
308 let (resp_tx, mut resp_rx) = oneshot::channel();
309 bus.send(EditorWriteRequest {
310 kind: InboundKind::OpenFile {
311 path: PathBuf::from("/a.rs"),
312 column: None,
313 },
314 response: resp_tx,
315 })
316 .expect("send ok");
317
318 let effects = drain();
319 assert_eq!(effects.len(), 1);
320 let reply = resp_rx.try_recv().expect("oneshot resolved");
321 assert!(reply.ok);
322 }
323}