pub fn make_inbound<T, H>(
wake: Arc<Notify>,
handler: H,
) -> (InboundBus<T>, TickCallback)Expand description
Build an inbound primitive.
Returns the InboundBus (the sender, whose send wakes wake) and a
TickCallback drain closure. The caller registers the drain with the
TickCallbackRegistry; each
tick the drain try_recvs every pending item, runs it through handler,
and returns the concatenated Effects for the host to apply.
handler maps one inbound item to zero or more Effects — the
crate-owned validate→map step (e.g. the I3 optimistic-ack: resolve the
request’s oneshot, return one effect on a valid target, none on an unknown
one). It is FnMut so it may carry mutable state (a read-state cache, a
counter) across drains, exactly like the existing make_drain closures.
Subsystems normally reach this through
SubsystemBoot::inbound, which supplies
the editor’s wake and registers the drain for them.
§Examples
use std::sync::Arc;
use lattice_grammar::effect::{EchoLevel, Effect};
use lattice_mode::inbound::make_inbound;
use tokio::sync::Notify;
let wake = Arc::new(Notify::new()); // the editor's `async_landed`
let (bus, mut drain) = make_inbound::<String, _>(wake, |text| {
vec![Effect::Echo { level: EchoLevel::Info, text }]
});
// Off-thread producer: sending wakes the editor.
std::thread::spawn(move || bus.send("indexed 120 files".into()).unwrap())
.join()
.unwrap();
// Editor actor, next tick: the drain maps every pending item to effects.
let effects = drain();
assert_eq!(effects.len(), 1);
assert!(drain().is_empty()); // nothing left