Skip to content
UrushiDocumentation

Message delivery and drawing

Urushi serializes accepted deliveries through one runtime-wide FIFO, but drawing is deliberately coalesced. Accepted deliveries keep their order relative to other deliveries. A draw that is ready at the same time as an ordinary asynchronous delivery has no public relative-order guarantee. Only an accepted surface delivery fences draw admission.

Event Guaranteed behavior
A delivery is accepted It receives one permanent position relative to other accepted deliveries
An ordinary delivery and a draw are both ready Either may be selected first; do not use their relative order as application logic
update returns normally The model is changed, drawing is marked dirty, the returned effect starts, then subscriptions are reconciled
Several updates arrive during a draw or cooldown All updates run; their invalidations coalesce into one later draw from the latest model
A surface observation is accepted It fences new draws until its entire synchronous delivery is applied
A frame fails It is not committed; the run ends unless terminal_errors is subscribed
update returns Effect::shutdown() No later view is admitted; runtime work stops and presentation cleanup completes before return

Effect completions, custom subscription values, built-in source messages, and handled terminal errors all enter the same queue. Acceptance—moving a value into that queue—is the ordering point. The queue does not inspect the application’s message enum.

use urushi_tui_app::Effect;
#[derive(Debug, Clone, Copy)]
enum Message { Add(i32) }
fn update(model: &mut Vec<i32>, message: Message) -> Effect<Message> {
let Message::Add(value) = message;
model.push(value);
Effect::none()
}
let mut model = Vec::new();
for accepted in [Message::Add(1), Message::Add(2), Message::Add(3)] {
let _ = update(&mut model, accepted);
}
assert_eq!(model, [1, 2, 3]);
Accepted-order trace

Acceptance fixes delivery order; each update observes the model left by the previous accepted message.

  1. dataAdd(1)model: [] → [1]
  2. next
    dataAdd(2)model: [1] → [1, 2]
  3. next
    dataAdd(3)model: [1, 2] → [1, 2, 3]

This does not promise which concurrent producer finishes first. Effect completion order and cross-source production order depend on the executor and external events. Once each value is accepted, however, deliveries follow FIFO insertion order relative to one another. Admission::latest() can replace only a source value that has not yet been accepted.

Draw admission is a separate runtime choice, not an item in that message FIFO. When an ordinary asynchronous delivery and a ready draw can both proceed, the public contract does not specify which one wins:

Ordinary delivery versus a ready draw

Both scheduler choices are valid. In either case, delivery N still precedes accepted delivery N+1.

Allowed A — draw wins first

  1. processDraw old model
  2. then
    processupdate Add(1)
  3. dirty
    stateDrawing dirty
  4. later
    processDraw new model

Allowed B — delivery wins first

  1. processupdate Add(1)
  2. dirty
    stateDrawing dirty
  3. then
    processDraw new model

Never use a frame as an acknowledgement that the preceding ordinary message has run. Use another message or an effect completion for application sequencing. Surface delivery is the one exception on the draw boundary: its accepted synchronous delivery raises a fence that prevents a new draw until the complete surface delivery has been applied.

Once the runtime selects an ordinary accepted delivery, it performs these steps in this order:

One accepted message lifecycle

The runtime finishes the delivery path before a later scheduler decision admits a draw.

  1. dataEarliest accepted message
  2. deliver
    processapplication.update(&mut model, message)
  3. after update
    stateDrawing dirty
  4. start work
    processStart the returned Effect
  5. reconcile
    processReconcile subscriptions(&model)Only when the application is continuing.
  6. admit draw
    processapplication.view(&model)Submit the frame when drawing is later admitted.

Effect start therefore precedes subscription reconciliation. Do not use an effect and a subscription as an implicit sequencing mechanism. If one action must follow another, make the first completion a message and return the second effect from that message’s update branch.

view is not called once per update. It is called only when the scheduler admits a draw, and it sees the latest model at that moment. A draw already ready for admission may be selected before an ordinary delivery; the lifecycle above starts after the delivery itself is selected.

The initial model invalidates drawing immediately. Each later update also sets one dirty flag. While a frame is in progress, or while the scheduler is in its post-frame cooldown, repeated invalidations do not queue repeated frames.

The built-in minimum interval is 1 / 30 s (approximately 33.33 ms) measured from draw completion. There is currently no public frame-rate builder setting.

Coalescing timeline

