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}