Skip to main content

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}