Run the library from local data when the server is unreachable
Also carries in-flight work that shared these files: the zoom structure-key fix in the adjust pipeline, nearest-neighbour filtering past 1:1, the timeline scrub marker correction, the 423-Locked retry in the metadata sweep, and the thumbnail size-class migration. # Offline mode (FR-CAT-9) The app previously assumed the server was reachable and treated its absence as a series of unrelated per-operation failures. A launch without a connection produced an empty grid, even with a complete catalog on disk and every thumbnail already in the shards. Reachability is now inferred from traffic the app was already making, rather than probed for. `RemoteError::indicates_offline` draws the line that makes this possible: a dead connection is offline, a 403 or a 500 is not — the server answered, so blanking the library over one forbidden file would be a worse error than the one being reported. `Reachability` turns those outcomes into a state, so a library browsing happily never issues a probe at all. Going offline takes one failure, because the user is already experiencing it. Coming back requires evidence — a completed scan or a fetched thumbnail — with a capped exponential backoff behind the manual retry, so twelve sweep lanes failing together do not schedule twelve immediate probes. What keeps working: the catalog opens even when the scan that normally provides it failed, so the grid fills from the last successful scan. Thumbnails come from the shards. Rating, flagging and collecting are catalog writes that never touched the network. What stops is opening an original that was never stored locally, and it now says so in those words instead of reporting "network error: connection refused" over a photograph. Work that is pure network is refused rather than left to fail slowly: the metadata sweep, derived sync, and sidecar writes. The sweep would otherwise spend a timeout per image across the whole library while the progress bar implied something was happening. Deferring sidecars is a real gap rather than a hidden one — a rating made offline reaches its sidecar only when that image is judged again while connected — and it is recorded as such at the call site. # The "On this device" filter A chip beside the rating filters, narrowing the grid to images whose original is held locally. It composes with the rating terms rather than replacing them, so "five-star frames I can actually edit on this train" is one filter. The predicate is SQL, like the rating terms and for the same reason: the count in the header has to agree with the cells drawn. It reads `image_cache.tier_actual`, which nothing writes yet — the next commit fills it. Until then the chip honestly reports zero. `Tier` gains an explicit on-disk encoding. The variants are ordered by generosity and the derived `Ord` invites reordering them, which would silently reinterpret every cached row; the round-trip test is what holds the two in agreement. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
@@ -0,0 +1,16 @@
|
||||
fn main() {
|
||||
for p in std::env::args().skip(1) {
|
||||
let Ok(d) = std::fs::read(&p) else { continue };
|
||||
let n = p.rsplit('/').next().unwrap();
|
||||
// Exactly what the sweep sees: the first HEADER_BYTES only.
|
||||
let head = &d[..d.len().min(dr_decode::HEADER_BYTES as usize)];
|
||||
match dr_decode::metadata(head) {
|
||||
Ok(m) => println!("{n}: header-only at={:?} model={:?}", m.captured_at, m.model),
|
||||
Err(e) => println!("{n}: header-only ERROR {e}"),
|
||||
}
|
||||
match dr_decode::metadata(&d) {
|
||||
Ok(m) => println!("{n}: whole-file at={:?}", m.captured_at),
|
||||
Err(e) => println!("{n}: whole-file ERROR {e}"),
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -695,15 +695,28 @@ mod tests {
|
||||
|
||||
#[test]
|
||||
fn zooming_does_not_recompile() {
|
||||
// The property that makes scroll-wheel zoom smooth: a new zoom level
|
||||
// is a uniform upload, never a pipeline build. If zoom reached the
|
||||
// structure hash, every wheel notch would stall on a shader compile.
|
||||
// The property that makes scroll-wheel zoom smooth: a new zoom *level*
|
||||
// is a uniform upload, never a pipeline build. If the magnitude reached
|
||||
// the structure hash, every wheel notch would stall on a compile.
|
||||
//
|
||||
// Entering the zoom at all is the one exception, and it is deliberate
|
||||
// — see `zooming_after_an_unzoomed_render_actually_zooms`. So the walk
|
||||
// below starts already zoomed, and the count is taken from there.
|
||||
let Some(ctx) = ctx() else { return };
|
||||
let mut pass = AdjustPass::new(&ctx);
|
||||
let img = split_image(&ctx, false);
|
||||
|
||||
let mut g = EditGraph::default_chain();
|
||||
for (i, extent) in [1.0f32, 0.5, 0.25, 0.125].iter().enumerate() {
|
||||
g.framing_mut().set_view(dr_pipeline::CropRect {
|
||||
x: 0.0,
|
||||
y: 0.0,
|
||||
width: 0.5,
|
||||
height: 0.5,
|
||||
});
|
||||
pass.render(&img, &g.compose(), 32, 32).expect("render");
|
||||
let baseline = pass.cached_pipelines();
|
||||
|
||||
for (i, extent) in [0.4f32, 0.25, 0.125].iter().enumerate() {
|
||||
g.framing_mut().set_view(dr_pipeline::CropRect {
|
||||
x: 0.0,
|
||||
y: 0.0,
|
||||
@@ -713,12 +726,59 @@ mod tests {
|
||||
pass.render(&img, &g.compose(), 32, 32).expect("render");
|
||||
assert_eq!(
|
||||
pass.cached_pipelines(),
|
||||
1,
|
||||
baseline,
|
||||
"zoom step {i} compiled a second pipeline"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn zooming_after_an_unzoomed_render_actually_zooms() {
|
||||
// The regression: every earlier zoom test set a view *before* the first
|
||||
// render, so the first pipeline compiled was already the one carrying
|
||||
// the crop mapping. Real use is the other way round — the image is
|
||||
// shown fitted, and only then does the wheel turn.
|
||||
//
|
||||
// A neutral framing emits a prologue that never reads `u.crop_rect`.
|
||||
// While zoom was excluded from the structure hash, that neutral
|
||||
// pipeline stayed cached under the same key once zoomed, so the view
|
||||
// uploaded on every frame was read by nobody and the canvas never
|
||||
// changed. This renders unzoomed first and asserts the pixels move.
|
||||
let Some(ctx) = ctx() else { return };
|
||||
let mut pass = AdjustPass::new(&ctx);
|
||||
let img = split_image(&ctx, false);
|
||||
|
||||
let mut g = EditGraph::default_chain();
|
||||
|
||||
// Fitted: the frame spans both halves, so the two edges differ.
|
||||
let tex = pass.render(&img, &g.compose(), 32, 32).expect("render");
|
||||
let fitted_left = read_pixel(&ctx, tex, 4, 16)[0];
|
||||
let fitted_right = read_pixel(&ctx, tex, 28, 16)[0];
|
||||
assert!(
|
||||
(i32::from(fitted_left) - i32::from(fitted_right)).abs() > 40,
|
||||
"the unzoomed frame should span both halves: \
|
||||
left={fitted_left} right={fitted_right}"
|
||||
);
|
||||
|
||||
// Now zoom into the bright half. Both edges must come up bright.
|
||||
g.framing_mut().set_view(dr_pipeline::CropRect {
|
||||
x: 0.0,
|
||||
y: 0.4,
|
||||
width: 0.2,
|
||||
height: 0.2,
|
||||
});
|
||||
let tex = pass.render(&img, &g.compose(), 32, 32).expect("render");
|
||||
let zoomed_left = read_pixel(&ctx, tex, 4, 16)[0];
|
||||
let zoomed_right = read_pixel(&ctx, tex, 28, 16)[0];
|
||||
|
||||
assert!(
|
||||
zoomed_left > 100 && zoomed_right > 100,
|
||||
"zooming into the bright half after an unzoomed render must show \
|
||||
it edge to edge — the neutral pipeline was reused and the view \
|
||||
was ignored: left={zoomed_left} right={zoomed_right}"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn cropping_to_one_half_shows_only_that_half() {
|
||||
// The property a crop exists for, checked against content rather than
|
||||
|
||||
@@ -608,12 +608,26 @@ impl Framing {
|
||||
/// the crop handles or the straighten slider must reuse the compiled
|
||||
/// pipeline and upload uniforms only. Only the presence of each
|
||||
/// transform, never its magnitude, may enter this.
|
||||
///
|
||||
/// The last bit is *whether the prologue is emitted at all*, which zoom
|
||||
/// reaches through [`Self::is_active`]. It has to be here even though zoom
|
||||
/// is not an edit: the neutral prologue never reads `u.crop_rect`, so a
|
||||
/// pipeline compiled while unzoomed ignores every later view upload. Two
|
||||
/// framings that generate different WGSL must not share a cache key — the
|
||||
/// symptom otherwise is scroll-to-zoom on an otherwise-unedited image
|
||||
/// doing nothing at all, because the first frame compiled the neutral
|
||||
/// prologue and the hash never moved off it.
|
||||
///
|
||||
/// What this must *not* do is vary with the zoom level: the bit is set by
|
||||
/// any zoom and cleared by none, so a wheel notch is still a uniform
|
||||
/// upload rather than a shader build.
|
||||
pub fn structure_key(&self) -> u64 {
|
||||
u64::from(!self.crop.is_full())
|
||||
| u64::from(self.angle != 0.0) << 1
|
||||
| u64::from(self.flip_h) << 2
|
||||
| u64::from(self.flip_v) << 3
|
||||
| u64::from(self.quarter_turns) << 4
|
||||
| u64::from(self.is_active()) << 6
|
||||
}
|
||||
}
|
||||
|
||||
@@ -639,11 +653,15 @@ mod tests {
|
||||
#[test]
|
||||
fn zooming_does_not_change_the_exported_image() {
|
||||
// The property that makes zoom a viewing tool rather than an edit: it
|
||||
// must not reach the output size, the structure hash, or the crop.
|
||||
// If it did, exporting while zoomed would write the zoomed view.
|
||||
// must not reach the output size or the crop. If it did, exporting
|
||||
// while zoomed would write the zoomed view.
|
||||
//
|
||||
// The structure key is deliberately not asserted here — see
|
||||
// `zooming_from_neutral_changes_the_structure_key` for why it must
|
||||
// move, and `zoom_level_does_not_change_the_structure_key` for the
|
||||
// part that must not.
|
||||
let mut f = Framing::new();
|
||||
let before_size = f.output_size(6000, 4000);
|
||||
let before_key = f.structure_key();
|
||||
|
||||
f.set_view(CropRect {
|
||||
x: 0.25,
|
||||
@@ -653,10 +671,62 @@ mod tests {
|
||||
});
|
||||
|
||||
assert_eq!(f.output_size(6000, 4000), before_size, "zoom resized output");
|
||||
assert_eq!(f.structure_key(), before_key, "zoom forced a recompile");
|
||||
assert!(f.crop().is_full(), "zoom altered the crop");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn zooming_from_neutral_changes_the_structure_key() {
|
||||
// The regression this guards: a neutral framing emits a prologue that
|
||||
// never reads `u.crop_rect`, so if zooming leaves the key alone the
|
||||
// GPU reuses that pipeline and the uploaded view is ignored — zoom
|
||||
// silently does nothing on an otherwise-unedited image.
|
||||
let mut f = Framing::new();
|
||||
let neutral = f.structure_key();
|
||||
|
||||
f.set_view(CropRect {
|
||||
x: 0.25,
|
||||
y: 0.25,
|
||||
width: 0.5,
|
||||
height: 0.5,
|
||||
});
|
||||
|
||||
assert_ne!(
|
||||
f.structure_key(),
|
||||
neutral,
|
||||
"a zoomed framing generates different WGSL and must not share the \
|
||||
neutral cache key"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn zoom_level_does_not_change_the_structure_key() {
|
||||
// The other half of the contract: crossing from unzoomed to zoomed is
|
||||
// a recompile, but every notch after that is a uniform upload. If the
|
||||
// magnitude reached the key, every wheel step would stall on a build.
|
||||
let mut f = Framing::new();
|
||||
f.set_view(CropRect {
|
||||
x: 0.25,
|
||||
y: 0.25,
|
||||
width: 0.5,
|
||||
height: 0.5,
|
||||
});
|
||||
let zoomed = f.structure_key();
|
||||
|
||||
for extent in [0.4, 0.3, 0.2, 0.1] {
|
||||
f.set_view(CropRect {
|
||||
x: 0.1,
|
||||
y: 0.1,
|
||||
width: extent,
|
||||
height: extent,
|
||||
});
|
||||
assert_eq!(
|
||||
f.structure_key(),
|
||||
zoomed,
|
||||
"zoom level {extent} forced a recompile"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn the_view_nests_inside_the_crop() {
|
||||
// Both rects live in the same normalised space and the shader applies
|
||||
|
||||
@@ -27,6 +27,45 @@ async fn main() {
|
||||
Err(e) => println!("READ FAILED: {e}"),
|
||||
}
|
||||
|
||||
// What is actually in the derived folder on the server?
|
||||
{
|
||||
let dir = RemotePath::new(format!("{}/.darkroom-derived", session.root));
|
||||
match backend.list(&dir, None).await {
|
||||
Ok(entries) => {
|
||||
println!("\n[derived] {} entry/entries on the server:", entries.len());
|
||||
for e in &entries {
|
||||
println!(" {:<28} {:>10} bytes", e.path.name(), e.size);
|
||||
}
|
||||
}
|
||||
Err(e) => println!("\n[derived] listing failed: {e}"),
|
||||
}
|
||||
}
|
||||
|
||||
// Probe a file the sweep reported as 423 Locked: is it the file, the
|
||||
// range request, or the folder?
|
||||
{
|
||||
let c = dr_sync_nextcloud::http_client("DarkRoom").unwrap();
|
||||
let base = format!("{}/remote.php/dav/files/{}",
|
||||
creds.server.trim_end_matches('/'), session.user_id);
|
||||
let f = "PhotosRaw/Darktable/20230629_no_name/20230629_0030.jpeg";
|
||||
let enc: String = f.split('/').map(|seg| {
|
||||
seg.bytes().map(|b| match b {
|
||||
b'A'..=b'Z'|b'a'..=b'z'|b'0'..=b'9'|b'-'|b'_'|b'.'|b'~' => (b as char).to_string(),
|
||||
_ => format!("%{b:02X}"),
|
||||
}).collect::<String>()
|
||||
}).collect::<Vec<_>>().join("/");
|
||||
let url = format!("{base}/{enc}");
|
||||
|
||||
for (what, range) in [("ranged 0-256k", Some("bytes=0-262143")), ("whole file", None)] {
|
||||
let mut rq = c.get(&url).basic_auth(&creds.login_name, Some(&creds.app_password));
|
||||
if let Some(r) = range { rq = rq.header("Range", r); }
|
||||
match rq.send().await {
|
||||
Ok(r) => println!("GET {what}: {}", r.status()),
|
||||
Err(e) => println!("GET {what}: transport {e}"),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Raw HTTP, to see the status the connector maps away.
|
||||
{
|
||||
let url = format!("{}/remote.php/dav/files/{}/{}/.darkroom-write-test",
|
||||
|
||||
@@ -92,6 +92,114 @@ impl NextcloudBackend {
|
||||
}
|
||||
}
|
||||
|
||||
/// TRACES: FR-NC-7
|
||||
/// Upload a large body with chunked upload v2.
|
||||
///
|
||||
/// `MKCOL` an upload directory, `PUT` each chunk into it under a numeric
|
||||
/// name, then `MOVE` the `.file` pseudo-entry to the destination, which is
|
||||
/// where the server assembles them.
|
||||
///
|
||||
/// The alternative — refusing anything over the single-shot threshold —
|
||||
/// is what blocked thumbnail shards from ever reaching the server: they
|
||||
/// are 25 MB by design.
|
||||
///
|
||||
/// `OC-Total-Length` is sent on every chunk so quota is checked up front
|
||||
/// rather than at assembly, when the bytes have already been transferred.
|
||||
async fn put_chunked(
|
||||
&self,
|
||||
path: &RemotePath,
|
||||
body: Vec<u8>,
|
||||
) -> Result<Validator, RemoteError> {
|
||||
let total = body.len() as u64;
|
||||
// Named from the destination so a resumed or abandoned upload is
|
||||
// identifiable, and so two uploads cannot collide in one directory.
|
||||
let token: String = path
|
||||
.as_str()
|
||||
.bytes()
|
||||
.map(|b| match b {
|
||||
b'A'..=b'Z' | b'a'..=b'z' | b'0'..=b'9' => (b as char).to_string(),
|
||||
_ => "-".to_string(),
|
||||
})
|
||||
.collect();
|
||||
let dir = format!(
|
||||
"{}/remote.php/dav/uploads/{}/{token}",
|
||||
self.server, self.login
|
||||
);
|
||||
|
||||
self.mkcol_url(&dir).await?;
|
||||
|
||||
// Chunks are numbered from 1 and must sort correctly as strings, which
|
||||
// is why they are zero-padded rather than bare integers.
|
||||
let chunk_size = CHUNKS.min_chunk as usize;
|
||||
for (i, chunk) in body.chunks(chunk_size).enumerate() {
|
||||
if i + 1 > CHUNKS.max_chunks as usize {
|
||||
return Err(RemoteError::Protocol(format!(
|
||||
"{total} bytes exceeds {} chunks", CHUNKS.max_chunks
|
||||
)));
|
||||
}
|
||||
let resp = self
|
||||
.client
|
||||
.put(format!("{dir}/{:05}", i + 1))
|
||||
.basic_auth(&self.login, Some(&self.password))
|
||||
.header("OC-Total-Length", total.to_string())
|
||||
.body(chunk.to_vec())
|
||||
.send()
|
||||
.await
|
||||
.map_err(|e| RemoteError::Network(e.to_string()))?;
|
||||
map_status(resp.status(), path.as_str())?;
|
||||
}
|
||||
|
||||
// Assemble. The destination is an absolute URL in the Destination
|
||||
// header, and `Overwrite: T` because a re-uploaded shard replaces the
|
||||
// one already there.
|
||||
let resp = self
|
||||
.client
|
||||
.request(
|
||||
reqwest::Method::from_bytes(b"MOVE").expect("valid method"),
|
||||
format!("{dir}/.file"),
|
||||
)
|
||||
.basic_auth(&self.login, Some(&self.password))
|
||||
.header("Destination", self.url_for(path))
|
||||
.header("Overwrite", "T")
|
||||
.header("OC-Total-Length", total.to_string())
|
||||
.send()
|
||||
.await
|
||||
.map_err(|e| RemoteError::Network(e.to_string()))?;
|
||||
map_status(resp.status(), path.as_str())?;
|
||||
|
||||
// The MOVE response carries the assembled file's ETag on Nextcloud,
|
||||
// but not on every version; fall back to asking rather than failing an
|
||||
// upload that in fact succeeded.
|
||||
if let Some(v) = resp
|
||||
.headers()
|
||||
.get(reqwest::header::ETAG)
|
||||
.and_then(|v| v.to_str().ok())
|
||||
{
|
||||
return Ok(Validator::new(v));
|
||||
}
|
||||
self.dir_validator(path).await
|
||||
}
|
||||
|
||||
/// `MKCOL` at an absolute URL, treating "already there" as success.
|
||||
async fn mkcol_url(&self, url: &str) -> Result<(), RemoteError> {
|
||||
let resp = self
|
||||
.client
|
||||
.request(
|
||||
reqwest::Method::from_bytes(b"MKCOL").expect("valid method"),
|
||||
url,
|
||||
)
|
||||
.basic_auth(&self.login, Some(&self.password))
|
||||
.send()
|
||||
.await
|
||||
.map_err(|e| RemoteError::Network(e.to_string()))?;
|
||||
|
||||
// 405 is "already exists", which is exactly what we want.
|
||||
if resp.status() == 405 {
|
||||
return Ok(());
|
||||
}
|
||||
map_status(resp.status(), url)
|
||||
}
|
||||
|
||||
async fn propfind(
|
||||
&self,
|
||||
path: &RemotePath,
|
||||
@@ -209,10 +317,14 @@ impl RemoteBackend for NextcloudBackend {
|
||||
) -> Result<Validator, RemoteError> {
|
||||
// Chunked upload is an implementation detail of put, chosen by size —
|
||||
// exposing it on the trait would leak this protocol (ARCH §8.3).
|
||||
if body.len() as u64 >= CHUNKS.single_shot_below {
|
||||
return Err(RemoteError::Unsupported(
|
||||
"chunked upload v2 not implemented yet",
|
||||
));
|
||||
//
|
||||
// A precondition cannot ride on a chunked upload: the guard belongs to
|
||||
// the assembling MOVE, not to the individual chunks, and Nextcloud
|
||||
// does not honour `If-Match` there. Large bodies are shards and
|
||||
// catalog snapshots, which are written whole and never merged, so
|
||||
// there is no conflict to guard against.
|
||||
if body.len() as u64 >= CHUNKS.single_shot_below && precond.is_none() {
|
||||
return self.put_chunked(path, body).await;
|
||||
}
|
||||
|
||||
let mut req = self
|
||||
|
||||
@@ -54,11 +54,35 @@ impl RemoteError {
|
||||
RemoteError::Network(_) => true,
|
||||
RemoteError::Server { status, .. } => {
|
||||
// 5xx and 429 are worth retrying; other 4xx are not.
|
||||
*status >= 500 || *status == 429
|
||||
//
|
||||
// 423 Locked is the exception, and it is not hypothetical:
|
||||
// Nextcloud's file locking returns it on a plain *read* under
|
||||
// concurrency, and the same range re-read seconds later
|
||||
// succeeds. Treating it as permanent marks an image
|
||||
// permanently undated over a lock that lasted moments.
|
||||
*status >= 500 || *status == 429 || *status == 423
|
||||
}
|
||||
_ => false,
|
||||
}
|
||||
}
|
||||
|
||||
/// Whether this failure means *the server could not be reached*, as
|
||||
/// opposed to the server answering and refusing.
|
||||
///
|
||||
/// The distinction is the whole basis of offline mode (FR-CAT-9). A 403
|
||||
/// and a dead connection are both "the operation failed", but only one of
|
||||
/// them is fixed by waiting, and only one of them should put the whole app
|
||||
/// into a degraded mode. Signing the user out — or showing "you are
|
||||
/// offline" — because a single file was forbidden would be a much worse
|
||||
/// error than the one it reported.
|
||||
///
|
||||
/// A 5xx is deliberately **not** offline: the server is up and talking, it
|
||||
/// is just failing, and a retry is the right response rather than a
|
||||
/// mode change. 429 and 423 likewise — those are the server working
|
||||
/// correctly under load.
|
||||
pub fn indicates_offline(&self) -> bool {
|
||||
matches!(self, RemoteError::Network(_))
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
@@ -107,4 +131,29 @@ mod tests {
|
||||
"the message must point at permissions, not the login: {denied}"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_lock_is_transient() {
|
||||
// Observed against a real server: 12 concurrent range reads produced
|
||||
// 423 on some files, and the identical request succeeded moments
|
||||
// later. Classing it with the permanent 4xx left those images
|
||||
// undated for good.
|
||||
assert!(RemoteError::Server {
|
||||
status: 423,
|
||||
detail: String::new()
|
||||
}
|
||||
.is_transient());
|
||||
|
||||
// Still permanent, so the exception stays narrow.
|
||||
assert!(!RemoteError::Server {
|
||||
status: 404,
|
||||
detail: String::new()
|
||||
}
|
||||
.is_transient());
|
||||
assert!(!RemoteError::Server {
|
||||
status: 400,
|
||||
detail: String::new()
|
||||
}
|
||||
.is_transient());
|
||||
}
|
||||
}
|
||||
|
||||
@@ -22,11 +22,13 @@ use async_trait::async_trait;
|
||||
|
||||
pub mod capability;
|
||||
pub mod error;
|
||||
pub mod reachability;
|
||||
pub mod scan;
|
||||
pub mod types;
|
||||
|
||||
pub use capability::{Capabilities, ChangeDetection, ChunkConstraints, ServerPreviews};
|
||||
pub use error::RemoteError;
|
||||
pub use reachability::{Connectivity, Reachability};
|
||||
pub use scan::{scan, ScanProgress, ScanResult};
|
||||
pub use types::{
|
||||
Cursor, EntryKind, Identity, Precondition, RemoteChange, RemoteEntry, RemoteId, RemotePath,
|
||||
|
||||
@@ -0,0 +1,335 @@
|
||||
//! TRACES: FR-CAT-9 | FR-NC-12
|
||||
//! Whether the remote is reachable, and what the app does while it is not.
|
||||
//!
|
||||
//! # Why this is a state machine rather than a boolean
|
||||
//!
|
||||
//! "Are we online?" cannot be answered by asking the operating system. A
|
||||
//! laptop with a live wifi association and no route, a captive portal that
|
||||
//! answers every request with a login page, a Nextcloud instance that is down
|
||||
//! while the internet is fine — all three report a working network and none of
|
||||
//! them can serve an image. The only evidence that counts is whether *this
|
||||
//! backend* answered, so reachability is inferred from the traffic the app was
|
||||
//! already making rather than probed for separately.
|
||||
//!
|
||||
//! That inversion is what keeps the cost at zero. Every remote call already
|
||||
//! returns a `Result`; [`Reachability::observe`] turns those results into the
|
||||
//! state, so a library that is browsing happily never issues a probe at all.
|
||||
//! A probe happens only when something failed and the app wants to know
|
||||
//! whether it has come back (ARCH §9.0).
|
||||
//!
|
||||
//! # Why leaving offline is harder than entering it
|
||||
//!
|
||||
//! One failed request is enough to go offline: the user is *already*
|
||||
//! experiencing the failure, and the honest thing is to say so immediately.
|
||||
//! But a single success is not enough to declare recovery, because the failure
|
||||
//! mode that matters — a flapping connection — produces exactly that. So
|
||||
//! recovery requires a deliberate probe, and the app backs off between
|
||||
//! attempts rather than retrying in a tight loop against a server that is
|
||||
//! plainly down.
|
||||
|
||||
use std::time::{Duration, Instant};
|
||||
|
||||
use crate::error::RemoteError;
|
||||
|
||||
/// How long to wait before the first reconnection probe.
|
||||
///
|
||||
/// Short enough that a brief drop — a laptop changing access points, a phone
|
||||
/// moving between cells — recovers before the user has finished noticing it.
|
||||
const FIRST_BACKOFF: Duration = Duration::from_secs(5);
|
||||
|
||||
/// The longest gap between probes.
|
||||
///
|
||||
/// Capped rather than growing without bound: a user who left the app open
|
||||
/// overnight on a dead connection should reconnect within a minute of the
|
||||
/// server returning, not hours later because the backoff had doubled its way
|
||||
/// into the distance.
|
||||
const MAX_BACKOFF: Duration = Duration::from_secs(60);
|
||||
|
||||
/// What the app currently believes about the remote.
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||
pub enum Connectivity {
|
||||
/// The remote answered the last time anything asked.
|
||||
Online,
|
||||
/// The remote could not be reached. The library runs from local data.
|
||||
Offline,
|
||||
}
|
||||
|
||||
impl Connectivity {
|
||||
pub fn is_online(self) -> bool {
|
||||
matches!(self, Connectivity::Online)
|
||||
}
|
||||
|
||||
pub fn is_offline(self) -> bool {
|
||||
matches!(self, Connectivity::Offline)
|
||||
}
|
||||
}
|
||||
|
||||
/// Tracks reachability from observed request outcomes.
|
||||
///
|
||||
/// Cheap to construct and `Clone`-free by design — one lives beside the
|
||||
/// session and every worker reports into it.
|
||||
#[derive(Debug)]
|
||||
pub struct Reachability {
|
||||
state: Connectivity,
|
||||
/// When the app went offline. Shown to the user, because "offline" without
|
||||
/// "since when" leaves them unable to tell a momentary drop from a
|
||||
/// connection that died an hour ago.
|
||||
since: Option<Instant>,
|
||||
/// How long to wait before the next probe, doubling per failure.
|
||||
backoff: Duration,
|
||||
/// When the next probe becomes worthwhile.
|
||||
next_probe: Option<Instant>,
|
||||
/// Why we think we are offline, for the banner. The underlying transport
|
||||
/// message, which is usually specific enough to be actionable ("dns error",
|
||||
/// "connection refused").
|
||||
reason: Option<String>,
|
||||
}
|
||||
|
||||
impl Default for Reachability {
|
||||
fn default() -> Self {
|
||||
Self::new()
|
||||
}
|
||||
}
|
||||
|
||||
impl Reachability {
|
||||
/// Start optimistic.
|
||||
///
|
||||
/// Assuming online until proven otherwise is deliberate: the alternative
|
||||
/// is a probe on every launch, which makes startup wait on the network for
|
||||
/// a library that may be entirely cached. The first real request settles
|
||||
/// it either way, and settles it with evidence.
|
||||
pub fn new() -> Self {
|
||||
Self {
|
||||
state: Connectivity::Online,
|
||||
since: None,
|
||||
backoff: FIRST_BACKOFF,
|
||||
next_probe: None,
|
||||
reason: None,
|
||||
}
|
||||
}
|
||||
|
||||
pub fn state(&self) -> Connectivity {
|
||||
self.state
|
||||
}
|
||||
|
||||
pub fn is_offline(&self) -> bool {
|
||||
self.state.is_offline()
|
||||
}
|
||||
|
||||
/// Why the app believes it is offline, if it does.
|
||||
pub fn reason(&self) -> Option<&str> {
|
||||
self.reason.as_deref()
|
||||
}
|
||||
|
||||
/// How long the app has been offline.
|
||||
pub fn offline_for(&self, now: Instant) -> Option<Duration> {
|
||||
self.since.map(|t| now.saturating_duration_since(t))
|
||||
}
|
||||
|
||||
/// Record the outcome of a remote call.
|
||||
///
|
||||
/// Returns `true` if the connectivity state *changed*, so the caller can
|
||||
/// repaint a banner or kick off a rescan without diffing the state itself.
|
||||
///
|
||||
/// Takes the result by reference so callers can report an outcome they are
|
||||
/// still going to use — this observes, it never consumes.
|
||||
pub fn observe<T>(&mut self, outcome: &Result<T, RemoteError>, now: Instant) -> bool {
|
||||
match outcome {
|
||||
Ok(_) => self.mark_reachable(now),
|
||||
Err(e) if e.indicates_offline() => self.mark_unreachable(e.to_string(), now),
|
||||
// The server answered. Whatever went wrong is not connectivity, so
|
||||
// it must not move this state — a 404 on one file says nothing
|
||||
// about the other 17,000.
|
||||
Err(_) => false,
|
||||
}
|
||||
}
|
||||
|
||||
/// Record that the remote answered.
|
||||
pub fn mark_reachable(&mut self, _now: Instant) -> bool {
|
||||
let changed = self.state.is_offline();
|
||||
self.state = Connectivity::Online;
|
||||
self.since = None;
|
||||
self.backoff = FIRST_BACKOFF;
|
||||
self.next_probe = None;
|
||||
self.reason = None;
|
||||
changed
|
||||
}
|
||||
|
||||
/// Record that the remote could not be reached.
|
||||
///
|
||||
/// Repeated calls while already offline extend the backoff rather than
|
||||
/// resetting it, so a library with twelve workers all failing at once
|
||||
/// does not schedule twelve immediate probes.
|
||||
pub fn mark_unreachable(&mut self, reason: String, now: Instant) -> bool {
|
||||
let changed = self.state.is_online();
|
||||
if changed {
|
||||
self.state = Connectivity::Offline;
|
||||
self.since = Some(now);
|
||||
self.backoff = FIRST_BACKOFF;
|
||||
} else {
|
||||
self.backoff = (self.backoff * 2).min(MAX_BACKOFF);
|
||||
}
|
||||
self.next_probe = Some(now + self.backoff);
|
||||
self.reason = Some(reason);
|
||||
changed
|
||||
}
|
||||
|
||||
/// Whether enough time has passed to be worth trying the remote again.
|
||||
///
|
||||
/// Always false while online — there is nothing to probe for.
|
||||
pub fn should_probe(&self, now: Instant) -> bool {
|
||||
self.state.is_offline() && self.next_probe.is_some_and(|t| now >= t)
|
||||
}
|
||||
|
||||
/// How long until the next probe is due, for a countdown in the banner.
|
||||
pub fn until_probe(&self, now: Instant) -> Option<Duration> {
|
||||
self.next_probe.map(|t| t.saturating_duration_since(now))
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
fn network() -> RemoteError {
|
||||
RemoteError::Network("connection refused".into())
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn starts_online_without_probing() {
|
||||
// Launch must not wait on the network: a fully cached library opens
|
||||
// with no request at all, and an optimistic start is what allows that.
|
||||
let r = Reachability::new();
|
||||
assert!(r.state().is_online());
|
||||
assert!(!r.should_probe(Instant::now()));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn one_network_failure_goes_offline() {
|
||||
let mut r = Reachability::new();
|
||||
let now = Instant::now();
|
||||
let changed = r.observe::<()>(&Err(network()), now);
|
||||
assert!(changed, "the first failure is a state change");
|
||||
assert!(r.is_offline());
|
||||
assert_eq!(r.reason(), Some("network error: connection refused"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_server_error_is_not_offline() {
|
||||
// The distinction the whole mode rests on: the server answered, so it
|
||||
// is reachable. Going offline here would blank the grid over one
|
||||
// forbidden file.
|
||||
let mut r = Reachability::new();
|
||||
let now = Instant::now();
|
||||
|
||||
for e in [
|
||||
RemoteError::PermissionDenied,
|
||||
RemoteError::NotFound("a.CR2".into()),
|
||||
RemoteError::AuthFailed,
|
||||
RemoteError::Server {
|
||||
status: 500,
|
||||
detail: String::new(),
|
||||
},
|
||||
RemoteError::Server {
|
||||
status: 423,
|
||||
detail: String::new(),
|
||||
},
|
||||
] {
|
||||
let mut probe = Reachability::new();
|
||||
assert!(!probe.observe::<()>(&Err(e), now));
|
||||
assert!(probe.state().is_online());
|
||||
}
|
||||
|
||||
assert!(!r.observe::<()>(&Err(RemoteError::PermissionDenied), now));
|
||||
assert!(r.state().is_online());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn success_brings_it_back() {
|
||||
let mut r = Reachability::new();
|
||||
let now = Instant::now();
|
||||
r.observe::<()>(&Err(network()), now);
|
||||
assert!(r.is_offline());
|
||||
|
||||
let changed = r.observe(&Ok(()), now);
|
||||
assert!(changed, "recovery is a state change");
|
||||
assert!(r.state().is_online());
|
||||
assert_eq!(r.reason(), None);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn repeated_failures_back_off_rather_than_reset() {
|
||||
// Twelve sweep lanes failing together must not schedule twelve
|
||||
// immediate probes against a server that is plainly down.
|
||||
let mut r = Reachability::new();
|
||||
let now = Instant::now();
|
||||
|
||||
assert!(r.observe::<()>(&Err(network()), now));
|
||||
let first = r.until_probe(now).unwrap();
|
||||
|
||||
assert!(!r.observe::<()>(&Err(network()), now), "already offline");
|
||||
let second = r.until_probe(now).unwrap();
|
||||
|
||||
assert!(
|
||||
second > first,
|
||||
"backoff must grow: {second:?} should exceed {first:?}"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn backoff_is_capped() {
|
||||
let mut r = Reachability::new();
|
||||
let now = Instant::now();
|
||||
for _ in 0..20 {
|
||||
r.observe::<()>(&Err(network()), now);
|
||||
}
|
||||
assert!(
|
||||
r.until_probe(now).unwrap() <= MAX_BACKOFF,
|
||||
"an app left overnight must still reconnect promptly"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn probing_waits_for_the_backoff() {
|
||||
let mut r = Reachability::new();
|
||||
let now = Instant::now();
|
||||
r.observe::<()>(&Err(network()), now);
|
||||
|
||||
assert!(!r.should_probe(now), "not immediately");
|
||||
assert!(r.should_probe(now + FIRST_BACKOFF));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn recovery_resets_the_backoff() {
|
||||
// Otherwise a connection that flaps all day arrives at the maximum
|
||||
// backoff and stays there, so the next real drop takes a minute to
|
||||
// notice recovery.
|
||||
let mut r = Reachability::new();
|
||||
let now = Instant::now();
|
||||
for _ in 0..5 {
|
||||
r.observe::<()>(&Err(network()), now);
|
||||
}
|
||||
r.observe(&Ok(()), now);
|
||||
r.observe::<()>(&Err(network()), now);
|
||||
|
||||
assert_eq!(r.until_probe(now), Some(FIRST_BACKOFF));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn offline_duration_is_measured_from_the_first_failure() {
|
||||
// Not from the most recent one: a connection that has been down for an
|
||||
// hour must not report five seconds because a worker retried.
|
||||
let mut r = Reachability::new();
|
||||
let start = Instant::now();
|
||||
r.observe::<()>(&Err(network()), start);
|
||||
|
||||
let later = start + Duration::from_secs(600);
|
||||
r.observe::<()>(&Err(network()), later);
|
||||
|
||||
let elapsed = r.offline_for(later).expect("offline since the first failure");
|
||||
assert!(
|
||||
elapsed >= Duration::from_secs(600),
|
||||
"measured from the first failure, got {elapsed:?}"
|
||||
);
|
||||
}
|
||||
}
|
||||
@@ -329,6 +329,12 @@ impl ThumbStore {
|
||||
if create {
|
||||
conn.execute_batch(SHARD_SCHEMA)?;
|
||||
}
|
||||
// Every shard carries its own `thumbs` table, so migrating the index
|
||||
// alone is not enough — and `CREATE TABLE IF NOT EXISTS` leaves an
|
||||
// existing one untouched, so a shard written before the size class
|
||||
// keeps the old shape and every write to it fails with "no column
|
||||
// named size".
|
||||
migrate_size_column(&conn, "thumbs")?;
|
||||
Ok(conn)
|
||||
}
|
||||
|
||||
@@ -794,6 +800,53 @@ mod tests {
|
||||
assert_eq!(ThumbSize::for_cell(400), ThumbSize::Large);
|
||||
}
|
||||
|
||||
|
||||
#[test]
|
||||
fn a_shard_written_before_the_size_class_accepts_new_thumbnails() {
|
||||
// The index is not the only table with a `size` column: every shard
|
||||
// carries its own `thumbs`. Migrating the index alone left existing
|
||||
// shards in the old shape, and `CREATE TABLE IF NOT EXISTS` will not
|
||||
// fix one — so every write failed with "no column named size" and the
|
||||
// store silently stopped accepting thumbnails.
|
||||
let dir = std::env::temp_dir().join(format!("dr-thumbs-shardmig-{}", std::process::id()));
|
||||
let _ = std::fs::remove_dir_all(&dir);
|
||||
std::fs::create_dir_all(&dir).unwrap();
|
||||
|
||||
// A shard in the old shape, holding one thumbnail.
|
||||
{
|
||||
let c = Connection::open(dir.join("shard-0000.sqlite")).unwrap();
|
||||
c.execute_batch(
|
||||
"CREATE TABLE thumbs (
|
||||
file_id INTEGER PRIMARY KEY, width INTEGER NOT NULL,
|
||||
height INTEGER NOT NULL, bytes BLOB NOT NULL);
|
||||
INSERT INTO thumbs VALUES (42, 8, 8, x'FFD8FFD9');",
|
||||
)
|
||||
.unwrap();
|
||||
}
|
||||
{
|
||||
let c = Connection::open(dir.join("index.sqlite")).unwrap();
|
||||
c.execute_batch(
|
||||
"CREATE TABLE entries (
|
||||
file_id INTEGER PRIMARY KEY, shard INTEGER NOT NULL, bytes INTEGER NOT NULL);
|
||||
CREATE TABLE shards (
|
||||
id INTEGER PRIMARY KEY, bytes INTEGER NOT NULL DEFAULT 0,
|
||||
sealed INTEGER NOT NULL DEFAULT 0);
|
||||
INSERT INTO shards(id, bytes, sealed) VALUES (0, 4, 0);
|
||||
INSERT INTO entries(file_id, shard, bytes) VALUES (42, 0, 4);",
|
||||
)
|
||||
.unwrap();
|
||||
}
|
||||
|
||||
let mut s = ThumbStore::open(&dir).unwrap();
|
||||
|
||||
// The old thumbnail survives as grid-sized...
|
||||
assert!(s.get(42, ThumbSize::Grid).unwrap().is_some());
|
||||
// ...and the shard now accepts new writes rather than rejecting them.
|
||||
s.put(43, ThumbSize::Grid, &thumb(64))
|
||||
.expect("a migrated shard must accept writes");
|
||||
assert!(s.contains(43, ThumbSize::Grid));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_store_written_before_the_size_class_keeps_its_thumbnails() {
|
||||
// The migration case: an existing library must not lose the thumbnails
|
||||
|
||||
@@ -154,6 +154,36 @@ pub enum Tier {
|
||||
Original,
|
||||
}
|
||||
|
||||
impl Tier {
|
||||
/// The value stored in `image_cache.tier_actual` / `tier_desired`.
|
||||
///
|
||||
/// Written out explicitly rather than derived from the discriminant: these
|
||||
/// integers are **on disk**, so reordering the variants — which the `Ord`
|
||||
/// derive above openly invites, since generosity ordering is the point —
|
||||
/// would silently reinterpret every existing row. The `from_stored` round
|
||||
/// trip below is what holds the two in agreement.
|
||||
pub fn stored(self) -> i64 {
|
||||
match self {
|
||||
Tier::Metadata => 0,
|
||||
Tier::Preview => 1,
|
||||
Tier::Original => 2,
|
||||
}
|
||||
}
|
||||
|
||||
/// Read a tier back from the catalog.
|
||||
///
|
||||
/// An unknown value reads as [`Tier::Metadata`] — the tier that promises
|
||||
/// nothing — so a catalog written by a newer version degrades to "not
|
||||
/// cached" rather than claiming to hold pixels it does not have.
|
||||
pub fn from_stored(v: i64) -> Self {
|
||||
match v {
|
||||
2 => Tier::Original,
|
||||
1 => Tier::Preview,
|
||||
_ => Tier::Metadata,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
@@ -202,6 +232,38 @@ mod tests {
|
||||
assert_eq!(found, vec![CollectionId(1), CollectionId(2)]);
|
||||
}
|
||||
|
||||
#[test]
|
||||
#[test]
|
||||
fn stored_tiers_round_trip() {
|
||||
// These integers are on disk. A reordering of the variants that broke
|
||||
// this test would silently reinterpret every cached row as a different
|
||||
// tier — an image recorded as holding its original would come back
|
||||
// claiming metadata, or worse, the reverse.
|
||||
for t in [Tier::Metadata, Tier::Preview, Tier::Original] {
|
||||
assert_eq!(Tier::from_stored(t.stored()), t);
|
||||
}
|
||||
assert_eq!(Tier::Metadata.stored(), 0);
|
||||
assert_eq!(Tier::Preview.stored(), 1);
|
||||
assert_eq!(Tier::Original.stored(), 2);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn an_unknown_stored_tier_promises_nothing() {
|
||||
// A catalog written by a newer version must not have its unknown tier
|
||||
// read as "the original is here"; the safe direction is downwards.
|
||||
assert_eq!(Tier::from_stored(99), Tier::Metadata);
|
||||
assert_eq!(Tier::from_stored(-1), Tier::Metadata);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn stored_order_matches_generosity_order() {
|
||||
// The SQL predicates compare `tier_actual >= n`, so the stored
|
||||
// integers must sort the same way the enum does or a ">= Original"
|
||||
// query would match a Preview row.
|
||||
assert!(Tier::Metadata.stored() < Tier::Preview.stored());
|
||||
assert!(Tier::Preview.stored() < Tier::Original.stored());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn tiers_order_by_generosity() {
|
||||
// ARCH §9.3: where rules disagree, the most generous wins, which is
|
||||
|
||||
Reference in New Issue
Block a user