Many imorovments
This commit is contained in:
@@ -87,10 +87,9 @@ pub fn spawn_sync(
|
||||
let (tx, rx) = std::sync::mpsc::channel();
|
||||
|
||||
std::thread::spawn(move || {
|
||||
let rt = match tokio::runtime::Builder::new_current_thread()
|
||||
.enable_all()
|
||||
.build()
|
||||
{
|
||||
// A current-thread runtime here is what made the library scan itself
|
||||
// rather than adopt the shards the server already had; see net_runtime.
|
||||
let rt = match crate::net_runtime::build() {
|
||||
Ok(rt) => rt,
|
||||
Err(e) => {
|
||||
let _ = tx.send(SyncMessage::Failed(e.to_string()));
|
||||
|
||||
@@ -253,6 +253,24 @@ impl LaunchModel {
|
||||
};
|
||||
}
|
||||
|
||||
/// TRACES: FR-NC-1
|
||||
/// Begin a sign-in with an app password the user supplied directly.
|
||||
///
|
||||
/// The browser flow is the recommended route because it never exposes a
|
||||
/// credential to us, but it needs a browser and a round trip through one.
|
||||
/// An app password is what Nextcloud offers instead: device-scoped,
|
||||
/// individually revocable, and created by the user under Settings →
|
||||
/// Security. That makes it the workable option where no browser can
|
||||
/// complete the handshake.
|
||||
pub fn begin_direct_sign_in(&mut self, server: impl Into<String>) {
|
||||
self.server_url = normalise_server(&server.into());
|
||||
self.error = None;
|
||||
self.state = LaunchState::Busy {
|
||||
message: "Checking the credentials…".into(),
|
||||
session: None,
|
||||
};
|
||||
}
|
||||
|
||||
pub fn await_approval(&mut self, login_url: impl Into<String>) {
|
||||
self.status = Some("Approve the sign-in in your browser.".into());
|
||||
self.state = LaunchState::AwaitingApproval {
|
||||
|
||||
+121
-7
@@ -106,6 +106,28 @@ where
|
||||
});
|
||||
}
|
||||
|
||||
// --- sign in with an app password -----------------------------------
|
||||
{
|
||||
let weak = window.as_weak();
|
||||
let ctl = controller.clone();
|
||||
window.on_launch_sign_in_direct(move |server, login, password| {
|
||||
let Some(w) = weak.upgrade() else { return };
|
||||
ctl.model
|
||||
.borrow_mut()
|
||||
.begin_direct_sign_in(server.to_string());
|
||||
render(&w, &ctl);
|
||||
|
||||
let server = ctl.model.borrow().server_url.clone();
|
||||
spawn_direct_login(
|
||||
w.as_weak(),
|
||||
ctl.clone(),
|
||||
server,
|
||||
login.to_string(),
|
||||
password.to_string(),
|
||||
);
|
||||
});
|
||||
}
|
||||
|
||||
// --- sign out ------------------------------------------------------
|
||||
{
|
||||
let weak = window.as_weak();
|
||||
@@ -279,11 +301,18 @@ fn spawn_login(weak: slint::Weak<AppWindow>, ctl: Rc<LaunchController>, server:
|
||||
};
|
||||
step("thread started");
|
||||
let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(move || {
|
||||
// `enable_all()` also enables the signal driver, which wants to own
|
||||
// process-wide signal handling and is not something a worker thread
|
||||
// inside an Android app can safely claim — the flow needs only IO
|
||||
// (reqwest) and time (the poll interval), so ask for just those.
|
||||
let rt = match tokio::runtime::Builder::new_current_thread()
|
||||
// Multi-thread, not current-thread: a current-thread runtime drives
|
||||
// its reactor only while the thread sits inside `block_on`, and on
|
||||
// Android that left reqwest's connection future never polled — the
|
||||
// await never resolved, so the worker neither failed nor returned.
|
||||
// A multi-thread runtime owns worker threads that poll the reactor
|
||||
// regardless. One worker is plenty for a single login flow.
|
||||
//
|
||||
// Only IO and time are enabled; `enable_all()` would also start the
|
||||
// signal driver, which wants process-wide signal handling that an
|
||||
// Android app's runtime already owns.
|
||||
let rt = match tokio::runtime::Builder::new_multi_thread()
|
||||
.worker_threads(1)
|
||||
.enable_io()
|
||||
.enable_time()
|
||||
.build()
|
||||
@@ -313,6 +342,85 @@ fn spawn_login(weak: slint::Weak<AppWindow>, ctl: Rc<LaunchController>, server:
|
||||
poll_channel(weak, ctl, rx);
|
||||
}
|
||||
|
||||
/// TRACES: FR-NC-1
|
||||
/// Verify an app password the user supplied, then keep it.
|
||||
///
|
||||
/// No browser and no polling: one authenticated request establishes both that
|
||||
/// the credential works and what the account's canonical user id is, which is
|
||||
/// what the DAV paths are built from. Reuses the same channel and drain loop as
|
||||
/// the browser flow, so success and failure land in the UI identically.
|
||||
fn spawn_direct_login(
|
||||
weak: slint::Weak<AppWindow>,
|
||||
ctl: Rc<LaunchController>,
|
||||
server: String,
|
||||
login: String,
|
||||
password: String,
|
||||
) {
|
||||
let (tx, rx) = std::sync::mpsc::channel::<LoginMessage>();
|
||||
|
||||
std::thread::spawn(move || {
|
||||
let panic_tx = tx.clone();
|
||||
let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(move || {
|
||||
// Multi-thread for the reason the browser flow is: a current-thread
|
||||
// runtime left reqwest's future unpolled on Android.
|
||||
let rt = match tokio::runtime::Builder::new_multi_thread()
|
||||
.worker_threads(1)
|
||||
.enable_io()
|
||||
.enable_time()
|
||||
.build()
|
||||
{
|
||||
Ok(rt) => rt,
|
||||
Err(e) => {
|
||||
let _ = tx.send(LoginMessage::Failed(format!("tokio runtime: {e}")));
|
||||
return;
|
||||
}
|
||||
};
|
||||
|
||||
rt.block_on(async {
|
||||
let client = match dr_sync_nextcloud::http_client("DarkRoom") {
|
||||
Ok(c) => c,
|
||||
Err(e) => {
|
||||
let _ = tx.send(LoginMessage::Failed(format!("http client: {e}")));
|
||||
return;
|
||||
}
|
||||
};
|
||||
|
||||
let creds = dr_sync_nextcloud::AppCredentials {
|
||||
server,
|
||||
login_name: login,
|
||||
app_password: password,
|
||||
};
|
||||
|
||||
// The credential is only worth storing if it actually works,
|
||||
// and this is the cheapest request that proves it. A 401 comes
|
||||
// back as an error here rather than as a puzzling failure on
|
||||
// the first listing.
|
||||
match fetch_user_id(&client, &creds).await {
|
||||
Ok(user_id) => {
|
||||
let _ = tx.send(LoginMessage::Success(Box::new((creds, user_id))));
|
||||
}
|
||||
Err(e) => {
|
||||
let _ = tx.send(LoginMessage::Failed(format!(
|
||||
"could not sign in with that app password: {e}"
|
||||
)));
|
||||
}
|
||||
}
|
||||
});
|
||||
}));
|
||||
|
||||
if let Err(panic) = result {
|
||||
let detail = panic
|
||||
.downcast_ref::<&str>()
|
||||
.map(|s| (*s).to_string())
|
||||
.or_else(|| panic.downcast_ref::<String>().cloned())
|
||||
.unwrap_or_else(|| "panicked with a non-string payload".to_string());
|
||||
let _ = panic_tx.send(LoginMessage::Failed(format!("internal error: {detail}")));
|
||||
}
|
||||
});
|
||||
|
||||
poll_channel(weak, ctl, rx);
|
||||
}
|
||||
|
||||
/// The body of the login flow, split out so the worker above can wrap it in
|
||||
/// `catch_unwind` without a deeply nested closure.
|
||||
fn run_login_flow(
|
||||
@@ -487,8 +595,14 @@ fn spawn_folder_list(weak: slint::Weak<AppWindow>, ctl: Rc<LaunchController>, pa
|
||||
let user_id = session.user_id.clone();
|
||||
|
||||
std::thread::spawn(move || {
|
||||
let rt = tokio::runtime::Builder::new_current_thread()
|
||||
.enable_all()
|
||||
// Multi-thread for the same reason as the login worker: a
|
||||
// current-thread runtime left reqwest's connection future unpolled on
|
||||
// Android, so the await never resolved and the thread stopped without
|
||||
// failing or returning.
|
||||
let rt = tokio::runtime::Builder::new_multi_thread()
|
||||
.worker_threads(1)
|
||||
.enable_io()
|
||||
.enable_time()
|
||||
.build();
|
||||
let Ok(rt) = rt else {
|
||||
let _ = tx.send(Err("runtime".into()));
|
||||
|
||||
@@ -16,6 +16,7 @@
|
||||
|
||||
mod collections_ui;
|
||||
mod derived_sync;
|
||||
mod net_runtime;
|
||||
mod develop;
|
||||
mod labels;
|
||||
mod library;
|
||||
|
||||
+72
-24
@@ -79,6 +79,13 @@ pub struct ThumbnailReady {
|
||||
pub width: u32,
|
||||
pub height: u32,
|
||||
pub rgba: Vec<u8>,
|
||||
/// Whether these pixels came off local disk rather than the server.
|
||||
///
|
||||
/// The grid paints both identically, so this exists solely for
|
||||
/// reachability: a store hit is evidence about the *cache*, not the
|
||||
/// network, and treating one as proof of connectivity clears offline mode
|
||||
/// before a single request has been attempted.
|
||||
pub from_cache: bool,
|
||||
}
|
||||
|
||||
/// Capture metadata read from the same header the thumbnail needed.
|
||||
@@ -330,9 +337,7 @@ pub fn spawn_sidecar_writes(
|
||||
let (tx, rx) = std::sync::mpsc::channel();
|
||||
|
||||
std::thread::spawn(move || {
|
||||
let rt = match tokio::runtime::Builder::new_current_thread()
|
||||
.enable_all()
|
||||
.build()
|
||||
let rt = match crate::net_runtime::build()
|
||||
{
|
||||
Ok(rt) => rt,
|
||||
Err(e) => {
|
||||
@@ -560,9 +565,7 @@ fn run_scan(
|
||||
}
|
||||
let catalog = Catalog::open(&catalog_path).map_err(ScanFailure::local)?;
|
||||
|
||||
let rt = tokio::runtime::Builder::new_current_thread()
|
||||
.enable_all()
|
||||
.build()
|
||||
let rt = crate::net_runtime::build()
|
||||
.map_err(ScanFailure::local)?;
|
||||
|
||||
rt.block_on(async {
|
||||
@@ -887,9 +890,7 @@ pub fn spawn_pin_fetch(
|
||||
return;
|
||||
}
|
||||
|
||||
let rt = match tokio::runtime::Builder::new_current_thread()
|
||||
.enable_all()
|
||||
.build()
|
||||
let rt = match crate::net_runtime::build()
|
||||
{
|
||||
Ok(rt) => rt,
|
||||
Err(e) => {
|
||||
@@ -1055,9 +1056,7 @@ pub fn spawn_full_fetch(
|
||||
}
|
||||
}
|
||||
|
||||
let rt = match tokio::runtime::Builder::new_current_thread()
|
||||
.enable_all()
|
||||
.build()
|
||||
let rt = match crate::net_runtime::build()
|
||||
{
|
||||
Ok(rt) => rt,
|
||||
Err(e) => {
|
||||
@@ -1160,6 +1159,7 @@ pub fn spawn_thumbnails(
|
||||
width,
|
||||
height,
|
||||
rgba,
|
||||
from_cache: true,
|
||||
});
|
||||
}
|
||||
// A corrupt stored blob is a miss, not a failure.
|
||||
@@ -1202,9 +1202,7 @@ pub fn spawn_thumbnails(
|
||||
return;
|
||||
}
|
||||
|
||||
let rt = match tokio::runtime::Builder::new_current_thread()
|
||||
.enable_all()
|
||||
.build()
|
||||
let rt = match crate::net_runtime::build()
|
||||
{
|
||||
Ok(rt) => rt,
|
||||
Err(e) => {
|
||||
@@ -1392,6 +1390,7 @@ async fn fetch_one(
|
||||
width: preview.width,
|
||||
height: preview.height,
|
||||
rgba: preview.rgba,
|
||||
from_cache: false,
|
||||
}))
|
||||
}
|
||||
|
||||
@@ -1675,9 +1674,7 @@ pub fn spawn_sweep(
|
||||
return;
|
||||
}
|
||||
|
||||
let rt = match tokio::runtime::Builder::new_current_thread()
|
||||
.enable_all()
|
||||
.build()
|
||||
let rt = match crate::net_runtime::build()
|
||||
{
|
||||
Ok(rt) => rt,
|
||||
Err(e) => {
|
||||
@@ -2169,9 +2166,7 @@ mod tests {
|
||||
// The ordering guarantee is what lets a caller pair results back to
|
||||
// their inputs; without it a lane's dates could be attributed to the
|
||||
// wrong images.
|
||||
let rt = tokio::runtime::Builder::new_current_thread()
|
||||
.enable_all()
|
||||
.build()
|
||||
let rt = crate::net_runtime::build()
|
||||
.unwrap();
|
||||
|
||||
let out = rt.block_on(async {
|
||||
@@ -2195,9 +2190,7 @@ mod tests {
|
||||
|
||||
#[test]
|
||||
fn join_all_of_nothing_completes() {
|
||||
let rt = tokio::runtime::Builder::new_current_thread()
|
||||
.enable_all()
|
||||
.build()
|
||||
let rt = crate::net_runtime::build()
|
||||
.unwrap();
|
||||
let out: Vec<i32> =
|
||||
rt.block_on(async { futures_join_all(Vec::<std::future::Ready<i32>>::new()).await });
|
||||
@@ -2636,6 +2629,61 @@ mod tests {
|
||||
assert!(out[2] > out[0], "channel order survived the round trip");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_store_hit_is_not_evidence_the_server_is_reachable() {
|
||||
// The regression this guards: store hits were delivered as the same
|
||||
// `Ready` the network path sends, and the UI took any `Ready` as proof
|
||||
// of connectivity. A mostly-cached window then declared "back online"
|
||||
// against a server that was down — clearing the banner and kicking off
|
||||
// a sweep that immediately failed, on every scroll.
|
||||
let dir = std::env::temp_dir().join(format!("dr-ui-provenance-{}", std::process::id()));
|
||||
let _ = std::fs::remove_dir_all(&dir);
|
||||
|
||||
let rgba: Vec<u8> = std::iter::repeat_n([10u8, 20, 30, 255], 8 * 8)
|
||||
.flatten()
|
||||
.collect();
|
||||
let mut store = ThumbStore::open(&dir).unwrap();
|
||||
let bytes = dr_thumbs::encode_rgba(8, 8, &rgba).unwrap();
|
||||
store
|
||||
.put(99, dr_thumbs::ThumbSize::Grid, &dr_thumbs::Thumbnail {
|
||||
width: 8,
|
||||
height: 8,
|
||||
bytes,
|
||||
})
|
||||
.unwrap();
|
||||
|
||||
// The split in `spawn_thumbnails`: a hit decodes off local disk and is
|
||||
// reported with `from_cache` set, which is what the reachability gate
|
||||
// keys on.
|
||||
let stored = store.get(99, dr_thumbs::ThumbSize::Grid).unwrap().expect("stored");
|
||||
let (width, height, rgba) = dr_thumbs::decode_rgba(&stored.bytes).unwrap();
|
||||
let hit = ThumbnailReady {
|
||||
row: 0,
|
||||
width,
|
||||
height,
|
||||
rgba,
|
||||
from_cache: true,
|
||||
};
|
||||
|
||||
assert!(
|
||||
hit.from_cache,
|
||||
"a thumbnail read from the store must not be mistaken for a fetch"
|
||||
);
|
||||
|
||||
// And the mechanism it feeds: an offline tracker must survive it.
|
||||
let mut reach = dr_sync::Reachability::new();
|
||||
let now = std::time::Instant::now();
|
||||
reach.mark_unreachable("network error".into(), now);
|
||||
|
||||
if !hit.from_cache {
|
||||
reach.mark_reachable(now);
|
||||
}
|
||||
assert!(
|
||||
reach.is_offline(),
|
||||
"replaying cached thumbnails must leave offline mode intact"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn thumbs_live_beside_the_catalog_not_in_the_cache() {
|
||||
// They sync to the server and are shared with other clients, so a
|
||||
|
||||
@@ -1650,14 +1650,22 @@ fn drain_thumbnails(
|
||||
// carry no embedded preview.
|
||||
ThumbnailMessage::Ready(t) => {
|
||||
w.set_library_thumbs_done(w.get_library_thumbs_done() + 1);
|
||||
// Bytes arrived, so the server is reachable. This is
|
||||
// what clears the banner when a connection returns
|
||||
// while the user is simply scrolling, without waiting
|
||||
// for a probe or a manual retry.
|
||||
if ctl_cb
|
||||
.reachability
|
||||
.borrow_mut()
|
||||
.mark_reachable(std::time::Instant::now())
|
||||
// Bytes arrived *from the server*, so it is reachable.
|
||||
// This is what clears the banner when a connection
|
||||
// returns while the user is simply scrolling, without
|
||||
// waiting for a probe or a manual retry.
|
||||
//
|
||||
// Store hits are excluded deliberately: they are read
|
||||
// from local disk and say nothing about the network. A
|
||||
// window that is mostly cached delivers a run of them
|
||||
// before the first request is even attempted, so
|
||||
// counting them declared "back online" against a server
|
||||
// that was plainly down.
|
||||
if !t.from_cache
|
||||
&& ctl_cb
|
||||
.reachability
|
||||
.borrow_mut()
|
||||
.mark_reachable(std::time::Instant::now())
|
||||
{
|
||||
log::info!("back online");
|
||||
refresh_offline(&w, &ctl_cb);
|
||||
|
||||
@@ -0,0 +1,34 @@
|
||||
//! The tokio runtime every network worker is built from.
|
||||
//!
|
||||
//! One function rather than a builder repeated at each call site, because the
|
||||
//! choice it encodes is not obvious and was wrong in eleven places at once.
|
||||
//!
|
||||
//! **Why not `new_current_thread`.** A current-thread runtime drives its
|
||||
//! reactor only while the thread is inside `block_on`. On Android that left
|
||||
//! reqwest's connection future unpolled: the await never resolved, so the
|
||||
//! worker neither failed nor returned. Every symptom was an absence — no error,
|
||||
//! no panic for `catch_unwind` to catch, no log line, and a channel that closed
|
||||
//! only when the thread was finally torn down. What the user saw was a sign-in
|
||||
//! that stopped, and a library that quietly scanned itself rather than adopting
|
||||
//! the shards the server already held.
|
||||
//!
|
||||
//! A multi-thread runtime owns worker threads that poll the reactor regardless.
|
||||
//! One worker is enough: these are single request-response flows, not
|
||||
//! throughput-bound work.
|
||||
//!
|
||||
//! **Why not `enable_all`.** That also starts the signal driver, which wants
|
||||
//! process-wide signal handling. Inside an Android app the runtime already owns
|
||||
//! that, and a worker thread claiming it is asking for trouble. These flows need
|
||||
//! IO (reqwest) and time (poll intervals, timeouts) and nothing else.
|
||||
|
||||
/// Build a runtime suitable for a network worker thread.
|
||||
///
|
||||
/// Call from the spawned thread, not from the caller: the runtime must live on
|
||||
/// the thread that blocks on it.
|
||||
pub fn build() -> std::io::Result<tokio::runtime::Runtime> {
|
||||
tokio::runtime::Builder::new_multi_thread()
|
||||
.worker_threads(1)
|
||||
.enable_io()
|
||||
.enable_time()
|
||||
.build()
|
||||
}
|
||||
@@ -175,9 +175,7 @@ pub fn spawn_move(
|
||||
|
||||
std::thread::spawn(move || {
|
||||
let total = moves.len();
|
||||
let rt = match tokio::runtime::Builder::new_current_thread()
|
||||
.enable_all()
|
||||
.build()
|
||||
let rt = match crate::net_runtime::build()
|
||||
{
|
||||
Ok(rt) => rt,
|
||||
Err(e) => {
|
||||
@@ -291,9 +289,7 @@ pub fn spawn_purge(
|
||||
|
||||
std::thread::spawn(move || {
|
||||
let total = paths.len();
|
||||
let rt = match tokio::runtime::Builder::new_current_thread()
|
||||
.enable_all()
|
||||
.build()
|
||||
let rt = match crate::net_runtime::build()
|
||||
{
|
||||
Ok(rt) => rt,
|
||||
Err(e) => {
|
||||
|
||||
Reference in New Issue
Block a user