Skip to content

Futures And Timers

ctx.future runs an async operation that automatically re-executes when its dependencies change. The returned FutureHandle exposes a reactive Signal<FutureState<T, E>>.

For requests whose data should be shared across components and retained across navigation, enable the query feature and use Queries. Queries add a shared cache, freshness, request deduplication, and invalidation to finite async reads.

Use ctx.future for finite async work that resolves to one result, such as loading a page, querying an endpoint, or submitting a request. Do not use it to manually chain a continuous subscription by changing a dependency after every completion. For watch receivers, sockets, event feeds, and other multi-item sources, use ctx.stream.

use lurq::{
app::{component::Component, ctx::{Ctx, FutureStatus}},
components::Text,
node::Element,
};
struct UserProfile;
impl Component for UserProfile {
type Props = ();
fn create(_ctx: &mut Ctx) -> Self { Self }
fn render(&self, ctx: &mut Ctx) -> impl Into<Element> {
let handle = ctx.future((), |_| async {
Ok::<_, String>("loaded data".to_owned())
});
let state = handle.state().get();
match state.status {
FutureStatus::Pending => Text::new("Loading..."),
FutureStatus::Fulfilled => Text::new(&state.data.unwrap()),
FutureStatus::Rejected => Text::new(&format!("Error: {}", state.error.unwrap())),
FutureStatus::Idle => Text::new(""),
}
}
}

The first argument is a dependency value. When it changes between renders, the future restarts.

fn render(&self, ctx: &mut Ctx) -> impl Into<Element> {
let page = self.page.get();
let handle = ctx.future(page, |page| async move {
fetch_page(page).await
});
// ...
}

The dependency must implement Clone + PartialEq + Send + Sync + 'static. Use a tuple to combine multiple deps.

ctx.future restarts when the dependency value is different on a later render. If a future result is consumed and that render then changes the dependency, the new future is only created by a following render. That is fine for finite request/retry flows, but it is the wrong shape for continuous streams.

FieldTypeDescription
statusFutureStatusIdle, Pending, Fulfilled, or Rejected.
dataOption<T>Present on Fulfilled; preserved during re-fetch.
errorOption<E>Present on Rejected.

Convenience methods: is_idle(), is_pending(), is_fulfilled(), is_rejected().

MethodDescription
.state()Returns a Signal<FutureState<T, E>> for reactive reads.
.cancel()Cancels the in-flight future.
.is_active()Whether a future is currently running.

ctx.stream runs a continuous async producer. The producer receives a StreamEmitter<T, E> and can call .emit(value) repeatedly. The returned StreamHandle exposes the same reactive Signal<FutureState<T, E>> shape as futures, but the task stays alive after each emitted item.

Use streams for watch::Receiver, websocket subscriptions, event feeds, filesystem watchers, and other sources that can produce more than one value.

The example below uses a Tokio watch channel. Add a direct tokio dependency with the sync feature. Wrap the receiver in props that satisfy PartialEq and DevTools inspection bounds; increment source_id when replacing its channel so the stream restarts.

use lurq::{
app::{
component::Component,
ctx::{Ctx, FutureStatus, StreamEmitter},
},
components::Text,
node::Element,
};
use tokio::sync::watch;
struct ServerEvents;
#[derive(Clone, lurq::DevtoolsInspectable)]
struct ServerEventsProps {
source_id: u64,
#[devtools_ignore]
receiver: watch::Receiver<String>,
}
impl PartialEq for ServerEventsProps {
fn eq(&self, other: &Self) -> bool {
self.source_id == other.source_id
}
}
impl Component for ServerEvents {
type Props = ServerEventsProps;
fn create(_ctx: &mut Ctx) -> Self { Self }
fn render(&self, ctx: &mut Ctx) -> impl Into<Element> {
let props = ctx.props::<Self::Props>().clone();
let receiver = props.receiver;
let handle = ctx.stream(props.source_id, move |_, emitter: StreamEmitter<String, String>| {
let mut receiver = receiver.clone();
async move {
loop {
if receiver.changed().await.is_err() {
break;
}
if !emitter.emit(receiver.borrow().clone()) {
break;
}
}
}
});
let state = handle.state().get();
match state.status {
FutureStatus::Fulfilled => Text::new(&state.data.unwrap()),
FutureStatus::Pending => Text::new("Waiting..."),
FutureStatus::Rejected => Text::new(&format!("Stream error: {}", state.error.unwrap())),
FutureStatus::Idle => Text::new(""),
}
}
}
MethodDescription
.state()Returns a Signal<FutureState<T, E>> containing the latest emitted item or error.
.cancel()Cancels the stream task.
.is_active()Whether the stream task is currently running.
MethodDescription
.emit(value)Publishes a fulfilled item to the UI state. Returns false if the receiver was dropped.
.reject(error)Publishes a rejected state while keeping the stream task alive. Returns false if the receiver was dropped.

