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.
This commit is contained in:
2026-09-27 07:08:37 -04:00
parent 7be1efff32
commit b1d1c47261
24 changed files with 87 additions and 63 deletions
+3 -1
View File
@@ -300,7 +300,9 @@ fn install_bundled_models(app: slint::android::AndroidApp) {
// on the worker because `AAssetManager` is thread-safe by contract and // on the worker because `AAssetManager` is thread-safe by contract and
// reading the pointer takes only the app's read lock, which `poll_events` // reading the pointer takes only the app's read lock, which `poll_events`
// also only ever holds shared. // also only ever holds shared.
std::thread::spawn(move || unpack_bundled_models(&app)); dr_ui::executors::spawn(dr_ui::executors::Executor::Io, "models", move || {
unpack_bundled_models(&app)
});
} }
/// The copy itself, on the worker [`install_bundled_models`] starts. /// The copy itself, on the worker [`install_bundled_models`] starts.
+8 -9
View File
@@ -45,6 +45,7 @@
//! //!
//! Waiting is the client's job (poll `locate`), which keeps this stateless. //! Waiting is the client's job (poll `locate`), which keeps this stateless.
use crate::executors::{self, Executor};
use std::io::{BufRead, BufReader, Write}; use std::io::{BufRead, BufReader, Write};
use std::os::unix::net::{UnixListener, UnixStream}; use std::os::unix::net::{UnixListener, UnixStream};
use std::sync::mpsc; use std::sync::mpsc;
@@ -71,15 +72,13 @@ pub(crate) fn attach(window: &AppWindow) {
}; };
log::info!("automation: listening on {path:?}"); log::info!("automation: listening on {path:?}");
let weak = window.as_weak(); let weak = window.as_weak();
std::thread::Builder::new() executors::try_spawn(Executor::Io, "automation", move || {
.name("automation".into()) for stream in listener.incoming().flatten() {
.spawn(move || { let weak = weak.clone();
for stream in listener.incoming().flatten() { executors::spawn(Executor::Io, "auto-conn", move || serve(stream, weak));
let weak = weak.clone(); }
std::thread::spawn(move || serve(stream, weak)); })
} .ok();
})
.ok();
} }
fn serve(stream: UnixStream, weak: slint::Weak<AppWindow>) { fn serve(stream: UnixStream, weak: slint::Weak<AppWindow>) {
+2 -1
View File
@@ -36,6 +36,7 @@
//! abandoned half way leaves the signatures it did compute — they are permanent //! abandoned half way leaves the signatures it did compute — they are permanent
//! and correct — and the previous grouping intact. //! and correct — and the previous grouping intact.
use crate::executors::{self, Executor};
use std::cell::{Cell, RefCell}; use std::cell::{Cell, RefCell};
use std::collections::HashMap; use std::collections::HashMap;
use std::path::PathBuf; use std::path::PathBuf;
@@ -82,7 +83,7 @@ pub enum BurstMessage {
pub fn spawn_grouping(catalog_path: PathBuf, thumbs_dir: PathBuf) -> Receiver<BurstMessage> { pub fn spawn_grouping(catalog_path: PathBuf, thumbs_dir: PathBuf) -> Receiver<BurstMessage> {
let (tx, rx) = std::sync::mpsc::channel(); let (tx, rx) = std::sync::mpsc::channel();
std::thread::spawn(move || { executors::spawn(Executor::Decode, "bursts", move || {
let catalog = match Catalog::open(&catalog_path) { let catalog = match Catalog::open(&catalog_path) {
Ok(c) => c, Ok(c) => c,
Err(e) => { Err(e) => {
+3 -2
View File
@@ -32,6 +32,7 @@
//! share boundary the account cannot write to. The scanner excludes it by the //! share boundary the account cannot write to. The scanner excludes it by the
//! same mechanism that excludes the trash. //! same mechanism that excludes the trash.
use crate::executors::{self, Executor};
use std::path::{Path, PathBuf}; use std::path::{Path, PathBuf};
use dr_sync::{Connection, RemoteBackend, RemoteError, RemoteId, RemotePath}; use dr_sync::{Connection, RemoteBackend, RemoteError, RemoteId, RemotePath};
@@ -132,7 +133,7 @@ pub fn spawn_sync(
) -> std::sync::mpsc::Receiver<SyncMessage> { ) -> std::sync::mpsc::Receiver<SyncMessage> {
let (tx, rx) = std::sync::mpsc::channel(); let (tx, rx) = std::sync::mpsc::channel();
std::thread::spawn(move || { executors::spawn(Executor::Network, "sync", move || {
// A current-thread runtime here is what made the library scan itself // A current-thread runtime here is what made the library scan itself
// rather than adopt the shards the server already had; see net_runtime. // rather than adopt the shards the server already had; see net_runtime.
let rt = match crate::net_runtime::build() { let rt = match crate::net_runtime::build() {
@@ -1018,7 +1019,7 @@ pub fn spawn_place_fetch(
) -> std::sync::mpsc::Receiver<Result<Option<dr_types::Place>, String>> { ) -> std::sync::mpsc::Receiver<Result<Option<dr_types::Place>, String>> {
let (tx, rx) = std::sync::mpsc::channel(); let (tx, rx) = std::sync::mpsc::channel();
std::thread::spawn(move || { executors::spawn(Executor::Network, "places", move || {
let rt = match crate::net_runtime::build() { let rt = match crate::net_runtime::build() {
Ok(rt) => rt, Ok(rt) => rt,
Err(e) => { Err(e) => {
+3 -2
View File
@@ -38,6 +38,7 @@
//! is a function of the image id, so the next review finds each such file //! is a function of the image id, so the next review finds each such file
//! there, and consolidating again treats it as already moved. //! there, and consolidating again treats it as already moved.
use crate::executors::{self, Executor};
use std::path::PathBuf; use std::path::PathBuf;
use std::sync::atomic::{AtomicBool, Ordering}; use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::mpsc::{Receiver, Sender}; use std::sync::mpsc::{Receiver, Sender};
@@ -492,7 +493,7 @@ pub fn spawn_check(
stop: Arc<AtomicBool>, stop: Arc<AtomicBool>,
) -> Receiver<DupMessage> { ) -> Receiver<DupMessage> {
let (tx, rx) = std::sync::mpsc::channel(); let (tx, rx) = std::sync::mpsc::channel();
std::thread::spawn(move || { executors::spawn(Executor::Network, "dupes", move || {
let stopped = run_check(&conn, &catalog_path, cache_dir, groups, &stop, &tx).err(); let stopped = run_check(&conn, &catalog_path, cache_dir, groups, &stop, &tx).err();
let _ = tx.send(DupMessage::Finished { stopped }); let _ = tx.send(DupMessage::Finished { stopped });
}); });
@@ -758,7 +759,7 @@ pub fn spawn_consolidate(
stop: Arc<AtomicBool>, stop: Arc<AtomicBool>,
) -> Receiver<DupMessage> { ) -> Receiver<DupMessage> {
let (tx, rx) = std::sync::mpsc::channel(); let (tx, rx) = std::sync::mpsc::channel();
std::thread::spawn(move || { executors::spawn(Executor::Network, "consolidate", move || {
let stopped = run_consolidate(&conn, &catalog_path, cache_dir, plans, &stop, &tx).err(); let stopped = run_consolidate(&conn, &catalog_path, cache_dir, plans, &stop, &tx).err();
let _ = tx.send(DupMessage::Finished { stopped }); let _ = tx.send(DupMessage::Finished { stopped });
}); });
+9 -6
View File
@@ -35,12 +35,15 @@
//! //!
//! # Which executor a job belongs to //! # Which executor a job belongs to
//! //!
//! By what it spends its time waiting on. A worker that builds a network //! By what it spends its time on. One that decodes images or runs a CPU pass
//! runtime and moves bytes to or from a server is [`Executor::Network`], even //! over their pixels or embeddings — the thumbnail and metadata sweeps, face
//! when it touches the catalog on the way; one that decodes images or runs a //! indexing, grouping, repairs — is [`Executor::Decode`], even when it fetches
//! CPU pass over their pixels or embeddings is [`Executor::Decode`]; one that //! the bytes first. One that drives the GPU — batch export, a merge — is
//! drives the GPU is [`Executor::GpuSubmit`]; the rest — catalog, files, local //! [`Executor::GpuSubmit`]. One whose work is moving files or metadata to or
//! sockets, relaying a channel to the log — is [`Executor::Io`]. //! 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::cell::Cell;
use std::thread::JoinHandle; use std::thread::JoinHandle;
+5 -2
View File
@@ -55,6 +55,7 @@
//! same file from the grid does — subject, both times, to what the settings //! same file from the grid does — subject, both times, to what the settings
//! allow (FR-EXP-8). //! allow (FR-EXP-8).
use crate::executors::{self, Executor};
use std::collections::HashSet; use std::collections::HashSet;
use std::path::{Path, PathBuf}; use std::path::{Path, PathBuf};
use std::rc::Rc; use std::rc::Rc;
@@ -414,7 +415,7 @@ pub fn spawn_upload(
) -> std::sync::mpsc::Receiver<UploadMessage> { ) -> std::sync::mpsc::Receiver<UploadMessage> {
let (tx, rx) = std::sync::mpsc::channel(); let (tx, rx) = std::sync::mpsc::channel();
std::thread::spawn(move || { executors::spawn(Executor::Network, "upload", move || {
let rt = match crate::net_runtime::build() { let rt = match crate::net_runtime::build() {
Ok(rt) => rt, Ok(rt) => rt,
Err(e) => { Err(e) => {
@@ -665,7 +666,9 @@ pub enum BatchMessage {
/// Export a selection on a thread of its own. /// Export a selection on a thread of its own.
pub fn spawn_batch(request: BatchRequest, cancel: Cancel) -> Receiver<BatchMessage> { pub fn spawn_batch(request: BatchRequest, cancel: Cancel) -> Receiver<BatchMessage> {
let (tx, rx) = std::sync::mpsc::channel(); let (tx, rx) = std::sync::mpsc::channel();
std::thread::spawn(move || run(request, &cancel, &tx)); executors::spawn(Executor::GpuSubmit, "export", move || {
run(request, &cancel, &tx)
});
rx rx
} }
+4 -3
View File
@@ -23,6 +23,7 @@
//! on every one. It therefore runs debounced, when detection has been idle and //! on every one. It therefore runs debounced, when detection has been idle and
//! the face count has moved materially (catalog.md §10.2). //! the face count has moved materially (catalog.md §10.2).
use crate::executors::{self, Executor};
use std::path::PathBuf; use std::path::PathBuf;
use std::sync::mpsc::{Receiver, Sender}; use std::sync::mpsc::{Receiver, Sender};
@@ -816,7 +817,7 @@ pub fn spawn_store_face_sweep(
) -> Receiver<FaceSweepMessage> { ) -> Receiver<FaceSweepMessage> {
let (tx, rx) = std::sync::mpsc::channel(); let (tx, rx) = std::sync::mpsc::channel();
std::thread::spawn(move || { executors::spawn(Executor::Decode, "faces", move || {
let finish_empty = |tx: &Sender<FaceSweepMessage>| { let finish_empty = |tx: &Sender<FaceSweepMessage>| {
let _ = tx.send(FaceSweepMessage::Finished { let _ = tx.send(FaceSweepMessage::Finished {
images: 0, images: 0,
@@ -1348,7 +1349,7 @@ pub fn spawn_grouping_preview(
) -> Receiver<PreviewMessage> { ) -> Receiver<PreviewMessage> {
let (tx, rx) = std::sync::mpsc::channel(); let (tx, rx) = std::sync::mpsc::channel();
std::thread::spawn(move || { executors::spawn(Executor::Decode, "grouping", move || {
let msg = match Catalog::open(&catalog_path) { let msg = match Catalog::open(&catalog_path) {
Ok(catalog) => match preview_grouping(&catalog, &model_id, &grouping) { Ok(catalog) => match preview_grouping(&catalog, &model_id, &grouping) {
Ok(p) => PreviewMessage::Ready(p), Ok(p) => PreviewMessage::Ready(p),
@@ -1397,7 +1398,7 @@ pub fn spawn_recluster(
) -> Receiver<ReclusterMessage> { ) -> Receiver<ReclusterMessage> {
let (tx, rx) = std::sync::mpsc::channel(); let (tx, rx) = std::sync::mpsc::channel();
std::thread::spawn(move || { executors::spawn(Executor::Decode, "recluster", move || {
let catalog = match Catalog::open(&catalog_path) { let catalog = match Catalog::open(&catalog_path) {
Ok(c) => c, Ok(c) => c,
Err(e) => { Err(e) => {
+7 -8
View File
@@ -48,6 +48,7 @@
//! record afterwards is each file's digest (`dr_catalog::dedup`), which the //! record afterwards is each file's digest (`dr_catalog::dedup`), which the
//! scan cannot know because it never reads a whole file. //! scan cannot know because it never reads a whole file.
use crate::executors::{self, Executor};
use std::path::{Path, PathBuf}; use std::path::{Path, PathBuf};
use std::sync::atomic::{AtomicBool, Ordering}; use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::mpsc::{Receiver, Sender}; use std::sync::mpsc::{Receiver, Sender};
@@ -228,14 +229,12 @@ pub struct Outcome {
/// does, and the caller holds it. /// does, and the caller holds it.
pub fn spawn(request: Request, cancel: Cancel) -> Receiver<Message> { pub fn spawn(request: Request, cancel: Cancel) -> Receiver<Message> {
let (tx, rx) = std::sync::mpsc::channel(); let (tx, rx) = std::sync::mpsc::channel();
std::thread::Builder::new() executors::try_spawn(Executor::Io, "import", move || {
.name("import".into()) if let Err(e) = run(request, &cancel, &tx) {
.spawn(move || { let _ = tx.send(Message::Failed(e));
if let Err(e) = run(request, &cancel, &tx) { }
let _ = tx.send(Message::Failed(e)); })
} .expect("spawning the import worker");
})
.expect("spawning the import worker");
rx rx
} }
+4 -3
View File
@@ -4,6 +4,7 @@
//! the develop window's wiring. The model holds the state machine and is //! the develop window's wiring. The model holds the state machine and is
//! tested headless; this module only moves values across the boundary. //! tested headless; this module only moves values across the boundary.
use crate::executors::{self, Executor};
use std::cell::RefCell; use std::cell::RefCell;
use std::rc::Rc; use std::rc::Rc;
@@ -464,7 +465,7 @@ fn spawn_login(weak: slint::Weak<AppWindow>, ctl: Rc<LaunchController>, server:
// the thread boundary. // the thread boundary.
let (tx, rx) = std::sync::mpsc::channel::<LoginMessage>(); let (tx, rx) = std::sync::mpsc::channel::<LoginMessage>();
std::thread::spawn(move || { executors::spawn(Executor::Network, "login", move || {
// A panic anywhere below would unwind the thread, drop `tx`, and leave // A panic anywhere below would unwind the thread, drop `tx`, and leave
// the UI with nothing but a closed channel — which it can only report // the UI with nothing but a closed channel — which it can only report
// as "failed unexpectedly", losing the one piece of information that // as "failed unexpectedly", losing the one piece of information that
@@ -533,7 +534,7 @@ fn spawn_direct_login(
) { ) {
let (tx, rx) = std::sync::mpsc::channel::<LoginMessage>(); let (tx, rx) = std::sync::mpsc::channel::<LoginMessage>();
std::thread::spawn(move || { executors::spawn(Executor::Network, "login", move || {
let panic_tx = tx.clone(); let panic_tx = tx.clone();
let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(move || { let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(move || {
// Multi-thread for the reason the browser flow is: a current-thread // Multi-thread for the reason the browser flow is: a current-thread
@@ -767,7 +768,7 @@ fn spawn_folder_list(weak: slint::Weak<AppWindow>, ctl: Rc<LaunchController>, pa
let (tx, rx) = std::sync::mpsc::channel::<Result<Vec<String>, String>>(); let (tx, rx) = std::sync::mpsc::channel::<Result<Vec<String>, String>>();
std::thread::spawn(move || { executors::spawn(Executor::Network, "folders", move || {
// Multi-thread for the same reason as the login worker: a // Multi-thread for the same reason as the login worker: a
// current-thread runtime left reqwest's connection future unpolled on // current-thread runtime left reqwest's connection future unpolled on
// Android, so the await never resolved and the thread stopped without // Android, so the await never resolved and the thread stopped without
+1 -1
View File
@@ -747,7 +747,7 @@ fn drain_outbox(library: &Rc<library_ui::LibraryController>) {
let root = conn.account.root.clone(); let root = conn.account.root.clone();
let rx = export::spawn_upload(conn, root, outbox); let rx = export::spawn_upload(conn, root, outbox);
std::thread::spawn(move || { executors::spawn(executors::Executor::Io, "upload-log", move || {
while let Ok(msg) = rx.recv() { while let Ok(msg) = rx.recv() {
match msg { match msg {
export::UploadMessage::Status(s) => log::info!("export: {s}"), export::UploadMessage::Status(s) => log::info!("export: {s}"),
+2 -1
View File
@@ -1,6 +1,7 @@
//! Walking a remote library into the catalog, and pulling other devices' //! Walking a remote library into the catalog, and pulling other devices'
//! judgements out of the sidecars the walk finds along the way. //! judgements out of the sidecars the walk finds along the way.
use crate::executors::{self, Executor};
use dr_catalog::Catalog; use dr_catalog::Catalog;
use dr_sync::{Connection, RemoteBackend, RemoteId, RemotePath}; use dr_sync::{Connection, RemoteBackend, RemoteId, RemotePath};
use dr_types::FormatFilter; use dr_types::FormatFilter;
@@ -77,7 +78,7 @@ pub fn spawn_scan(
) -> Receiver<ScanMessage> { ) -> Receiver<ScanMessage> {
let (tx, rx) = std::sync::mpsc::channel(); let (tx, rx) = std::sync::mpsc::channel();
std::thread::spawn(move || { executors::spawn(Executor::Network, "scan", move || {
let started = std::time::Instant::now(); let started = std::time::Instant::now();
if let Err(e) = run_scan(&tx, conn, root, filter, catalog_path, started) { if let Err(e) = run_scan(&tx, conn, root, filter, catalog_path, started) {
let _ = tx.send(ScanMessage::Failed { let _ = tx.send(ScanMessage::Failed {
+3 -2
View File
@@ -1,6 +1,7 @@
//! Writing local edits out to the catalog's sidecar outbox, and draining //! Writing local edits out to the catalog's sidecar outbox, and draining
//! that outbox to the remote once a connection is available. //! that outbox to the remote once a connection is available.
use crate::executors::{self, Executor};
use crate::sidecar_cache::SidecarCache; use crate::sidecar_cache::SidecarCache;
use dr_sync::{Connection, RemoteBackend, RemoteError, RemoteId, RemotePath}; use dr_sync::{Connection, RemoteBackend, RemoteError, RemoteId, RemotePath};
use std::path::PathBuf; use std::path::PathBuf;
@@ -138,7 +139,7 @@ pub fn spawn_sidecar_writes(
) -> Receiver<SidecarMessage> { ) -> Receiver<SidecarMessage> {
let (tx, rx) = std::sync::mpsc::channel(); let (tx, rx) = std::sync::mpsc::channel();
std::thread::spawn(move || { executors::spawn(Executor::Network, "sc-write", move || {
let cache = SidecarCache::open(cache_dir); let cache = SidecarCache::open(cache_dir);
// The runtime and the backend are only needed to *upload*. Offline, // The runtime and the backend are only needed to *upload*. Offline,
@@ -451,7 +452,7 @@ pub(super) async fn write_one_sidecar_online(
pub fn spawn_outbox_drain(conn: Connection, cache_dir: PathBuf) -> Receiver<SidecarMessage> { pub fn spawn_outbox_drain(conn: Connection, cache_dir: PathBuf) -> Receiver<SidecarMessage> {
let (tx, rx) = std::sync::mpsc::channel(); let (tx, rx) = std::sync::mpsc::channel();
std::thread::spawn(move || { executors::spawn(Executor::Network, "outbox", move || {
let cache = SidecarCache::open(cache_dir); let cache = SidecarCache::open(cache_dir);
let queued = cache.pending(); let queued = cache.pending();
if queued.is_empty() { if queued.is_empty() {
+3 -2
View File
@@ -1,6 +1,7 @@
//! The background sweeps: metadata extraction and thumbnail generation //! The background sweeps: metadata extraction and thumbnail generation
//! for whatever the catalog still owes, a chunk at a time. //! for whatever the catalog still owes, a chunk at a time.
use crate::executors::{self, Executor};
use dr_catalog::Catalog; use dr_catalog::Catalog;
use dr_sync::{Connection, RemoteBackend, RemoteId, RemotePath}; use dr_sync::{Connection, RemoteBackend, RemoteId, RemotePath};
use dr_thumbs::ThumbStore; use dr_thumbs::ThumbStore;
@@ -280,7 +281,7 @@ pub fn spawn_sweep(conn: Connection, catalog_path: PathBuf) -> Receiver<SweepMes
// The one place this job names a decoder; everything below takes it. // The one place this job names a decoder; everything below takes it.
let decoder = dr_decode::default(); let decoder = dr_decode::default();
std::thread::spawn(move || { executors::spawn(Executor::Decode, "metadata", move || {
let catalog = match Catalog::open(&catalog_path) { let catalog = match Catalog::open(&catalog_path) {
Ok(c) => c, Ok(c) => c,
Err(e) => { Err(e) => {
@@ -602,7 +603,7 @@ pub fn spawn_thumbnail_sweep(
// The one place this job names a decoder; everything below takes it. // The one place this job names a decoder; everything below takes it.
let decoder = dr_decode::default(); let decoder = dr_decode::default();
std::thread::spawn(move || { executors::spawn(Executor::Decode, "thumb-sweep", move || {
let finish_empty = |tx: &Sender<ThumbSweepMessage>| { let finish_empty = |tx: &Sender<ThumbSweepMessage>| {
let _ = tx.send(ThumbSweepMessage::Finished { let _ = tx.send(ThumbSweepMessage::Finished {
stored: 0, stored: 0,
+9 -8
View File
@@ -2,6 +2,7 @@
//! cache, dehydration, and the prefetcher that keeps the grid ahead of //! cache, dehydration, and the prefetcher that keeps the grid ahead of
//! scrolling. //! scrolling.
use crate::executors::{self, Executor};
use crate::sidecar_cache::SidecarCache; use crate::sidecar_cache::SidecarCache;
use dr_catalog::Catalog; use dr_catalog::Catalog;
#[cfg(test)] #[cfg(test)]
@@ -136,7 +137,7 @@ pub fn spawn_pin_fetch(
) -> Receiver<PinMessage> { ) -> Receiver<PinMessage> {
let (tx, rx) = std::sync::mpsc::channel(); let (tx, rx) = std::sync::mpsc::channel();
std::thread::spawn(move || { executors::spawn(Executor::Network, "pin", move || {
let (store, catalog) = match ( let (store, catalog) = match (
dr_catalog::Cache::open(&cache_dir, budget), dr_catalog::Cache::open(&cache_dir, budget),
Catalog::open(&catalog_path), Catalog::open(&catalog_path),
@@ -324,7 +325,7 @@ pub fn spawn_dehydrate(
) -> Receiver<usize> { ) -> Receiver<usize> {
let (tx, rx) = std::sync::mpsc::channel(); let (tx, rx) = std::sync::mpsc::channel();
std::thread::spawn(move || { executors::spawn(Executor::Network, "dehydrate", move || {
let Ok(catalog) = Catalog::open(&catalog_path) else { let Ok(catalog) = Catalog::open(&catalog_path) else {
return; return;
}; };
@@ -442,7 +443,7 @@ pub fn spawn_sidecar_fetch(
) -> Receiver<Option<dr_pipeline::Sidecar>> { ) -> Receiver<Option<dr_pipeline::Sidecar>> {
let (tx, rx) = std::sync::mpsc::channel(); let (tx, rx) = std::sync::mpsc::channel();
std::thread::spawn(move || { executors::spawn(Executor::Network, "sidecars", move || {
let cache = SidecarCache::open(cache_dir); let cache = SidecarCache::open(cache_dir);
let path_str = sidecar_path(&image_path); let path_str = sidecar_path(&image_path);
@@ -565,7 +566,7 @@ pub fn spawn_full_fetch(
cache: Option<CacheContext>, cache: Option<CacheContext>,
) -> Receiver<Result<Vec<u8>, FetchFailure>> { ) -> Receiver<Result<Vec<u8>, FetchFailure>> {
let (tx, rx) = std::sync::mpsc::channel(); let (tx, rx) = std::sync::mpsc::channel();
std::thread::spawn(move || { executors::spawn(Executor::Network, "fetch", move || {
let _ = tx.send(fetch_original(conn, &path, cache.as_ref())); let _ = tx.send(fetch_original(conn, &path, cache.as_ref()));
}); });
rx rx
@@ -845,10 +846,10 @@ impl Prefetcher {
let shared = std::sync::Arc::new(PrefetchShared::default()); let shared = std::sync::Arc::new(PrefetchShared::default());
let (tx, events) = std::sync::mpsc::channel(); let (tx, events) = std::sync::mpsc::channel();
let worker = shared.clone(); let worker = shared.clone();
std::thread::Builder::new() executors::try_spawn(Executor::Network, "prefetch", move || {
.name("prefetch".into()) serve_prefetches(&worker, &tx)
.spawn(move || serve_prefetches(&worker, &tx)) })
.expect("spawning the prefetch worker"); .expect("spawning the prefetch worker");
Self { shared, events } Self { shared, events }
} }
+2 -1
View File
@@ -1,6 +1,7 @@
//! Generating thumbnails locally from a decoded preview or original, and //! Generating thumbnails locally from a decoded preview or original, and
//! the metadata that comes along for the ride. //! the metadata that comes along for the ride.
use crate::executors::{self, Executor};
#[cfg(test)] #[cfg(test)]
use dr_catalog::Catalog; use dr_catalog::Catalog;
use dr_sync::{Connection, RemoteBackend, RemoteId, RemotePath}; use dr_sync::{Connection, RemoteBackend, RemoteId, RemotePath};
@@ -101,7 +102,7 @@ pub fn spawn_thumbnails(
// The one place this job names a decoder; everything below takes it. // The one place this job names a decoder; everything below takes it.
let decoder = dr_decode::default(); let decoder = dr_decode::default();
std::thread::spawn(move || { executors::spawn(Executor::Decode, "thumbs", move || {
let mut store = match ThumbStore::open(&store_dir) { let mut store = match ThumbStore::open(&store_dir) {
Ok(s) => Some(s), Ok(s) => Some(s),
Err(e) => { Err(e) => {
+3 -2
View File
@@ -1,6 +1,7 @@
//! Pushing local judgements out to XMP sidecars, and reloading sidecars a //! Pushing local judgements out to XMP sidecars, and reloading sidecars a
//! person chose to trust by hand. //! person chose to trust by hand.
use crate::executors::{self, Executor};
use dr_catalog::Catalog; use dr_catalog::Catalog;
use dr_sync::{Connection, RemoteBackend, RemoteId, RemotePath}; use dr_sync::{Connection, RemoteBackend, RemoteId, RemotePath};
use std::path::PathBuf; use std::path::PathBuf;
@@ -32,7 +33,7 @@ pub struct XmpWrite {
/// name. Reported once at the end, as the judgement writes are. /// name. Reported once at the end, as the judgement writes are.
pub fn spawn_xmp_writes(conn: Connection, writes: Vec<XmpWrite>) -> Receiver<XmpMessage> { pub fn spawn_xmp_writes(conn: Connection, writes: Vec<XmpWrite>) -> Receiver<XmpMessage> {
let (tx, rx) = std::sync::mpsc::channel(); let (tx, rx) = std::sync::mpsc::channel();
std::thread::spawn(move || { executors::spawn(Executor::Network, "xmp-write", move || {
let mut written = 0usize; let mut written = 0usize;
let mut failed = 0usize; let mut failed = 0usize;
let mut last_error = None; let mut last_error = None;
@@ -140,7 +141,7 @@ pub fn spawn_xmp_reload(
paths: Vec<String>, paths: Vec<String>,
) -> Receiver<XmpMessage> { ) -> Receiver<XmpMessage> {
let (tx, rx) = std::sync::mpsc::channel(); let (tx, rx) = std::sync::mpsc::channel();
std::thread::spawn(move || { executors::spawn(Executor::Network, "xmp-reload", move || {
let mut written = 0usize; let mut written = 0usize;
let mut failed = 0usize; let mut failed = 0usize;
let mut last_error = None; let mut last_error = None;
+2 -1
View File
@@ -8,6 +8,7 @@
//! (`offline`) it hands off to on a failure. See `docs/dev/code-health.md` //! (`offline`) it hands off to on a failure. See `docs/dev/code-health.md`
//! CH-1. //! CH-1.
use crate::executors::{self, Executor};
use std::path::PathBuf; use std::path::PathBuf;
use std::rc::Rc; use std::rc::Rc;
use std::sync::mpsc::Receiver; use std::sync::mpsc::Receiver;
@@ -271,7 +272,7 @@ enum CatalogOpen {
fn spawn_catalog_open(path: PathBuf) -> Receiver<CatalogOpen> { fn spawn_catalog_open(path: PathBuf) -> Receiver<CatalogOpen> {
let (tx, rx) = std::sync::mpsc::channel(); let (tx, rx) = std::sync::mpsc::channel();
std::thread::spawn(move || { executors::spawn(Executor::Io, "catalog-open", move || {
let started = std::time::Instant::now(); let started = std::time::Instant::now();
let message = match Catalog::open_verified(&path) { let message = match Catalog::open_verified(&path) {
Ok(cat) => { Ok(cat) => {
+3 -2
View File
@@ -7,6 +7,7 @@
//! library rather than in response to what the grid is showing. //! library rather than in response to what the grid is showing.
//! See `docs/dev/code-health.md` CH-1. //! See `docs/dev/code-health.md` CH-1.
use crate::executors::{self, Executor};
use std::rc::Rc; use std::rc::Rc;
use slint::ComponentHandle; use slint::ComponentHandle;
@@ -33,7 +34,7 @@ fn spawn_scheduled_backup(ctl: &Rc<LibraryController>) {
if !dr_catalog::recovery::backup_due(&catalog_path) { if !dr_catalog::recovery::backup_due(&catalog_path) {
return; return;
} }
std::thread::spawn(move || { executors::spawn(Executor::Io, "backup", move || {
let catalog = match dr_catalog::Catalog::open(&catalog_path) { let catalog = match dr_catalog::Catalog::open(&catalog_path) {
Ok(c) => c, Ok(c) => c,
Err(e) => { Err(e) => {
@@ -111,7 +112,7 @@ pub(super) fn start_derived_sync(window: &AppWindow, ctl: &Rc<LibraryController>
let outbox = crate::export::outbox_dir(&conn.account); let outbox = crate::export::outbox_dir(&conn.account);
if crate::export::pending_count(&outbox) > 0 { if crate::export::pending_count(&outbox) > 0 {
let rx = crate::export::spawn_upload(conn.clone(), conn.account.root.clone(), outbox); let rx = crate::export::spawn_upload(conn.clone(), conn.account.root.clone(), outbox);
std::thread::spawn(move || { executors::spawn(Executor::Io, "upload-log", move || {
while let Ok(msg) = rx.recv() { while let Ok(msg) = rx.recv() {
match msg { match msg {
crate::export::UploadMessage::Status(s) => log::info!("export: {s}"), crate::export::UploadMessage::Status(s) => log::info!("export: {s}"),
+2 -1
View File
@@ -27,6 +27,7 @@
//! sensor data stays on the CPU at a fifth of the size — and it is what //! sensor data stays on the CPU at a fifth of the size — and it is what
//! makes a twelve-frame set fit. //! makes a twelve-frame set fit.
use crate::executors::{self, Executor};
use std::collections::VecDeque; use std::collections::VecDeque;
use std::path::{Path, PathBuf}; use std::path::{Path, PathBuf};
use std::sync::mpsc::{Receiver, Sender}; use std::sync::mpsc::{Receiver, Sender};
@@ -727,7 +728,7 @@ fn run_inner(
let fill_for_writer = fill_cam.is_some(); let fill_for_writer = fill_cam.is_some();
let file = let file =
std::fs::File::create(&out_path).map_err(|e| format!("{}: {e}", out_path.display()))?; std::fs::File::create(&out_path).map_err(|e| format!("{}: {e}", out_path.display()))?;
let writer = std::thread::spawn(move || -> Result<(), String> { let writer = executors::spawn(Executor::Io, "dng-write", move || -> Result<(), String> {
let mut file = std::io::BufWriter::new(file); let mut file = std::io::BufWriter::new(file);
dr_export::write_linear_dng( dr_export::write_linear_dng(
&mut file, &mut file,
+2 -1
View File
@@ -11,6 +11,7 @@
//! staged, the outbox drains and the library rescans, and the composite //! staged, the outbox drains and the library rescans, and the composite
//! appears in the grid beside its sources. //! appears in the grid beside its sources.
use crate::executors::{self, Executor};
use std::cell::{Cell, RefCell}; use std::cell::{Cell, RefCell};
use std::rc::Rc; use std::rc::Rc;
use std::sync::mpsc::{Receiver, Sender}; use std::sync::mpsc::{Receiver, Sender};
@@ -446,7 +447,7 @@ fn start<Fetch>(
{ {
let cancel = cancel.clone(); let cancel = cancel.clone();
let tx = tx.clone(); let tx = tx.clone();
std::thread::spawn(move || { executors::spawn(Executor::GpuSubmit, "merge", move || {
let Some(frames) = fetch(&tx, &cancel) else { let Some(frames) = fetch(&tx, &cancel) else {
return; return;
}; };
+2 -1
View File
@@ -9,6 +9,7 @@
//! locally: the server may have normalised it, or refused it, and the list //! locally: the server may have normalised it, or refused it, and the list
//! the user picks from should be the server's, not this app's guess at it. //! the user picks from should be the server's, not this app's guess at it.
use crate::executors::{self, Executor};
use std::cell::RefCell; use std::cell::RefCell;
use std::rc::Rc; use std::rc::Rc;
use std::sync::mpsc::TryRecvError; use std::sync::mpsc::TryRecvError;
@@ -105,7 +106,7 @@ where
{ {
let work: Work = Box::new(work); let work: Work = Box::new(work);
let (tx, rx) = std::sync::mpsc::channel::<Result<Vec<String>, String>>(); let (tx, rx) = std::sync::mpsc::channel::<Result<Vec<String>, String>>();
std::thread::spawn(move || { executors::spawn(Executor::Network, "remote-dirs", move || {
// Multi-thread, for the reason the login worker records: a // Multi-thread, for the reason the login worker records: a
// current-thread runtime left reqwest's connection future unpolled on // current-thread runtime left reqwest's connection future unpolled on
// Android, and the await never resolved. // Android, and the await never resolved.
+2 -1
View File
@@ -41,6 +41,7 @@
//! runs over everything the *chosen detector* has not been over, whatever a //! runs over everything the *chosen detector* has not been over, whatever a
//! weaker one found there. Same job, one predicate differs. //! weaker one found there. Same job, one predicate differs.
use crate::executors::{self, Executor};
use std::collections::HashSet; use std::collections::HashSet;
use std::path::PathBuf; use std::path::PathBuf;
use std::sync::mpsc::{Receiver, Sender}; use std::sync::mpsc::{Receiver, Sender};
@@ -1008,7 +1009,7 @@ pub fn spawn(
) -> Receiver<FaceSweepMessage> { ) -> Receiver<FaceSweepMessage> {
let (tx, rx) = std::sync::mpsc::channel(); let (tx, rx) = std::sync::mpsc::channel();
std::thread::spawn(move || { executors::spawn(Executor::Decode, "repairs", move || {
let finish_empty = |tx: &Sender<FaceSweepMessage>| { let finish_empty = |tx: &Sender<FaceSweepMessage>| {
let _ = tx.send(FaceSweepMessage::Finished { let _ = tx.send(FaceSweepMessage::Finished {
images: 0, images: 0,
+3 -2
View File
@@ -27,6 +27,7 @@
//! halfway with no record of where it stopped is worse than one that reports //! halfway with no record of where it stopped is worse than one that reports
//! "38 of 40". //! "38 of 40".
use crate::executors::{self, Executor};
use std::path::PathBuf; use std::path::PathBuf;
use std::sync::mpsc::Receiver; use std::sync::mpsc::Receiver;
@@ -169,7 +170,7 @@ pub fn spawn_move(
) -> Receiver<TrashMessage> { ) -> Receiver<TrashMessage> {
let (tx, rx) = std::sync::mpsc::channel(); let (tx, rx) = std::sync::mpsc::channel();
std::thread::spawn(move || { executors::spawn(Executor::Network, "trash", move || {
let total = moves.len(); let total = moves.len();
let rt = match crate::net_runtime::build() { let rt = match crate::net_runtime::build() {
Ok(rt) => rt, Ok(rt) => rt,
@@ -281,7 +282,7 @@ pub fn spawn_purge(
) -> Receiver<TrashMessage> { ) -> Receiver<TrashMessage> {
let (tx, rx) = std::sync::mpsc::channel(); let (tx, rx) = std::sync::mpsc::channel();
std::thread::spawn(move || { executors::spawn(Executor::Network, "purge", move || {
let total = paths.len(); let total = paths.len();
let rt = match crate::net_runtime::build() { let rt = match crate::net_runtime::build() {
Ok(rt) => rt, Ok(rt) => rt,