Thirty-nine spawn sites in dr-ui, and one in the Android entry point, called std::thread::spawn or a Builder of their own, and most of the threads they started were <unnamed> in a panic message or a profiler. Each now calls executors::spawn with its executor and a role, so the thread is named <executor>:<role> — net:sync, decode:thumbs, io:catalog-open — and knows which executor it is on. The three that already set a name (automation, import, prefetch) keep their name as the role. Behaviour is unchanged: each job still gets a thread of its own when it starts, and spawn panics where std::thread::spawn did. The module's documentation now says how a job is assigned: by what it spends its time on, so a sweep that fetches bytes and then decodes them is Decode, and a sidecar write that touches the catalog is Network. Left as they were: the segmentation and refine workers in masks_ui.rs, which another change is reworking, and test-only threads.
255 lines
10 KiB
Rust
255 lines
10 KiB
Rust
//! 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 on. One that decodes images or runs a CPU pass
|
||
//! over their pixels or embeddings — the thumbnail and metadata sweeps, face
|
||
//! indexing, grouping, repairs — is [`Executor::Decode`], even when it fetches
|
||
//! the bytes first. One that drives the GPU — batch export, a merge — is
|
||
//! [`Executor::GpuSubmit`]. One whose work is moving files or metadata to or
|
||
//! from a server — sync, sidecars, fetches, trash, login — is
|
||
//! [`Executor::Network`], even when it touches the catalog on the way. The
|
||
//! rest — opening the catalog, a backup, a local import, 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);
|
||
}
|
||
}
|