//! TRACES: FR-MRG-6 //! A merged panorama in the catalog the moment it is written, rather than //! whenever a scan next happens to find it. //! //! # Why the merge catalogues its own file //! //! The composite reaches the library through the outbox, like an export, //! and the grid used to learn of it the way it learns of anything: a rescan, //! fired as the merge finished. That scan raced the upload it followed. On a //! folder library the copy of an 800 MB file was still running when the //! folder was listed; on Nextcloud the upload takes minutes. Either way the //! listing did not have the file, the folder's validator was recorded as it //! stood, and nothing looked again until the next sync pass — the panorama //! the photographer had just watched being made was not in the grid. //! //! So the job hands the library what it knows — the name the file will have, //! the size of the picture, when it was taken — and the row is written at //! once, in one transaction. The upload later supplies what only the server //! knows, its file id, through [`record_uploaded`]; the scan that follows //! finds a row already there and updates it in place, on the same //! `(root_id, source_ref)` key it would have inserted under. use dr_catalog::Catalog; use super::scan::now_secs; /// What the catalog is told about a composite before any scan has seen it. #[derive(Debug, Clone, PartialEq)] pub struct CompositeRow { /// Where it will be in the library, spelled as the scan spells it. pub source_ref: String, /// The picture it opens on — the crop, where the merge cropped — upright. pub width: u32, pub height: u32, /// When it was taken: the middle of its sweep, as the DNG says. `None` /// where no frame carried a time, and then the earliest of `sources` as /// the catalog has them. pub captured_at: Option, pub captured_offset: Option, pub camera: Option, pub lens: Option, pub iso: Option, pub file_size: u64, /// The frames it was merged from, as source refs. pub sources: Vec, } /// Write the composite's row, and return its image id. /// /// One transaction: the root, the capture time the sources lend it if it has /// none of its own, and the row. At `metadata_state = 2`, since everything a /// header read would add is already here — the sweep has no reason to fetch /// the tail of an 800 MB file for a date the merge wrote into it. /// /// A row already under the name is brought up to date rather than /// duplicated. The merge picks a name the catalog does not hold, so that is /// a safety net: the file at that name is the one just written. pub fn catalogue_composite( catalog: &Catalog, root: &str, row: &CompositeRow, ) -> Result { let conn = catalog.connection(); let tx = conn.unchecked_transaction()?; tx.execute( "INSERT INTO roots(kind, label, last_seen) VALUES ('remote', ?1, ?2) ON CONFLICT DO NOTHING", rusqlite::params![root, now_secs()], )?; let root_id: i64 = tx.query_row( "SELECT id FROM roots WHERE label = ?1 AND kind = 'remote'", [root], |r| r.get(0), )?; // The folder, if a scan has recorded it; the sources are in it, so it // almost always has. `None` is what the scan itself writes for a file at // the root, and the next scan fills it either way. let folder_id: Option = match row.source_ref.rsplit_once('/') { Some((parent, _)) => tx .query_row( "SELECT id FROM folders WHERE root_id = ?1 AND path = ?2", rusqlite::params![root_id, parent], |r| r.get(0), ) .ok(), None => None, }; // The sources' earliest, where the composite carries no time. One // aggregate over the handful of frames, by their key. let captured_at = match row.captured_at { Some(t) => Some(t), None if row.sources.is_empty() => None, None => { let placeholders = std::iter::repeat_n("?", row.sources.len()) .collect::>() .join(","); let mut params: Vec = vec![rusqlite::types::Value::Integer(root_id)]; params.extend( row.sources .iter() .map(|s| rusqlite::types::Value::Text(s.clone())), ); tx.query_row( &format!( "SELECT min(captured_at) FROM images WHERE root_id = ? AND source_ref IN ({placeholders})" ), rusqlite::params_from_iter(params.iter()), |r| r.get::<_, Option>(0), )? } }; let format = row .source_ref .rsplit_once('.') .map(|(_, e)| e.to_ascii_lowercase()); tx.execute( "INSERT INTO images(root_id, folder_id, source_ref, format, w, h, captured_at, captured_offset, camera, lens, iso, file_size, availability, metadata_state, added_at) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, 0, 2, ?13) ON CONFLICT(root_id, source_ref) DO UPDATE SET folder_id = coalesce(excluded.folder_id, images.folder_id), w = excluded.w, h = excluded.h, captured_at = coalesce(excluded.captured_at, images.captured_at), captured_offset = coalesce(excluded.captured_offset, images.captured_offset), camera = coalesce(excluded.camera, images.camera), lens = coalesce(excluded.lens, images.lens), iso = coalesce(excluded.iso, images.iso), file_size = excluded.file_size, metadata_state = 2", rusqlite::params![ root_id, folder_id, row.source_ref, format, row.width as i64, row.height as i64, captured_at, row.captured_offset, row.camera, row.lens, row.iso.map(i64::from), row.file_size as i64, now_secs(), ], )?; let image_id: i64 = tx.query_row( "SELECT id FROM images WHERE root_id = ?1 AND source_ref = ?2", rusqlite::params![root_id, row.source_ref], |r| r.get(0), )?; tx.commit()?; Ok(image_id) } /// The image a library path names, if the catalog has it. What the upload /// asks before it spends a listing on a file: a file the catalog has no row /// for is the scan's to find, and only a row written ahead of the scan is /// waiting for its server identity. pub fn image_at(catalog: &Catalog, root: &str, source_ref: &str) -> Option { catalog .connection() .query_row( "SELECT i.id FROM images i JOIN roots r ON r.id = i.root_id WHERE r.label = ?1 AND r.kind = 'remote' AND i.source_ref = ?2", rusqlite::params![root, source_ref], |r| r.get(0), ) .ok() } /// The names the catalog holds directly in `folder` of the library at /// `root`, whatever their state — what a new file there must not be called. /// /// A range over the `(root_id, source_ref)` key, `/` to `0` being the byte /// after it, so a folder of a thousand costs a seek and a thousand index /// entries, not the library. pub fn names_in_folder( catalog: &Catalog, root: &str, folder: &str, ) -> std::collections::HashSet { let folder = folder.trim_matches('/'); let (from, to) = if folder.is_empty() { (String::new(), "\u{10ffff}".to_string()) } else { (format!("{folder}/"), format!("{folder}0")) }; let Ok(mut stmt) = catalog.connection().prepare( "SELECT i.source_ref FROM images i JOIN roots r ON r.id = i.root_id WHERE r.label = ?1 AND r.kind = 'remote' AND i.source_ref >= ?2 AND i.source_ref < ?3", ) else { return Default::default(); }; let Ok(rows) = stmt.query_map(rusqlite::params![root, from, to], |r| r.get::<_, String>(0)) else { return Default::default(); }; rows.flatten() .filter_map(|p| { let name = p.get(from.len()..)?; (!name.contains('/')).then(|| name.to_string()) }) .collect() } /// Give a catalogued file the identity the server assigned it on upload: /// its file id, which the thumbnail store keys on, and its validator. /// /// The same row the scan's `remote` upsert writes, so whichever of the two /// runs second agrees with the first. pub fn record_uploaded( catalog: &Catalog, image_id: i64, entry: &dr_sync::RemoteEntry, ) -> Result<(), dr_catalog::CatalogError> { let conn = catalog.connection(); let tx = conn.unchecked_transaction()?; if let dr_sync::RemoteId::Stable(file_id) = entry.id { tx.execute( "INSERT INTO remote(image_id, file_id, etag, remote_path) VALUES (?1, ?2, ?3, ?4) ON CONFLICT(image_id) DO UPDATE SET file_id = excluded.file_id, etag = excluded.etag, remote_path = excluded.remote_path", rusqlite::params![ image_id, file_id as i64, entry.validator.as_str(), entry.path.as_str() ], )?; } tx.execute( "UPDATE images SET file_size = ?2 WHERE id = ?1 AND file_size IS NOT ?2", rusqlite::params![image_id, entry.size as i64], )?; tx.commit()?; Ok(()) } #[cfg(test)] mod tests { use super::*; use crate::library::read_cells; use crate::library::scan::persist; use crate::library::test_support::entry; use dr_sync::RemotePath; /// Twelve frames of a sweep a second apart, scanned and dated. fn a_sweep() -> Catalog { let catalog = Catalog::in_memory().unwrap(); let images = (0..12) .map(|i| { entry( &format!("Alps/_MG_{:04}.CR2", 8320 + i), 100 + i, 30_000_000, ) }) .chain(std::iter::once(entry("Alps/later.CR2", 200, 30_000_000))) .collect(); persist( &catalog, "", &dr_sync::ScanResult { images, directories: vec![(RemotePath::new("Alps"), dr_sync::Validator::new("e1"))], progress: Default::default(), sidecars: Vec::new(), }, ) .unwrap(); for i in 0..12 { catalog .connection() .execute( "UPDATE images SET captured_at = ?2, metadata_state = 2 WHERE source_ref = ?1", rusqlite::params![format!("Alps/_MG_{:04}.CR2", 8320 + i), 1_000 + i], ) .unwrap(); } catalog .connection() .execute( "UPDATE images SET captured_at = 5000 WHERE source_ref = 'Alps/later.CR2'", [], ) .unwrap(); catalog } fn composite(captured_at: Option) -> CompositeRow { CompositeRow { source_ref: "Alps/_MG_8320-pano.dng".into(), width: 22_000, height: 5_600, captured_at, captured_offset: None, camera: Some("Canon EOS 6D".into()), lens: None, iso: Some(100), file_size: 800_000_000, sources: (0..12) .map(|i| format!("Alps/_MG_{:04}.CR2", 8320 + i)) .collect(), } } fn names(catalog: &Catalog) -> Vec { read_cells(catalog, 0, 100) .unwrap() .into_iter() .map(|c| c.name) .collect() } #[test] fn a_composite_sits_among_its_sources_the_moment_it_is_written() { let catalog = a_sweep(); catalogue_composite(&catalog, "", &composite(Some(1_005))).unwrap(); let names = names(&catalog); let at = names.iter().position(|n| n == "_MG_8320-pano.dng").unwrap(); // Dated at the middle of the sweep, so it is among the frames and not // after the photograph taken an hour later. assert!(at > 0 && at < 12, "{names:?}"); assert_eq!(names.last().unwrap(), "later.CR2"); } #[test] fn a_composite_with_no_time_takes_its_earliest_source() { let catalog = a_sweep(); let id = catalogue_composite(&catalog, "", &composite(None)).unwrap(); let (t, w, h, state): (Option, i64, i64, i64) = catalog .connection() .query_row( "SELECT captured_at, w, h, metadata_state FROM images WHERE id = ?1", [id], |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?, r.get(3)?)), ) .unwrap(); assert_eq!(t, Some(1_000)); assert_eq!((w, h, state), (22_000, 5_600, 2)); } #[test] fn a_scan_that_ran_before_the_upload_does_not_lose_the_row() { // Nextcloud: the merge finishes, the row is written, and a sync pass // lists the folder before the upload has finished — the file is not // on the server yet. let catalog = a_sweep(); let id = catalogue_composite(&catalog, "", &composite(Some(1_005))).unwrap(); persist( &catalog, "", &dr_sync::ScanResult { images: vec![entry("Alps/later.CR2", 200, 30_000_000)], directories: vec![(RemotePath::new("Alps"), dr_sync::Validator::new("e2"))], progress: Default::default(), sidecars: Vec::new(), }, ) .unwrap(); assert!(names(&catalog).contains(&"_MG_8320-pano.dng".to_string())); // The upload lands and the server names it; then the scan finds it. let listed = entry("Alps/_MG_8320-pano.dng", 9_001, 812_345_678); record_uploaded(&catalog, id, &listed).unwrap(); persist( &catalog, "", &dr_sync::ScanResult { images: vec![listed], directories: vec![(RemotePath::new("Alps"), dr_sync::Validator::new("e3"))], progress: Default::default(), sidecars: Vec::new(), }, ) .unwrap(); let rows: Vec<(i64, Option, i64, Option)> = catalog .connection() .prepare( "SELECT i.id, r.file_id, i.file_size, i.captured_at FROM images i LEFT JOIN remote r ON r.image_id = i.id WHERE i.source_ref = 'Alps/_MG_8320-pano.dng'", ) .unwrap() .query_map([], |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?, r.get(3)?))) .unwrap() .map(Result::unwrap) .collect(); assert_eq!(rows, vec![(id, Some(9_001), 812_345_678, Some(1_005))]); let cell = read_cells(&catalog, 0, 100) .unwrap() .into_iter() .find(|c| c.image_id == id) .unwrap(); assert_eq!(cell.file_id, Some(9_001)); } #[test] fn merging_again_under_the_same_name_updates_the_one_row() { let catalog = a_sweep(); let first = catalogue_composite(&catalog, "", &composite(Some(1_005))).unwrap(); let mut again = composite(Some(1_006)); again.width = 21_000; let second = catalogue_composite(&catalog, "", &again).unwrap(); assert_eq!(first, second); let w: i64 = catalog .connection() .query_row("SELECT w FROM images WHERE id = ?1", [first], |r| r.get(0)) .unwrap(); assert_eq!(w, 21_000); } #[test] fn the_names_in_a_folder_are_its_own_and_not_its_subfolders() { let catalog = a_sweep(); catalogue_composite(&catalog, "", &composite(Some(1_005))).unwrap(); persist( &catalog, "", &dr_sync::ScanResult { images: vec![ entry("Alps/2019/deep.CR2", 300, 1), entry("Alpsine/near.CR2", 301, 1), ], directories: Vec::new(), progress: Default::default(), sidecars: Vec::new(), }, ) .unwrap(); let names = names_in_folder(&catalog, "", "Alps"); assert!(names.contains("_MG_8320-pano.dng")); assert!(names.contains("_MG_8320.CR2")); assert!(!names .iter() .any(|n| n.contains("deep") || n.contains("near"))); assert_eq!(names.len(), 14); } #[test] fn only_a_row_the_merge_wrote_is_found_for_the_upload() { let catalog = a_sweep(); assert!(image_at(&catalog, "", "Alps/_MG_8320-pano.dng").is_none()); let id = catalogue_composite(&catalog, "", &composite(Some(1_005))).unwrap(); assert_eq!(image_at(&catalog, "", "Alps/_MG_8320-pano.dng"), Some(id)); assert!(image_at(&catalog, "elsewhere", "Alps/_MG_8320-pano.dng").is_none()); } }