The dependency argument works the same way as ctx.future: changing it between renders cancels the existing stream task and starts a new one.

ctx.future_action creates a future that does not run automatically. Instead, you call .run(args) to start it — useful for form submissions, button-triggered requests, or any operation that should not run on every render.

fn render(&self, ctx: &mut Ctx) -> impl Into<Element> {
let action = ctx.future_action(|query: String| async move {
search(query).await
});
let state = action.state().get();
Column::new()
.child(Button::new("Search").on_click({
let action = action.clone();
move |_| action.run("lurq".to_owned())
}))
.child(match state.status {
FutureStatus::Idle => Text::new("Press search"),
FutureStatus::Pending => Text::new("Searching..."),
FutureStatus::Fulfilled => Text::new(&state.data.unwrap()),
FutureStatus::Rejected => Text::new("Failed"),
})
}

FutureAction has the same .state(), .cancel(), and .is_active() methods as FutureHandle, plus .run(args).

A ctx.watch on action.state() may call .run(args) again, for example to retry when the state becomes Rejected. The new run sets the state to Pending after the current notification. See writes from callbacks.

When using the form feature, FormProps::submit_action(action) wires a FutureAction<FormValues, _, FormErrors> into a mounted form. It validates before running the action, exposes form.submitting(), blocks duplicate submits while pending, and maps rejected FormErrors back into field errors.

Futures, streams, and future actions belong to the component that renders them. A task is cancelled when:

  • its dependency changes (the new task replaces it),
  • .cancel() is called,
  • a render no longer reaches its call: the component makes fewer ctx.future, ctx.stream, and ctx.future_action calls than before, or a different one at that position,
  • its component unmounts, which includes mounting another root and dropping the Tree.

A cancelled task’s result is never applied. With a Tokio handle, cancelling aborts the Tokio task: it stops the next time it yields to the runtime, at an .await whose future is not ready yet, and the runtime drops its future, so a stream does not have to reach its next emit to stop. Code between such points keeps running until it yields: CPU-bound work, or a loop whose .awaits are always ready, is not interrupted. .run(args) on a FutureAction whose component has unmounted does nothing.

An offstage component (ctx.mount_offstage, Router::mount_offstage) is still mounted, so its tasks are not cancelled, but its futures are not polled until it is active again. A future or stream polled cooperatively (without a Tokio handle) waits. A task already running on Tokio keeps running, and what it produced while offstage is applied, in order, once the component is active again. Removing an offstage component unmounts it and cancels its tasks.

Enable tokio and configure a live runtime handle to spawn futures and streams on Tokio. Add a direct Tokio dependency when application code uses its APIs:

lurq = { version = "0.41.1", features = ["tokio"] }
tokio = { version = "1", features = ["rt-multi-thread", "sync", "time", "net"] }

Pass a tokio handle when creating the App:

let tokio_rt = tokio::runtime::Runtime::new().unwrap();
let app = App::new().with_tokio_handle(tokio_rt.handle().clone());

With a tokio handle, futures and streams spawn onto the tokio runtime and complete independently. Results are delivered back to the UI thread on the next tree.tick_futures() call.

Keep tokio_rt alive for as long as the UI uses its tasks. Without a configured Tokio handle, futures and streams are polled cooperatively during tree.tick_futures(), even if the Cargo feature is enabled. Tokio I/O and timers require an appropriately configured runtime.

A one-shot timer that fires once after a duration.

use std::time::Duration;
fn create(ctx: &mut Ctx) -> Self {
let count = ctx.signal(0);
let timeout = ctx.create_timeout(Duration::from_secs(2), {
let count = count.clone();
move || count.update(|n| *n += 1)
});
timeout.start();
Self { count, _timeout: timeout }
}
MethodDescription
.start()Arm the timer. If already armed, does nothing.
.restart()Reset the deadline to now + duration.
.cancel()Stop the timer without firing.
.is_active()Whether the timer is armed.

A timeout fires at most once. After firing it becomes inactive.

A repeating timer that fires every duration.

let interval = ctx.create_interval(Duration::from_millis(500), {
let tick = tick.clone();
move || tick.update(|n| *n += 1)
});
interval.start();
MethodDescription
.start()Arm the interval.
.restart()Reset the next fire to now + duration.
.stop()Stop repeating.
.is_active()Whether the interval is running.

Create timers in Component::create and store them in the component struct. The runtime ticks timers each frame via tree.tick_timers(). When a timer fires it invokes its callback, which typically updates a signal, causing a re-render.