Files
DarkRoom/ui/dr-ui/src/executors.rs
T
dtourolle b1d1c47261 Start every worker thread through the executors module
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.
2026-09-27 07:08:37 -04:00

255 lines
10 KiB
Rust
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
//! 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);
}
}