Models 1 and 2 are applied but never presented; one dirty flag carries the latest model to draw B.

  1. statet=0 ms · drawing model 0Draw A is admitted; phase is drawing.
  2. +4 ms
    statet=4 ms · model 1 dirtyupdate(+1); no second draw starts.
  3. +3 ms
    statet=7 ms · model 2 dirtyupdate(+1); still one dirty flag.
  4. +5 ms
    statet=12 ms · cooldownDraw A completes; the post-frame cooldown begins.
  5. +8 ms
    statet=20 ms · model 3 dirtyupdate(+1); still one dirty flag.
  6. +25 ms
    statet=45 ms · drawing model 3The cooldown expires; draw B is admitted. Presented frames: 0 and 3.

This keeps message handling responsive and prevents a slow terminal from building an unbounded frame backlog. Put durable work in update or an effect, not in an assumption that every intermediate model will be presented.

Subscribe with Subscription::surface when layout depends on terminal size or cell-pixel information. Surface observations are latest-only while waiting for acceptance. Once accepted, the runtime gives the observation a global queue position and treats its messages as one synchronous delivery.

use urushi_tui_app::{Effect, Subscription, Surface};
#[derive(Debug)]
enum Message { SurfaceChanged(Surface) }
fn subscriptions() -> Subscription<Message> {
Subscription::surface(Message::SurfaceChanged)
}
fn update(surface: &mut Option<Surface>, message: Message) -> Effect<Message> {
let Message::SurfaceChanged(new_surface) = message;
*surface = Some(new_surface);
Effect::none()
}
let _declared = subscriptions();
Resize barrier timeline

An accepted surface observation fences only later draw admission; it does not cancel an already-admitted frame.

  1. processOld-size draw admitted
  2. resize
    dataResize accepted at FIFO position N
  3. raise fence
    stateDraw fence active
  4. deliver
    processApply the synchronous surface deliveryApply every message in order.
  5. complete
    processRelease the fence
  6. admit draw
    processNext view uses the resized surfaceThe updated model and presentation geometry now agree.

A resize accepted after a draw was already admitted follows that draw and fences the next one. The barrier prevents a new frame from mixing updated surface state with stale presentation geometry; it does not retroactively cancel an already-admitted frame.

A terminal frame is transactional: a failed frame is not committed. Without a Subscription::terminal_errors declaration, the runtime stops and run returns a terminal error. With the subscription, the failure becomes an ordinary message and the application decides whether to continue or shut down.

use std::io;
use urushi_tui_app::{Effect, Subscription};
#[derive(Debug)]
enum Message { DrawFailed(io::Error) }
#[derive(Default)]
struct Model { last_error: Option<String> }
fn subscriptions() -> Subscription<Message> {
Subscription::terminal_errors(Message::DrawFailed)
}
fn update(model: &mut Model, message: Message) -> Effect<Message> {
match message {
Message::DrawFailed(error) => {
model.last_error = Some(error.to_string());
Effect::shutdown()
}
}
}
let mut model = Model::default();
let effect = update(
&mut model,
Message::DrawFailed(io::Error::other("terminal disconnected")),
);
assert_eq!(model.last_error.as_deref(), Some("terminal disconnected"));
let _shutdown = effect;
Failure state transition

A failed frame never becomes the committed baseline; the subscribed application receives the error as a message.

  1. processSubmit frame
  2. write error
    stateFrame uncommittedTerminal drawing failed.
  3. error event
    dataDrawFailed(error) accepted
  4. deliver
    processupdate records the error
  5. shutdown
    processShutdown and terminal cleanuprun returns the final model.

Failure of the runtime’s presentation worker itself is not a recoverable draw message; it ends the run as a runtime error.

Effect::shutdown() is a control request, not a delivered message. When an update returns it, the runtime:

  1. marks the model dirty, then detects shutdown before any new draw is admitted;
  2. starts no sibling in the same effect tree;
  3. stops subscriptions and suppresses effect completions not yet accepted;
  4. shuts down presentation and restores the terminal session; and
  5. returns the final model.

An effect completion accepted before shutdown keeps its earlier queue position, but once the shutdown-producing message is processed the event loop stops; it does not continue draining later queued messages. Work that must complete before exit should finish as an ordinary effect and return shutdown from the update that receives its completion.

