Hand the photographer's place between devices
A place recorded on the tablet should be where the desktop opens. Exchanged through `.darkroom-derived/place.json`, beside the thumbnail shards and the catalog snapshot. Newest timestamp wins outright: unlike the catalog this is replaced rather than merged, because two devices cannot both be where the photographer is and so there is nothing of theirs inside ours to preserve. It still refuses to upload over a copy it could not read, for a smaller version of the reason `sync_catalog` does: a record we have not compared against may be the newer one, and overwriting it would move the other device's photographer without ever having seen where they were. Last in the pass, and its failures are logged rather than reported. Everything else in that folder is *derived* -- a faster way to learn what the device could work out for itself -- so losing it costs time. A place is a fact only the other device knew, and losing it costs a scroll. A sync that ran out of connectivity should spend what it had on the shards. The full pass runs after a thumbnail sweep or when Sync is pressed, neither of which happens on an ordinary launch -- so a handover would arrive one launch late, which is one too many for a feature whose whole claim is picking up where you stopped. `spawn_place_fetch` is the small half: one GET of a few hundred bytes, started beside the scan. And it can still be refused. A handover is welcome on the way in and unwelcome once the photographer has started: a grid that jumped elsewhere mid-scroll because a round trip finally landed would have lost their place to the feature meant to keep it. Any scroll, scrub, scope change, filter or opened photograph closes the latch, and a record arriving after that is written to disk and takes effect next launch. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
This commit is contained in:
@@ -70,6 +70,16 @@ pub struct SyncReport {
|
||||
pub face_shards_downloaded: usize,
|
||||
/// Images whose faces this device took from a peer instead of detecting.
|
||||
pub faces_adopted: usize,
|
||||
|
||||
/// TRACES: FR-UI-8
|
||||
/// Whether the exchange moved where this device thinks the photographer is.
|
||||
///
|
||||
/// Counted apart from everything above because it is the one thing here
|
||||
/// that is not derived state: a shard or a catalog snapshot is a faster way
|
||||
/// to learn what this device could have worked out for itself, and a place
|
||||
/// is a fact only the other device knew.
|
||||
pub place_adopted: bool,
|
||||
pub place_uploaded: bool,
|
||||
}
|
||||
|
||||
impl SyncReport {
|
||||
@@ -80,6 +90,7 @@ impl SyncReport {
|
||||
|| self.catalog_merged
|
||||
|| self.face_shards_uploaded > 0
|
||||
|| self.face_shards_downloaded > 0
|
||||
|| self.place_adopted
|
||||
}
|
||||
}
|
||||
|
||||
@@ -100,6 +111,7 @@ pub fn spawn_sync(
|
||||
root: String,
|
||||
thumbs_dir: PathBuf,
|
||||
catalog_path: PathBuf,
|
||||
place_path: PathBuf,
|
||||
scratch: PathBuf,
|
||||
) -> std::sync::mpsc::Receiver<SyncMessage> {
|
||||
let (tx, rx) = std::sync::mpsc::channel();
|
||||
@@ -124,7 +136,17 @@ pub fn spawn_sync(
|
||||
}
|
||||
};
|
||||
|
||||
match run(&*backend, &root, &thumbs_dir, &catalog_path, &scratch, &tx).await {
|
||||
match run(
|
||||
&*backend,
|
||||
&root,
|
||||
&thumbs_dir,
|
||||
&catalog_path,
|
||||
&place_path,
|
||||
&scratch,
|
||||
&tx,
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(report) => {
|
||||
let _ = tx.send(SyncMessage::Finished(Box::new(report)));
|
||||
}
|
||||
@@ -143,6 +165,7 @@ async fn run(
|
||||
root: &str,
|
||||
thumbs_dir: &Path,
|
||||
catalog_path: &Path,
|
||||
place_path: &Path,
|
||||
scratch: &Path,
|
||||
tx: &std::sync::mpsc::Sender<SyncMessage>,
|
||||
) -> Result<SyncReport, String> {
|
||||
@@ -162,6 +185,13 @@ async fn run(
|
||||
let _ = tx.send(SyncMessage::Status("checking collections…".into()));
|
||||
sync_catalog(backend, &base, catalog_path, scratch, &mut report).await?;
|
||||
|
||||
// TRACES: FR-UI-8
|
||||
// Last, and it costs one small GET plus at most one small PUT. Last because
|
||||
// it is the only thing here that is not derived state and so the only thing
|
||||
// whose loss the user would not notice: a pass that ran out of connectivity
|
||||
// should spend what it had on the shards and the catalog.
|
||||
sync_place(backend, &base, place_path, &mut report).await;
|
||||
|
||||
Ok(report)
|
||||
}
|
||||
|
||||
@@ -624,6 +654,154 @@ async fn sync_catalog(
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// TRACES: FR-UI-8
|
||||
/// The place's name inside the derived folder.
|
||||
///
|
||||
/// The same name the local copy has, so `.darkroom-derived/place.json` and the
|
||||
/// file beside the catalog are visibly the same thing. There is one of these per
|
||||
/// library, not one per device: the question it answers — where is the
|
||||
/// photographer — has one answer, and a folder of per-device files would need a
|
||||
/// listing and N fetches to work out which of them was current.
|
||||
const PLACE_NAME: &str = "place.json";
|
||||
|
||||
/// TRACES: FR-UI-8
|
||||
/// Exchange the place with the server. The newer record wins, in both
|
||||
/// directions.
|
||||
///
|
||||
/// # Why this cannot clobber the way the catalog could
|
||||
///
|
||||
/// [`sync_catalog`] has to refuse to upload when it cannot read the server's
|
||||
/// copy, because its upload is a read-modify-write: writing without merging
|
||||
/// discards the other device's collections. This is not that. A place is
|
||||
/// *replaced*, never merged, so there is nothing of theirs inside ours to lose.
|
||||
///
|
||||
/// It still gives up on an unreadable read rather than uploading over it, for a
|
||||
/// smaller reason: a server that will not answer a GET is not one to spend a PUT
|
||||
/// on, and a record we could not compare against might be newer than ours —
|
||||
/// overwriting it would move the other device's photographer without ever having
|
||||
/// seen where they were.
|
||||
///
|
||||
/// # Failures are not propagated
|
||||
///
|
||||
/// This returns nothing and takes no `?`. Every other step in [`run`] carries
|
||||
/// state that has to arrive; this one carries a scroll position, and a sync that
|
||||
/// reported itself failed — putting an error in front of the user and skipping
|
||||
/// nothing, since it runs last — because a position file could not be written
|
||||
/// would be reporting the wrong thing entirely.
|
||||
async fn sync_place(
|
||||
backend: &dyn RemoteBackend,
|
||||
base: &RemotePath,
|
||||
place_path: &Path,
|
||||
report: &mut SyncReport,
|
||||
) {
|
||||
let target = RemotePath::new(format!("{}/{PLACE_NAME}", base.as_str()));
|
||||
|
||||
let theirs = match read_derived(backend, &target).await {
|
||||
Ok(bytes) => crate::place::from_bytes(&bytes, "the place on the server"),
|
||||
// Nobody has recorded one for this library yet.
|
||||
Err(RemoteError::NotFound(_)) => None,
|
||||
Err(e) => {
|
||||
log::debug!("not exchanging the place: reading the server's copy: {e}");
|
||||
return;
|
||||
}
|
||||
};
|
||||
|
||||
let store = crate::place::PlaceStore::open_at(place_path.to_path_buf());
|
||||
|
||||
// Take theirs first, so what goes back up is whichever of the two is
|
||||
// current rather than always ours.
|
||||
if let Some(theirs) = theirs.as_ref() {
|
||||
if store.adopt(theirs) {
|
||||
log::info!(
|
||||
"the place moved to where {} left off",
|
||||
if theirs.device.is_empty() {
|
||||
"another device"
|
||||
} else {
|
||||
&theirs.device
|
||||
}
|
||||
);
|
||||
report.place_adopted = true;
|
||||
}
|
||||
}
|
||||
|
||||
// Then push, if there is anything to say. Re-read rather than reusing what
|
||||
// was loaded above: `adopt` may have just replaced it, and uploading a
|
||||
// record the server already has is a PUT for nothing.
|
||||
let Some(ours) = store.load() else {
|
||||
return;
|
||||
};
|
||||
if !ours.supersedes(theirs.as_ref()) {
|
||||
return;
|
||||
}
|
||||
let bytes = match crate::place::to_bytes(&ours) {
|
||||
Ok(b) => b,
|
||||
Err(e) => {
|
||||
log::debug!("serialising the place: {e}");
|
||||
return;
|
||||
}
|
||||
};
|
||||
match backend.put(&target, bytes, None).await {
|
||||
Ok(_) => report.place_uploaded = true,
|
||||
Err(e) => log::debug!("uploading the place: {e}"),
|
||||
}
|
||||
}
|
||||
|
||||
/// TRACES: FR-UI-8
|
||||
/// Fetch just the place, for the handover at launch.
|
||||
///
|
||||
/// # Why this is not simply [`sync_place`]
|
||||
///
|
||||
/// The full pass runs after a thumbnail sweep or when the Sync button is
|
||||
/// pressed, neither of which happens on an ordinary launch — so a place left on
|
||||
/// the tablet would reach the desktop one launch late, which is one launch too
|
||||
/// many for a feature whose whole claim is that you pick up where you stopped.
|
||||
///
|
||||
/// This is the small half: one GET of a few hundred bytes, started beside the
|
||||
/// scan rather than after it, so the answer is usually in hand before the user
|
||||
/// has decided what to look at. It writes nothing and uploads nothing — the
|
||||
/// exchange proper still happens in [`sync_place`], which is where "ours is
|
||||
/// newer, push it" belongs.
|
||||
///
|
||||
/// `Ok(None)` is "there is not one there", which on a first launch against a
|
||||
/// fresh library is the normal answer and not worth a word to the user.
|
||||
pub fn spawn_place_fetch(
|
||||
conn: Connection,
|
||||
root: String,
|
||||
) -> std::sync::mpsc::Receiver<Result<Option<dr_types::Place>, String>> {
|
||||
let (tx, rx) = std::sync::mpsc::channel();
|
||||
|
||||
std::thread::spawn(move || {
|
||||
let rt = match crate::net_runtime::build() {
|
||||
Ok(rt) => rt,
|
||||
Err(e) => {
|
||||
let _ = tx.send(Err(e.to_string()));
|
||||
return;
|
||||
}
|
||||
};
|
||||
|
||||
rt.block_on(async {
|
||||
let backend = match crate::remote::connect(&conn) {
|
||||
Ok(b) => b,
|
||||
Err(e) => {
|
||||
let _ = tx.send(Err(e.to_string()));
|
||||
return;
|
||||
}
|
||||
};
|
||||
|
||||
let base = derived_path(&root);
|
||||
let target = RemotePath::new(format!("{}/{PLACE_NAME}", base.as_str()));
|
||||
let got = match read_derived(&*backend, &target).await {
|
||||
Ok(bytes) => Ok(crate::place::from_bytes(&bytes, "the place on the server")),
|
||||
Err(RemoteError::NotFound(_)) => Ok(None),
|
||||
Err(e) => Err(e.to_string()),
|
||||
};
|
||||
let _ = tx.send(got);
|
||||
});
|
||||
});
|
||||
|
||||
rx
|
||||
}
|
||||
|
||||
/// TRACES: FR-NC-6c
|
||||
/// Read a derived file, fetching its content first if only a placeholder is
|
||||
/// here.
|
||||
@@ -755,8 +933,9 @@ mod tests {
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod catalog_guard_tests {
|
||||
//! What `sync_catalog` does when it cannot read the server's copy.
|
||||
mod derived_guard_tests {
|
||||
//! What the two read-then-write steps do when they cannot read the server's
|
||||
//! copy — the catalog snapshot, and the place.
|
||||
//!
|
||||
//! The bug these exist for was a control-flow one — `if let Ok(bytes)`
|
||||
//! folding every failure into "there is none yet" and falling through to
|
||||
@@ -770,20 +949,39 @@ mod catalog_guard_tests {
|
||||
/// A backend whose read fails in a chosen way, counting writes.
|
||||
struct Fussy {
|
||||
fail_with: Option<RemoteError>,
|
||||
/// What a successful read returns. Empty for the catalog tests, which
|
||||
/// only ever exercise failures; the place tests need real content,
|
||||
/// because "theirs is newer than ours" is the branch under test.
|
||||
body: Vec<u8>,
|
||||
puts: Arc<AtomicUsize>,
|
||||
/// What the last successful write carried, so a test can assert *which*
|
||||
/// record went up rather than only that one did.
|
||||
last_put: Arc<std::sync::Mutex<Vec<u8>>>,
|
||||
caps: dr_sync::Capabilities,
|
||||
}
|
||||
|
||||
impl Fussy {
|
||||
fn reading(fail_with: Option<RemoteError>) -> (Self, Arc<AtomicUsize>) {
|
||||
let (f, puts, _) = Self::serving(fail_with, Vec::new());
|
||||
(f, puts)
|
||||
}
|
||||
|
||||
fn serving(
|
||||
fail_with: Option<RemoteError>,
|
||||
body: Vec<u8>,
|
||||
) -> (Self, Arc<AtomicUsize>, Arc<std::sync::Mutex<Vec<u8>>>) {
|
||||
let puts = Arc::new(AtomicUsize::new(0));
|
||||
let last_put = Arc::new(std::sync::Mutex::new(Vec::new()));
|
||||
(
|
||||
Self {
|
||||
fail_with,
|
||||
body,
|
||||
puts: puts.clone(),
|
||||
last_put: last_put.clone(),
|
||||
caps: dr_sync::Capabilities::minimal(),
|
||||
},
|
||||
puts,
|
||||
last_put,
|
||||
)
|
||||
}
|
||||
}
|
||||
@@ -826,7 +1024,7 @@ mod catalog_guard_tests {
|
||||
Err(RemoteError::NotMaterialised(s.clone()))
|
||||
}
|
||||
Some(_) => Err(RemoteError::PermissionDenied),
|
||||
None => Ok(Vec::new()),
|
||||
None => Ok(self.body.clone()),
|
||||
}
|
||||
}
|
||||
async fn put(
|
||||
@@ -836,6 +1034,7 @@ mod catalog_guard_tests {
|
||||
_pc: Option<dr_sync::Precondition>,
|
||||
) -> Result<dr_sync::Validator, RemoteError> {
|
||||
self.puts.fetch_add(1, Ordering::SeqCst);
|
||||
*self.last_put.lock().unwrap() = _b;
|
||||
Ok(dr_sync::Validator::new("v"))
|
||||
}
|
||||
async fn delete(
|
||||
@@ -914,4 +1113,165 @@ mod catalog_guard_tests {
|
||||
assert_eq!(puts, 1, "nothing to merge, so ours goes up");
|
||||
assert!(report.catalog_uploaded);
|
||||
}
|
||||
|
||||
// --- the place (FR-UI-8) ---------------------------------------------
|
||||
//
|
||||
// The same "do not write over what you could not read" rule as above, for a
|
||||
// file with different stakes. A place is replaced rather than merged, so an
|
||||
// unread copy costs no data — but it may be the *newer* one, and pushing
|
||||
// ours over it would move the other device's photographer without ever
|
||||
// having seen where they were.
|
||||
|
||||
fn a_place(at: i64, device: &str) -> dr_types::Place {
|
||||
dr_types::Place {
|
||||
at,
|
||||
device: device.into(),
|
||||
path: "2019/a.CR2".into(),
|
||||
..Default::default()
|
||||
}
|
||||
}
|
||||
|
||||
/// A local place file with a chosen record already in it, or none.
|
||||
fn local(name: &str, ours: Option<&dr_types::Place>) -> std::path::PathBuf {
|
||||
let dir = std::env::temp_dir().join(format!("dr-place-guard-{name}"));
|
||||
let _ = std::fs::remove_dir_all(&dir);
|
||||
std::fs::create_dir_all(&dir).unwrap();
|
||||
let path = dir.join("place.json");
|
||||
if let Some(p) = ours {
|
||||
crate::place::PlaceStore::open_at(path.clone()).save(p);
|
||||
}
|
||||
path
|
||||
}
|
||||
|
||||
async fn exchange(
|
||||
fail_with: Option<RemoteError>,
|
||||
theirs: Option<&dr_types::Place>,
|
||||
ours: Option<&dr_types::Place>,
|
||||
name: &str,
|
||||
) -> (usize, SyncReport, Option<dr_types::Place>, Vec<u8>) {
|
||||
let body = theirs
|
||||
.map(|p| crate::place::to_bytes(p).unwrap())
|
||||
.unwrap_or_default();
|
||||
let (backend, puts, last) = Fussy::serving(fail_with, body);
|
||||
let path = local(name, ours);
|
||||
let mut report = SyncReport::default();
|
||||
sync_place(
|
||||
&backend,
|
||||
&RemotePath::new(".darkroom-derived"),
|
||||
&path,
|
||||
&mut report,
|
||||
)
|
||||
.await;
|
||||
let after = crate::place::PlaceStore::open_at(path.clone()).load();
|
||||
let sent = last.lock().unwrap().clone();
|
||||
let _ = std::fs::remove_dir_all(path.parent().unwrap());
|
||||
(puts.load(Ordering::SeqCst), report, after, sent)
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn a_place_that_could_not_be_read_is_never_written_over() {
|
||||
// A dehydrated placeholder, and the record on the server may well be
|
||||
// newer than ours. Uploading blind would discard it.
|
||||
let ours = a_place(100, "desktop");
|
||||
let (puts, report, after, _) = exchange(
|
||||
Some(RemoteError::NotMaterialised("place.json".into())),
|
||||
None,
|
||||
Some(&ours),
|
||||
"notmaterialised",
|
||||
)
|
||||
.await;
|
||||
assert_eq!(puts, 0, "must not upload over a place it could not read");
|
||||
assert!(!report.place_uploaded);
|
||||
assert!(!report.place_adopted);
|
||||
assert_eq!(after.unwrap().at, 100, "and ours is untouched");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn a_first_exchange_uploads_ours() {
|
||||
// Genuinely nothing there: ours is the whole truth, and refusing to
|
||||
// push it would mean the place never travels at all.
|
||||
let ours = a_place(100, "desktop");
|
||||
let (puts, report, _, sent) = exchange(
|
||||
Some(RemoteError::NotFound("nope".into())),
|
||||
None,
|
||||
Some(&ours),
|
||||
"firstrun",
|
||||
)
|
||||
.await;
|
||||
assert_eq!(puts, 1);
|
||||
assert!(report.place_uploaded);
|
||||
assert_eq!(
|
||||
crate::place::from_bytes(&sent, "sent").unwrap().device,
|
||||
"desktop"
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn a_newer_record_from_another_device_is_adopted_and_not_pushed_back() {
|
||||
// The handover, and the half that stops it oscillating: having taken
|
||||
// theirs, ours *is* theirs, so there is nothing left to send.
|
||||
let ours = a_place(100, "desktop");
|
||||
let theirs = a_place(200, "tablet");
|
||||
let (puts, report, after, _) = exchange(None, Some(&theirs), Some(&ours), "adopt").await;
|
||||
assert!(report.place_adopted);
|
||||
assert_eq!(after.unwrap().device, "tablet");
|
||||
assert_eq!(puts, 0, "nothing to say that the server does not know");
|
||||
assert!(!report.place_uploaded);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn an_older_record_on_the_server_is_replaced_by_ours() {
|
||||
let ours = a_place(300, "desktop");
|
||||
let theirs = a_place(200, "tablet");
|
||||
let (puts, report, after, sent) = exchange(None, Some(&theirs), Some(&ours), "push").await;
|
||||
assert!(!report.place_adopted);
|
||||
assert_eq!(after.unwrap().device, "desktop", "ours stands");
|
||||
assert_eq!(puts, 1);
|
||||
assert!(report.place_uploaded);
|
||||
assert_eq!(crate::place::from_bytes(&sent, "sent").unwrap().at, 300);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn two_idle_devices_settle_rather_than_ping_pong() {
|
||||
// Equal timestamps mean equal records. A pass that decided it had
|
||||
// something to say here would put a PUT on every sync of every device,
|
||||
// for ever.
|
||||
let same = a_place(100, "desktop");
|
||||
let (puts, report, _, _) = exchange(None, Some(&same), Some(&same), "settle").await;
|
||||
assert_eq!(puts, 0);
|
||||
assert!(!report.place_adopted);
|
||||
assert!(!report.place_uploaded);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn a_device_with_no_place_of_its_own_takes_theirs() {
|
||||
// A second machine signing into an existing library. There is nothing
|
||||
// local to compare against, so anything on the server wins.
|
||||
let theirs = a_place(200, "tablet");
|
||||
let (puts, report, after, _) = exchange(None, Some(&theirs), None, "fresh").await;
|
||||
assert!(report.place_adopted);
|
||||
assert_eq!(after.unwrap().device, "tablet");
|
||||
assert_eq!(puts, 0);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn a_torn_record_on_the_server_is_replaced_rather_than_obeyed() {
|
||||
// Unlike the catalog, an unparseable place is not a reason to hold
|
||||
// back: there is nothing inside it to lose, and leaving a corrupt file
|
||||
// in place would mean the exchange never recovers.
|
||||
let ours = a_place(100, "desktop");
|
||||
let (backend, puts, _) = Fussy::serving(None, b"{\"at\": 12, \"scr".to_vec());
|
||||
let path = local("torn", Some(&ours));
|
||||
let mut report = SyncReport::default();
|
||||
sync_place(
|
||||
&backend,
|
||||
&RemotePath::new(".darkroom-derived"),
|
||||
&path,
|
||||
&mut report,
|
||||
)
|
||||
.await;
|
||||
assert_eq!(puts.load(Ordering::SeqCst), 1, "ours goes over the rubble");
|
||||
assert!(report.place_uploaded);
|
||||
let _ = std::fs::remove_dir_all(path.parent().unwrap());
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user