Skip to main content

lattice_mode/
inbound.rs

1//! Boot-composition BC.1: the generic *inbound* primitive.
2//!
3//! Generalizes the I3 `ClaudeCodeInboundBus` and LSP's hand-rolled inbound
4//! buses into ONE reusable primitive: a channel whose `send` **wakes the
5//! editor** (`async_landed.notify_one`) so off-keystroke work reaches the
6//! screen WITHOUT a keypress, and whose items are drained per-tick through a
7//! handler that maps each to existing [`Effect`]s.
8//!
9//! CRITICAL (paramount goal #4 — async-correct *by construction*, not by
10//! discipline): the wake is baked into [`InboundBus::send`], so it is
11//! structurally impossible to forget. This is the bug class
12//! `boot-composition.md` §3 designs out: "forget the wake and a refresh only
13//! reaches the screen on the next keypress."
14//!
15//! Pairs with [`TickCallbackRegistry`](crate::tick_callback::TickCallbackRegistry):
16//! [`make_inbound`] returns the bus PLUS a [`TickCallback`] drain closure the
17//! caller registers with the registry. Per `feedback_mode_owns_its_surface`,
18//! the `handler` (the map from request → `Effect`) lives in the owning crate;
19//! the host only runs the generic drain + applies the returned effects.
20//!
21//! The host-facing bundle (`lattice_host::BootContext`) exposes this as
22//! `BootContext::inbound::<T>(handler)`, which calls [`make_inbound`] against
23//! the editor's `async_landed` and registers the drain on the shared
24//! tick-callback registry — so a subsystem gets the bus back and never
25//! touches the wake or the registry directly.
26
27use std::sync::Arc;
28
29use lattice_grammar::effect::Effect;
30use tokio::sync::{Notify, mpsc};
31
32use crate::tick_callback::TickCallback;
33
34/// Sender half of the inbound primitive. Held by the off-thread producer (a
35/// WS task, an LSP forwarder, …); `send` hands an item to the editor thread
36/// and wakes it.
37///
38/// `Clone` is implemented manually so `T` need *not* be `Clone` — only the
39/// channel sender and the `Arc<Notify>` are cloned (a `#[derive(Clone)]`
40/// would spuriously bound `T: Clone`).
41pub struct InboundBus<T> {
42    tx: mpsc::UnboundedSender<T>,
43    wake: Arc<Notify>,
44}
45
46impl<T> Clone for InboundBus<T> {
47    fn clone(&self) -> Self {
48        Self {
49            tx: self.tx.clone(),
50            wake: Arc::clone(&self.wake),
51        }
52    }
53}
54
55impl<T> InboundBus<T> {
56    /// Send an item and **wake the editor** so the per-tick drain runs
57    /// off-keystroke. The wake fires only on a successful send.
58    ///
59    /// Returns the item back on failure (receiver dropped — the subsystem
60    /// stopped) so the caller reports a graceful error instead of awaiting a
61    /// reply that will never resolve. Mirrors `ClaudeCodeInboundBus::send`.
62    pub fn send(&self, item: T) -> Result<(), T> {
63        self.tx.send(item).map_err(|e| e.0)?;
64        self.wake.notify_one();
65        Ok(())
66    }
67}
68
69/// Build an inbound primitive.
70///
71/// Returns the [`InboundBus`] (the sender, whose `send` wakes `wake`) and a
72/// [`TickCallback`] drain closure. The caller registers the drain with the
73/// [`TickCallbackRegistry`](crate::tick_callback::TickCallbackRegistry); each
74/// tick the drain `try_recv`s every pending item, runs it through `handler`,
75/// and returns the concatenated `Effect`s for the host to apply.
76///
77/// `handler` maps one inbound item to zero or more `Effect`s — the
78/// crate-owned validate→map step (e.g. the I3 optimistic-ack: resolve the
79/// request's oneshot, return one effect on a valid target, none on an unknown
80/// one). It is `FnMut` so it may carry mutable state (a read-state cache, a
81/// counter) across drains, exactly like the existing `make_drain` closures.
82///
83/// Subsystems normally reach this through
84/// [`SubsystemBoot::inbound`](crate::SubsystemBoot::inbound), which supplies
85/// the editor's wake and registers the drain for them.
86///
87/// # Examples
88///
89/// ```
90/// use std::sync::Arc;
91/// use lattice_grammar::effect::{EchoLevel, Effect};
92/// use lattice_mode::inbound::make_inbound;
93/// use tokio::sync::Notify;
94///
95/// let wake = Arc::new(Notify::new()); // the editor's `async_landed`
96/// let (bus, mut drain) = make_inbound::<String, _>(wake, |text| {
97///     vec![Effect::Echo { level: EchoLevel::Info, text }]
98/// });
99///
100/// // Off-thread producer: sending wakes the editor.
101/// std::thread::spawn(move || bus.send("indexed 120 files".into()).unwrap())
102///     .join()
103///     .unwrap();
104///
105/// // Editor actor, next tick: the drain maps every pending item to effects.
106/// let effects = drain();
107/// assert_eq!(effects.len(), 1);
108/// assert!(drain().is_empty()); // nothing left
109/// ```
110pub fn make_inbound<T, H>(wake: Arc<Notify>, mut handler: H) -> (InboundBus<T>, TickCallback)
111where
112    T: Send + 'static,
113    H: FnMut(T) -> Vec<Effect> + Send + 'static,
114{
115    let (tx, mut rx) = mpsc::unbounded_channel::<T>();
116    let bus = InboundBus { tx, wake };
117    let drain: TickCallback = Box::new(move || {
118        let mut effects = Vec::new();
119        while let Ok(item) = rx.try_recv() {
120            effects.extend(handler(item));
121        }
122        effects
123    });
124    (bus, drain)
125}
126
127/// Build a *host-drained* inbound primitive: the wake-baked [`InboundBus`]
128/// sender PLUS the raw receiver, with NO per-tick handler.
129///
130/// Use this (instead of [`make_inbound`]) when the per-item work needs mutable
131/// host state a `FnMut(T) -> Vec<Effect>` handler closure can't capture — e.g.
132/// LSP `workspace/applyEdit`, whose apply is irreducibly `&mut Editor` and
133/// carries `lsp_types` that can't cross the [`Effect`] boundary. The caller
134/// holds the receiver and drains it from its own tick (with whatever `&mut`
135/// access it needs); the wake still fires on every [`InboundBus::send`], so the
136/// work is answered off-keystroke (paramount #4 — the wake lives in the sender,
137/// so it can't be forgotten), exactly like [`make_inbound`]. Keeping the drain
138/// host-side this way avoids introducing an internal-pump variant into the
139/// [`Effect`] vocabulary.
140pub fn make_inbound_raw<T>(wake: Arc<Notify>) -> (InboundBus<T>, mpsc::UnboundedReceiver<T>)
141where
142    T: Send + 'static,
143{
144    let (tx, rx) = mpsc::unbounded_channel::<T>();
145    (InboundBus { tx, wake }, rx)
146}
147
148#[cfg(test)]
149mod tests {
150    #![allow(clippy::unwrap_used)]
151    use super::*;
152    use std::sync::Mutex;
153    use std::time::Duration;
154
155    /// A non-`Clone` payload, to prove `InboundBus<T>: Clone` does not require
156    /// `T: Clone`. (`Debug` only so `send(..).unwrap()` can format the `Err`.)
157    #[derive(Debug)]
158    struct NotClone(u32);
159
160    #[test]
161    fn drain_runs_handler_over_each_item_in_order() {
162        let wake = Arc::new(Notify::new());
163        let seen = Arc::new(Mutex::new(Vec::<u32>::new()));
164        let seen_in = Arc::clone(&seen);
165        let (bus, mut drain) = make_inbound::<u32, _>(wake, move |n| {
166            seen_in.lock().unwrap().push(n);
167            // Return one effect per item so ordering / concatenation is also
168            // observable at the effect level.
169            vec![Effect::None]
170        });
171
172        bus.send(1).unwrap();
173        bus.send(2).unwrap();
174        bus.send(3).unwrap();
175
176        let effects = drain();
177        assert_eq!(effects.len(), 3, "one effect per drained item");
178        assert_eq!(
179            *seen.lock().unwrap(),
180            vec![1, 2, 3],
181            "handler runs over items in send order"
182        );
183    }
184
185    #[test]
186    fn drain_with_nothing_pending_is_empty() {
187        let wake = Arc::new(Notify::new());
188        let (_bus, mut drain) = make_inbound::<u32, _>(wake, |_| vec![Effect::None]);
189        assert!(drain().is_empty(), "no items pending → no effects");
190    }
191
192    #[test]
193    fn handler_may_drop_items_returning_no_effect() {
194        // The I3 optimistic-ack maps unknown targets to *no* effect. Prove a
195        // handler returning an empty vec contributes nothing.
196        let wake = Arc::new(Notify::new());
197        let (bus, mut drain) = make_inbound::<u32, _>(wake, |n| {
198            if n % 2 == 0 {
199                vec![Effect::None]
200            } else {
201                Vec::new()
202            }
203        });
204        bus.send(1).unwrap(); // dropped
205        bus.send(2).unwrap(); // kept
206        bus.send(3).unwrap(); // dropped
207        assert_eq!(drain().len(), 1, "only the even item yields an effect");
208    }
209
210    #[tokio::test]
211    async fn send_wakes_the_editor() {
212        let wake = Arc::new(Notify::new());
213        let (bus, _drain) = make_inbound::<u32, _>(Arc::clone(&wake), |_| vec![Effect::None]);
214        bus.send(7).unwrap();
215        // The permit stored by `notify_one` must let a `notified()` resolve
216        // promptly — i.e. the actor would wake off-keystroke.
217        let woke = tokio::time::timeout(Duration::from_millis(200), wake.notified()).await;
218        assert!(woke.is_ok(), "send must wake the editor");
219    }
220
221    #[test]
222    fn dropped_receiver_makes_send_fail_gracefully() {
223        let wake = Arc::new(Notify::new());
224        let (bus, drain) = make_inbound::<u32, _>(wake, |_| vec![Effect::None]);
225        drop(drain); // the drain owns the receiver — dropping it = subsystem stopped
226        let result = bus.send(7);
227        assert_eq!(
228            result.err(),
229            Some(7),
230            "dropped receiver → send returns the item back"
231        );
232    }
233
234    #[tokio::test]
235    async fn raw_send_wakes_and_receiver_gets_item() {
236        // make_inbound_raw: the wake still fires on send (off-keystroke), but
237        // the caller drains the receiver itself (no handler). Used by LSP
238        // apply-edit, whose apply is irreducibly &mut Editor.
239        let wake = Arc::new(Notify::new());
240        let (bus, mut rx) = make_inbound_raw::<u32>(Arc::clone(&wake));
241        bus.send(42).unwrap();
242        let woke = tokio::time::timeout(Duration::from_millis(200), wake.notified()).await;
243        assert!(woke.is_ok(), "raw send must wake the editor");
244        assert_eq!(
245            rx.try_recv().ok(),
246            Some(42),
247            "host drains the raw receiver itself"
248        );
249    }
250
251    #[test]
252    fn bus_is_clone_without_t_clone() {
253        let wake = Arc::new(Notify::new());
254        let (bus, mut drain) = make_inbound::<NotClone, _>(wake, |v| {
255            let _ = v.0;
256            vec![Effect::None]
257        });
258        let bus2 = bus.clone();
259        bus.send(NotClone(1)).unwrap();
260        bus2.send(NotClone(2)).unwrap();
261        assert_eq!(drain().len(), 2, "both clones feed the same drain");
262    }
263}