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.
The guarantees at a glance
Section titled “The guarantees at a glance”| 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 |
Accepted deliveries have one FIFO order
Section titled “Accepted deliveries have one FIFO order”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]);Acceptance fixes delivery order; each update observes the model left by the previous accepted message.
- dataAdd(1)model: [] → [1]
- nextdataAdd(2)model: [1] → [1, 2]
- nextdataAdd(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:
Both scheduler choices are valid. In either case, delivery N still precedes accepted delivery N+1.
Allowed A — draw wins first
- processDraw old model
- thenprocessupdate Add(1)
- dirtystateDrawing dirty
- laterprocessDraw new model
Allowed B — delivery wins first
- processupdate Add(1)
- dirtystateDrawing dirty
- thenprocessDraw 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.
One message lifecycle
Section titled “One message lifecycle”Once the runtime selects an ordinary accepted delivery, it performs these steps in this order:
The runtime finishes the delivery path before a later scheduler decision admits a draw.
- dataEarliest accepted message
- deliverprocessapplication.update(&mut model, message)
- after updatestateDrawing dirty
- start workprocessStart the returned Effect
- reconcileprocessReconcile subscriptions(&model)Only when the application is continuing.
- admit drawprocessapplication.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.
Draws coalesce around the latest model
Section titled “Draws coalesce around the latest model”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.
Models 1 and 2 are applied but never presented; one dirty flag carries the latest model to draw B.
- statet=0 ms · drawing model 0Draw A is admitted; phase is drawing.
- +4 msstatet=4 ms · model 1 dirtyupdate(+1); no second draw starts.
- +3 msstatet=7 ms · model 2 dirtyupdate(+1); still one dirty flag.
- +5 msstatet=12 ms · cooldownDraw A completes; the post-frame cooldown begins.
- +8 msstatet=20 ms · model 3 dirtyupdate(+1); still one dirty flag.
- +25 msstatet=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.
Surface changes form a draw barrier
Section titled “Surface changes form a draw barrier”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();An accepted surface observation fences only later draw admission; it does not cancel an already-admitted frame.
- processOld-size draw admitted
- resizedataResize accepted at FIFO position N
- raise fencestateDraw fence active
- deliverprocessApply the synchronous surface deliveryApply every message in order.
- completeprocessRelease the fence
- admit drawprocessNext 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.
Handle draw failures explicitly
Section titled “Handle draw failures explicitly”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;A failed frame never becomes the committed baseline; the subscribed application receives the error as a message.
- processSubmit frame
- write errorstateFrame uncommittedTerminal drawing failed.
- error eventdataDrawFailed(error) accepted
- deliverprocessupdate records the error
- shutdownprocessShutdown 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.
Shutdown wins before another draw
Section titled “Shutdown wins before another draw”Effect::shutdown() is a control request, not a delivered message. When an
update returns it, the runtime:
- marks the model dirty, then detects shutdown before any new draw is admitted;
- starts no sibling in the same effect tree;
- stops subscriptions and suppresses effect completions not yet accepted;
- shuts down presentation and restores the terminal session; and
- 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.
Configure the runtime boundaries
Section titled “Configure the runtime boundaries”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.
Test deterministically
Section titled “Test deterministically”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.
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 time_and_cancellation_move_only_when_the_test_moves_them ... okUse 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:
[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.
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);}test runtime_draws_delivers_input_and_restores_the_session ... okExpose the two support files to one integration test:
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.
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);}test controlled_boundaries_drive_one_runtime_timeline ... okPass 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.
Reference checklist
Section titled “Reference checklist”- 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,
viewfrom the latest model. - Surface synchronization: latest before acceptance, globally ordered and draw-fenced after acceptance.
- Errors: failed frame is uncommitted;
terminal_errorsconverts 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.