diff --git a/ui/dr-ui/src/executors.rs b/ui/dr-ui/src/executors.rs new file mode 100644 index 0000000..fe99364 --- /dev/null +++ b/ui/dr-ui/src/executors.rs @@ -0,0 +1,251 @@ +//! TRACES: NFR-ARCH-1 | R4 +//! The executors: the one place that names them, states their thread counts, +//! and starts the threads that belong to them. +//! +//! `architecture.md §7.1` is the design; this is where the code says the same +//! thing. There are five, and every long-lived thread the app starts belongs +//! to one of them: +//! +//! | Executor | Threads | Work | +//! |---|---|---| +//! | [`Executor::Ui`] | 1 | Slint's event loop. Never blocks. | +//! | [`Executor::GpuSubmit`] | 1 | Command encoding and queue submission | +//! | [`Executor::Decode`] | cores − 2 | RAW decode, previews, and the CPU passes over them | +//! | [`Executor::Io`] | 4 | Catalog, sidecar, cache, filesystem | +//! | [`Executor::Network`] | 2 | Nextcloud and other remote transfer | +//! +//! # What this module does and does not enforce +//! +//! It **names**: [`spawn`] gives every thread `:` — `net:sync`, +//! `decode:thumbs` — which is what a panic message, `top -H` and a profiler +//! show. Before it, most workers were ``, and a panic on one said +//! nothing about which of thirty jobs it was. +//! +//! It **marks**: each thread knows which executor it belongs to +//! ([`current`]), and the UI thread is marked where the event loop starts +//! ([`mark_ui_thread`], called at the top of `run`). [`assert_not_ui`] is the +//! guard every blocking helper calls; see there. +//! +//! It does **not yet bound** the counts. A job still gets a thread of its own +//! when it starts, as it did before this module existed, so the counts in the +//! table are the budget the design states rather than a limit the code holds +//! — [`Executor::threads`] is what a pool will be sized from when one exists. +//! Pooling, and the priority between executors that NFR-ARCH-2 asks for, is a +//! change of behaviour; naming them first is what makes that change visible. +//! +//! # Which executor a job belongs to +//! +//! By what it spends its time waiting on. A worker that builds a network +//! runtime and moves bytes to or from a server is [`Executor::Network`], even +//! when it touches the catalog on the way; one that decodes images or runs a +//! CPU pass over their pixels or embeddings is [`Executor::Decode`]; one that +//! drives the GPU is [`Executor::GpuSubmit`]; the rest — catalog, files, local +//! sockets, relaying a channel to the log — is [`Executor::Io`]. + +use std::cell::Cell; +use std::thread::JoinHandle; + +/// One of the app's executors. See the module documentation for the table. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum Executor { + /// The Slint event loop. One thread, and it never blocks (NFR-P9). + Ui, + /// Command encoding and queue submission. + GpuSubmit, + /// RAW decode, preview extraction, and CPU-bound passes over the result. + Decode, + /// Catalog, sidecars, caches, the filesystem. + Io, + /// Transfer to and from a remote library. + Network, +} + +impl Executor { + /// Every executor, in the order §7.1 lists them. + pub const ALL: [Executor; 5] = [ + Executor::Ui, + Executor::GpuSubmit, + Executor::Decode, + Executor::Io, + Executor::Network, + ]; + + /// The prefix of every thread name on this executor. + /// + /// Short because Linux keeps fifteen bytes of a thread name, and the role + /// after the colon is the part worth reading. + pub const fn name(self) -> &'static str { + match self { + Executor::Ui => "ui", + Executor::GpuSubmit => "gpu", + Executor::Decode => "decode", + Executor::Io => "io", + Executor::Network => "net", + } + } + + /// The thread count §7.1 states for this executor, on this machine. + /// + /// - **UI, 1**: Slint has one event loop, on the thread that created the + /// window. + /// - **GPU submit, 1**: one queue; more submitters only contend for it. + /// - **Decode, cores − 2**: CPU-bound, so as many as there are cores, + /// less one for the UI and one for the GPU submitter, so neither is + /// starved by a sweep. Never fewer than one. + /// - **I/O, 4**: enough to overlap a catalog write with file reads; more + /// only queue on the same disk. + /// - **Network, 2**: one transfer and one metadata request in flight; a + /// WebDAV server serialises most of the rest per account. + pub fn threads(self) -> usize { + match self { + Executor::Ui | Executor::GpuSubmit => 1, + Executor::Decode => std::thread::available_parallelism() + .map_or(1, |n| n.get()) + .saturating_sub(2) + .max(1), + Executor::Io => 4, + Executor::Network => 2, + } + } +} + +thread_local! { + /// Which executor this thread belongs to; `None` for threads this module + /// did not start — the test harness, a library's own workers. + static CURRENT: Cell> = const { Cell::new(None) }; +} + +/// The executor the calling thread belongs to, if it was started here or is +/// the marked UI thread. +pub fn current() -> Option { + CURRENT.with(Cell::get) +} + +/// Mark the calling thread as the UI executor. +/// +/// Called once, at the top of `run`, on the thread that goes on to create the +/// window and enter Slint's event loop — on desktop the main thread, on +/// Android the one `android_main` runs on. +pub fn mark_ui_thread() { + CURRENT.with(|c| c.set(Some(Executor::Ui))); +} + +/// TRACES: NFR-ARCH-1 | NFR-P9 +/// Fail loudly, in debug and test builds, if the calling thread is the UI +/// thread. `what` names the blocking call, for the message. +/// +/// Every blocking helper calls this before it blocks: a `block_on` on the UI +/// thread is a frozen window for as long as the future takes, which for a +/// network call is unbounded. A debug assertion rather than a check in release +/// because the fault is in the code, not the input — a release build that met +/// it would still freeze rather than crash, which is no worse than before, and +/// every test and debug run that reaches the call finds it. +#[track_caller] +pub fn assert_not_ui(what: &str) { + debug_assert!( + current() != Some(Executor::Ui), + "{what} on the UI thread: it would block the event loop (NFR-ARCH-1). \ + Move it onto a worker with executors::spawn." + ); +} + +/// Start a thread on `executor`, named `:`. +/// +/// A drop-in for `std::thread::spawn`, and like it panics if the OS will not +/// give a thread. `role` is what the job is — `scan`, `sync`, `thumbs` — and +/// should be short: the name is cut to fifteen bytes where the OS keeps it. +pub fn spawn(executor: Executor, role: &str, f: F) -> JoinHandle +where + F: FnOnce() -> T + Send + 'static, + T: Send + 'static, +{ + try_spawn(executor, role, f) + .unwrap_or_else(|e| panic!("spawning {}:{role}: {e}", executor.name())) +} + +/// [`spawn`], for a caller that can carry on without the thread. +pub fn try_spawn(executor: Executor, role: &str, f: F) -> std::io::Result> +where + F: FnOnce() -> T + Send + 'static, + T: Send + 'static, +{ + debug_assert!( + executor != Executor::Ui, + "the UI executor is the event loop's thread; it is marked, not spawned" + ); + std::thread::Builder::new() + .name(format!("{}:{role}", executor.name())) + .spawn(move || { + CURRENT.with(|c| c.set(Some(executor))); + f() + }) +} + +#[cfg(test)] +mod tests { + use super::*; + + /// TRACES: NFR-ARCH-1 + /// Every executor has a name and at least one thread, and the fixed + /// counts are the ones §7.1 states. + #[test] + fn the_executors_and_their_counts_are_the_ones_the_design_states() { + let names: Vec<_> = Executor::ALL.iter().map(|e| e.name()).collect(); + assert_eq!(names, ["ui", "gpu", "decode", "io", "net"]); + assert!(Executor::ALL.iter().all(|e| e.threads() >= 1)); + assert_eq!(Executor::Ui.threads(), 1); + assert_eq!(Executor::GpuSubmit.threads(), 1); + assert_eq!(Executor::Io.threads(), 4); + assert_eq!(Executor::Network.threads(), 2); + } + + /// TRACES: NFR-ARCH-1 + /// A spawned thread carries its executor's name and knows which executor + /// it is on; the spawning thread is unaffected. + #[test] + fn a_spawned_thread_is_named_and_marked() { + let (name, on) = spawn(Executor::Network, "sync", || { + (std::thread::current().name().map(str::to_owned), current()) + }) + .join() + .unwrap(); + assert_eq!(name.as_deref(), Some("net:sync")); + assert_eq!(on, Some(Executor::Network)); + assert_ne!(current(), Some(Executor::Network)); + } + + /// TRACES: NFR-ARCH-1 | NFR-P9 + /// A `block_on` on the thread marked as the UI executor panics, and says + /// why. On a thread of its own, because the mark is thread-local and with + /// `--test-threads=1` every test runs on the harness's main thread. + #[cfg(debug_assertions)] + #[test] + fn block_on_the_ui_thread_panics() { + let outcome = std::thread::spawn(|| { + mark_ui_thread(); + let rt = crate::net_runtime::build().expect("runtime"); + rt.block_on(async { 1 }) + }) + .join(); + let payload = outcome.expect_err("a block_on on the UI thread must panic"); + let message = payload + .downcast_ref::() + .cloned() + .or_else(|| payload.downcast_ref::<&str>().map(|s| s.to_string())) + .unwrap_or_default(); + assert!(message.contains("UI thread"), "message: {message}"); + } + + /// TRACES: NFR-ARCH-1 + /// The same `block_on` on a worker is what the network runtime is for. + #[test] + fn block_on_a_worker_is_allowed() { + let answer = spawn(Executor::Network, "test", || { + let rt = crate::net_runtime::build().expect("runtime"); + rt.block_on(async { 42 }) + }) + .join() + .expect("a worker may block"); + assert_eq!(answer, 42); + } +} diff --git a/ui/dr-ui/src/launch_ui.rs b/ui/dr-ui/src/launch_ui.rs index 9c9d88e..52ef41e 100644 --- a/ui/dr-ui/src/launch_ui.rs +++ b/ui/dr-ui/src/launch_ui.rs @@ -491,12 +491,7 @@ fn spawn_login(weak: slint::Weak, ctl: Rc, server: // Only IO and time are enabled; `enable_all()` would also start the // signal driver, which wants process-wide signal handling that an // Android app's runtime already owns. - let rt = match tokio::runtime::Builder::new_multi_thread() - .worker_threads(1) - .enable_io() - .enable_time() - .build() - { + let rt = match crate::net_runtime::build() { Ok(rt) => rt, Err(e) => { let _ = tx.send(LoginMessage::Failed(format!("tokio runtime: {e}"))); @@ -543,12 +538,7 @@ fn spawn_direct_login( let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(move || { // Multi-thread for the reason the browser flow is: a current-thread // runtime left reqwest's future unpolled on Android. - let rt = match tokio::runtime::Builder::new_multi_thread() - .worker_threads(1) - .enable_io() - .enable_time() - .build() - { + let rt = match crate::net_runtime::build() { Ok(rt) => rt, Err(e) => { let _ = tx.send(LoginMessage::Failed(format!("tokio runtime: {e}"))); @@ -604,7 +594,7 @@ fn spawn_direct_login( /// The body of the login flow, split out so the worker above can wrap it in /// `catch_unwind` without a deeply nested closure. fn run_login_flow( - rt: tokio::runtime::Runtime, + rt: crate::net_runtime::NetRuntime, tx: std::sync::mpsc::Sender, server: String, step: &dyn Fn(&str), @@ -782,11 +772,7 @@ fn spawn_folder_list(weak: slint::Weak, ctl: Rc, pa // current-thread runtime left reqwest's connection future unpolled on // Android, so the await never resolved and the thread stopped without // failing or returning. - let rt = tokio::runtime::Builder::new_multi_thread() - .worker_threads(1) - .enable_io() - .enable_time() - .build(); + let rt = crate::net_runtime::build(); let Ok(rt) = rt else { let _ = tx.send(Err("runtime".into())); return; diff --git a/ui/dr-ui/src/lib.rs b/ui/dr-ui/src/lib.rs index 150eb1c..95f4fcf 100644 --- a/ui/dr-ui/src/lib.rs +++ b/ui/dr-ui/src/lib.rs @@ -33,6 +33,7 @@ mod develop_ui; mod display_ui; mod duplicates; mod duplicates_ui; +pub mod executors; mod export; pub mod faces; mod folder_dialog; @@ -1142,6 +1143,10 @@ fn shared_gpu() -> Option { /// TRACES: M-13 | M-14 /// Build and run the viewer. pub fn run(paths: Vec) -> Result<()> { + // This thread creates the window and enters Slint's event loop: it is the + // UI executor, and every blocking helper asserts it is not (NFR-ARCH-1). + executors::mark_ui_thread(); + // Mutable because the browsing list has two sources: the command line at // startup, and whatever the library grid is showing when a cell is // clicked. Opening from the grid replaces this so next/previous walk the diff --git a/ui/dr-ui/src/net_runtime.rs b/ui/dr-ui/src/net_runtime.rs index ac77e81..17664ac 100644 --- a/ui/dr-ui/src/net_runtime.rs +++ b/ui/dr-ui/src/net_runtime.rs @@ -21,14 +21,43 @@ //! that, and a worker thread claiming it is asking for trouble. These flows need //! IO (reqwest) and time (poll intervals, timeouts) and nothing else. +//! +//! **Why a wrapper.** [`NetRuntime::block_on`] is tokio's, behind the guard +//! that fails a debug or test build when it is called on the UI thread +//! (NFR-ARCH-1; see `executors::assert_not_ui`). Everything else derefs to the +//! runtime. + /// Build a runtime suitable for a network worker thread. /// /// Call from the spawned thread, not from the caller: the runtime must live on -/// the thread that blocks on it. -pub fn build() -> std::io::Result { +/// the thread that blocks on it — a thread from `executors::spawn`, normally +/// on `Executor::Network`. +pub fn build() -> std::io::Result { tokio::runtime::Builder::new_multi_thread() .worker_threads(1) .enable_io() .enable_time() .build() + .map(NetRuntime) +} + +/// A network worker's runtime, whose `block_on` refuses the UI thread. +pub struct NetRuntime(tokio::runtime::Runtime); + +impl NetRuntime { + /// TRACES: NFR-ARCH-1 | NFR-P9 + /// Tokio's `block_on`, after asserting the caller is not the UI thread. + #[track_caller] + pub fn block_on(&self, future: F) -> F::Output { + crate::executors::assert_not_ui("block_on"); + self.0.block_on(future) + } +} + +impl std::ops::Deref for NetRuntime { + type Target = tokio::runtime::Runtime; + + fn deref(&self) -> &Self::Target { + &self.0 + } } diff --git a/ui/dr-ui/src/remote_folders.rs b/ui/dr-ui/src/remote_folders.rs index 8395e79..dc3823d 100644 --- a/ui/dr-ui/src/remote_folders.rs +++ b/ui/dr-ui/src/remote_folders.rs @@ -109,11 +109,7 @@ where // Multi-thread, for the reason the login worker records: a // current-thread runtime left reqwest's connection future unpolled on // Android, and the await never resolved. - let rt = tokio::runtime::Builder::new_multi_thread() - .worker_threads(1) - .enable_io() - .enable_time() - .build(); + let rt = crate::net_runtime::build(); let Ok(rt) = rt else { let _ = tx.send(Err("could not start the network runtime".into())); return;