Skip to main content

make_inbound

Function make_inbound 

Source
pub fn make_inbound<T, H>(
    wake: Arc<Notify>,
    handler: H,
) -> (InboundBus<T>, TickCallback)
where T: Send + 'static, H: FnMut(T) -> Vec<Effect> + Send + 'static,
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