use std::sync::Arc; use axum::body::Body; use axum::http::{Method, Request, StatusCode}; use http_body_util::BodyExt; use jacquard_common::IntoStatic; use jacquard_oauth::{atproto::AtprotoClientMetadata, session::ClientData}; use serde::Deserialize; use serde_json::{Value, json}; use sqlx::SqlitePool; use sqlx::sqlite::SqlitePoolOptions; use tower::ServiceExt; use crate::{oauth_store::SqliteAuthStore, routes, state::AppState}; // ---- Helpers ---- /// KOSync clients send MD5(password) as x-auth-key. Mirror that here so tests /// exercise the same code path as real clients. fn md5_hex(s: &str) -> String { format!("{:x}", md5::compute(s)) } struct TestApp { router: axum::Router, db: SqlitePool, } async fn make_app_inner(email_client: Option>) -> TestApp { let db = SqlitePoolOptions::new() .connect("sqlite::memory:") .await .unwrap(); sqlx::migrate!().run(&db).await.unwrap(); let metadata_json = json!({ "client_id": "https://example.com/oauth/client-metadata.json", "client_name": "test", "client_uri": "https://example.com", "redirect_uris": ["https://example.com/oauth/callback"], "scopes": ["atproto", "transition:email", "repo:social.popfeed.book.progress?action=create", "repo:social.popfeed.feed.listItem?action=create&action=update", "repo:page.headway.sync.link?action=create&action=update", "repo:page.headway.sync.progress?action=create&action=update"], "grant_types": ["authorization_code", "refresh_token"], "jwks_uri": null }); let metadata_str = metadata_json.to_string(); let mut meta_de = serde_json::Deserializer::from_str(&metadata_str); let metadata = AtprotoClientMetadata::deserialize(&mut meta_de).unwrap(); let client_data = ClientData::new_public(metadata.into_static()); let oauth_store = SqliteAuthStore { db: db.clone() }; let oauth_client = jacquard_oauth::client::OAuthClient::new(oauth_store, client_data); let state = AppState { db: db.clone(), base_url: "https://example.com".to_owned(), oauth_client: Arc::new(oauth_client), client_metadata_json: Arc::new(metadata_json), email_client, sync_throttle: Arc::new(crate::sync_throttle::Throttle::new( std::time::Duration::from_secs(300), )), }; TestApp { router: routes::router(state), db } } async fn make_app() -> TestApp { make_app_inner(None).await } async fn make_app_with_email_client() -> TestApp { make_app_inner(Some(Arc::new(crate::email::EmailClient::new( "test-api-key".to_string(), "notifications@example.com".to_string(), )))) .await } async fn insert_user(db: &SqlitePool, username: &str, password: &str) { let hash = bcrypt::hash(md5_hex(password), 4).unwrap(); sqlx::query( "INSERT INTO users (atproto_did, atproto_handle, kosync_username, kosync_key_hash) VALUES (?, ?, ?, ?)", ) .bind(format!("did:plc:{username}test")) .bind(format!("{username}.bsky.social")) .bind(username) .bind(&hash) .execute(db) .await .unwrap(); } async fn send(router: axum::Router, req: Request) -> axum::response::Response { router.oneshot(req).await.unwrap() } async fn body_json(resp: axum::response::Response) -> Value { let bytes = resp.into_body().collect().await.unwrap().to_bytes(); serde_json::from_slice(&bytes).unwrap() } fn get_req(uri: &str) -> Request { Request::builder() .method(Method::GET) .uri(uri) .body(Body::empty()) .unwrap() } fn authed_get(uri: &str, user: &str, password: &str) -> Request { Request::builder() .method(Method::GET) .uri(uri) .header("x-auth-user", user) .header("x-auth-key", md5_hex(password)) .body(Body::empty()) .unwrap() } fn opds_get(uri: &str, user: &str, password: &str) -> Request { use base64::Engine; let creds = base64::engine::general_purpose::STANDARD.encode(format!("{user}:{password}")); Request::builder() .method(Method::GET) .uri(uri) .header("authorization", format!("Basic {creds}")) .body(Body::empty()) .unwrap() } fn put_progress_req(user: &str, password: &str, body: Value) -> Request { Request::builder() .method(Method::PUT) .uri("/1/syncs/progress") .header("content-type", "application/json") .header("x-auth-user", user) .header("x-auth-key", md5_hex(password)) .body(Body::from(serde_json::to_vec(&body).unwrap())) .unwrap() } async fn insert_user_returning_id(db: &SqlitePool, username: &str, password: &str) -> i64 { let hash = bcrypt::hash(md5_hex(password), 4).unwrap(); sqlx::query_scalar::<_, i64>( "INSERT INTO users (atproto_did, atproto_handle, kosync_username, kosync_key_hash) VALUES (?, ?, ?, ?) RETURNING id", ) .bind(format!("did:plc:{username}test")) .bind(format!("{username}.bsky.social")) .bind(username) .bind(&hash) .fetch_one(db) .await .unwrap() } async fn insert_web_session(db: &SqlitePool, user_id: i64) -> String { let token = crate::web_auth::generate_session_token(); sqlx::query("INSERT INTO web_sessions (user_id, token) VALUES (?, ?)") .bind(user_id) .bind(&token) .execute(db) .await .unwrap(); token } fn web_get(uri: &str, token: &str) -> Request { Request::builder() .method(Method::GET) .uri(uri) .header("cookie", format!("headway_session={token}")) .body(Body::empty()) .unwrap() } fn web_post_form(uri: &str, token: &str, body: &str) -> Request { Request::builder() .method(Method::POST) .uri(uri) .header("cookie", format!("headway_session={token}")) .header("content-type", "application/x-www-form-urlencoded") .body(Body::from(body.to_owned())) .unwrap() } async fn body_string(resp: axum::response::Response) -> String { let bytes = resp.into_body().collect().await.unwrap().to_bytes(); String::from_utf8(bytes.to_vec()).unwrap() } // ---- Health ---- #[tokio::test] async fn test_health_ok() { let app = make_app().await; let resp = send(app.router, get_req("/health")).await; assert_eq!(resp.status(), StatusCode::OK); let body = body_json(resp).await; assert_eq!(body["status"], "ok"); } // ---- OAuth static endpoints ---- #[tokio::test] async fn test_client_metadata_structure() { let app = make_app().await; let resp = send(app.router, get_req("/oauth/client-metadata.json")).await; assert_eq!(resp.status(), StatusCode::OK); let body = body_json(resp).await; assert!(body["client_id"].is_string()); assert!(body["redirect_uris"].is_array()); assert!(body["grant_types"].is_array()); // Scope guarantees: the published metadata MUST match the narrowed scopes // the app actually needs (see src/main.rs). If anyone widens or shifts // these without updating both call sites, the published surface drifts // out of sync with what the OAuth flow asks for at consent time. // // The served document may carry scopes either as a `scope` space-separated // string (RFC 6749 OAuth client metadata) or a `scopes` array // (ATProto-flavored client metadata). Accept either, then assert on the // set of distinct scope strings. let scopes: std::collections::HashSet = if let Some(s) = body["scope"].as_str() { s.split_whitespace().map(|t| t.to_owned()).collect() } else if let Some(arr) = body["scopes"].as_array() { arr.iter().filter_map(|v| v.as_str().map(str::to_owned)).collect() } else { panic!("expected either 'scope' (string) or 'scopes' (array) in client metadata"); }; assert!(scopes.contains("atproto"), "missing atproto scope"); assert!(scopes.contains("transition:email"), "missing transition:email scope"); assert!( scopes.contains("repo:social.popfeed.book.progress?action=create"), "book.progress should be create-only (append-only log model); got: {scopes:?}" ); assert!( scopes.contains("repo:social.popfeed.feed.listItem?action=create&action=update"), "listItem must permit both create and update; got: {scopes:?}" ); assert!( scopes.contains("repo:page.headway.sync.link?action=create&action=update"), "sync.link must permit both create and update (putRecord creates or merges); got: {scopes:?}" ); assert!( scopes.contains("repo:page.headway.sync.progress?action=create&action=update"), "sync.progress must permit both create and update (putRecord overwrites); got: {scopes:?}" ); // Defenses against accidental widening. let serialized = serde_json::to_string(&body).unwrap(); assert!( !serialized.contains("social.popfeed.book.progress?action=create&action=update"), "book.progress must not carry action=update (records are append-only)" ); assert!( !serialized.contains("action=delete"), "no scope should carry action=delete" ); } #[tokio::test] async fn test_jwks_returns_json() { let app = make_app().await; let resp = send(app.router, get_req("/oauth/jwks.json")).await; assert_eq!(resp.status(), StatusCode::OK); let body = body_json(resp).await; assert!(body.is_object()); } // ---- KOSync: create user ---- #[tokio::test] async fn test_create_user_is_forbidden() { let app = make_app().await; let req = Request::builder() .method(Method::POST) .uri("/1/users/create") .header("content-type", "application/json") .body(Body::from(r#"{"username":"x","password":"y"}"#)) .unwrap(); let resp = send(app.router, req).await; assert_eq!(resp.status(), StatusCode::FORBIDDEN); } // ---- KOSync: auth ---- #[tokio::test] async fn test_auth_no_credentials_returns_401() { let app = make_app().await; let resp = send(app.router, get_req("/1/users/auth")).await; assert_eq!(resp.status(), StatusCode::UNAUTHORIZED); } #[tokio::test] async fn test_auth_unknown_user_returns_401() { let app = make_app().await; let resp = send(app.router, authed_get("/1/users/auth", "nobody", "anykey")).await; assert_eq!(resp.status(), StatusCode::UNAUTHORIZED); } #[tokio::test] async fn test_auth_wrong_password_returns_401() { let app = make_app().await; insert_user(&app.db, "alice", "correct_password").await; let resp = send( app.router, authed_get("/1/users/auth", "alice", "wrong_password"), ) .await; assert_eq!(resp.status(), StatusCode::UNAUTHORIZED); } #[tokio::test] async fn test_auth_valid_credentials_returns_200() { let app = make_app().await; insert_user(&app.db, "bob", "bobspassword").await; let resp = send( app.router, authed_get("/1/users/auth", "bob", "bobspassword"), ) .await; assert_eq!(resp.status(), StatusCode::OK); let body = body_json(resp).await; assert_eq!(body["authorized"], "OK"); } // ---- KOSync: put progress ---- #[tokio::test] async fn test_put_progress_no_auth_returns_401() { let app = make_app().await; let req = Request::builder() .method(Method::PUT) .uri("/1/syncs/progress") .header("content-type", "application/json") .body(Body::from( r#"{"document":"book.epub","progress":"3/100","percentage":0.03,"device":"Kindle"}"#, )) .unwrap(); let resp = send(app.router, req).await; assert_eq!(resp.status(), StatusCode::UNAUTHORIZED); } #[tokio::test] async fn test_put_progress_creates_record() { let app = make_app().await; insert_user(&app.db, "carol", "carolpass").await; let resp = send( app.router, put_progress_req( "carol", "carolpass", json!({ "document": "mybook.epub", "progress": "42/200", "percentage": 0.21, "device": "KOReader" }), ), ) .await; assert_eq!(resp.status(), StatusCode::OK); let body = body_json(resp).await; assert_eq!(body["document"], "mybook.epub"); assert!(body["timestamp"].is_number()); } #[tokio::test] async fn test_put_progress_returns_latest_local_value() { // KOSync semantics: GET /1/syncs/progress returns whatever the last // PUT was. The local SQLite row is upserted on every PUT (this is // unrelated to the append-only `book.progress` log we write to the // user's PDS — that side genuinely creates a new record per emit). let app = make_app().await; insert_user(&app.db, "dave", "davepass").await; for pct in [0.10, 0.50, 0.90] { let resp = send( app.router.clone(), put_progress_req( "dave", "davepass", json!({ "document": "samebook.epub", "progress": format!("{}/100", (pct * 100.0) as u32), "percentage": pct, "device": "KOReader" }), ), ) .await; assert_eq!(resp.status(), StatusCode::OK); } let resp = send( app.router, authed_get("/1/syncs/progress/samebook.epub", "dave", "davepass"), ) .await; assert_eq!(resp.status(), StatusCode::OK); let body = body_json(resp).await; let got: f64 = body["percentage"].as_f64().unwrap(); assert!((got - 0.90).abs() < 1e-9, "expected 0.90, got {got}"); } #[tokio::test] async fn test_put_progress_round_trips_client_timestamp() { // The server must echo KOReader's page-turn timestamp rather than stamp its // own wall-clock; otherwise its record always looks like the newest writer // and KOReader raises a spurious "server vs device" conflict prompt. let app = make_app().await; insert_user(&app.db, "tina", "tinapass").await; let resp = send( app.router.clone(), put_progress_req( "tina", "tinapass", json!({ "document": "clock.epub", "progress": "10/100", "percentage": 0.10, "device": "KOReader", "timestamp": 1_600_000_000i64 }), ), ) .await; assert_eq!(resp.status(), StatusCode::OK); assert_eq!(body_json(resp).await["timestamp"], 1_600_000_000i64); let resp = send( app.router, authed_get("/1/syncs/progress/clock.epub", "tina", "tinapass"), ) .await; assert_eq!(body_json(resp).await["timestamp"], 1_600_000_000i64); } #[tokio::test] async fn test_put_progress_stale_timestamp_does_not_clobber() { // Monotonic write guard (optimistic locking): an older-timestamped write // must not overwrite a newer position, and its response reports the // authoritative newer timestamp so the client learns it is behind. let app = make_app().await; insert_user(&app.db, "una", "unapass").await; send( app.router.clone(), put_progress_req( "una", "unapass", json!({ "document": "guard.epub", "progress": "90/100", "percentage": 0.90, "device": "KOReader", "timestamp": 2000i64 }), ), ) .await; let resp = send( app.router.clone(), put_progress_req( "una", "unapass", json!({ "document": "guard.epub", "progress": "10/100", "percentage": 0.10, "device": "KOReader", "timestamp": 1000i64 }), ), ) .await; assert_eq!(resp.status(), StatusCode::OK); assert_eq!(body_json(resp).await["timestamp"], 2000i64); let resp = send( app.router, authed_get("/1/syncs/progress/guard.epub", "una", "unapass"), ) .await; let body = body_json(resp).await; let got = body["percentage"].as_f64().unwrap(); assert!((got - 0.90).abs() < 1e-9, "stale write clobbered position: got {got}"); assert_eq!(body["timestamp"], 2000i64); } #[tokio::test] async fn test_put_progress_newer_timestamp_applies() { let app = make_app().await; insert_user(&app.db, "vic", "vicpass").await; for (progress, pct, ts) in [("10/100", 0.10, 1000i64), ("50/100", 0.50, 2000i64)] { send( app.router.clone(), put_progress_req( "vic", "vicpass", json!({ "document": "fwd.epub", "progress": progress, "percentage": pct, "device": "KOReader", "timestamp": ts }), ), ) .await; } let resp = send( app.router, authed_get("/1/syncs/progress/fwd.epub", "vic", "vicpass"), ) .await; let body = body_json(resp).await; let got = body["percentage"].as_f64().unwrap(); assert!((got - 0.50).abs() < 1e-9, "expected 0.50, got {got}"); assert_eq!(body["timestamp"], 2000i64); } #[tokio::test] async fn test_put_progress_far_future_timestamp_falls_back_to_server_clock() { // A far-future client timestamp (broken clock, or milliseconds sent as // seconds) must not be stored: it would win the monotonic guard against // every legitimate write for that document, forever. let app = make_app().await; insert_user(&app.db, "wes", "wespass").await; let millis_not_secs = 1_700_000_000_000i64; let resp = send( app.router.clone(), put_progress_req( "wes", "wespass", json!({ "document": "skew.epub", "progress": "10/100", "percentage": 0.10, "device": "KOReader", "timestamp": millis_not_secs }), ), ) .await; assert_eq!(resp.status(), StatusCode::OK); let ts = body_json(resp).await["timestamp"].as_i64().unwrap(); assert!(ts < millis_not_secs, "far-future timestamp was stored: {ts}"); // The document is not poisoned: a later timestamp-less write still applies. send( app.router.clone(), put_progress_req( "wes", "wespass", json!({ "document": "skew.epub", "progress": "20/100", "percentage": 0.20, "device": "KOReader" }), ), ) .await; let resp = send( app.router, authed_get("/1/syncs/progress/skew.epub", "wes", "wespass"), ) .await; let got = body_json(resp).await["percentage"].as_f64().unwrap(); assert!((got - 0.20).abs() < 1e-9, "expected 0.20, got {got}"); } // ---- KOSync: get progress ---- #[tokio::test] async fn test_get_progress_unknown_document_returns_empty() { let app = make_app().await; insert_user(&app.db, "eve", "evepass").await; let resp = send( app.router, authed_get("/1/syncs/progress/nonexistent.epub", "eve", "evepass"), ) .await; assert_eq!(resp.status(), StatusCode::OK); let body = body_json(resp).await; assert_eq!(body, json!({})); } #[tokio::test] async fn test_get_progress_after_put_returns_correct_data() { let app = make_app().await; insert_user(&app.db, "frank", "frankpass").await; send( app.router.clone(), put_progress_req( "frank", "frankpass", json!({ "document": "novel.epub", "progress": "75/300", "percentage": 0.25, "device": "PocketBook" }), ), ) .await; let resp = send( app.router, authed_get("/1/syncs/progress/novel.epub", "frank", "frankpass"), ) .await; assert_eq!(resp.status(), StatusCode::OK); let body = body_json(resp).await; assert_eq!(body["document"], "novel.epub"); assert_eq!(body["progress"], "75/300"); assert_eq!(body["device"], "PocketBook"); let pct: f64 = body["percentage"].as_f64().unwrap(); assert!((pct - 0.25).abs() < 1e-9, "expected 0.25, got {pct}"); assert!(body["timestamp"].as_i64().unwrap() > 0); } #[tokio::test] async fn test_get_progress_isolated_per_user() { let app = make_app().await; insert_user(&app.db, "grace", "gracepass").await; insert_user(&app.db, "henry", "henrypass").await; send( app.router.clone(), put_progress_req( "grace", "gracepass", json!({ "document": "shared.epub", "progress": "10/100", "percentage": 0.10, "device": "KOReader" }), ), ) .await; let resp = send( app.router, authed_get("/1/syncs/progress/shared.epub", "henry", "henrypass"), ) .await; assert_eq!(resp.status(), StatusCode::OK); let body = body_json(resp).await; assert_eq!(body, json!({}), "henry should not see grace's progress"); } // ---- KOSync: get progress — error paths ---- #[tokio::test] async fn test_get_progress_no_auth_returns_401() { let app = make_app().await; let resp = send(app.router, get_req("/1/syncs/progress/somebook.epub")).await; assert_eq!(resp.status(), StatusCode::UNAUTHORIZED); } // ---- KOSync: put progress — boundary and error paths ---- #[tokio::test] async fn test_put_progress_malformed_json_returns_422() { let app = make_app().await; insert_user(&app.db, "ivan", "ivanpass").await; let req = Request::builder() .method(Method::PUT) .uri("/1/syncs/progress") .header("content-type", "application/json") .header("x-auth-user", "ivan") .header("x-auth-key", md5_hex("ivanpass")) .body(Body::from("not json at all")) .unwrap(); let resp = send(app.router, req).await; assert_eq!(resp.status(), StatusCode::BAD_REQUEST); } #[tokio::test] async fn test_put_progress_missing_required_field_returns_422() { let app = make_app().await; insert_user(&app.db, "julia", "juliapass").await; // Missing "device" field let req = Request::builder() .method(Method::PUT) .uri("/1/syncs/progress") .header("content-type", "application/json") .header("x-auth-user", "julia") .header("x-auth-key", md5_hex("juliapass")) .body(Body::from(r#"{"document":"book.epub","progress":"1/10","percentage":0.1}"#)) .unwrap(); let resp = send(app.router, req).await; assert_eq!(resp.status(), StatusCode::UNPROCESSABLE_ENTITY); } #[tokio::test] async fn test_put_progress_zero_percentage() { let app = make_app().await; insert_user(&app.db, "karen", "karenpass").await; let resp = send( app.router.clone(), put_progress_req("karen", "karenpass", json!({ "document": "start.epub", "progress": "0/100", "percentage": 0.0, "device": "KOReader" })), ) .await; assert_eq!(resp.status(), StatusCode::OK); let resp = send( app.router, authed_get("/1/syncs/progress/start.epub", "karen", "karenpass"), ) .await; let body = body_json(resp).await; assert!((body["percentage"].as_f64().unwrap() - 0.0).abs() < 1e-9); } #[tokio::test] async fn test_put_progress_full_percentage() { let app = make_app().await; insert_user(&app.db, "liam", "liampass").await; let resp = send( app.router.clone(), put_progress_req("liam", "liampass", json!({ "document": "done.epub", "progress": "100/100", "percentage": 1.0, "device": "KOReader" })), ) .await; assert_eq!(resp.status(), StatusCode::OK); let resp = send( app.router, authed_get("/1/syncs/progress/done.epub", "liam", "liampass"), ) .await; let body = body_json(resp).await; assert!((body["percentage"].as_f64().unwrap() - 1.0).abs() < 1e-9); } // ---- OAuth: callback error path ---- #[tokio::test] async fn test_oauth_callback_error_param_returns_400() { let app = make_app().await; let resp = send( app.router, get_req("/oauth/callback?error=access_denied&error_description=User+denied+access"), ) .await; assert_eq!(resp.status(), StatusCode::BAD_REQUEST); } // ---- Settings: GET ---- #[tokio::test] async fn test_settings_get_unauthenticated_redirects() { let app = make_app().await; let resp = send(app.router, get_req("/settings")).await; assert_eq!(resp.status(), StatusCode::SEE_OTHER); } #[tokio::test] async fn test_settings_get_shows_form() { let app = make_app().await; let user_id = insert_user_returning_id(&app.db, "sue", "suepass").await; let token = insert_web_session(&app.db, user_id).await; let resp = send(app.router, web_get("/settings", &token)).await; assert_eq!(resp.status(), StatusCode::OK); let body = body_string(resp).await; assert!(body.contains("Notification email"), "form should contain email field label"); } #[tokio::test] async fn test_settings_get_shows_saved_banner() { let app = make_app().await; let user_id = insert_user_returning_id(&app.db, "sue2", "suepass2").await; let token = insert_web_session(&app.db, user_id).await; let resp = send(app.router, web_get("/settings?saved=1", &token)).await; assert_eq!(resp.status(), StatusCode::OK); let body = body_string(resp).await; assert!(body.contains("Settings saved."), "saved banner should be present"); } // ---- Settings: POST ---- #[tokio::test] async fn test_settings_post_saves_email_and_redirects() { let app = make_app().await; let user_id = insert_user_returning_id(&app.db, "tina", "tinapass").await; let token = insert_web_session(&app.db, user_id).await; let resp = send( app.router, web_post_form("/settings", &token, "email=tina%40example.com¬ify_new_doc=on"), ) .await; assert_eq!(resp.status(), StatusCode::SEE_OTHER); let location = resp.headers().get("location").and_then(|v| v.to_str().ok()).unwrap_or(""); assert!(location.contains("saved=1"), "should redirect to ?saved=1"); let email: Option = sqlx::query_scalar("SELECT email FROM users WHERE id = ?") .bind(user_id) .fetch_one(&app.db) .await .unwrap(); assert_eq!(email.as_deref(), Some("tina@example.com")); } #[tokio::test] async fn test_settings_post_clears_email_when_blank() { let app = make_app().await; let user_id = insert_user_returning_id(&app.db, "uma", "umapass").await; sqlx::query("UPDATE users SET email = 'uma@example.com' WHERE id = ?") .bind(user_id) .execute(&app.db) .await .unwrap(); let token = insert_web_session(&app.db, user_id).await; let resp = send( app.router, web_post_form("/settings", &token, "email=¬ify_new_doc=on"), ) .await; assert_eq!(resp.status(), StatusCode::SEE_OTHER); let email: Option = sqlx::query_scalar("SELECT email FROM users WHERE id = ?") .bind(user_id) .fetch_one(&app.db) .await .unwrap(); assert!(email.is_none(), "blank email should clear the stored address"); } #[tokio::test] async fn test_settings_post_invalid_email_returns_error_page() { let app = make_app().await; let user_id = insert_user_returning_id(&app.db, "vera", "verapass").await; let token = insert_web_session(&app.db, user_id).await; let resp = send( app.router, web_post_form("/settings", &token, "email=notanemail¬ify_new_doc=on"), ) .await; assert_eq!(resp.status(), StatusCode::OK); let body = body_string(resp).await; assert!(body.contains("Invalid email address"), "should show validation error"); } #[tokio::test] async fn test_settings_post_unchecked_notify_sets_zero() { let app = make_app().await; let user_id = insert_user_returning_id(&app.db, "will", "willpass").await; let token = insert_web_session(&app.db, user_id).await; // Omit notify_new_doc entirely (unchecked HTML checkbox sends nothing) let resp = send( app.router, web_post_form("/settings", &token, "email=will%40example.com"), ) .await; assert_eq!(resp.status(), StatusCode::SEE_OTHER); let notify: i64 = sqlx::query_scalar("SELECT notify_new_doc FROM users WHERE id = ?") .bind(user_id) .fetch_one(&app.db) .await .unwrap(); assert_eq!(notify, 0, "absent checkbox should store notify_new_doc = 0"); } // ---- Documents: email nudge banner ---- #[tokio::test] async fn test_documents_nudge_shown_when_no_email() { let app = make_app_with_email_client().await; let user_id = insert_user_returning_id(&app.db, "noel", "noelpass").await; let token = insert_web_session(&app.db, user_id).await; let resp = send(app.router, web_get("/documents", &token)).await; assert_eq!(resp.status(), StatusCode::OK); let body = body_string(resp).await; assert!( body.contains("add your email in Settings"), "nudge should appear when no email is stored" ); } #[tokio::test] async fn test_documents_nudge_hidden_when_email_set() { let app = make_app_with_email_client().await; let user_id = insert_user_returning_id(&app.db, "nora", "norapass").await; sqlx::query("UPDATE users SET email = 'nora@example.com' WHERE id = ?") .bind(user_id) .execute(&app.db) .await .unwrap(); let token = insert_web_session(&app.db, user_id).await; let resp = send(app.router, web_get("/documents", &token)).await; assert_eq!(resp.status(), StatusCode::OK); let body = body_string(resp).await; assert!( !body.contains("add your email in Settings"), "nudge should not appear when email is already set" ); } #[tokio::test] async fn test_documents_nudge_hidden_without_email_client() { let app = make_app().await; // email_client is None let user_id = insert_user_returning_id(&app.db, "ned", "nedpass").await; let token = insert_web_session(&app.db, user_id).await; let resp = send(app.router, web_get("/documents", &token)).await; assert_eq!(resp.status(), StatusCode::OK); let body = body_string(resp).await; assert!( !body.contains("add your email in Settings"), "nudge should not appear when email delivery is not configured" ); } // ---- Home page: logged-in redirect ---- #[tokio::test] async fn test_index_logged_in_redirects_to_documents() { let app = make_app().await; let user_id = insert_user_returning_id(&app.db, "otto", "ottopass").await; let token = insert_web_session(&app.db, user_id).await; let resp = send(app.router, web_get("/", &token)).await; assert_eq!(resp.status(), StatusCode::SEE_OTHER); let location = resp.headers().get("location").and_then(|v| v.to_str().ok()).unwrap_or(""); assert_eq!(location, "/documents", "logged-in user should be redirected to /documents"); } #[tokio::test] async fn test_index_unauthenticated_shows_form() { let app = make_app().await; let resp = send(app.router, get_req("/")).await; assert_eq!(resp.status(), StatusCode::OK); let body = body_string(resp).await; assert!(body.contains("ATProto handle"), "home page should show the sign-in form"); } // ---- Document merge ---- async fn insert_document( db: &SqlitePool, user_id: i64, kosync_id: &str, percent: f64, status: &str, ) -> (i64, String) { let public_id = crate::web_auth::generate_public_id(); let id = sqlx::query_scalar::<_, i64>( "INSERT INTO documents (user_id, kosync_document_id, percent, association_status, public_id) VALUES (?, ?, ?, ?, ?) RETURNING id", ) .bind(user_id) .bind(kosync_id) .bind(percent) .bind(status) .bind(&public_id) .fetch_one(db) .await .unwrap(); (id, public_id) } #[tokio::test] async fn test_merge_sets_merged_into_id_and_redirects() { let app = make_app().await; let user_id = insert_user_returning_id(&app.db, "mergeuser", "mergepass").await; let token = insert_web_session(&app.db, user_id).await; let (dup_id, dup_pub_id) = insert_document(&app.db, user_id, "hash-device-a", 42.0, "unassociated").await; let (canon_id, canon_pub_id) = insert_document(&app.db, user_id, "hash-device-b", 55.0, "associated").await; let resp = send( app.router, web_post_form( &format!("/documents/{dup_pub_id}/merge"), &token, &format!("canonical_id={canon_pub_id}"), ), ) .await; // Should redirect to the canonical doc assert_eq!(resp.status(), StatusCode::SEE_OTHER); let location = resp.headers().get("location").and_then(|v| v.to_str().ok()).unwrap_or(""); assert_eq!(location, format!("/documents/{canon_pub_id}")); // Duplicate row should now have merged_into_id set let merged: Option = sqlx::query_scalar( "SELECT merged_into_id FROM documents WHERE id = ?", ) .bind(dup_id) .fetch_one(&app.db) .await .unwrap(); assert_eq!(merged, Some(canon_id)); } #[tokio::test] async fn test_merge_promotes_higher_progress_to_canonical() { let app = make_app().await; let user_id = insert_user_returning_id(&app.db, "mergeuser2", "pass").await; let token = insert_web_session(&app.db, user_id).await; // Duplicate has higher progress than canonical let (_, dup_pub_id) = insert_document(&app.db, user_id, "hash-a", 80.0, "unassociated").await; let (canon_id, canon_pub_id) = insert_document(&app.db, user_id, "hash-b", 30.0, "associated").await; let resp = send( app.router, web_post_form( &format!("/documents/{dup_pub_id}/merge"), &token, &format!("canonical_id={canon_pub_id}"), ), ) .await; assert_eq!(resp.status(), StatusCode::SEE_OTHER); let canon_pct: f64 = sqlx::query_scalar( "SELECT percent FROM documents WHERE id = ?", ) .bind(canon_id) .fetch_one(&app.db) .await .unwrap(); assert_eq!(canon_pct, 80.0, "canonical should inherit the higher progress"); } #[tokio::test] async fn test_put_progress_on_merged_doc_updates_canonical() { let app = make_app().await; insert_user(&app.db, "syncuser", "syncpass").await; let user_id: i64 = sqlx::query_scalar("SELECT id FROM users WHERE kosync_username = 'syncuser'") .fetch_one(&app.db) .await .unwrap(); let (_, _) = insert_document(&app.db, user_id, "old-hash", 10.0, "unassociated").await; let (canon_id, _) = insert_document(&app.db, user_id, "new-hash", 10.0, "associated").await; // Merge old-hash into new-hash directly in DB sqlx::query("UPDATE documents SET merged_into_id = ? WHERE user_id = ? AND kosync_document_id = 'old-hash'") .bind(canon_id) .bind(user_id) .execute(&app.db) .await .unwrap(); // KOSync PUT using the old (merged) ID let resp = send( app.router, put_progress_req("syncuser", "syncpass", json!({ "document": "old-hash", "progress": "100/200", "percentage": 0.75, "device": "device-a", "device_id": "dev-a" })), ) .await; assert_eq!(resp.status(), StatusCode::OK); // Canonical doc's percent should be updated let canon_pct: f64 = sqlx::query_scalar("SELECT percent FROM documents WHERE id = ?") .bind(canon_id) .fetch_one(&app.db) .await .unwrap(); assert!((canon_pct - 75.0).abs() < 0.01, "canonical percent should be 75, got {canon_pct}"); // Merged doc's percent should NOT have changed let dup_pct: f64 = sqlx::query_scalar( "SELECT percent FROM documents WHERE user_id = ? AND kosync_document_id = 'old-hash'", ) .bind(user_id) .fetch_one(&app.db) .await .unwrap(); assert!((dup_pct - 10.0).abs() < 0.01, "dup percent should be unchanged at 10, got {dup_pct}"); } #[tokio::test] async fn test_get_progress_on_merged_doc_returns_canonical_progress() { let app = make_app().await; insert_user(&app.db, "getuser", "getpass").await; let user_id: i64 = sqlx::query_scalar("SELECT id FROM users WHERE kosync_username = 'getuser'") .fetch_one(&app.db) .await .unwrap(); let (_, _) = insert_document(&app.db, user_id, "old-id", 20.0, "unassociated").await; let (canon_id, _) = insert_document(&app.db, user_id, "new-id", 65.0, "associated").await; sqlx::query("UPDATE documents SET merged_into_id = ? WHERE user_id = ? AND kosync_document_id = 'old-id'") .bind(canon_id) .bind(user_id) .execute(&app.db) .await .unwrap(); // GET progress using the old (merged) ID — should return canonical's 65% let resp = send(app.router, authed_get("/1/syncs/progress/old-id", "getuser", "getpass")).await; assert_eq!(resp.status(), StatusCode::OK); let body = body_json(resp).await; let pct = body["percentage"].as_f64().unwrap_or(0.0); assert!((pct - 0.65).abs() < 0.001, "should return canonical percentage 0.65, got {pct}"); } // ---- Settings: regenerate password ---- #[tokio::test] async fn test_regenerate_password_shows_credentials() { let app = make_app().await; let user_id = insert_user_returning_id(&app.db, "petra", "petrapass").await; let token = insert_web_session(&app.db, user_id).await; let resp = send(app.router, web_post_form("/settings/regenerate", &token, "")).await; assert_eq!(resp.status(), StatusCode::OK); let body = body_string(resp).await; assert!(body.contains("Connected"), "credentials page should be shown after regeneration"); assert!(body.contains("KOReader"), "credentials page should contain KOReader setup instructions"); } #[tokio::test] async fn test_regenerate_password_updates_hash_in_db() { let app = make_app().await; let user_id = insert_user_returning_id(&app.db, "quinn", "quinnpass").await; let token = insert_web_session(&app.db, user_id).await; let old_hash: String = sqlx::query_scalar("SELECT kosync_key_hash FROM users WHERE id = ?") .bind(user_id) .fetch_one(&app.db) .await .unwrap(); send(app.router, web_post_form("/settings/regenerate", &token, "")).await; let new_hash: String = sqlx::query_scalar("SELECT kosync_key_hash FROM users WHERE id = ?") .bind(user_id) .fetch_one(&app.db) .await .unwrap(); assert_ne!(old_hash, new_hash, "kosync_key_hash should change after regeneration"); } // ---- Lexicon schema ---- /// The `page.headway.sync.link` lexicon must be a well-formed Lexicon 1 document /// per jacquard's own schema model (the same model the codegen path consumes). /// This guards the hand-authored JSON as the schema evolves. #[test] fn link_lexicon_is_valid_lexicon_doc() { use jacquard_lexicon::lexicon::{ LexObjectProperty, LexRecordRecord, LexUserType, Lexicon, LexiconDoc, }; const LINK_LEXICON: &str = include_str!("../lexicons/page/headway/sync/link.json"); let doc: LexiconDoc = serde_json::from_str(LINK_LEXICON) .expect("link.json must parse as a jacquard LexiconDoc"); assert_eq!(doc.lexicon, Lexicon::Lexicon1); assert_eq!(doc.id.as_ref(), "page.headway.sync.link"); // `main` is a record whose rkey is application-chosen (`any`) so it can reuse // the referenced listItem's rkey — the one-link-record-per-book invariant. let main = doc.defs.get("main").expect("missing main def"); let LexUserType::Record(record) = main else { panic!("main def must be a record, got {main:?}"); }; assert_eq!( record.key.as_ref().map(|k| k.as_ref()), Some("any"), "record key must be `any` (app reuses the listItem rkey)" ); // documentIds must be an array of (bare-string) identifiers, allowing many // KOSync hashes to resolve to one book. LexRecordRecord has a single Object // variant, so this destructuring is irrefutable. let LexRecordRecord::Object(obj) = &record.record; assert!( matches!(obj.properties.get("documentIds"), Some(LexObjectProperty::Array(_))), "documentIds must be an array property" ); assert!(obj.properties.contains_key("target"), "missing target property"); } #[test] fn progress_lexicon_is_valid_lexicon_doc() { use jacquard_lexicon::lexicon::{LexRecordRecord, LexUserType, Lexicon, LexiconDoc}; const PROGRESS_LEXICON: &str = include_str!("../lexicons/page/headway/sync/progress.json"); let doc: LexiconDoc = serde_json::from_str(PROGRESS_LEXICON) .expect("progress.json must parse as a jacquard LexiconDoc"); assert_eq!(doc.lexicon, Lexicon::Lexicon1); assert_eq!(doc.id.as_ref(), "page.headway.sync.progress"); let main = doc.defs.get("main").expect("missing main def"); let LexUserType::Record(record) = main else { panic!("main def must be a record, got {main:?}"); }; assert_eq!(record.key.as_ref().map(|k| k.as_ref()), Some("any")); let LexRecordRecord::Object(obj) = &record.record; for field in ["document", "progress", "percentage", "updatedAt"] { assert!(obj.properties.contains_key(field), "missing {field} property"); } // device fields must NOT be present (privacy decision). assert!(!obj.properties.contains_key("device"), "device must not be stored"); assert!(!obj.properties.contains_key("device_id"), "device_id must not be stored"); } // ---- Rebuild from PDS link records (apply_link_to_documents) ---- #[tokio::test] async fn apply_link_to_documents_associates_only_unassociated() { use sqlx::Row; let app = make_app().await; let user_id = insert_user_returning_id(&app.db, "rebuilder", "pw").await; let other_id = insert_user_returning_id(&app.db, "stranger", "pw").await; insert_document(&app.db, user_id, "md5a", 10.0, "unassociated").await; insert_document(&app.db, user_id, "md5b", 20.0, "unassociated").await; // Same hash under a different user must not be touched. insert_document(&app.db, other_id, "md5a", 30.0, "unassociated").await; // An already-associated document must not be overwritten. insert_document(&app.db, user_id, "md5c", 40.0, "associated").await; sqlx::query( "UPDATE documents SET popfeed_list_item_uri = 'at://keep/social.popfeed.feed.listItem/keep' WHERE user_id = ? AND kosync_document_id = 'md5c'", ) .bind(user_id) .execute(&app.db) .await .unwrap(); let uri = "at://did:plc:testdid/social.popfeed.feed.listItem/r1"; let n = crate::popfeed::apply_link_to_documents( &app.db, user_id, uri, &["md5a".to_owned(), "md5b".to_owned(), "md5c".to_owned()], Some("Dune"), ) .await .unwrap(); assert_eq!(n, 2, "only the two unassociated docs should be newly associated"); for id in ["md5a", "md5b"] { let row = sqlx::query( "SELECT popfeed_list_item_uri, association_status, title, sync_link_synced_at FROM documents WHERE user_id = ? AND kosync_document_id = ?", ) .bind(user_id) .bind(id) .fetch_one(&app.db) .await .unwrap(); let u: Option = row.try_get("popfeed_list_item_uri").unwrap(); let st: String = row.try_get("association_status").unwrap(); let title: Option = row.try_get("title").unwrap(); let ts: Option = row.try_get("sync_link_synced_at").unwrap(); assert_eq!(u.as_deref(), Some(uri), "{id} should be associated with the uri"); assert_eq!(st, "associated"); assert_eq!(title.as_deref(), Some("Dune"), "{id} title backfilled from link"); assert!(ts.is_some(), "{id} should have sync_link_synced_at stamped"); } // md5c keeps its original association. let kept: Option = sqlx::query_scalar( "SELECT popfeed_list_item_uri FROM documents WHERE user_id = ? AND kosync_document_id = 'md5c'", ) .bind(user_id) .fetch_one(&app.db) .await .unwrap(); assert_eq!( kept.as_deref(), Some("at://keep/social.popfeed.feed.listItem/keep"), "an existing association must not be overwritten" ); // The other user's identical hash is untouched. let other: Option = sqlx::query_scalar( "SELECT popfeed_list_item_uri FROM documents WHERE user_id = ? AND kosync_document_id = 'md5a'", ) .bind(other_id) .fetch_one(&app.db) .await .unwrap(); assert_eq!(other, None, "another user's document must not be associated"); } #[tokio::test] async fn reconstruct_apply_helpers_rebuild_documents_from_scratch() { use sqlx::Row; let app = make_app().await; let user_id = insert_user_returning_id(&app.db, "wiped", "pw").await; // Simulates a wiped DB: no documents rows exist. Reconstruct one book from // its link record (association + title) and its progress record (position). crate::popfeed::apply_reconstructed_link( &app.db, user_id, "md5a", "at://did:plc:x/site.standard.publication/r1", Some("Dune"), ) .await .unwrap(); crate::popfeed::apply_reconstructed_position( &app.db, user_id, "md5a", &crate::popfeed::RestoredPosition { progress: "/body/DocFragment[3]/p[7]".to_owned(), percentage: 0.4212, timestamp: 1_700_000_000, }, ) .await .unwrap(); let row = sqlx::query( "SELECT popfeed_list_item_uri, association_status, title, progress, percent, synced_at, public_id FROM documents WHERE user_id = ? AND kosync_document_id = 'md5a'", ) .bind(user_id) .fetch_one(&app.db) .await .unwrap(); let uri: Option = row.try_get("popfeed_list_item_uri").unwrap(); let status: String = row.try_get("association_status").unwrap(); let title: Option = row.try_get("title").unwrap(); let progress: Option = row.try_get("progress").unwrap(); let percent: Option = row.try_get("percent").unwrap(); let synced: Option = row.try_get("synced_at").unwrap(); let public_id: Option = row.try_get("public_id").unwrap(); assert_eq!(uri.as_deref(), Some("at://did:plc:x/site.standard.publication/r1")); assert_eq!(status, "associated", "reconstructed book is associated (Gap 1/2)"); assert_eq!(title.as_deref(), Some("Dune"), "title restored from link (Gap 3)"); assert_eq!(progress.as_deref(), Some("/body/DocFragment[3]/p[7]")); assert!((percent.unwrap() - 42.12).abs() < 1e-6, "ten-thousandths → percent"); assert_eq!(synced, Some(1_700_000_000)); assert!(public_id.is_some(), "a reconstructed row gets a public_id"); // A progress record with no link → an unassociated row (still surfaced). crate::popfeed::apply_reconstructed_position( &app.db, user_id, "md5b", &crate::popfeed::RestoredPosition { progress: "p1".to_owned(), percentage: 0.1, timestamp: 1, }, ) .await .unwrap(); let st: String = sqlx::query_scalar( "SELECT association_status FROM documents WHERE user_id = ? AND kosync_document_id = 'md5b'", ) .bind(user_id) .fetch_one(&app.db) .await .unwrap(); assert_eq!(st, "unassociated"); } // ---- OPDS root (Basic Auth) ---- #[tokio::test] async fn test_opds_root_requires_auth_with_challenge() { let app = make_app().await; let resp = send(app.router, get_req("/opds")).await; assert_eq!(resp.status(), StatusCode::UNAUTHORIZED); let challenge = resp .headers() .get("www-authenticate") .and_then(|v| v.to_str().ok()) .unwrap_or_default(); assert!( challenge.starts_with("Basic"), "401 must carry a Basic challenge so OPDS clients prompt; got {challenge:?}" ); } #[tokio::test] async fn test_opds_root_rejects_bad_password() { let app = make_app().await; insert_user_returning_id(&app.db, "olivia", "opdspass").await; let resp = send(app.router, opds_get("/opds", "olivia", "wrongpass")).await; assert_eq!(resp.status(), StatusCode::UNAUTHORIZED); } #[tokio::test] async fn test_opds_root_with_basic_auth_returns_atom_nav_feed() { let app = make_app().await; insert_user_returning_id(&app.db, "olivia", "opdspass").await; let resp = send(app.router, opds_get("/opds", "olivia", "opdspass")).await; assert_eq!(resp.status(), StatusCode::OK); let ctype = resp .headers() .get("content-type") .and_then(|v| v.to_str().ok()) .unwrap_or_default() .to_owned(); assert!( ctype.contains("application/atom+xml") && ctype.contains("kind=navigation"), "expected an OPDS navigation Atom content type; got {ctype:?}" ); let body = body_string(resp).await; assert!(body.contains("headway.page library")); assert!(body.contains("/opds/publications"), "root should link to subscriptions"); } #[tokio::test] async fn test_opds_publications_requires_auth() { let app = make_app().await; let resp = send(app.router, get_req("/opds/publications")).await; assert_eq!(resp.status(), StatusCode::UNAUTHORIZED); } #[tokio::test] async fn test_opds_basic_auth_with_did_username_succeeds() { // The KOSync username is the user's DID, which contains colons — the Basic // Auth "user:pass" split must not break on them. (Earlier tests used a // colon-free username and missed this.) let app = make_app().await; sqlx::query( "INSERT INTO users (atproto_did, atproto_handle, kosync_username, kosync_key_hash) VALUES (?, ?, ?, ?)", ) .bind("did:plc:colonuser") .bind("colon.user.test") .bind("did:plc:colonuser") .bind(bcrypt::hash(md5_hex("opdspass"), 4).unwrap()) .execute(&app.db) .await .unwrap(); let resp = send(app.router, opds_get("/opds", "did:plc:colonuser", "opdspass")).await; assert_eq!(resp.status(), StatusCode::OK, "DID username must authenticate over Basic Auth"); } #[tokio::test] async fn test_opds_documents_requires_auth() { // Unauthenticated request is rejected before any network resolution happens. let app = make_app().await; let resp = send(app.router, get_req("/opds/publications/anything")).await; assert_eq!(resp.status(), StatusCode::UNAUTHORIZED); } #[tokio::test] async fn test_opds_documents_malformed_segment_is_404() { // Authenticated, but the path segment isn't a valid base64url publication // at-uri — rejected with 404 before any resolution/fetch. let app = make_app().await; insert_user_returning_id(&app.db, "olivia", "opdspass").await; let resp = send(app.router, opds_get("/opds/publications/not!base64", "olivia", "opdspass")).await; assert_eq!(resp.status(), StatusCode::NOT_FOUND); } #[tokio::test] async fn test_opds_publications_without_session_is_empty_feed() { // Authenticated, but the test user has no ATProto OAuth session, so the // catalog is reachable and well-formed — just empty. let app = make_app().await; insert_user_returning_id(&app.db, "olivia", "opdspass").await; let resp = send(app.router, opds_get("/opds/publications", "olivia", "opdspass")).await; assert_eq!(resp.status(), StatusCode::OK); let body = body_string(resp).await; assert!(body.contains("Subscriptions")); assert_eq!(body.matches("").count(), 0, "no session → no publications"); }