Effects and subscriptions
An Effect<Message> is one-shot work returned by init or update. A
Subscription<Message> is a source that remains active for as long as the
current model declares it. Both produce messages; only update mutates the
model.
Choose the right value
Section titled “Choose the right value”| Need | Use | Lifetime and delivery |
|---|---|---|
| No work | Effect::none() |
Produces no message |
| Blocking filesystem, database, or CPU work | Effect::perform |
Runs off the application thread; completion becomes one message |
| Existing async operation | Effect::future |
The runtime polls the future; completion becomes one message |
| One delayed message | Effect::after |
Sleeps on the runtime clock, then produces one message |
| Replace unfinished work for one purpose | perform_latest, future_latest, or after_latest |
A newer effect with the same key replaces the older unfinished effect |
| Several independent one-shot operations | Effect::batch |
Starts members in source order; runs them concurrently; completion order is unspecified |
| Stop the application | Effect::shutdown() |
Starts no sibling effect, stops subscriptions, suppresses in-flight completions, and restores the terminal |
| Terminal input, resize, interval, signal, or draw errors | Built-in Subscription constructor |
Active while declared by the current model |
| Existing asynchronous stream | Subscription::stream[_with] |
The stream is polled while its key remains declared |
| Custom async producer | Subscription::run[_with] |
Receives a Sender; use send(...).await |
| Custom blocking producer | Subscription::run_blocking[_with] |
Runs on a blocking thread; use blocking_send |
Perform blocking, asynchronous, and delayed work
Section titled “Perform blocking, asynchronous, and delayed work”Effect construction is inert: work starts only when the runtime interprets the
effect returned from init or update.
use std::time::Duration;use urushi_tui_app::Effect;
#[derive(Debug, PartialEq)]enum Message { FileRead(String), RequestFinished(u16), Timeout,}
fn read_file() -> Effect<Message> { Effect::perform(|| Message::FileRead("settings = loaded".into()))}
fn request() -> Effect<Message> { Effect::future(async { Message::RequestFinished(200) })}
fn timeout() -> Effect<Message> { Effect::after(Duration::from_secs(2), |_| Message::Timeout)}
let (_read, _request, _timeout) = (read_file(), request(), timeout());Every effect form eventually rejoins the application as an accepted message.
Blocking work
- processupdate returns perform
- startsactorBlocking executor
- donedataFile read accepted
- deliversprocessupdate
Asynchronous work
- processupdate returns future
- startsactorAsync executor
- donedataRequest 200 accepted
- deliversprocessupdate
Delayed work
- processupdate returns after
- waitsactorRuntime clock
- firesdataTimeout accepted
- deliversprocessupdate
Use perform only for an operation that may block its calling thread. Wrapping
an async operation in perform occupies a blocking worker unnecessarily;
running blocking work inside future can stall an async executor worker.
What latest does—and does not do
Section titled “What latest does—and does not do”The key identifies one replaceable operation. Replacement cancels an unstarted or pending future/timer through the executor’s cancellation handle. Tokio’s handle also aborts an async task when it is dropped.
Already-running blocking work is different: dropping Tokio’s spawn_blocking
handle cannot stop its closure. Urushi still marks its completion stale, so it
cannot enter the delivery queue if replacement won the race. Any external side
effect performed by that closure can still happen.
An old completion that was accepted before replacement already owns a place
in the runtime-wide queue and is not removed. When that distinction matters,
check domain identity in update before applying the result.
The following complete program starts a future that can never finish, then
replaces it under the same key. The replacement completes with "new"; the
runtime drops the first future’s execution handle and only the new completion
reaches update.
use std::future::pending;use std::time::Duration;
use urushi::{TextStyle, View};use urushi_tui_app::{Application, Effect, Runtime, Subscription};
struct LatestDemo;
enum Message { StartOld, Replace, Finished(&'static str),}
impl Application for LatestDemo { type Model = Vec<&'static str>; type Message = Message;
fn init(&self) -> (Self::Model, Effect<Self::Message>) { ( Vec::new(), Effect::future(async { Message::StartOld }), ) }
fn update( &self, model: &mut Self::Model, message: Self::Message, ) -> Effect<Self::Message> { match message { Message::StartOld => Effect::batch([ Effect::future_latest("search", async { pending::<()>().await; Message::Finished("old") }), Effect::after(Duration::from_millis(10), |_| Message::Replace), ]), Message::Replace => Effect::future_latest("search", async { Message::Finished("new") }), Message::Finished(label) => { model.push(label); Effect::shutdown() } } }
fn view(&self, model: &Self::Model) -> View { View::text(format!("completions: {model:?}"), TextStyle::new()) }
fn subscriptions(&self, _model: &Self::Model) -> Subscription<Self::Message> { Subscription::none() }}
fn main() -> Result<(), urushi_tui_app::Error> { let completions = Runtime::new(LatestDemo).run()?; assert_eq!(completions, ["new"]); println!("completions: {completions:?}"); Ok(())}After the full-screen session is restored, it prints:
completions: ["new"]A newer effect with the same key replaces unfinished work before its completion is accepted.
- stateFuture A pendingStartOld creates work under key “search”.
- 10 msdataReplace arrives10 ms later.
- same keyprocessReplace A with BDrop A’s execution handle and start future B under the same key.
- B donedataFinished("new") is acceptedupdate records it and shuts down; final model is ["new"].
The deliberately pending A demonstrates cancellation before completion. In a
real search, A can instead win the race and be accepted just before Replace;
that accepted result stays queued, so the domain-identity check described above
is still necessary.
after_latest follows the same rule and restarts the wait, making it suitable
for debounce. It is not a repeating timer; use Subscription::interval for
that.
Combine, sequence, map, and shut down effects
Section titled “Combine, sequence, map, and shut down effects”batch is concurrency, not sequencing. Return the next effect from the update
that receives the prior result when order is required.
use urushi_tui_app::Effect;
#[derive(Debug)]enum ChildMessage { Loaded(usize) }
#[derive(Debug)]enum Message { Child(ChildMessage), Saved }
let child: Effect<ChildMessage> = Effect::perform(|| ChildMessage::Loaded(3));let mapped: Effect<Message> = child.map(Message::Child);
let concurrent = Effect::batch([ mapped, Effect::perform(|| Message::Saved),]);
let _ = concurrent;Members are started in array order, but Saved may be delivered before or
after Child(Loaded(3)). map preserves delays, latest keys, batches, and
shutdown while transforming produced messages.
If any member of an effect tree is Effect::shutdown(), the runtime detects it
before starting the tree. No sibling starts. Existing effects are stopped,
their not-yet-accepted completions are suppressed, subscriptions stop, and the
terminal is cleaned up before run returns the final model.
Declare and reconcile subscriptions
Section titled “Declare and reconcile subscriptions”After a non-shutdown update, the runtime evaluates subscriptions(&model) and
reconciles the declaration by source identity:
| Source constructor | Identity |
|---|---|
input |
One runtime-owned input singleton |
surface |
One runtime-owned surface singleton |
terminal_errors |
One runtime-owned terminal-error singleton |
signal(signal, ...) |
The Signal value; each signal is a separate source |
interval(key, period, ...) |
The application key and period |
stream, run, run_blocking and their _with forms |
The application key; constructor kind and admission are not additional identity |
Reconciliation applies these rules:
| Declaration change | Runtime action |
|---|---|
| New identity appears | Start the source |
| Same identity remains | Keep the source running and refresh its message mapper |
| Identity disappears | Stop the source; waiting sends return SendError |
| Same interval key, different period | Restart it, because the period is part of interval identity |
An identity occurs more than once in one batch |
Keep only the last declaration and start or refresh it once |
The mapper closure is not identity. This lets a long-lived source keep running while later values use the mapper from the newest declaration. For a custom source, changing constructor kind or admission while retaining the same key does not restart its producer or inbox; give the changed source a new key (or remove it in an intervening model) when those properties must change.
use std::time::Duration;use urushi_tui_app::{Subscription, Surface};
#[derive(Debug)]enum Message { FastTick, SlowTick, Surface(Surface) }
fn subscriptions(fast: bool) -> Subscription<Message> { let period = if fast { Duration::from_millis(100) } else { Duration::from_secs(1) };
Subscription::batch([ Subscription::surface(Message::Surface), Subscription::interval("refresh", period, move |_| { if fast { Message::FastTick } else { Message::SlowTick } }), ])}
let _fast = subscriptions(true);let _slow = subscriptions(false);The runtime compares each new subscription declaration with the active source identities.
- statefast = trueStart Surface and (refresh, 100 ms).
- same identitiesstatefast = true againKeep both sources and refresh their mappers.
- new identitystatefast = falseKeep Surface; stop the 100 ms interval and start the 1 s interval.
Subscription::batch flattens several declarations. Subscription::map
preserves every source identity and transforms its later messages, which is
useful when embedding a child application.
Build a custom source
Section titled “Build a custom source”Choose the constructor from how values are produced, then choose admission
from what may happen before the runtime accepts them. This complete application
runs both policies. Its bounded producer has capacity one, so later sends wait
for space and all three values arrive in FIFO order. Its latest producer sends
three values without yielding; each successful send replaces the still-pending
slot, so only 3 reaches update.
use urushi::{TextStyle, View};use urushi_tui_app::{ Admission, Application, Effect, Runtime, Subscription,};
struct AdmissionDemo;
enum Message { Bounded(u8), Latest(u8),}
#[derive(Default)]struct Model { bounded: Vec<u8>, latest: Vec<u8>,}
impl Application for AdmissionDemo { type Model = Model; type Message = Message;
fn init(&self) -> (Self::Model, Effect<Self::Message>) { (Model::default(), Effect::none()) }
fn update( &self, model: &mut Self::Model, message: Self::Message, ) -> Effect<Self::Message> { match message { Message::Bounded(value) => model.bounded.push(value), Message::Latest(value) => model.latest.push(value), }
if model.bounded == [1, 2, 3] && model.latest == [3] { Effect::shutdown() } else { Effect::none() } }
fn view(&self, model: &Self::Model) -> View { View::text( format!("bounded={:?} latest={:?}", model.bounded, model.latest), TextStyle::new(), ) }
fn subscriptions(&self, _model: &Self::Model) -> Subscription<Self::Message> { let bounded = Subscription::run_with( "bounded-readings", Admission::bounded(1), |sender| async move { for value in 1..=3 { if sender.send(Message::Bounded(value)).await.is_err() { return; } } }, );
let latest = Subscription::run_with( "latest-reading", Admission::latest(), |sender| async move { for value in 1..=3 { if sender.send(Message::Latest(value)).await.is_err() { return; } } }, );
Subscription::batch([bounded, latest]) }}
fn main() -> Result<(), urushi_tui_app::Error> { let model = Runtime::new(AdmissionDemo).run()?; assert_eq!(model.bounded, [1, 2, 3]); assert_eq!(model.latest, [3]); println!("bounded: {:?}", model.bounded); println!("latest: {:?}", model.latest); Ok(())}bounded: [1, 2, 3]latest: [3]This exact trace uses Runtime’s production default, whose async work is polled
on one runtime thread. The three latest sends are immediately ready and finish
in one producer poll before the acceptance task runs. A custom multi-threaded
executor may accept an earlier value concurrently; latest guarantees only
that the one value still unaccepted is replaced, not that every burst reduces
to one delivered message.
run_blocking_with follows the same policy but runs the producer on a blocking
thread; call blocking_send there instead of awaiting send.
Use stream when a library already gives you a Stream<Item = Message>:
use std::pin::Pin;use std::task::{Context, Poll};use futures_core::Stream;use urushi_tui_app::{Admission, Subscription};
struct OneReading(Option<u64>);
impl Stream for OneReading { type Item = u64;
fn poll_next( mut self: Pin<&mut Self>, _context: &mut Context<'_>, ) -> Poll<Option<Self::Item>> { Poll::Ready(self.0.take()) }}
let source = Subscription::stream_with( "history", Admission::bounded(8), OneReading(Some(1)),);
let _mapped = source.map(|value| format!("reading={value}"));The stream example needs futures-core in the application’s dependencies;
most stream-producing libraries already expose a compatible Stream.
The non-_with forms use Admission::default(), which is
Admission::bounded(64).
Admission and backpressure
Section titled “Admission and backpressure”Admission governs values still waiting in a source-local inbox. A successful
send or blocking_send means only that the value entered that local inbox;
it does not mean the runtime accepted it into the global queue or delivered
it to update. Under latest, a successful value can therefore still be
replaced by the next send. Once the runtime accepts a value, it has a permanent
position in the runtime-wide FIFO and can no longer be replaced.
| Policy | Pending capacity | When full | Suitable for |
|---|---|---|---|
Admission::bounded(n) |
max(n, 1) |
Async send waits; blocking_send blocks |
Logs, commands, or events where every value matters |
Admission::latest() |
1 | New value replaces the one unaccepted value; sending never waits | Progress, telemetry, or snapshots where only the newest pending value matters |
Bounded admission preserves pending values with backpressure; latest admission replaces only the one value that has not been accepted.
bounded(2)
- statepending [1]
- send 2statepending [1, 2]
- send 3statesend 3 waitsThe pending inbox is full.
- accept 1stateglobal FIFO [1]; pending [2, 3]Accepting 1 frees space and resumes send 3.
latest()
- statepending [old]
- send newstatepending [new]The unaccepted old value is replaced.
- acceptstateglobal FIFO [new]
- send nextstatepending [next]The accepted new value is never replaced.
Sender::send and Sender::blocking_send return SendError after the source
is no longer declared or the runtime begins shutdown. Treat that as the normal
stop signal and return from the producer. Stopping also clears values that were
never accepted and releases senders waiting for bounded capacity.
Reference checklist
Section titled “Reference checklist”| Type | Complete public operations covered here |
|---|---|
Effect |
none, shutdown, perform, perform_latest, future, future_latest, after, after_latest, batch, map |
Subscription |
none, input, surface, interval, signal, terminal_errors, stream, stream_with, run, run_with, run_blocking, run_blocking_with, batch, map |
Admission |
bounded, latest, default (bounded(64)) |
Sender |
asynchronous send, blocking blocking_send, Clone |
SendError |
Signals that the runtime stopped taking that source’s messages |
For the order in which accepted messages, updates, reconciliation, and frames interact, continue to Message delivery and drawing.