use crate::io_timeout::{wait_with_timeout, WaitOutcome};
use crate::semaphore::Semaphore;
use anyhow::Context as _;
use image::DynamicImage;
use std::fs::File;
use std::path::{Path, PathBuf};
use std::sync::OnceLock;
use std::time::Duration;
const QLMANAGE_TIMEOUT: Duration = Duration::from_secs(20);
const QLMANAGE_MAX_CONCURRENT_DEFAULT: usize = 6;
static QLMANAGE_CONCURRENCY_OVERRIDE: OnceLock<usize> = OnceLock::new();
pub fn set_qlmanage_concurrency(n: usize) {
let _ = QLMANAGE_CONCURRENCY_OVERRIDE.set(n);
}
fn resolve_qlmanage_concurrency(override_val: Option<usize>) -> usize {
override_val.unwrap_or(QLMANAGE_MAX_CONCURRENT_DEFAULT)
}
pub fn qlmanage_semaphore() -> &'static Semaphore {
static SEM: OnceLock<Semaphore> = OnceLock::new();
SEM.get_or_init(|| {
let max = resolve_qlmanage_concurrency(QLMANAGE_CONCURRENCY_OVERRIDE.get().copied());
Semaphore::new(max)
})
}
static QUICKLOOK_SLOT_DIR: OnceLock<PathBuf> = OnceLock::new();
const SLOT_POLL: Duration = Duration::from_millis(50);
pub fn set_quicklook_slot_dir(dir: PathBuf) {
let _ = QUICKLOOK_SLOT_DIR.set(dir);
}
struct QuicklookSlot(#[allow(dead_code)] File);
fn acquire_slot(dir: &Path, n: usize, poll: Duration) -> std::io::Result<QuicklookSlot> {
acquire_slot_with(dir, n, poll, fs2::FileExt::try_lock_exclusive)
}
fn acquire_slot_with(
dir: &Path,
n: usize,
poll: Duration,
try_lock: impl Fn(&File) -> std::io::Result<()>,
) -> std::io::Result<QuicklookSlot> {
std::fs::create_dir_all(dir)?;
let slots = (0..n.max(1))
.map(|i| {
std::fs::OpenOptions::new()
.create(true)
.write(true)
.truncate(false)
.open(dir.join(format!("slot-{i}.lock")))
})
.collect::<std::io::Result<Vec<File>>>()?;
loop {
for slot in &slots {
match try_lock(slot) {
Ok(()) => return Ok(QuicklookSlot(slot.try_clone()?)),
Err(error) if error.kind() == std::io::ErrorKind::WouldBlock => {}
Err(error) => return Err(error),
}
}
std::thread::sleep(poll);
}
}
fn machine_slot() -> Option<QuicklookSlot> {
let dir = QUICKLOOK_SLOT_DIR.get()?;
match acquire_slot(dir, qlmanage_semaphore().max(), SLOT_POLL) {
Ok(slot) => Some(slot),
Err(error) => {
static WARNED: std::sync::Once = std::sync::Once::new();
WARNED.call_once(|| {
tracing::warn!(
"QuickLook slots in {} are unusable, limiting conversions in this process only: {error}",
dir.display()
)
});
None
}
}
}
pub fn decode_via_quicklook(
path: &Path,
tag: &str,
max_size: Option<u32>,
) -> anyhow::Result<DynamicImage> {
if !cfg!(target_os = "macos") {
warn_quicklook_unavailable_once();
return Err(anyhow::anyhow!("needs macOS QuickLook (`qlmanage`)")
.context(crate::error_kind::ErrorKind::QuicklookUnavailable));
}
let probe_ext = path
.extension()
.and_then(|e| e.to_str())
.map(|e| e.to_lowercase());
if matches!(probe_ext.as_deref(), Some("mov") | Some("mp4"))
&& !crate::video_probe::has_video_track(path)
{
anyhow::bail!("no video track (audio-only file); skipped without calling QuickLook");
}
let scratch = tempfile::Builder::new()
.prefix(&format!("videre_ql_{tag}_"))
.tempdir()
.context("create qlmanage temp dir")?;
let out_dir = scratch.path();
let _permit = qlmanage_semaphore().acquire();
let _slot = machine_slot();
let size_arg = max_size.unwrap_or(10000).to_string();
let mut child = std::process::Command::new("qlmanage")
.args(["-t", "-s", &size_arg, "-o"])
.arg(out_dir)
.arg(path)
.stdout(std::process::Stdio::null())
.stderr(std::process::Stdio::null())
.spawn()
.context("run qlmanage (requires macOS)")?;
let outcome = wait_with_timeout(&mut child, QLMANAGE_TIMEOUT);
if outcome == WaitOutcome::TimedOut {
return Err(anyhow::anyhow!(
"qlmanage timed out after {}s decoding {} (file may be unreachable - is its drive connected?)",
QLMANAGE_TIMEOUT.as_secs(),
path.display()
)
.context(crate::error_kind::ErrorKind::SourceUnavailable));
}
anyhow::ensure!(
outcome == WaitOutcome::Success,
"qlmanage failed for {}",
path.display()
);
let file_name = path.file_name().context("path has no file name")?;
let out_file = out_dir.join(quicklook_output_name(file_name));
image::open(&out_file).with_context(|| format!("decode qlmanage output for {}", path.display()))
}
pub fn publish_cached_original(
cache: &crate::library::CachePaths,
hash: &str,
jpeg: &[u8],
) -> bool {
let destination = crate::thumb_cache::original_path_in(cache, hash);
let written = crate::atomic_file::publish(&destination, |file| {
use std::io::Write as _;
file.write_all(jpeg).map_err(Into::into)
});
if let Err(error) = written {
tracing::warn!("could not cache the full-resolution original for {hash}: {error:#}");
return false;
}
true
}
pub fn decode_fullres_cached(
path: &Path,
tag: &str,
cache: Option<(&crate::library::CachePaths, &str)>,
) -> anyhow::Result<DynamicImage> {
if let Some((cache, hash)) = cache {
let cached = crate::thumb_cache::original_path_in(cache, hash);
let opened =
crate::io_timeout::run_with_timeout(crate::io_timeout::DEFAULT_IO_TIMEOUT, move || {
image::open(&cached)
});
if let Ok(Ok(img)) = opened {
return Ok(img);
}
}
let img = decode_via_quicklook(path, tag, None)?;
if let Some((cache, hash)) = cache {
let mut jpeg = Vec::new();
if img
.write_to(
&mut std::io::Cursor::new(&mut jpeg),
image::ImageFormat::Jpeg,
)
.is_ok()
{
publish_cached_original(cache, hash, &jpeg);
}
}
Ok(img)
}
fn quicklook_output_name(file_name: &std::ffi::OsStr) -> std::ffi::OsString {
let mut name = file_name.to_os_string();
name.push(".png");
name
}
pub fn warn_if_timeout(error: &anyhow::Error) {
if crate::error_kind::ErrorKind::in_chain(error)
== Some(crate::error_kind::ErrorKind::SourceUnavailable)
{
crate::error_log::report(tracing::Level::WARN, error, None);
}
}
pub const QUICKLOOK_UNAVAILABLE: &str =
"HEIC images and video frames are decoded via macOS QuickLook (`qlmanage`), \
which has no equivalent on this platform - those files are skipped. \
Scanning, dedupe, and search still work for jpg/jpeg/png/gif/webp/bmp/tiff.";
fn warn_quicklook_unavailable_once() {
static WARNED: std::sync::Once = std::sync::Once::new();
WARNED.call_once(|| {
crate::error_log::report(tracing::Level::WARN, &quicklook_unavailable(), None)
});
}
fn quicklook_unavailable() -> anyhow::Error {
anyhow::anyhow!(QUICKLOOK_UNAVAILABLE)
.context(crate::error_kind::ErrorKind::QuicklookUnavailable)
}
#[cfg(test)]
mod tests {
use super::*;
fn hold(dir: &Path, i: usize) -> std::fs::File {
use fs2::FileExt;
std::fs::create_dir_all(dir).unwrap();
let file = std::fs::OpenOptions::new()
.create(true)
.write(true)
.truncate(false)
.open(dir.join(format!("slot-{i}.lock")))
.unwrap();
file.try_lock_exclusive().unwrap();
file
}
fn is_held(dir: &Path, i: usize) -> bool {
use fs2::FileExt;
let file = std::fs::File::open(dir.join(format!("slot-{i}.lock"))).unwrap();
file.try_lock_exclusive().is_err()
}
#[test]
fn a_free_slot_is_taken() {
let temp = tempfile::tempdir().unwrap();
let dir = temp.path().join("quicklook");
let _other = hold(&dir, 0);
let _slot = acquire_slot(&dir, 2, Duration::from_millis(10)).unwrap();
assert!(is_held(&dir, 1), "the free slot is the one taken");
}
#[test]
fn a_full_pool_waits_for_a_release() {
let temp = tempfile::tempdir().unwrap();
let dir = temp.path().join("quicklook");
let first = hold(&dir, 0);
let _second = hold(&dir, 1);
let release = std::thread::spawn(move || {
std::thread::sleep(Duration::from_millis(300));
drop(first);
});
let started = std::time::Instant::now();
let _slot = acquire_slot(&dir, 2, Duration::from_millis(10)).unwrap();
assert!(started.elapsed() >= Duration::from_millis(300));
release.join().unwrap();
}
#[test]
fn a_released_slot_is_free_again() {
let temp = tempfile::tempdir().unwrap();
let dir = temp.path().join("quicklook");
drop(acquire_slot(&dir, 1, Duration::from_millis(10)).unwrap());
assert!(!is_held(&dir, 0));
}
#[test]
fn a_filesystem_that_rejects_locking_is_an_error_not_a_wait() {
let temp = tempfile::tempdir().unwrap();
let dir = temp.path().join("quicklook");
let (tx, rx) = std::sync::mpsc::channel();
std::thread::spawn(move || {
let rejected = |_: &File| Err(std::io::Error::from(std::io::ErrorKind::Unsupported));
let _ = tx.send(acquire_slot_with(&dir, 2, Duration::from_millis(10), rejected).err());
});
let error = rx
.recv_timeout(Duration::from_secs(5))
.expect("acquire returned instead of waiting")
.expect("a lock that is refused outright is an error");
assert_eq!(error.kind(), std::io::ErrorKind::Unsupported);
}
#[test]
fn an_unusable_directory_is_an_error() {
let temp = tempfile::tempdir().unwrap();
let dir = temp.path().join("quicklook");
std::fs::write(&dir, b"not a directory").unwrap();
assert!(acquire_slot(&dir, 2, Duration::from_millis(10)).is_err());
}
#[test]
fn a_published_original_opens_as_a_jpeg() {
let temp = tempfile::tempdir().unwrap();
let ctx =
crate::library::LibraryContext::new(temp.path(), &temp.path().join("cache")).unwrap();
let img = image::RgbImage::from_pixel(2, 2, image::Rgb([255, 0, 0]));
let mut jpeg = Vec::new();
image::DynamicImage::ImageRgb8(img)
.write_to(
&mut std::io::Cursor::new(&mut jpeg),
image::ImageFormat::Jpeg,
)
.unwrap();
assert!(publish_cached_original(&ctx.cache, "pub-test", &jpeg));
let opened =
image::open(crate::thumb_cache::original_path_in(&ctx.cache, "pub-test")).unwrap();
assert_eq!((opened.width(), opened.height()), (2, 2));
}
#[test]
fn the_quicklook_warning_carries_its_kind_and_the_full_explanation() {
let err = quicklook_unavailable();
assert_eq!(
crate::error_kind::ErrorKind::in_chain(&err),
Some(crate::error_kind::ErrorKind::QuicklookUnavailable)
);
assert!(format!("{err:#}").contains("qlmanage"));
}
#[test]
fn simultaneous_conversions_of_one_file_all_succeed() {
if !cfg!(target_os = "macos") {
return;
}
let video = concat!(
env!("CARGO_MANIFEST_DIR"),
"/../videre/tests/fixtures/red_1s.mp4"
);
let handles: Vec<_> = (0..4)
.map(|_| {
std::thread::spawn(move || {
decode_via_quicklook(Path::new(video), "vposter480", Some(480))
})
})
.collect();
let ok = handles
.into_iter()
.map(|h| h.join().unwrap().is_ok())
.filter(|ok| *ok)
.count();
assert_eq!(ok, 4, "every simultaneous conversion must produce an image");
}
#[test]
#[cfg(target_os = "macos")]
fn an_audio_only_video_is_rejected_before_quicklook_is_called() {
let audio_only = concat!(
env!("CARGO_MANIFEST_DIR"),
"/../videre/tests/fixtures/audio_only.mov"
);
let err = decode_via_quicklook(Path::new(audio_only), "audio-only-test", None).unwrap_err();
assert!(
err.to_string().contains("no video track"),
"unexpected error: {err:#}"
);
}
#[test]
fn without_quicklook_the_decode_fails_with_the_quicklook_unavailable_kind() {
if cfg!(target_os = "macos") {
return;
}
let err =
decode_via_quicklook(Path::new("/nonexistent.mov"), "kind-test", None).unwrap_err();
assert_eq!(
crate::error_kind::ErrorKind::in_chain(&err),
Some(crate::error_kind::ErrorKind::QuicklookUnavailable)
);
}
#[test]
#[cfg(unix)]
fn a_non_utf8_file_name_keeps_its_bytes_in_the_output_name() {
use std::os::unix::ffi::{OsStrExt, OsStringExt};
let name = std::ffi::OsStr::from_bytes(b"caf\xe9.heic");
assert_eq!(
quicklook_output_name(name).into_vec(),
b"caf\xe9.heic.png".to_vec()
);
}
#[test]
fn only_a_timeout_is_warned_for_callers_that_swallow_the_error() {
#[derive(Clone, Default)]
struct Buf(std::sync::Arc<std::sync::Mutex<Vec<u8>>>);
impl std::io::Write for Buf {
fn write(&mut self, b: &[u8]) -> std::io::Result<usize> {
self.0.lock().unwrap().extend_from_slice(b);
Ok(b.len())
}
fn flush(&mut self) -> std::io::Result<()> {
Ok(())
}
}
let buf = Buf::default();
let w = buf.clone();
let sub = tracing_subscriber::fmt()
.with_writer(move || w.clone())
.finish();
tracing::subscriber::with_default(sub, || {
warn_if_timeout(&anyhow::anyhow!("decode failed"));
warn_if_timeout(
&anyhow::anyhow!("qlmanage timed out")
.context(crate::error_kind::ErrorKind::SourceUnavailable),
);
});
let logged = String::from_utf8_lossy(&buf.0.lock().unwrap()).to_string();
assert!(logged.contains("qlmanage timed out"), "{logged}");
assert!(!logged.contains("decode failed"), "{logged}");
}
#[test]
fn resolve_qlmanage_concurrency_uses_override_when_present() {
assert_eq!(resolve_qlmanage_concurrency(Some(10)), 10);
}
#[test]
fn resolve_qlmanage_concurrency_falls_back_to_default_when_absent() {
assert_eq!(
resolve_qlmanage_concurrency(None),
QLMANAGE_MAX_CONCURRENT_DEFAULT
);
}
#[test]
fn resolve_qlmanage_concurrency_override_of_zero_is_honored_literally() {
assert_eq!(resolve_qlmanage_concurrency(Some(0)), 0);
}
}