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:
@@ -1,6 +1,7 @@
|
||||
//! Walking a remote library into the catalog, and pulling other devices'
|
||||
//! judgements out of the sidecars the walk finds along the way.
|
||||
|
||||
use crate::executors::{self, Executor};
|
||||
use dr_catalog::Catalog;
|
||||
use dr_sync::{Connection, RemoteBackend, RemoteId, RemotePath};
|
||||
use dr_types::FormatFilter;
|
||||
@@ -77,7 +78,7 @@ pub fn spawn_scan(
|
||||
) -> Receiver<ScanMessage> {
|
||||
let (tx, rx) = std::sync::mpsc::channel();
|
||||
|
||||
std::thread::spawn(move || {
|
||||
executors::spawn(Executor::Network, "scan", move || {
|
||||
let started = std::time::Instant::now();
|
||||
if let Err(e) = run_scan(&tx, conn, root, filter, catalog_path, started) {
|
||||
let _ = tx.send(ScanMessage::Failed {
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
//! Writing local edits out to the catalog's sidecar outbox, and draining
|
||||
//! that outbox to the remote once a connection is available.
|
||||
|
||||
use crate::executors::{self, Executor};
|
||||
use crate::sidecar_cache::SidecarCache;
|
||||
use dr_sync::{Connection, RemoteBackend, RemoteError, RemoteId, RemotePath};
|
||||
use std::path::PathBuf;
|
||||
@@ -138,7 +139,7 @@ pub fn spawn_sidecar_writes(
|
||||
) -> Receiver<SidecarMessage> {
|
||||
let (tx, rx) = std::sync::mpsc::channel();
|
||||
|
||||
std::thread::spawn(move || {
|
||||
executors::spawn(Executor::Network, "sc-write", move || {
|
||||
let cache = SidecarCache::open(cache_dir);
|
||||
|
||||
// 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> {
|
||||
let (tx, rx) = std::sync::mpsc::channel();
|
||||
|
||||
std::thread::spawn(move || {
|
||||
executors::spawn(Executor::Network, "outbox", move || {
|
||||
let cache = SidecarCache::open(cache_dir);
|
||||
let queued = cache.pending();
|
||||
if queued.is_empty() {
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
//! The background sweeps: metadata extraction and thumbnail generation
|
||||
//! for whatever the catalog still owes, a chunk at a time.
|
||||
|
||||
use crate::executors::{self, Executor};
|
||||
use dr_catalog::Catalog;
|
||||
use dr_sync::{Connection, RemoteBackend, RemoteId, RemotePath};
|
||||
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.
|
||||
let decoder = dr_decode::default();
|
||||
|
||||
std::thread::spawn(move || {
|
||||
executors::spawn(Executor::Decode, "metadata", move || {
|
||||
let catalog = match Catalog::open(&catalog_path) {
|
||||
Ok(c) => c,
|
||||
Err(e) => {
|
||||
@@ -602,7 +603,7 @@ pub fn spawn_thumbnail_sweep(
|
||||
// The one place this job names a decoder; everything below takes it.
|
||||
let decoder = dr_decode::default();
|
||||
|
||||
std::thread::spawn(move || {
|
||||
executors::spawn(Executor::Decode, "thumb-sweep", move || {
|
||||
let finish_empty = |tx: &Sender<ThumbSweepMessage>| {
|
||||
let _ = tx.send(ThumbSweepMessage::Finished {
|
||||
stored: 0,
|
||||
|
||||
@@ -2,6 +2,7 @@
|
||||
//! cache, dehydration, and the prefetcher that keeps the grid ahead of
|
||||
//! scrolling.
|
||||
|
||||
use crate::executors::{self, Executor};
|
||||
use crate::sidecar_cache::SidecarCache;
|
||||
use dr_catalog::Catalog;
|
||||
#[cfg(test)]
|
||||
@@ -136,7 +137,7 @@ pub fn spawn_pin_fetch(
|
||||
) -> Receiver<PinMessage> {
|
||||
let (tx, rx) = std::sync::mpsc::channel();
|
||||
|
||||
std::thread::spawn(move || {
|
||||
executors::spawn(Executor::Network, "pin", move || {
|
||||
let (store, catalog) = match (
|
||||
dr_catalog::Cache::open(&cache_dir, budget),
|
||||
Catalog::open(&catalog_path),
|
||||
@@ -324,7 +325,7 @@ pub fn spawn_dehydrate(
|
||||
) -> Receiver<usize> {
|
||||
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 {
|
||||
return;
|
||||
};
|
||||
@@ -442,7 +443,7 @@ pub fn spawn_sidecar_fetch(
|
||||
) -> Receiver<Option<dr_pipeline::Sidecar>> {
|
||||
let (tx, rx) = std::sync::mpsc::channel();
|
||||
|
||||
std::thread::spawn(move || {
|
||||
executors::spawn(Executor::Network, "sidecars", move || {
|
||||
let cache = SidecarCache::open(cache_dir);
|
||||
let path_str = sidecar_path(&image_path);
|
||||
|
||||
@@ -565,7 +566,7 @@ pub fn spawn_full_fetch(
|
||||
cache: Option<CacheContext>,
|
||||
) -> Receiver<Result<Vec<u8>, FetchFailure>> {
|
||||
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()));
|
||||
});
|
||||
rx
|
||||
@@ -845,10 +846,10 @@ impl Prefetcher {
|
||||
let shared = std::sync::Arc::new(PrefetchShared::default());
|
||||
let (tx, events) = std::sync::mpsc::channel();
|
||||
let worker = shared.clone();
|
||||
std::thread::Builder::new()
|
||||
.name("prefetch".into())
|
||||
.spawn(move || serve_prefetches(&worker, &tx))
|
||||
.expect("spawning the prefetch worker");
|
||||
executors::try_spawn(Executor::Network, "prefetch", move || {
|
||||
serve_prefetches(&worker, &tx)
|
||||
})
|
||||
.expect("spawning the prefetch worker");
|
||||
Self { shared, events }
|
||||
}
|
||||
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
//! Generating thumbnails locally from a decoded preview or original, and
|
||||
//! the metadata that comes along for the ride.
|
||||
|
||||
use crate::executors::{self, Executor};
|
||||
#[cfg(test)]
|
||||
use dr_catalog::Catalog;
|
||||
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.
|
||||
let decoder = dr_decode::default();
|
||||
|
||||
std::thread::spawn(move || {
|
||||
executors::spawn(Executor::Decode, "thumbs", move || {
|
||||
let mut store = match ThumbStore::open(&store_dir) {
|
||||
Ok(s) => Some(s),
|
||||
Err(e) => {
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
//! Pushing local judgements out to XMP sidecars, and reloading sidecars a
|
||||
//! person chose to trust by hand.
|
||||
|
||||
use crate::executors::{self, Executor};
|
||||
use dr_catalog::Catalog;
|
||||
use dr_sync::{Connection, RemoteBackend, RemoteId, RemotePath};
|
||||
use std::path::PathBuf;
|
||||
@@ -32,7 +33,7 @@ pub struct XmpWrite {
|
||||
/// name. Reported once at the end, as the judgement writes are.
|
||||
pub fn spawn_xmp_writes(conn: Connection, writes: Vec<XmpWrite>) -> Receiver<XmpMessage> {
|
||||
let (tx, rx) = std::sync::mpsc::channel();
|
||||
std::thread::spawn(move || {
|
||||
executors::spawn(Executor::Network, "xmp-write", move || {
|
||||
let mut written = 0usize;
|
||||
let mut failed = 0usize;
|
||||
let mut last_error = None;
|
||||
@@ -140,7 +141,7 @@ pub fn spawn_xmp_reload(
|
||||
paths: Vec<String>,
|
||||
) -> Receiver<XmpMessage> {
|
||||
let (tx, rx) = std::sync::mpsc::channel();
|
||||
std::thread::spawn(move || {
|
||||
executors::spawn(Executor::Network, "xmp-reload", move || {
|
||||
let mut written = 0usize;
|
||||
let mut failed = 0usize;
|
||||
let mut last_error = None;
|
||||
|
||||
Reference in New Issue
Block a user