Catalogue a merged panorama the moment it is written

A finished merge drained the outbox and started a rescan beside it. The
scan raced the upload: on a folder library the 800 MB copy was still
running when the folder was listed, on Nextcloud the upload takes
minutes, and either way the listing lacked the composite, recorded the
folder's validator, and nothing looked again until the next sync pass.

The job now reports what the catalog needs (the name, the size of the
picture it opens on, the capture time it wrote into the DNG) and the
library writes the row at once, in one transaction, keyed where the scan
will list the file; a composite with no time of its own takes its
sources' earliest. The grid reloads and shows it beside its sources.

The upload then gives the row what only the server knows: after sending
a file the catalog already has a row for, it lists the folder once, takes
the file id the server assigned, and records it (and, once the merge
makes them, the thumbnails waiting beside the payload) under that id.
Every drain now rescans when something landed in the library, after the
upload rather than beside it. A second merge of the same frames is no
longer named over the first: the name is checked against the catalog's
names in that folder, which is all that knows it once the outbox is empty.
This commit is contained in:
2026-09-28 19:57:37 -04:00
parent 37136f7377
commit 98a67393d9
11 changed files with 1335 additions and 89 deletions
+576 -22
View File
@@ -146,6 +146,20 @@ impl Pending {
}
}
/// TRACES: FR-MRG-6
/// Where a file staged for `remote_dir` of the library at `root`, called
/// `name`, will be once it is uploaded — spelled as the scan will list it,
/// because the catalog row written ahead of the scan is keyed on it.
pub(crate) fn staged_remote_path(root: &str, remote_dir: &str, name: &str) -> RemotePath {
Pending {
local: PathBuf::new(),
remote_dir: remote_dir.to_string(),
account: false,
name: name.to_string(),
}
.remote_path(root)
}
/// Where an export was put, for the interface to report.
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum Placed {
@@ -332,6 +346,75 @@ pub(crate) fn destination_record(payload: &Path) -> PathBuf {
))
}
/// TRACES: FR-MRG-6
/// Where the thumbnails made with a staged payload wait for its upload:
/// `name.ext.thumbs` beside it.
///
/// Not a record: [`pending`] reads only `.dest` files, so this never goes to
/// the server. The upload reads it once the server has named the file
/// ([`take_thumbnails`]) and puts its contents in the thumbnail store under
/// that name; until then there is no key to store them under.
pub(crate) fn thumbnails_record(payload: &Path) -> PathBuf {
let mut name = payload.as_os_str().to_owned();
name.push(".thumbs");
PathBuf::from(name)
}
/// The magic line a thumbnails record starts with.
const THUMBS_MAGIC: &[u8] = b"darkroom-thumbs 1\n";
/// Write the thumbnails for a staged payload. A size-class byte and three
/// little-endian `u32`s — width, height, length — before each JPEG.
pub(crate) fn write_thumbnails(
payload: &Path,
thumbnails: &[(dr_thumbs::ThumbSize, dr_thumbs::Thumbnail)],
) -> std::io::Result<()> {
let mut out = THUMBS_MAGIC.to_vec();
for (size, t) in thumbnails {
out.push(*size as i64 as u8);
for v in [t.width, t.height, t.bytes.len() as u32] {
out.extend_from_slice(&v.to_le_bytes());
}
out.extend_from_slice(&t.bytes);
}
std::fs::write(thumbnails_record(payload), out)
}
/// Read and remove the thumbnails waiting beside a payload. Nothing, for a
/// payload that has none or a record that does not parse — a composite whose
/// thumbnails are lost is thumbnailed the ordinary way.
pub(crate) fn take_thumbnails(payload: &Path) -> Vec<(dr_thumbs::ThumbSize, dr_thumbs::Thumbnail)> {
let path = thumbnails_record(payload);
let Ok(bytes) = std::fs::read(&path) else {
return Vec::new();
};
let _ = std::fs::remove_file(&path);
parse_thumbnails(&bytes).unwrap_or_default()
}
fn parse_thumbnails(bytes: &[u8]) -> Option<Vec<(dr_thumbs::ThumbSize, dr_thumbs::Thumbnail)>> {
let mut rest = bytes.strip_prefix(THUMBS_MAGIC)?;
let mut out = Vec::new();
while let Some((&class, tail)) = rest.split_first() {
let word = |at: usize| -> Option<u32> {
Some(u32::from_le_bytes(tail.get(at..at + 4)?.try_into().ok()?))
};
let (width, height, len) = (word(0)?, word(4)?, word(8)? as usize);
let data = tail.get(12..12 + len)?;
let size = dr_thumbs::ThumbSize::from_stored(i64::from(class))?;
out.push((
size,
dr_thumbs::Thumbnail {
width,
height,
bytes: data.to_vec(),
},
));
rest = &tail[12 + len..];
}
Some(out)
}
/// Everything currently waiting in the outbox.
///
/// A payload with no record is skipped rather than guessed at — see the write
@@ -381,11 +464,12 @@ pub fn pending_count(outbox: &Path) -> usize {
pending(outbox).len()
}
/// Remove an entry and its record, once it is safely on the server.
/// Remove an entry and its records, once it is safely on the server.
fn clear(entry: &Pending) {
let record = PathBuf::from(format!("{}.dest", entry.local.display()));
let _ = std::fs::remove_file(&entry.local);
let _ = std::fs::remove_file(&record);
let _ = std::fs::remove_file(thumbnails_record(&entry.local));
}
/// Progress from the upload worker.
@@ -393,13 +477,170 @@ fn clear(entry: &Pending) {
pub enum UploadMessage {
Status(String),
/// Uploaded, still pending, and the first error if there was one.
/// `landed` counts the uploads that went into the library itself rather
/// than an album outside it: those are what a scan has something to
/// find.
Finished {
uploaded: usize,
remaining: usize,
error: Option<String>,
landed: usize,
},
}
/// Log a drain as it goes and, once files have landed in the library, ask
/// the grid to rescan — the scan that finds them, run after the upload
/// rather than beside it (FR-MRG-6).
///
/// On a thread of its own, since the drain reports only when it finishes and
/// that can be minutes for a composite; the window is reached through the
/// event loop.
pub(crate) fn watch_upload(rx: Receiver<UploadMessage>, window: slint::Weak<AppWindow>) {
executors::spawn(Executor::Io, "upload-log", move || {
while let Ok(msg) = rx.recv() {
match msg {
UploadMessage::Status(s) => log::info!("export: {s}"),
UploadMessage::Finished {
uploaded,
remaining,
error,
landed,
} => {
log::info!("export: {uploaded} uploaded, {remaining} still queued");
if let Some(e) = error {
log::warn!("export upload stopped: {e}");
}
if landed > 0 {
let _ = window.upgrade_in_event_loop(|w| {
w.global::<Library>().invoke_library_rescan();
});
}
}
}
}
});
}
/// Where a library's catalog and thumbnail store are, for the upload to
/// record what the server made of a file.
pub(crate) struct LibraryFiles {
pub catalog: PathBuf,
pub thumbs: PathBuf,
}
impl LibraryFiles {
fn of(account: &Account) -> Self {
LibraryFiles {
catalog: crate::library::catalog_path(account),
thumbs: crate::library::thumbs_dir(account),
}
}
}
/// What became of one outbox entry.
#[derive(Debug, PartialEq, Eq)]
enum Sent {
Uploaded,
/// Its payload had gone; the record went with it.
Gone,
}
/// Send one outbox entry, and clear it once it is on the server.
async fn send(
backend: &dyn dr_sync::RemoteBackend,
library: &LibraryFiles,
root: &str,
entry: &Pending,
) -> Result<Sent, String> {
let Ok(bytes) = std::fs::read(&entry.local) else {
// The payload vanished under us. Drop the record too; retrying
// forever against a file that is gone helps nobody.
clear(entry);
return Ok(Sent::Gone);
};
// The folder may not exist — this is the first export into it — and
// `create_dir` treats "already there" as success, so it is unconditional
// rather than guarded by a check that would cost a request every time.
backend
.create_dir(&entry.remote_folder(root))
.await
.map_err(|e| e.to_string())?;
backend
.put(&entry.remote_path(root), bytes, None)
.await
.map_err(|e| e.to_string())?;
let thumbnails = take_thumbnails(&entry.local);
register_upload(backend, library, root, entry, thumbnails).await;
clear(entry);
Ok(Sent::Uploaded)
}
/// TRACES: FR-MRG-6
/// Tell the catalog what the server made of a file it has just been given.
///
/// Only for a file the catalog already has a row for — a composite the merge
/// catalogued before its upload. Anything else is the scan's to find, and
/// this spends no request on it. The listing is one request for the folder,
/// and it is the only way to learn the id a server assigns on upload: the
/// thumbnail store keys on it, and the thumbnails made during the merge wait
/// beside the payload until it is known.
async fn register_upload(
backend: &dyn dr_sync::RemoteBackend,
library: &LibraryFiles,
root: &str,
entry: &Pending,
thumbnails: Vec<(dr_thumbs::ThumbSize, dr_thumbs::Thumbnail)>,
) {
if entry.account {
return;
}
let path = entry.remote_path(root);
let catalog = match dr_catalog::Catalog::open(&library.catalog) {
Ok(c) => c,
Err(e) => {
log::warn!("upload: opening the catalog to record {}: {e}", entry.name);
return;
}
};
let Some(image) = crate::library::image_at(&catalog, root, path.as_str()) else {
return;
};
let listed = match backend.list(&entry.remote_folder(root), None).await {
Ok(entries) => entries.into_iter().find(|e| e.path == path),
Err(e) => {
log::warn!("upload: listing {} after sending it: {e}", entry.name);
None
}
};
let Some(listed) = listed else {
log::warn!("upload: {} is not listed where it was sent", entry.name);
return;
};
if let Err(e) = crate::library::record_uploaded(&catalog, image, &listed) {
log::warn!("upload: recording {}: {e}", entry.name);
}
let dr_sync::RemoteId::Stable(file_id) = listed.id else {
return;
};
if thumbnails.is_empty() {
return;
}
match dr_thumbs::ThumbStore::open(&library.thumbs) {
Ok(mut store) => {
for (size, thumb) in &thumbnails {
crate::library::store_thumbnail(&mut store, file_id, *size, thumb);
}
log::info!(
"upload: {} thumbnail(s) of {} stored under file {file_id}",
thumbnails.len(),
entry.name
);
}
Err(e) => log::warn!("upload: opening the thumbnail store: {e}"),
}
}
/// Held by whichever drain is sending the outbox. See [`spawn_upload`].
static OUTBOX_DRAIN: std::sync::Mutex<()> = std::sync::Mutex::new(());
@@ -436,6 +677,7 @@ pub fn spawn_upload(
uploaded: 0,
remaining: pending_count(&outbox),
error: Some(e.to_string()),
landed: 0,
});
return;
}
@@ -449,14 +691,17 @@ pub fn spawn_upload(
uploaded: 0,
remaining: pending_count(&outbox),
error: Some(e.to_string()),
landed: 0,
});
return;
}
};
let library = LibraryFiles::of(&conn.account);
let queue = pending(&outbox);
let total = queue.len();
let mut uploaded = 0;
let mut landed = 0;
let mut error = None;
for (i, entry) in queue.iter().enumerate() {
@@ -466,30 +711,16 @@ pub fn spawn_upload(
i + 1
)));
let Ok(bytes) = std::fs::read(&entry.local) else {
// The payload vanished under us. Drop the record too;
// retrying forever against a file that is gone helps
// nobody.
clear(entry);
continue;
};
// The folder may not exist — this is the first export into it
// — and `create_dir` treats "already there" as success, so it
// is unconditional rather than guarded by a check that would
// cost a request every time.
if let Err(e) = backend.create_dir(&entry.remote_folder(&root)).await {
error = Some(e.to_string());
break;
}
match backend.put(&entry.remote_path(&root), bytes, None).await {
Ok(_) => {
clear(entry);
match send(&*backend, &library, &root, entry).await {
Ok(Sent::Gone) => {}
Ok(Sent::Uploaded) => {
uploaded += 1;
if !entry.account {
landed += 1;
}
}
Err(e) => {
error = Some(e.to_string());
error = Some(e);
break;
}
}
@@ -499,6 +730,7 @@ pub fn spawn_upload(
uploaded,
remaining: pending_count(&outbox),
error,
landed,
});
});
});
@@ -1295,6 +1527,328 @@ fn stop_timer(slot: &Rc<std::cell::RefCell<Option<slint::Timer>>>) {
mod tests {
use super::*;
/// A server that names a file when it is given one, as Nextcloud does:
/// the id exists only once the upload has landed, and a listing is the
/// way to learn it.
#[derive(Default)]
struct Assigning {
files: std::sync::Mutex<std::collections::BTreeMap<String, (u64, Vec<u8>)>>,
lists: std::sync::atomic::AtomicUsize,
caps: std::sync::OnceLock<dr_sync::Capabilities>,
}
impl Assigning {
fn listing(&self, dir: &str) -> Vec<dr_sync::RemoteEntry> {
self.files
.lock()
.unwrap()
.iter()
.filter(|(p, _)| p.rsplit_once('/').map_or("", |(d, _)| d) == dir)
.map(|(p, (id, body))| dr_sync::RemoteEntry {
id: dr_sync::RemoteId::Stable(*id),
path: RemotePath::new(p.clone()),
kind: dr_sync::EntryKind::File,
validator: dr_sync::Validator::new(format!("etag-{id}")),
size: body.len() as u64,
modified: None,
has_preview: false,
materialised: true,
})
.collect()
}
}
#[async_trait::async_trait]
impl dr_sync::RemoteBackend for Assigning {
fn capabilities(&self) -> &dr_sync::Capabilities {
self.caps.get_or_init(dr_sync::Capabilities::minimal)
}
fn name(&self) -> &str {
"assigning"
}
async fn list(
&self,
dir: &RemotePath,
_since: Option<&dr_sync::Validator>,
) -> Result<Vec<dr_sync::RemoteEntry>, dr_sync::RemoteError> {
self.lists.fetch_add(1, Ordering::SeqCst);
Ok(self.listing(dir.as_str()))
}
async fn dir_validator(
&self,
_dir: &RemotePath,
) -> Result<dr_sync::Validator, dr_sync::RemoteError> {
Err(dr_sync::RemoteError::Unsupported("test"))
}
async fn delta(
&self,
_c: &dr_sync::Cursor,
) -> Result<(Vec<dr_sync::RemoteChange>, dr_sync::Cursor), dr_sync::RemoteError> {
Err(dr_sync::RemoteError::Unsupported("test"))
}
async fn get(
&self,
_id: &dr_sync::RemoteId,
_r: Option<std::ops::Range<u64>>,
) -> Result<Vec<u8>, dr_sync::RemoteError> {
Err(dr_sync::RemoteError::Unsupported("test"))
}
async fn put(
&self,
path: &RemotePath,
body: Vec<u8>,
_pc: Option<dr_sync::Precondition>,
) -> Result<dr_sync::Validator, dr_sync::RemoteError> {
let mut files = self.files.lock().unwrap();
let id = files
.get(path.as_str())
.map(|(id, _)| *id)
.unwrap_or(40_000 + files.len() as u64);
files.insert(path.as_str().to_string(), (id, body));
Ok(dr_sync::Validator::new(format!("etag-{id}")))
}
async fn delete(
&self,
_id: &dr_sync::RemoteId,
_pc: Option<dr_sync::Precondition>,
) -> Result<(), dr_sync::RemoteError> {
Ok(())
}
async fn move_to(
&self,
_f: &dr_sync::RemoteId,
_t: &RemotePath,
) -> Result<(), dr_sync::RemoteError> {
Ok(())
}
async fn create_dir(&self, _p: &RemotePath) -> Result<(), dr_sync::RemoteError> {
Ok(())
}
}
/// A library of one sweep, its catalog on disk, a store beside it, and
/// an outbox holding the composite the merge staged: the row written
/// ahead of the scan, the payload, its record, and its thumbnails.
struct Staged {
dir: PathBuf,
library: LibraryFiles,
entry: Pending,
image: i64,
}
fn staged(name: &str, root: &str) -> Staged {
let dir =
std::env::temp_dir().join(format!("dr-composite-upload-{name}-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
std::fs::create_dir_all(dir.join("outbox")).unwrap();
let library = LibraryFiles {
catalog: dir.join("catalog.sqlite"),
thumbs: dir.join("thumbs"),
};
let at = |n: &str| {
if root.is_empty() {
format!("Alps/{n}")
} else {
format!("{root}/Alps/{n}")
}
};
let catalog = dr_catalog::Catalog::open(&library.catalog).unwrap();
crate::library::test_support::scanned_listing(
&catalog,
root,
vec![crate::library::test_support::entry(
&at("_MG_8320.CR2"),
1,
30,
)],
);
let image = crate::library::catalogue_composite(
&catalog,
root,
&crate::library::CompositeRow {
source_ref: staged_remote_path(root, "Alps", "_MG_8320-pano.dng")
.as_str()
.to_string(),
width: 4000,
height: 1000,
captured_at: Some(1_000),
captured_offset: None,
camera: None,
lens: None,
iso: None,
file_size: 5,
sources: vec![at("_MG_8320.CR2")],
},
)
.unwrap();
// The scan that ran while the upload was still going: the folder is
// listed without the composite, and the row is still there after it.
crate::library::test_support::scanned_listing(
&catalog,
root,
vec![crate::library::test_support::entry(
&at("_MG_8320.CR2"),
1,
30,
)],
);
let local = dir.join("outbox").join("_MG_8320-pano.dng");
std::fs::write(&local, b"pano!").unwrap();
std::fs::write(destination_record(&local), "Alps\n_MG_8320-pano.dng\n").unwrap();
let jpeg = dr_thumbs::encode_rgba(4, 1, &[200u8; 16]).unwrap();
write_thumbnails(
&local,
&[
(
dr_thumbs::ThumbSize::Grid,
dr_thumbs::Thumbnail {
width: 4,
height: 1,
bytes: jpeg.clone(),
},
),
(
dr_thumbs::ThumbSize::Large,
dr_thumbs::Thumbnail {
width: 4,
height: 1,
bytes: jpeg,
},
),
],
)
.unwrap();
let entry = pending(&dir.join("outbox")).pop().unwrap();
Staged {
dir,
library,
entry,
image,
}
}
fn file_id_of(library: &LibraryFiles, image: i64) -> Option<i64> {
dr_catalog::Catalog::open(&library.catalog)
.unwrap()
.connection()
.query_row(
"SELECT file_id FROM remote WHERE image_id = ?1",
[image],
|r| r.get(0),
)
.ok()
}
#[tokio::test]
async fn an_upload_that_lands_after_the_scan_gives_the_row_its_server_id() {
let s = staged("nextcloud", "PhotosRaw");
let server = Assigning::default();
assert_eq!(file_id_of(&s.library, s.image), None);
let sent = send(&server, &s.library, "PhotosRaw", &s.entry)
.await
.unwrap();
assert_eq!(sent, Sent::Uploaded);
// The id the server assigned, on the row the merge wrote, and the
// merge's thumbnails in the store under it.
let listed = server.listing("PhotosRaw/Alps");
let dr_sync::RemoteId::Stable(id) = listed[0].id else {
panic!("the server names its files");
};
assert_eq!(file_id_of(&s.library, s.image), Some(id as i64));
let store = dr_thumbs::ThumbStore::open(&s.library.thumbs).unwrap();
assert!(store.contains(id, dr_thumbs::ThumbSize::Grid));
assert!(store.contains(id, dr_thumbs::ThumbSize::Large));
// Nothing left in the outbox: not the payload, not its records.
assert_eq!(
std::fs::read_dir(s.dir.join("outbox")).unwrap().count(),
0,
"the outbox is empty"
);
// The scan that follows the upload finds the row already there.
let catalog = dr_catalog::Catalog::open(&s.library.catalog).unwrap();
crate::library::test_support::scanned_listing(&catalog, "PhotosRaw", listed);
let n: i64 = catalog
.connection()
.query_row(
"SELECT count(*) FROM images WHERE source_ref LIKE '%pano%'",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(n, 1);
let _ = std::fs::remove_dir_all(&s.dir);
}
#[tokio::test]
async fn a_folder_library_learns_the_id_the_folder_gives_its_path() {
let s = staged("folder", "");
let library_dir = s.dir.join("library");
std::fs::create_dir_all(library_dir.join("Alps")).unwrap();
let folder = dr_sync_folder::FolderBackend::new(&library_dir).unwrap();
send(&folder, &s.library, "", &s.entry).await.unwrap();
assert_eq!(
std::fs::read(library_dir.join("Alps/_MG_8320-pano.dng")).unwrap(),
b"pano!"
);
use dr_sync::RemoteBackend as _;
let listed = folder
.list(&RemotePath::new("Alps"), None)
.await
.unwrap()
.into_iter()
.find(|e| e.path.name() == "_MG_8320-pano.dng")
.unwrap();
let dr_sync::RemoteId::Stable(id) = listed.id else {
panic!("a folder names a file by its path");
};
assert_eq!(file_id_of(&s.library, s.image), Some(id as i64));
let store = dr_thumbs::ThumbStore::open(&s.library.thumbs).unwrap();
assert!(store.contains(id, dr_thumbs::ThumbSize::Grid));
let _ = std::fs::remove_dir_all(&s.dir);
}
#[tokio::test]
async fn an_ordinary_export_costs_no_listing() {
let s = staged("export", "PhotosRaw");
let other = s.dir.join("outbox").join("print.jpg");
std::fs::write(&other, b"jpeg").unwrap();
std::fs::write(destination_record(&other), "Alps\nprint.jpg\n").unwrap();
let entry = pending(&s.dir.join("outbox"))
.into_iter()
.find(|p| p.name == "print.jpg")
.unwrap();
let server = Assigning::default();
send(&server, &s.library, "PhotosRaw", &entry)
.await
.unwrap();
assert_eq!(server.lists.load(Ordering::SeqCst), 0);
let _ = std::fs::remove_dir_all(&s.dir);
}
#[test]
fn thumbnails_survive_the_outbox_and_a_damaged_record_is_nothing() {
let dir = tmp();
let payload = dir.join("x.dng");
let t = dr_thumbs::Thumbnail {
width: 3,
height: 2,
bytes: vec![1, 2, 3, 4, 5],
};
write_thumbnails(&payload, &[(dr_thumbs::ThumbSize::Large, t.clone())]).unwrap();
assert_eq!(
take_thumbnails(&payload),
vec![(dr_thumbs::ThumbSize::Large, t)]
);
assert!(!thumbnails_record(&payload).exists(), "taken, not copied");
std::fs::write(thumbnails_record(&payload), b"darkroom-thumbs 1\n\x00\x01").unwrap();
assert!(take_thumbnails(&payload).is_empty());
}
fn encoded(name: &str, bytes: &[u8]) -> Encoded {
Encoded {
name: name.to_string(),