Runtime::new(application) uses the production executor, clock, and physical terminal selected by the default crossterm feature. Cell-only builds use the portable Crossterm backend; a graphics-enabled Unix build uses the native bidirectional connection so it can query protocol capabilities. Replace boundaries before calling run:

use urushi_tui_app::{Application, Clock, Executor, Runtime};
pub fn configured_runtime<A, E, C>(
application: A,
executor: E,
clock: C,
) -> Runtime<A>
where
A: Application,
E: Executor,
C: Clock,
{
Runtime::new(application)
.executor(executor)
.clock(clock)
}
Boundary Builder Runtime responsibility affected
Executor Runtime::executor Async effects, blocking effects, and custom async/blocking sources; dropping an Execution cancels unfinished work
Clock Runtime::clock Effect::after, interval subscriptions, and the draw cooldown
TerminalBackend Runtime::backend Physical session commands, input/events, terminal queries, output, and terminal size

The executor and clock must agree operationally: a clock sleep returns a future that the configured executor can poll. A custom backend must implement the urushi_terminal::TerminalBackend contract and is required when urushi-tui-app is built without its default crossterm feature.

Keep most tests at the application boundary: call init, deliver messages to update, and inspect the model and view without a terminal. This verifies all domain transitions and which effect/subscription values a state declares.

use urushi::{TextStyle, View};
use urushi_tui_app::{Application, Effect, Subscription};
struct Counter;
enum Message { Increment, Saved }
impl Application for Counter {
type Model = i32;
type Message = Message;
fn init(&self) -> (i32, Effect<Message>) {
(0, Effect::none())
}
fn update(&self, model: &mut i32, message: Message) -> Effect<Message> {
match message {
Message::Increment => {
*model += 1;
Effect::perform(|| Message::Saved)
}
Message::Saved => Effect::none(),
}
}
fn view(&self, model: &i32) -> View {
View::text(model.to_string(), TextStyle::new())
}
fn subscriptions(&self, _model: &i32) -> Subscription<Message> {
Subscription::none()
}
}
let app = Counter;
let (mut model, _initial) = app.init();
let _save = app.update(&mut model, Message::Increment);
assert_eq!(model, 1);
let _ = app.update(&mut model, Message::Saved);
assert_eq!(app.view(&model), View::text("1", TextStyle::new()));

For a runtime-level test, inject three controlled boundaries together:

Test control What to record or advance Assertions it enables
Manual Executor Spawn order, completion order, dropped Execution handles latest replacement, shutdown cancellation, concurrent batch completion
Manual Clock Current Instant; explicitly release sleeps delayed effects, intervals, and the 30 fps cooldown without wall-clock sleeps
Recording TerminalBackend Events supplied and commands/output written surface barriers, committed frames, draw errors, and session cleanup

The following integration test is a complete, executable implementation of the first two controls. The executor polls one selected job only when the test calls run_one; dropping its Execution handle cancels an unfinished job. The clock moves only when the test calls advance.

