Name the executors and fail a block_on on the UI thread
architecture.md §7.1 stated five executors and their thread counts, and no code named them. dr_ui::executors now does: the Executor enum with each one's thread name and the count §7.1 gives with its reason, and spawn, which starts a thread named <executor>:<role> and marks it with the executor it belongs to. The counts are the stated budget, not yet a bound: a job still gets a thread of its own when it starts. run marks its own thread as the UI executor before it builds the window. net_runtime::build now returns a NetRuntime whose block_on asserts, in debug and test builds, that the caller is not that thread; everything else derefs to the tokio runtime. The login, folder-list and remote-folder workers built the same runtime by hand and now take it from net_runtime, so their block_on is guarded too. Tests: a block_on on a thread marked as the UI executor panics naming the UI thread; the same call on a worker returns; a spawned thread carries its name and executor.
This commit is contained in:
@@ -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 `<executor>:<role>` — `net:sync`,
|
||||
//! `decode:thumbs` — which is what a panic message, `top -H` and a profiler
|
||||
//! show. Before it, most workers were `<unnamed>`, 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<Option<Executor>> = 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<Executor> {
|
||||
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 `<executor>:<role>`.
|
||||
///
|
||||
/// 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<F, T>(executor: Executor, role: &str, f: F) -> JoinHandle<T>
|
||||
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<F, T>(executor: Executor, role: &str, f: F) -> std::io::Result<JoinHandle<T>>
|
||||
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::<String>()
|
||||
.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);
|
||||
}
|
||||
}
|
||||
@@ -491,12 +491,7 @@ fn spawn_login(weak: slint::Weak<AppWindow>, ctl: Rc<LaunchController>, 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<LoginMessage>,
|
||||
server: String,
|
||||
step: &dyn Fn(&str),
|
||||
@@ -782,11 +772,7 @@ fn spawn_folder_list(weak: slint::Weak<AppWindow>, ctl: Rc<LaunchController>, 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;
|
||||
|
||||
@@ -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<dr_gpu::GpuContext> {
|
||||
/// TRACES: M-13 | M-14
|
||||
/// Build and run the viewer.
|
||||
pub fn run(paths: Vec<PathBuf>) -> 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
|
||||
|
||||
@@ -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<tokio::runtime::Runtime> {
|
||||
/// the thread that blocks on it — a thread from `executors::spawn`, normally
|
||||
/// on `Executor::Network`.
|
||||
pub fn build() -> std::io::Result<NetRuntime> {
|
||||
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<F: std::future::Future>(&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
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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;
|
||||
|
||||
Reference in New Issue
Block a user