mod common;
use cooklang_sync_client::chunker::{Chunker, InMemoryCache};
use cooklang_sync_client::connection::get_connection;
use cooklang_sync_client::errors::SyncError;
use cooklang_sync_client::models::{CreateForm, DeleteForm, FileRecord};
use cooklang_sync_client::registry;
use cooklang_sync_client::remote::Remote;
use cooklang_sync_client::syncer::{check_download_once, check_upload_once};
use std::sync::Arc;
use time::OffsetDateTime;
use tokio::sync::Mutex;
use wiremock::matchers::{body_string_contains, method, path};
use wiremock::{Mock, MockServer, ResponseTemplate};
const NS: i32 = 1;
const TOKEN: &str = "test-token";
fn sample_create(path: &str, size: i64) -> CreateForm {
CreateForm {
jid: None,
path: path.to_string(),
deleted: false,
size,
modified_at: OffsetDateTime::from_unix_timestamp(1_700_000_000).unwrap(),
namespace_id: NS,
}
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn check_upload_once_commits_success_and_marks_jid() {
let server = MockServer::start().await;
Mock::given(method("POST"))
.and(path("/metadata/commit"))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({ "Success": 100 })))
.expect(1)
.mount(&server)
.await;
let base = common::client_base();
tokio::fs::write(base.dir.path().join("a.cook"), b"Eggs\n").await.expect("write file");
{
let conn = &mut get_connection(&base.pool).expect("checkout");
registry::create(conn, &[sample_create("a.cook", 5)]).expect("create");
}
let remote = Remote::new(&server.uri(), TOKEN);
let chunker_arc = Arc::new(Mutex::new(base.chunker));
let all_committed = check_upload_once(&base.pool, Arc::clone(&chunker_arc), &remote, NS)
.await
.expect("check_upload_once");
assert!(all_committed, "all rows should commit in one pass");
let conn = &mut get_connection(&base.pool).expect("checkout");
let after: Vec<FileRecord> = registry::updated_locally(conn, NS).expect("updated_locally");
assert!(after.is_empty(), "no rows should remain unsynced after Success");
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn check_upload_once_triggers_upload_batch_when_server_asks_for_chunks() {
let server = MockServer::start().await;
let mut base = common::client_base();
tokio::fs::write(base.dir.path().join("a.cook"), b"Eggs\n").await.expect("write file");
let computed_ids = base.chunker.hashify("a.cook").await.expect("hashify");
assert!(!computed_ids.is_empty(), "text chunker must produce ids");
let chunk_id = computed_ids.first().expect("first id").clone();
Mock::given(method("POST"))
.and(path("/metadata/commit"))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({
"NeedChunks": chunk_id.clone()
})))
.expect(1)
.mount(&server)
.await;
Mock::given(method("POST"))
.and(path("/chunks/upload"))
.respond_with(ResponseTemplate::new(200))
.expect(1)
.mount(&server)
.await;
{
let conn = &mut get_connection(&base.pool).expect("checkout");
registry::create(conn, &[sample_create("a.cook", 5)]).expect("create");
}
let remote = Remote::new(&server.uri(), TOKEN);
let chunker_arc = Arc::new(Mutex::new(base.chunker));
let all_committed = check_upload_once(&base.pool, Arc::clone(&chunker_arc), &remote, NS)
.await
.expect("check_upload_once");
assert!(!all_committed, "NeedChunks path should return false");
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn check_upload_once_commits_tombstone_without_hashifying() {
let server = MockServer::start().await;
Mock::given(method("POST"))
.and(path("/metadata/commit"))
.and(body_string_contains("deleted=true"))
.and(body_string_contains("chunk_ids=&")) .respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({ "Success": 7 })))
.expect(1)
.mount(&server)
.await;
let base = common::client_base();
{
let conn = &mut get_connection(&base.pool).expect("checkout");
registry::create(conn, &[sample_create("gone.cook", 10)]).expect("create");
let live: Vec<FileRecord> = registry::non_deleted(conn, NS).expect("non_deleted");
let latest = live.first().expect("live row").clone();
let tombstone = DeleteForm {
path: latest.path.clone(),
jid: None,
size: 0,
modified_at: latest.modified_at,
deleted: true,
namespace_id: NS,
};
registry::delete(conn, &[tombstone]).expect("delete");
}
let remote = Remote::new(&server.uri(), TOKEN);
let chunker_arc = Arc::new(Mutex::new(base.chunker));
let ok = check_upload_once(&base.pool, Arc::clone(&chunker_arc), &remote, NS)
.await
.expect("check_upload_once");
assert!(ok, "tombstone commit is a Success => all_commited stays true");
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn check_download_once_writes_file_and_inserts_registry_row() {
let server = MockServer::start().await;
let scratch_dir = tempfile::TempDir::new().unwrap();
tokio::fs::write(scratch_dir.path().join("a.cook"), b"Eggs\n").await.unwrap();
let scratch_cache = InMemoryCache::new(10, 1_000_000);
let mut scratch = Chunker::new(scratch_cache, scratch_dir.path().to_path_buf());
let ids = scratch.hashify("a.cook").await.unwrap();
assert_eq!(ids.len(), 1, "text file produces one line-per-chunk id");
let chunk_id = ids[0].clone();
Mock::given(method("GET"))
.and(path("/metadata/list"))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!([
{ "id": 11, "path": "a.cook", "deleted": false, "chunk_ids": chunk_id }
])))
.expect(1)
.mount(&server)
.await;
let boundary = "downloadbound";
let body = format!(
"--{b}\r\nX-Chunk-ID: {id}\r\nContent-Type: application/octet-stream\r\n\r\nEggs\n\r\n--{b}--\r\n",
b = boundary,
id = chunk_id
);
Mock::given(method("POST"))
.and(path("/chunks/download"))
.respond_with(
ResponseTemplate::new(200)
.insert_header("content-type", format!("multipart/form-data; boundary={}", boundary).as_str())
.set_body_bytes(body.into_bytes()),
)
.expect(1)
.mount(&server)
.await;
let base = common::client_base();
let remote = Remote::new(&server.uri(), TOKEN);
let chunker_arc = Arc::new(Mutex::new(base.chunker));
let downloaded = check_download_once(
&base.pool,
Arc::clone(&chunker_arc),
&remote,
base.dir.path(),
NS,
)
.await
.expect("check_download_once");
assert!(downloaded, "non-empty remote list => returns true");
let bytes = tokio::fs::read(base.dir.path().join("a.cook")).await.expect("read");
assert_eq!(bytes, b"Eggs\n");
let conn = &mut get_connection(&base.pool).expect("checkout");
let rows = registry::non_deleted(conn, NS).expect("non_deleted");
assert_eq!(rows.len(), 1);
assert_eq!(rows[0].path, "a.cook");
assert_eq!(rows[0].jid, Some(11));
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn check_download_once_removes_local_file_and_appends_tombstone() {
let server = MockServer::start().await;
Mock::given(method("GET"))
.and(path("/metadata/list"))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!([
{ "id": 22, "path": "gone.cook", "deleted": true, "chunk_ids": "" }
])))
.expect(1)
.mount(&server)
.await;
let base = common::client_base();
tokio::fs::write(base.dir.path().join("gone.cook"), b"bye\n").await.unwrap();
{
let conn = &mut get_connection(&base.pool).expect("checkout");
registry::create(conn, &[sample_create("gone.cook", 4)]).expect("create");
}
let remote = Remote::new(&server.uri(), TOKEN);
let chunker_arc = Arc::new(Mutex::new(base.chunker));
check_download_once(
&base.pool,
Arc::clone(&chunker_arc),
&remote,
base.dir.path(),
NS,
)
.await
.expect("check_download_once");
assert!(
!base.dir.path().join("gone.cook").exists(),
"local file should be deleted when remote says deleted"
);
let conn = &mut get_connection(&base.pool).expect("checkout");
let live = registry::non_deleted(conn, NS).expect("non_deleted");
assert!(
live.iter().all(|r| r.path != "gone.cook"),
"gone.cook should not be in non_deleted after tombstone applied"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn check_download_once_empty_remote_list_is_noop() {
let server = MockServer::start().await;
Mock::given(method("GET"))
.and(path("/metadata/list"))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!([])))
.expect(1)
.mount(&server)
.await;
let base = common::client_base();
let remote = Remote::new(&server.uri(), TOKEN);
let chunker_arc = Arc::new(Mutex::new(base.chunker));
let downloaded = check_download_once(
&base.pool,
Arc::clone(&chunker_arc),
&remote,
base.dir.path(),
NS,
)
.await
.expect("check_download_once");
assert!(!downloaded, "empty list => returns false");
let conn = &mut get_connection(&base.pool).expect("checkout");
let rows = registry::non_deleted(conn, NS).expect("non_deleted");
assert!(rows.is_empty());
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn check_download_once_propagates_unauthorized_from_list() {
let server = MockServer::start().await;
Mock::given(method("GET"))
.and(path("/metadata/list"))
.respond_with(ResponseTemplate::new(401))
.mount(&server)
.await;
let base = common::client_base();
let remote = Remote::new(&server.uri(), TOKEN);
let chunker_arc = Arc::new(Mutex::new(base.chunker));
let err = check_download_once(
&base.pool,
Arc::clone(&chunker_arc),
&remote,
base.dir.path(),
NS,
)
.await
.unwrap_err();
assert!(
matches!(err, SyncError::Unauthorized),
"expected SyncError::Unauthorized, got {:?}",
err
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn check_upload_once_propagates_unauthorized_from_commit() {
let server = MockServer::start().await;
Mock::given(method("POST"))
.and(path("/metadata/commit"))
.respond_with(ResponseTemplate::new(401))
.mount(&server)
.await;
let base = common::client_base();
tokio::fs::write(base.dir.path().join("a.cook"), b"Eggs\n").await.expect("write file");
{
let conn = &mut get_connection(&base.pool).expect("checkout");
registry::create(conn, &[sample_create("a.cook", 5)]).expect("create");
}
let remote = Remote::new(&server.uri(), TOKEN);
let chunker_arc = Arc::new(Mutex::new(base.chunker));
let err = check_upload_once(&base.pool, Arc::clone(&chunker_arc), &remote, NS)
.await
.unwrap_err();
assert!(
matches!(err, SyncError::Unauthorized),
"expected SyncError::Unauthorized, got {:?}",
err
);
}