tests/support/manual.rs
use std::collections::VecDeque;
use std::future::Future;
use std::pin::Pin;
use std::sync::atomic::{AtomicU8, Ordering};
use std::sync::{Arc, Condvar, Mutex};
use std::task::{Context, Poll, Waker};
use std::time::{Duration, Instant};
use urushi_tui_app::{BlockingTask, Clock, Execution, Executor, Task};
const WAITING: u8 = 0;
const COMPLETED: u8 = 1;
const CANCELED: u8 = 2;
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum Step {
Completed,
Pending,
Canceled,
Idle,
}
#[derive(Clone, Default)]
pub struct ManualExecutor {
jobs: Arc<Mutex<VecDeque<Job>>>,
}
struct Job {
work: Work,
state: Arc<AtomicU8>,
}
enum Work {
Async(Task),
Blocking(Option<BlockingTask>),
}
struct ManualExecution {
state: Arc<AtomicU8>,
}
impl Execution for ManualExecution {}
impl Drop for ManualExecution {
fn drop(&mut self) {
let _ = self.state.compare_exchange(
WAITING,
CANCELED,
Ordering::AcqRel,
Ordering::Acquire,
);
}
}
impl ManualExecutor {
fn schedule(&self, work: Work) -> Box<dyn Execution> {
let state = Arc::new(AtomicU8::new(WAITING));
self.jobs
.lock()
.expect("job queue lock is healthy")
.push_back(Job {
work,
state: Arc::clone(&state),
});
Box::new(ManualExecution { state })
}
pub fn run_one(&self) -> Step {
let Some(mut job) = self
.jobs
.lock()
.expect("job queue lock is healthy")
.pop_front()
else {
return Step::Idle;
};
if job.state.load(Ordering::Acquire) == CANCELED {
return Step::Canceled;
}
let step = match &mut job.work {
Work::Async(task) => {
let mut context = Context::from_waker(Waker::noop());
match task.as_mut().poll(&mut context) {
Poll::Ready(()) => Step::Completed,
Poll::Pending => Step::Pending,
}
}
Work::Blocking(task) => {
task.take().expect("a blocking job runs once")();
Step::Completed
}
};
if step == Step::Pending {
self.jobs
.lock()
.expect("job queue lock is healthy")
.push_back(job);
} else {
job.state.store(COMPLETED, Ordering::Release);
}
step
}
}
impl Executor for ManualExecutor {
fn spawn(&self, task: Task) -> Box<dyn Execution> {
self.schedule(Work::Async(task))
}
fn spawn_blocking(&self, task: BlockingTask) -> Box<dyn Execution> {
self.schedule(Work::Blocking(Some(task)))
}
}
#[derive(Clone)]
pub struct ManualClock {
shared: Arc<ClockShared>,
}
struct ClockShared {
state: Mutex<ClockState>,
registered: Condvar,
}
struct ClockState {
now: Instant,
sleepers: Vec<Waker>,
registrations: usize,
}
impl ManualClock {
pub fn new() -> Self {
Self {
shared: Arc::new(ClockShared {
state: Mutex::new(ClockState {
now: Instant::now(),
sleepers: Vec::new(),
registrations: 0,
}),
registered: Condvar::new(),
}),
}
}
pub fn advance(&self, duration: Duration) {
let sleepers = {
let mut state = self.shared.state.lock().expect("clock lock is healthy");
state.now += duration;
std::mem::take(&mut state.sleepers)
};
sleepers.into_iter().for_each(Waker::wake);
}
pub fn wait_for_sleepers(&self, count: usize) {
let mut state = self.shared.state.lock().expect("clock lock is healthy");
while state.registrations < count {
state = self
.shared
.registered
.wait(state)
.expect("clock lock is healthy");
}
}
}
impl Clock for ManualClock {
fn now(&self) -> Instant {
self.shared.state.lock().expect("clock lock is healthy").now
}
fn sleep(
&self,
duration: Duration,
) -> Pin<Box<dyn Future<Output = Instant> + Send + 'static>> {
let deadline = self.now() + duration;
Box::pin(Sleep {
shared: Arc::clone(&self.shared),
deadline,
registered: false,
})
}
}
struct Sleep {
shared: Arc<ClockShared>,
deadline: Instant,
registered: bool,
}
impl Future for Sleep {
type Output = Instant;
fn poll(self: Pin<&mut Self>, context: &mut Context<'_>) -> Poll<Instant> {
let this = self.get_mut();
let mut state = this.shared.state.lock().expect("clock lock is healthy");
if state.now >= this.deadline {
Poll::Ready(this.deadline)
} else {
if !state
.sleepers
.iter()
.any(|waker| waker.will_wake(context.waker()))
{
state.sleepers.push(context.waker().clone());
}
if !this.registered {
state.registrations += 1;
this.registered = true;
this.shared.registered.notify_all();
}
Poll::Pending
}
}
}
#[test]
fn time_and_cancellation_move_only_when_the_test_moves_them() {
let executor = ManualExecutor::default();
let clock = ManualClock::new();
let fired_at = Arc::new(Mutex::new(None));
let recorded = Arc::clone(&fired_at);
let task_clock = clock.clone();
let _timer = executor.spawn(Box::pin(async move {
let at = task_clock.sleep(Duration::from_millis(33)).await;
*recorded.lock().expect("result lock is healthy") = Some(at);
}));
assert_eq!(executor.run_one(), Step::Pending);
assert_eq!(*fired_at.lock().expect("result lock is healthy"), None);
clock.advance(Duration::from_millis(33));
assert_eq!(executor.run_one(), Step::Completed);
assert_eq!(*fired_at.lock().expect("result lock is healthy"), Some(clock.now()));
let canceled = executor.spawn(Box::pin(std::future::pending()));
drop(canceled);
assert_eq!(executor.run_one(), Step::Canceled);
assert_eq!(executor.run_one(), Step::Idle);
}
Test result
test time_and_cancellation_move_only_when_the_test_moves_them ... ok

