Skip to content
UrushiDocumentation

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.

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());
Runtime trace

Every effect form eventually rejoins the application as an accepted message.

Blocking work

  1. processupdate returns perform
  2. starts
    actorBlocking executor
  3. done
    dataFile read accepted
  4. delivers
    processupdate

Asynchronous work

  1. processupdate returns future
  2. starts
    actorAsync executor
  3. done
    dataRequest 200 accepted
  4. delivers
    processupdate

Delayed work

  1. processupdate returns after
  2. waits
    actorRuntime clock
  3. fires
    dataTimeout accepted
  4. delivers
    processupdate

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.

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.

src/main.rs
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:

Output
completions: ["new"]
Replacement trace

A newer effect with the same key replaces unfinished work before its completion is accepted.

  1. stateFuture A pendingStartOld creates work under key “search”.
  2. 10 ms
    dataReplace arrives10 ms later.
  3. same key
    processReplace A with BDrop A’s execution handle and start future B under the same key.
  4. B done
    dataFinished("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.

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);
Reconciliation trace

The runtime compares each new subscription declaration with the active source identities.

  1. statefast = trueStart Surface and (refresh, 100 ms).
  2. same identities
    statefast = true againKeep both sources and refresh their mappers.
  3. new identity
    statefast = 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.

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.

src/main.rs
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(())
}
Output after the full-screen session closes
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 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
Admission behavior

Bounded admission preserves pending values with backpressure; latest admission replaces only the one value that has not been accepted.

bounded(2)

  1. statepending [1]
  2. send 2
    statepending [1, 2]
  3. send 3
    statesend 3 waitsThe pending inbox is full.
  4. accept 1
    stateglobal FIFO [1]; pending [2, 3]Accepting 1 frees space and resumes send 3.

latest()

  1. statepending [old]
  2. send new
    statepending [new]The unaccepted old value is replaced.
  3. accept
    stateglobal FIFO [new]
  4. send next
    statepending [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.

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.