Use clones of these values with Runtime::executor and Runtime::clock. A runtime-level test also needs a backend. TerminalBackend is implemented automatically by a type that implements its five direct public supertraits. CommandWriter itself also requires TerminalOutput:

Required trait Recording backend responsibility
TerminalOutput record flush boundaries
CommandWriter record each Command, including printable cells and session modes
EventSource return test-owned Event values and block or time out when none is available
RawModeControl record raw-mode state and whether the backend is interactive
TerminalQuery return fixed size, cursor, raw-mode, and capability observations
KeyboardEnhancementQuery return whether enhanced-keyboard negotiation is supported

The backend is a separate crate boundary, so add it directly for an integration test:

Cargo.toml
[dev-dependencies]
urushi = "0.1.0"
urushi-terminal = "0.1.0"
urushi-tui-app = "0.1.0"

This complete recording backend waits on a condition variable instead of busy polling. Its first rendered ready cell releases one queued q event. The application therefore cannot quit before a frame reaches the backend, and the test can observe drawing, input delivery, shutdown, and restoration end to end.

tests/support/backend.rs
use std::collections::VecDeque;
use std::io;
use std::sync::{Arc, Condvar, Mutex};
use std::time::Duration;
use urushi::{TextStyle, View};
use urushi_terminal::{
Command, CommandWriter, Event, EventSource, KeyCode, KeyEvent,
KeyboardEnhancementQuery, Position, RawModeControl, TerminalOutput,
TerminalQuery, TerminalSize, WindowSize,
};
use urushi_tui_app::{Application, Effect, Input, Runtime, Subscription};
#[derive(Clone, Debug)]
pub struct Snapshot {
pub printed: String,
pub alternate_screen: Vec<bool>,
pub cursor_visible: Vec<bool>,
pub raw_mode: bool,
pub flushes: usize,
}
struct State {
events: VecDeque<Event>,
printed: String,
alternate_screen: Vec<bool>,
cursor_visible: Vec<bool>,
raw_mode: bool,
flushes: usize,
quit_released: bool,
}
struct Shared {
state: Mutex<State>,
event_ready: Condvar,
}
pub struct RecordingBackend {
shared: Arc<Shared>,
}
#[derive(Clone)]
pub struct RecordingHandle {
shared: Arc<Shared>,
}
impl RecordingBackend {
pub fn new() -> (Self, RecordingHandle) {
let shared = Arc::new(Shared {
state: Mutex::new(State {
events: VecDeque::new(),
printed: String::new(),
alternate_screen: Vec::new(),
cursor_visible: Vec::new(),
raw_mode: false,
flushes: 0,
quit_released: false,
}),
event_ready: Condvar::new(),
});
(
Self {
shared: Arc::clone(&shared),
},
RecordingHandle { shared },
)
}
}
impl RecordingHandle {
pub fn snapshot(&self) -> Snapshot {
let state = self.shared.state.lock().expect("backend lock is healthy");
Snapshot {
printed: state.printed.clone(),
alternate_screen: state.alternate_screen.clone(),
cursor_visible: state.cursor_visible.clone(),
raw_mode: state.raw_mode,
flushes: state.flushes,
}
}
pub fn wait_for_text(&self, expected: &str) {
let mut state = self.shared.state.lock().expect("backend lock is healthy");
while !state.printed.contains(expected) {
state = self
.shared
.event_ready
.wait(state)
.expect("backend lock is healthy");
}
}
}
impl TerminalOutput for RecordingBackend {
fn flush(&mut self) -> io::Result<()> {
self.shared
.state
.lock()
.expect("backend lock is healthy")
.flushes += 1;
Ok(())
}
}
impl CommandWriter for RecordingBackend {
fn write_command(&mut self, command: Command<'_>) -> io::Result<()> {
let mut state = self.shared.state.lock().expect("backend lock is healthy");
match command {
Command::Print(text) => {
state.printed.push_str(text.as_str());
if state.printed.contains("ready") && !state.quit_released {
state.quit_released = true;
state
.events
.push_back(Event::Key(KeyEvent::new(KeyCode::Char('q'))));
self.shared.event_ready.notify_all();
}
self.shared.event_ready.notify_all();
}
Command::SetAlternateScreen(enabled) => {
state.alternate_screen.push(enabled);
}
Command::SetCursorVisible(visible) => {
state.cursor_visible.push(visible);
}
_ => {}
}
Ok(())
}
}
impl EventSource for RecordingBackend {
fn read_event(&mut self) -> io::Result<Event> {
let mut state = self.shared.state.lock().expect("backend lock is healthy");
loop {
if let Some(event) = state.events.pop_front() {
return Ok(event);
}
state = self
.shared
.event_ready
.wait(state)
.expect("backend lock is healthy");
}
}
fn poll_event(&mut self) -> io::Result<Option<Event>> {
Ok(self
.shared
.state
.lock()
.expect("backend lock is healthy")
.events
.pop_front())
}
fn poll_event_timeout(&mut self, timeout: Duration) -> io::Result<Option<Event>> {
let state = self.shared.state.lock().expect("backend lock is healthy");
let (mut state, _) = self
.shared
.event_ready
.wait_timeout_while(state, timeout, |state| state.events.is_empty())
.expect("backend lock is healthy");
Ok(state.events.pop_front())
}
}
impl RawModeControl for RecordingBackend {
fn is_interactive(&self) -> bool {
true
}
fn enable_raw_mode(&mut self) -> io::Result<()> {
self.shared
.state
.lock()
.expect("backend lock is healthy")
.raw_mode = true;
Ok(())
}
fn disable_raw_mode(&mut self) -> io::Result<()> {
self.shared
.state
.lock()
.expect("backend lock is healthy")
.raw_mode = false;
Ok(())
}
}
impl TerminalQuery for RecordingBackend {
fn terminal_size(&mut self) -> io::Result<TerminalSize> {
Ok(TerminalSize::new(20, 4))
}
fn cursor_position(&mut self) -> io::Result<Position> {
Ok(Position::new(0, 0))
}
fn window_size(&mut self) -> io::Result<WindowSize> {
Ok(WindowSize::new(TerminalSize::new(20, 4), None))
}
fn raw_mode_enabled(&mut self) -> io::Result<bool> {
Ok(self
.shared
.state
.lock()
.expect("backend lock is healthy")
.raw_mode)
}
}
impl KeyboardEnhancementQuery for RecordingBackend {
fn supports_keyboard_enhancement(&mut self) -> io::Result<bool> {
Ok(false)
}
}
struct QuitAfterFirstInput;
impl Application for QuitAfterFirstInput {
type Model = usize;
type Message = Input;
fn init(&self) -> (usize, Effect<Input>) {
(0, Effect::none())
}
fn update(&self, model: &mut usize, _message: Input) -> Effect<Input> {
*model += 1;
Effect::shutdown()
}
fn view(&self, _model: &usize) -> View {
View::text("ready", TextStyle::new())
}
fn subscriptions(&self, _model: &usize) -> Subscription<Input> {
Subscription::input(|input| input)
}
}
#[test]
fn runtime_draws_delivers_input_and_restores_the_session() {
let (backend, recording) = RecordingBackend::new();
let final_model = Runtime::new(QuitAfterFirstInput)
.backend(backend)
.run()
.expect("the in-memory session succeeds");
let snapshot = recording.snapshot();
assert_eq!(final_model, 1);
assert!(snapshot.printed.contains("ready"));
assert_eq!(snapshot.alternate_screen, [true, false]);
assert_eq!(snapshot.cursor_visible.first(), Some(&false));
assert_eq!(snapshot.cursor_visible.last(), Some(&true));
assert!(!snapshot.raw_mode);
assert!(snapshot.flushes >= 2);
}
End-to-end test result
test runtime_draws_delivers_input_and_restores_the_session ... ok

Expose the two support files to one integration test:

tests/support/mod.rs
pub mod backend;
pub mod manual;

The final test uses all three boundaries in the same Runtime. run is blocking, so a worker thread owns it. The recording backend is the first synchronization point: rendering ready-0 releases the input event. The test then polls the runtime’s queued source and timer work itself, waits until both the draw cooldown and delayed effect are sleeping, advances virtual time, and requires the updated frame before completing the shutdown effect.

tests/runtime_boundaries.rs
mod support;
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use std::time::Duration;
use support::backend::RecordingBackend;
use support::manual::{ManualClock, ManualExecutor, Step};
use urushi::{TextStyle, View};
use urushi_tui_app::{Application, Effect, Input, Runtime, Subscription};
struct TimedQuit;
enum Message {
Input(Input),
Quit,
}
impl Application for TimedQuit {
type Model = usize;
type Message = Message;
fn init(&self) -> (usize, Effect<Message>) {
(0, Effect::none())
}
fn update(&self, model: &mut usize, message: Message) -> Effect<Message> {
match message {
Message::Input(input) => {
let _ = input;
*model = 1;
Effect::after(Duration::from_millis(100), |_| Message::Quit)
}
Message::Quit => Effect::shutdown(),
}
}
fn view(&self, model: &usize) -> View {
View::text(format!("ready-{model}"), TextStyle::new())
}
fn subscriptions(&self, _model: &usize) -> Subscription<Message> {
Subscription::input(Message::Input)
}
}
#[test]
fn controlled_boundaries_drive_one_runtime_timeline() {
let executor = ManualExecutor::default();
let clock = ManualClock::new();
let (backend, recording) = RecordingBackend::new();
let runtime_executor = executor.clone();
let runtime_clock = clock.clone();
let worker = std::thread::spawn(move || {
Runtime::new(TimedQuit)
.executor(runtime_executor)
.clock(runtime_clock)
.backend(backend)
.run()
});
// The backend releases its q event only after this frame is written.
recording.wait_for_text("ready-0");
// Drive source acceptance and effect scheduling beside the blocking
// runtime. Stop the driver before advancing time so the delayed Quit
// cannot overtake the updated frame assertion below.
let driving = Arc::new(AtomicBool::new(true));
let driver_flag = Arc::clone(&driving);
let driver_executor = executor.clone();
let driver = std::thread::spawn(move || {
while driver_flag.load(Ordering::Acquire) {
let _ = driver_executor.run_one();
std::thread::yield_now();
}
});
// Wait until both the delayed effect and frame cooldown have been polled
// and registered their wakers. Advancing before this point would lose a
// wake and make the test itself racy.
clock.wait_for_sleepers(2);
driving.store(false, Ordering::Release);
driver.join().expect("executor driver does not panic");
// Releasing the cooldown permits a frame from model 1. Do not complete the
// delayed Quit effect until that updated frame is observable.
// The draw cooldown is 1/30 s (about 33.33 ms), so advance past it
// without completing the 100 ms shutdown timer.
clock.advance(Duration::from_millis(34));
// The terminal output is a cell diff, so only the changed final cell `1`
// follows the complete initial `ready-0` frame in the command log.
recording.wait_for_text("ready-01");
clock.advance(Duration::from_millis(66));
for _ in 0..1_000 {
if executor.run_one() == Step::Completed {
break;
}
std::thread::yield_now();
}
let final_model = worker
.join()
.expect("runtime thread does not panic")
.expect("recorded runtime succeeds");
let snapshot = recording.snapshot();
assert_eq!(final_model, 1);
assert!(snapshot.printed.contains("ready-0"));
assert!(snapshot.printed.contains("ready-01"));
assert_eq!(snapshot.alternate_screen, [true, false]);
assert!(!snapshot.raw_mode);
}
Coordinated runtime test result
test controlled_boundaries_drive_one_runtime_timeline ... ok

Pass a recording backend through Runtime::backend and keep its observations in a shared handle because run takes ownership of the backend. The condition variable matters: a backend that returns None immediately from every timed event poll creates a busy loop, while one that never wakes cannot deliver the event that lets the runtime stop.

Advance one controlled boundary at a time and assert the observable timeline. For example: accept two messages, complete the current draw, advance the clock by 1 / 30 s, then verify that exactly one new frame contains the latest model. Avoid real sleeps: they make effect, admission, and draw races nondeterministic rather than testing the ordering contract.

  • Delivery order: accepted deliveries are FIFO relative to one another; concurrent production order and the order versus a ready ordinary draw are not specified.
  • Update order: mutate model, invalidate drawing, start effect, then reconcile subscriptions when continuing.
  • Drawing: initial dirty state, one draw at a time, invalidation coalescing, fixed approximately 30 fps cooldown, view from the latest model.
  • Surface synchronization: latest before acceptance, globally ordered and draw-fenced after acceptance.
  • Errors: failed frame is uncommitted; terminal_errors converts it to a message; presentation-worker failure ends the runtime.
  • Shutdown: no later view, stop effects and subscriptions, await presentation cleanup, return final model.

For effect replacement, subscription reconciliation, and source admission, see Effects and subscriptions.