mod config;
mod error;
mod payload;
mod transport;
use bytes::Bytes;
use futures::StreamExt as _;
use loonfs_api::{
v0::{
AbortUploadResponse, BeginDownloadRequest, BeginDownloadResponse, BeginUploadRequest,
BeginUploadResponse, ChangesResponse, CommitResponse as ApiCommitResponse, CommittedChange,
CompleteUploadRequest, CompleteUploadResponse, CompletedUploadPart,
DirectMultipartContentClaim, DirectMultipartUploadOptions, DirectPutContentClaim,
DisableGrepIndexResponse, EnableGrepIndexResponse, FilesystemChange, GrepGcRequest,
GrepGcResponse, GrepIndexStatusResponse, ObjectTransferAccess, SignUploadPartsRequest,
SignUploadPartsResponse, SignedUploadPart, StoreProbeRequest, StoreProbeResponse,
UploadContentResponse, UploadPartChecksumClaim, UploadStatusResponse,
ValidatedContentToken,
},
AbsolutePath, AuthoritativePathEntry, CapabilityDocument, ChangeSeq, CheckpointId,
ChecksumAlgorithm, CommitId, CommitRequest, ContentRef, Crc64Nvme, CreateCheckpointRequest,
CreateCheckpointResponse, CreateNamespaceRequest, DeleteNamespaceResponse, ErrorCode,
FilesystemOperation, ForkNamespaceRequest, GrepRequest, GrepResponse, InodeId,
ListCheckpointsResponse, ListFileRevisionsResponse, ListPathEntriesResponse, ListTrashResponse,
MaintenanceStepRequest, MaintenanceStepResponse, NamespaceId, NamespaceStatusResponse,
NamespaceSummary, ReleaseCheckpointResponse, RevisionNo, Sha256, StorageChecksum, UploadId,
FEATURE_DOWNLOADS_DIRECT_GET, FEATURE_UPLOADS_DIRECT_MULTIPART,
LIMIT_DOWNLOAD_MAX_CONTENT_BYTES,
};
use payload::PartReader;
use std::sync::{Arc, OnceLock};
pub const STREAMING_PUT_MIN_BYTES: u64 = 8 * 1024 * 1024;
const DIRECT_MULTIPART_PARTS_IN_FLIGHT: usize = 4;
const DIRECT_MULTIPART_PART_ATTEMPTS: usize = 3;
#[derive(Debug, Clone, Copy)]
enum UploadedContent<'a> {
Bytes(&'a [u8]),
Streamed(&'a ContentRef),
}
impl UploadedContent<'_> {
fn matches(&self, expected: &StorageChecksum) -> Option<bool> {
match self {
Self::Bytes(bytes) => expected.matches(bytes),
Self::Streamed(content_ref) => {
let observed = digest_of(content_ref, expected.algorithm)?;
Some(observed == expected.value)
}
}
}
}
fn digest_of(content_ref: &ContentRef, algorithm: ChecksumAlgorithm) -> Option<&str> {
if content_ref.storage_checksum.algorithm == algorithm {
return Some(&content_ref.storage_checksum.value);
}
match algorithm {
ChecksumAlgorithm::Sha256 => content_ref.whole_file_sha256.as_deref(),
_ => None,
}
}
fn uploaded_matches_committed(uploaded: &UploadedContent<'_>, content_ref: &ContentRef) -> bool {
let evidence = match &content_ref.whole_file_sha256 {
Some(digest) => StorageChecksum {
algorithm: ChecksumAlgorithm::Sha256,
value: digest.clone(),
},
None => content_ref.storage_checksum.clone(),
};
uploaded.matches(&evidence) == Some(true)
}
fn reported_commit_receipt(error: &ClientError) -> Option<(ChangeSeq, String)> {
match error {
ClientError::Api { details, .. } => {
let details = details.as_ref()?;
Some((
details.committed_seq?,
details.committed_fingerprint.clone()?,
))
}
_ => None,
}
}
fn sole_committed_content_ref(change: &CommittedChange) -> Option<&ContentRef> {
let mut content = change.events.iter().filter_map(|event| match event {
FilesystemChange::Created { content_ref, .. } => content_ref.as_ref(),
FilesystemChange::ContentChanged { content_ref, .. } => Some(content_ref),
_ => None,
});
let only = content.next()?;
content.next().is_none().then_some(only)
}
pub use config::ClientConfig;
pub use error::ClientError;
pub use payload::{PayloadSource, PayloadStream};
use transport::{WireRequest, IO_INACTIVITY_TIMEOUT};
pub use ClientError as Error;
pub use loonfs_api::options::{
CopyOptions, CreateDirectoryOptions, DeleteOptions, MoveOptions, PutFileOptions,
RestoreRevisionOptions, UndeleteOptions,
};
pub type Result<T> = std::result::Result<T, ClientError>;
#[derive(Debug, Clone)]
pub struct Client {
base_url: String,
auth_token: Option<String>,
http: reqwest::Client,
transient_retry: bool,
capabilities: Arc<OnceLock<CapabilityDocument>>,
}
pub struct DirectDownloadStream {
body: payload::PayloadStream,
expected: ContentRef,
path: AbsolutePath,
sha256: Option<Sha256>,
size_bytes: u64,
resumed_from: u64,
prefix_folded: u64,
finished: bool,
}
impl DirectDownloadStream {
pub fn fold_resumed_prefix(&mut self, bytes: &[u8]) {
if let Some(sha256) = self.sha256.as_mut() {
sha256.update(bytes);
}
self.prefix_folded = self.prefix_folded.saturating_add(bytes.len() as u64);
}
pub async fn next_chunk(&mut self) -> Result<Option<Bytes>> {
if self.prefix_folded != self.resumed_from {
return Err(ClientError::Http(format!(
"a download of `{}` resumed at offset {} was given {} bytes of what it \
skipped; verification covers the whole object, so all of them are needed \
first",
self.path, self.resumed_from, self.prefix_folded
)));
}
if self.finished {
return Ok(None);
}
match self.body.next().await {
Some(Ok(chunk)) => {
self.size_bytes = self.size_bytes.saturating_add(chunk.len() as u64);
if self.size_bytes > self.expected.size_bytes {
self.finished = true;
return Err(ClientError::Http(format!(
"direct download of `{}` sent more than the {} bytes the grant named",
self.path, self.expected.size_bytes
)));
}
if let Some(sha256) = self.sha256.as_mut() {
sha256.update(&chunk);
}
Ok(Some(chunk))
}
Some(Err(error)) => {
self.finished = true;
Err(ClientError::Io(format!(
"read of `{}` failed: {error}",
self.path
)))
}
None => {
self.finished = true;
if self.size_bytes != self.expected.size_bytes {
return Err(ClientError::Http(format!(
"direct download of `{}` ended after {} bytes, not the {} the grant named",
self.path, self.size_bytes, self.expected.size_bytes
)));
}
if let (Some(sha256), Some(expected_sha256)) = (
self.sha256.take(),
self.expected.whole_file_sha256.as_deref(),
) {
let observed = sha256.finish().value;
if observed != expected_sha256 {
return Err(ClientError::Http(format!(
"direct download of `{}` hashed to {observed}, not the \
{expected_sha256} the grant named",
self.path
)));
}
}
Ok(None)
}
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct MultipartUploadResume {
pub upload_id: UploadId,
pub part_size_bytes: u64,
pub parts: Vec<CompletedUploadPart>,
}
pub trait MultipartUploadJournal: Send + Sync {
fn began(&self, upload_id: &UploadId, part_size_bytes: u64);
fn part_completed(&self, part: &CompletedUploadPart);
}
#[derive(Clone, Copy, Default)]
struct UploadContinuity<'a> {
resume: Option<&'a MultipartUploadResume>,
journal: Option<&'a dyn MultipartUploadJournal>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
struct StagedContent {
content_ref: ContentRef,
validated_content_token: Option<ValidatedContentToken>,
}
enum MultipartPayload<'a> {
Held(&'a [u8]),
Streamed(PayloadStream),
}
impl<'a> MultipartPayload<'a> {
fn into_parts_of(self, part_bytes: usize) -> MultipartParts<'a> {
match self {
Self::Held(bytes) => MultipartParts::Held {
bytes,
offset: 0,
part_bytes: part_bytes.max(1),
},
Self::Streamed(stream) => MultipartParts::Streamed(PartReader::new(stream, part_bytes)),
}
}
}
enum MultipartParts<'a> {
Held {
bytes: &'a [u8],
offset: usize,
part_bytes: usize,
},
Streamed(PartReader),
}
impl MultipartParts<'_> {
async fn next_part(&mut self) -> Result<Option<Bytes>> {
match self {
Self::Held {
bytes,
offset,
part_bytes,
} => {
if *offset >= bytes.len() {
return Ok(None);
}
let end = bytes.len().min(*offset + *part_bytes);
let part = Bytes::copy_from_slice(&bytes[*offset..end]);
*offset = end;
Ok(Some(part))
}
Self::Streamed(reader) => reader
.next_part()
.await
.map_err(|error| ClientError::Http(format!("reading the payload failed: {error}"))),
}
}
}
struct PendingPart {
claim: UploadPartChecksumClaim,
bytes: Bytes,
}
struct UploadedObject {
size_bytes: u64,
crc64nvme: StorageChecksum,
parts: Vec<CompletedUploadPart>,
}
fn signed_access(signed: &[SignedUploadPart], part_number: u32) -> Result<ObjectTransferAccess> {
signed
.iter()
.find(|part| part.part_number == part_number)
.map(|part| part.access.clone())
.ok_or_else(|| {
ClientError::Http(format!(
"server authorized no upload for part {part_number}"
))
})
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct NamespacePath {
namespace: NamespaceId,
absolute_path: AbsolutePath,
}
impl Client {
pub fn new(config: ClientConfig) -> Result<Self> {
config.validate()?;
let mut builder = reqwest::Client::builder()
.read_timeout(IO_INACTIVITY_TIMEOUT)
.connect_timeout(IO_INACTIVITY_TIMEOUT);
if let Some(timeout_ms) = config.request_timeout_ms {
builder = builder.timeout(std::time::Duration::from_millis(timeout_ms));
}
for certificate in config.extra_root_certificates()? {
builder = builder.add_root_certificate(certificate);
}
Ok(Self {
base_url: config.server_url.trim().trim_end_matches('/').to_owned(),
auth_token: config.auth_token,
http: builder
.build()
.map_err(|err| ClientError::Http(err.to_string()))?,
transient_retry: !config.disable_transient_retry,
capabilities: Arc::new(OnceLock::new()),
})
}
pub async fn capabilities(&self) -> Result<CapabilityDocument> {
if let Some(document) = self.capabilities.get() {
return Ok(document.clone());
}
let url = format!("{}/v0/capabilities", self.base_url);
let mut document: CapabilityDocument =
self.request_json::<(), _>(self.get(&url), None).await?;
document.retain_well_formed();
let _ = self.capabilities.set(document);
Ok(self
.capabilities
.get()
.expect("capability cache was just filled")
.clone())
}
pub async fn create_namespace(&self, namespace_id: &NamespaceId) -> Result<NamespaceSummary> {
let url = format!("{}/v0/namespaces", self.base_url);
self.request_json_once::<_, NamespaceSummary>(
self.post(&url),
Some(&CreateNamespaceRequest {
namespace_id: namespace_id.clone(),
}),
)
.await
}
pub async fn namespace_status(
&self,
namespace_id: &NamespaceId,
) -> Result<NamespaceStatusResponse> {
let url = format!("{}/v0/namespaces/{namespace_id}", self.base_url);
self.request_json::<(), NamespaceStatusResponse>(self.get(&url), None)
.await
}
pub async fn delete_namespace(
&self,
namespace_id: &NamespaceId,
expected_head_seq: Option<ChangeSeq>,
) -> Result<DeleteNamespaceResponse> {
let mut url = format!("{}/v0/namespaces/{namespace_id}", self.base_url);
if let Some(expected) = expected_head_seq {
url.push_str(&format!("?expected_head_seq={}", expected.0));
}
self.request_json_once::<(), DeleteNamespaceResponse>(self.delete(&url), None)
.await
}
pub async fn fork_namespace(
&self,
source_namespace_id: &NamespaceId,
new_namespace_id: &NamespaceId,
) -> Result<NamespaceSummary> {
let url = format!(
"{}/v0/namespaces/{source_namespace_id}/forks",
self.base_url
);
self.request_json_once::<_, NamespaceSummary>(
self.post(&url),
Some(&ForkNamespaceRequest {
new_namespace_id: new_namespace_id.clone(),
}),
)
.await
}
pub async fn list_path_entries_all(
&self,
spec: &NamespacePath,
) -> Result<ListPathEntriesResponse> {
let mut entries = Vec::new();
let mut envelope = None;
let mut cursor = None;
loop {
let page = self
.list_path_entries_page(spec, None, cursor.as_deref())
.await?;
let envelope_ref = envelope.get_or_insert_with(|| ListPathEntriesResponse {
namespace_id: page.namespace_id.clone(),
absolute_path: page.absolute_path.clone(),
head_seq: page.head_seq,
entries: Vec::new(),
next_cursor: None,
});
envelope_ref.head_seq = envelope_ref.head_seq.max(page.head_seq);
entries.extend(page.entries);
cursor = page.next_cursor;
if cursor.is_none() {
envelope_ref.entries = entries;
return Ok(envelope.expect("first page initializes response envelope"));
}
}
}
pub async fn list_path_entries_page(
&self,
spec: &NamespacePath,
limit: Option<u32>,
cursor: Option<&str>,
) -> Result<ListPathEntriesResponse> {
let mut url = format!(
"{}/v0/namespaces/{}/filesystem/list?path={}",
self.base_url,
spec.namespace().as_str(),
urlencoding::encode(spec.absolute_path().as_str())
);
let has_query = true;
append_optional_pagination_query(&mut url, has_query, limit, cursor);
self.request_json::<(), ListPathEntriesResponse>(self.get(&url), None)
.await
}
pub async fn stat_path(&self, spec: &NamespacePath) -> Result<AuthoritativePathEntry> {
let url = format!(
"{}/v0/namespaces/{}/filesystem/stat?path={}",
self.base_url,
spec.namespace().as_str(),
urlencoding::encode(spec.absolute_path().as_str())
);
self.request_json::<(), AuthoritativePathEntry>(self.get(&url), None)
.await
}
pub async fn get_file_bytes(&self, spec: &NamespacePath) -> Result<Vec<u8>> {
let url = format!(
"{}/v0/namespaces/{}/filesystem/content?path={}",
self.base_url,
spec.namespace().as_str(),
urlencoding::encode(spec.absolute_path().as_str())
);
self.request_bytes(&url).await
}
pub async fn get_file_revision_bytes(
&self,
spec: &NamespacePath,
revision_no: RevisionNo,
) -> Result<Vec<u8>> {
let url = format!(
"{}/v0/namespaces/{}/filesystem/content?path={}&revision_no={}",
self.base_url,
spec.namespace().as_str(),
urlencoding::encode(spec.absolute_path().as_str()),
revision_no.0
);
self.request_bytes(&url).await
}
pub async fn offers_direct_download(&self, size_bytes: u64) -> bool {
let Ok(capabilities) = self.capabilities().await else {
return false;
};
capabilities.supports(FEATURE_DOWNLOADS_DIRECT_GET)
&& capabilities
.limits
.get(LIMIT_DOWNLOAD_MAX_CONTENT_BYTES)
.is_some_and(|proxy_cap| size_bytes > *proxy_cap)
}
pub async fn begin_download(
&self,
spec: &NamespacePath,
revision_no: Option<RevisionNo>,
) -> Result<BeginDownloadResponse> {
let url = format!(
"{}/v0/namespaces/{}/filesystem/downloads",
self.base_url,
spec.namespace().as_str()
);
let request = match revision_no {
Some(revision_no) => {
BeginDownloadRequest::for_revision(spec.absolute_path().clone(), revision_no)
}
None => BeginDownloadRequest::for_path(spec.absolute_path().clone()),
};
self.request_json::<_, BeginDownloadResponse>(self.post(&url), Some(&request))
.await
}
pub async fn open_direct_download(
&self,
download: &BeginDownloadResponse,
) -> Result<DirectDownloadStream> {
self.open_direct_download_at(download, 0).await
}
pub async fn open_direct_download_at(
&self,
download: &BeginDownloadResponse,
start_offset: u64,
) -> Result<DirectDownloadStream> {
let ObjectTransferAccess::PresignedUrl {
method,
url,
headers,
..
} = &download.access;
if method != "GET" {
return Err(ClientError::Http(format!(
"unsupported presigned download method `{method}`"
)));
}
if start_offset > download.content_ref.size_bytes {
return Err(ClientError::Http(format!(
"cannot resume a download of `{}` at offset {start_offset} of {} bytes",
download.absolute_path, download.content_ref.size_bytes
)));
}
let mut request = WireRequest::presigned(reqwest::Method::GET, url);
for (name, value) in headers {
request = request.header(name, value);
}
if start_offset > 0 {
request = request.header("range", format!("bytes={start_offset}-"));
}
let body = self.call_for_response_stream(&request).await?;
Ok(DirectDownloadStream {
body,
expected: download.content_ref.clone(),
path: download.absolute_path.clone(),
sha256: download
.content_ref
.whole_file_sha256
.as_ref()
.map(|_| Sha256::new()),
size_bytes: start_offset,
resumed_from: start_offset,
prefix_folded: 0,
finished: false,
})
}
pub async fn download_via_presigned_url<W>(
&self,
download: &BeginDownloadResponse,
sink: &mut W,
) -> Result<u64>
where
W: tokio::io::AsyncWrite + Unpin,
{
use tokio::io::AsyncWriteExt as _;
let path = &download.absolute_path;
let mut download = self.open_direct_download(download).await?;
let mut size_bytes = 0u64;
while let Some(chunk) = download.next_chunk().await? {
size_bytes += chunk.len() as u64;
sink.write_all(&chunk)
.await
.map_err(|err| ClientError::Io(format!("write of `{path}` failed: {err}")))?;
}
sink.flush()
.await
.map_err(|err| ClientError::Io(format!("write of `{path}` failed: {err}")))?;
Ok(size_bytes)
}
pub async fn list_file_revisions_page(
&self,
spec: &NamespacePath,
limit: Option<u32>,
cursor: Option<&str>,
) -> Result<ListFileRevisionsResponse> {
let mut url = format!(
"{}/v0/namespaces/{}/filesystem/revisions?path={}",
self.base_url,
spec.namespace().as_str(),
urlencoding::encode(spec.absolute_path().as_str())
);
let has_query = true;
append_optional_pagination_query(&mut url, has_query, limit, cursor);
self.request_json::<(), ListFileRevisionsResponse>(self.get(&url), None)
.await
}
pub async fn list_trash_page(
&self,
namespace_id: &NamespaceId,
limit: Option<u32>,
cursor: Option<&str>,
) -> Result<ListTrashResponse> {
let mut url = format!(
"{}/v0/namespaces/{}/filesystem/trash",
self.base_url,
namespace_id.as_str()
);
let has_query = false;
append_optional_pagination_query(&mut url, has_query, limit, cursor);
self.request_json::<(), ListTrashResponse>(self.get(&url), None)
.await
}
pub async fn health(&self) -> Result<()> {
let url = format!("{}/health", self.base_url);
self.call_with_transient_retry(&self.get(&url), None)
.await?;
Ok(())
}
pub async fn begin_upload(
&self,
namespace_id: &NamespaceId,
request: &BeginUploadRequest,
) -> Result<BeginUploadResponse> {
let url = format!("{}/v0/namespaces/{namespace_id}/uploads", self.base_url);
self.request_json_once::<_, BeginUploadResponse>(self.post(&url), Some(request))
.await
}
pub async fn begin_direct_put(
&self,
namespace_id: &NamespaceId,
claim: DirectPutContentClaim,
) -> Result<BeginUploadResponse> {
self.begin_upload(
namespace_id,
&BeginUploadRequest::DirectPut { content: claim },
)
.await
}
pub async fn begin_direct_multipart(
&self,
namespace_id: &NamespaceId,
options: DirectMultipartUploadOptions,
) -> Result<BeginUploadResponse> {
self.begin_upload(
namespace_id,
&BeginUploadRequest::DirectMultipart {
multipart: Some(options),
},
)
.await
}
pub async fn sign_upload_parts(
&self,
namespace_id: &NamespaceId,
upload_id: &UploadId,
parts: Vec<UploadPartChecksumClaim>,
) -> Result<SignUploadPartsResponse> {
let url = format!(
"{}/v0/namespaces/{namespace_id}/uploads/{upload_id}/parts",
self.base_url
);
self.request_json::<_, SignUploadPartsResponse>(
self.post(&url),
Some(&SignUploadPartsRequest { parts }),
)
.await
}
pub async fn upload_part_via_presigned_url(
&self,
part_number: u32,
access: &ObjectTransferAccess,
crc64nvme: String,
bytes: Bytes,
) -> Result<CompletedUploadPart> {
let ObjectTransferAccess::PresignedUrl {
method,
url,
headers,
..
} = access;
if method != "PUT" {
return Err(ClientError::Http(format!(
"unsupported presigned part method `{method}`"
)));
}
let mut request = WireRequest::presigned(reqwest::Method::PUT, url);
for (name, value) in headers {
request = request.header(name, value);
}
let response = self
.call_with_transient_retry_headers(&request, Some(&bytes))
.await?;
let etag = response
.get(http::header::ETAG)
.and_then(|value| value.to_str().ok())
.ok_or_else(|| {
ClientError::Http(format!("part {part_number} upload returned no etag"))
})?
.to_owned();
Ok(CompletedUploadPart {
part_number,
etag,
crc64nvme,
})
}
pub async fn upload_via_presigned_url(
&self,
access: &ObjectTransferAccess,
bytes: &[u8],
) -> Result<()> {
let (method, url, headers) = match access {
ObjectTransferAccess::PresignedUrl {
method,
url,
headers,
..
} => (method, url, headers),
};
if method != "PUT" {
return Err(ClientError::Http(format!(
"unsupported presigned upload method `{method}`"
)));
}
let mut request = WireRequest::presigned(reqwest::Method::PUT, url);
for (name, value) in headers {
request = request.header(name, value);
}
self.call_once(&request, Some(&Bytes::copy_from_slice(bytes)))
.await
.map(|_| ())
}
pub async fn upload_content(
&self,
namespace_id: &NamespaceId,
upload_id: &UploadId,
bytes: &[u8],
) -> Result<UploadContentResponse> {
let request = self.upload_content_request(namespace_id, upload_id);
let response = self
.call_with_transient_retry(&request, Some(&Bytes::copy_from_slice(bytes)))
.await?;
serde_json::from_slice(&response).map_err(|err| ClientError::Json(err.to_string()))
}
pub async fn upload_streamed_content(
&self,
namespace_id: &NamespaceId,
upload_id: &UploadId,
source: PayloadSource,
) -> Result<UploadContentResponse> {
let request = self.upload_content_request(namespace_id, upload_id);
let (stream, size_bytes) = source.into_stream();
let response = self
.call_streamed_once(&request, stream, size_bytes)
.await?;
serde_json::from_slice(&response).map_err(|err| ClientError::Json(err.to_string()))
}
fn upload_content_request(
&self,
namespace_id: &NamespaceId,
upload_id: &UploadId,
) -> WireRequest {
let url = format!(
"{}/v0/namespaces/{namespace_id}/uploads/{upload_id}/content",
self.base_url
);
self.put(&url)
.header("content-type", "application/octet-stream")
}
pub async fn abort_upload(
&self,
namespace_id: &NamespaceId,
upload_id: &UploadId,
) -> Result<AbortUploadResponse> {
let url = format!(
"{}/v0/namespaces/{namespace_id}/uploads/{upload_id}/abort",
self.base_url
);
self.request_json::<(), AbortUploadResponse>(self.post(&url), None)
.await
}
pub async fn read_upload_status(
&self,
namespace_id: &NamespaceId,
upload_id: &UploadId,
) -> Result<UploadStatusResponse> {
let url = format!(
"{}/v0/namespaces/{namespace_id}/uploads/{upload_id}",
self.base_url
);
self.request_json::<(), UploadStatusResponse>(self.get(&url), None)
.await
}
pub async fn complete_upload(
&self,
namespace_id: &NamespaceId,
upload_id: &UploadId,
request: &CompleteUploadRequest,
) -> Result<CompleteUploadResponse> {
let url = format!(
"{}/v0/namespaces/{namespace_id}/uploads/{upload_id}/complete",
self.base_url
);
self.request_json::<_, CompleteUploadResponse>(self.post(&url), Some(request))
.await
}
pub async fn list_changes(
&self,
namespace_id: &NamespaceId,
after_seq: ChangeSeq,
limit: Option<u32>,
) -> Result<ChangesResponse> {
let mut url = format!(
"{}/v0/namespaces/{namespace_id}/changes?after_seq={}",
self.base_url, after_seq.0
);
if let Some(limit) = limit {
url.push_str(&format!("&limit={limit}"));
}
self.request_json::<(), ChangesResponse>(self.get(&url), None)
.await
}
pub async fn create_checkpoint(
&self,
namespace_id: &NamespaceId,
request: &CreateCheckpointRequest,
) -> Result<CreateCheckpointResponse> {
let url = format!(
"{}/v0/admin/namespaces/{namespace_id}/checkpoints",
self.base_url
);
self.request_json(self.post(&url), Some(request)).await
}
pub async fn list_checkpoints(
&self,
namespace_id: &NamespaceId,
) -> Result<ListCheckpointsResponse> {
let url = format!(
"{}/v0/admin/namespaces/{namespace_id}/checkpoints",
self.base_url
);
self.request_json::<(), ListCheckpointsResponse>(self.get(&url), None)
.await
}
pub async fn release_checkpoint(
&self,
namespace_id: &NamespaceId,
checkpoint_id: &CheckpointId,
) -> Result<ReleaseCheckpointResponse> {
let url = format!(
"{}/v0/admin/namespaces/{namespace_id}/checkpoints/{checkpoint_id}/release",
self.base_url
);
self.request_json::<(), ReleaseCheckpointResponse>(self.post(&url), None)
.await
}
pub async fn maintenance_step(
&self,
namespace_id: &NamespaceId,
request: &MaintenanceStepRequest,
) -> Result<MaintenanceStepResponse> {
let url = format!(
"{}/v0/admin/namespaces/{namespace_id}/maintenance/step",
self.base_url
);
self.request_json(self.post(&url), Some(request)).await
}
pub async fn probe_store(&self, request: &StoreProbeRequest) -> Result<StoreProbeResponse> {
let url = format!("{}/v0/admin/store/probe", self.base_url);
self.request_json(self.post(&url), Some(request)).await
}
pub async fn grep(
&self,
namespace_id: &NamespaceId,
request: &GrepRequest,
) -> Result<GrepResponse> {
let url = format!("{}/v0/namespaces/{namespace_id}/query/grep", self.base_url);
self.request_json(self.post(&url), Some(request)).await
}
pub async fn grep_index_status(
&self,
namespace_id: &NamespaceId,
) -> Result<GrepIndexStatusResponse> {
let url = format!(
"{}/v0/admin/namespaces/{namespace_id}/grep/index",
self.base_url
);
self.request_json::<(), GrepIndexStatusResponse>(self.get(&url), None)
.await
}
pub async fn enable_grep_index(
&self,
namespace_id: &NamespaceId,
) -> Result<EnableGrepIndexResponse> {
let url = format!(
"{}/v0/admin/namespaces/{namespace_id}/grep/index/enable",
self.base_url
);
self.request_json::<(), EnableGrepIndexResponse>(self.post(&url), None)
.await
}
pub async fn disable_grep_index(
&self,
namespace_id: &NamespaceId,
) -> Result<DisableGrepIndexResponse> {
let url = format!(
"{}/v0/admin/namespaces/{namespace_id}/grep/index/disable",
self.base_url
);
self.request_json::<(), DisableGrepIndexResponse>(self.post(&url), None)
.await
}
pub async fn gc_grep_index(
&self,
namespace_id: &NamespaceId,
request: &GrepGcRequest,
) -> Result<GrepGcResponse> {
let url = format!(
"{}/v0/admin/namespaces/{namespace_id}/grep/index/gc",
self.base_url
);
self.request_json(self.post(&url), Some(request)).await
}
pub async fn commit(
&self,
namespace_id: &NamespaceId,
request: &CommitRequest,
) -> Result<ApiCommitResponse> {
let url = format!("{}/v0/namespaces/{namespace_id}/commits", self.base_url);
self.request_json::<_, ApiCommitResponse>(self.post(&url), Some(request))
.await
}
async fn stage_bytes_as_content_ref(
&self,
namespace_id: &NamespaceId,
bytes: &[u8],
) -> Result<StagedContent> {
if bytes.len() as u64 >= STREAMING_PUT_MIN_BYTES && self.offers_direct_multipart().await {
return self
.stage_via_multipart(
namespace_id,
MultipartPayload::Held(bytes),
UploadContinuity::default(),
)
.await;
}
self.stage_bytes_via_server(namespace_id, bytes).await
}
async fn stage_source_as_content_ref(
&self,
namespace_id: &NamespaceId,
source: PayloadSource,
continuity: UploadContinuity<'_>,
) -> Result<StagedContent> {
let known_small = source
.size_bytes()
.is_some_and(|size_bytes| size_bytes < STREAMING_PUT_MIN_BYTES);
if !known_small && self.offers_direct_multipart().await {
let (stream, _) = source.into_stream();
return self
.stage_via_multipart(namespace_id, MultipartPayload::Streamed(stream), continuity)
.await;
}
self.stage_source_via_server(namespace_id, source).await
}
async fn offers_direct_multipart(&self) -> bool {
self.capabilities()
.await
.is_ok_and(|capabilities| capabilities.supports(FEATURE_UPLOADS_DIRECT_MULTIPART))
}
async fn stage_via_multipart(
&self,
namespace_id: &NamespaceId,
payload: MultipartPayload<'_>,
continuity: UploadContinuity<'_>,
) -> Result<StagedContent> {
let (upload_id, part_size_bytes) = match continuity.resume {
Some(resume) => (resume.upload_id.clone(), resume.part_size_bytes),
None => {
let begin = self
.begin_direct_multipart(namespace_id, DirectMultipartUploadOptions::default())
.await?;
let Some(multipart) = begin.direct_multipart else {
return Err(ClientError::Http(
"server accepted direct_multipart without part geometry".to_owned(),
));
};
if let Some(journal) = continuity.journal {
journal.began(&begin.upload_id, multipart.part_size_bytes);
}
(begin.upload_id, multipart.part_size_bytes)
}
};
let uploaded = self
.upload_every_part(
namespace_id,
&upload_id,
payload,
part_size_bytes,
continuity,
)
.await;
let uploaded = match uploaded {
Ok(uploaded) => uploaded,
Err(error) => {
let _ = self.abort_upload(namespace_id, &upload_id).await;
return Err(error);
}
};
if uploaded.parts.is_empty() {
let _ = self.abort_upload(namespace_id, &upload_id).await;
return self.stage_bytes_via_server(namespace_id, &[]).await;
}
let response = self
.complete_upload(
namespace_id,
&upload_id,
&CompleteUploadRequest::for_multipart(
DirectMultipartContentClaim {
size_bytes: uploaded.size_bytes,
crc64nvme: uploaded.crc64nvme.value,
},
uploaded.parts,
),
)
.await?;
Ok(Self::staged_from_completion(response))
}
async fn upload_every_part(
&self,
namespace_id: &NamespaceId,
upload_id: &UploadId,
payload: MultipartPayload<'_>,
part_size_bytes: u64,
continuity: UploadContinuity<'_>,
) -> Result<UploadedObject> {
let part_size = usize::try_from(part_size_bytes)
.map_err(|_| ClientError::Http("part size does not fit this platform".to_owned()))?;
let landed = continuity
.resume
.map_or::<&[CompletedUploadPart], _>(&[], |resume| &resume.parts);
let mut source = payload.into_parts_of(part_size);
let mut whole_object = Crc64Nvme::new();
let mut size_bytes = 0u64;
let mut parts = Vec::new();
let mut next_part_number = 1u32;
loop {
let mut wave = Vec::with_capacity(DIRECT_MULTIPART_PARTS_IN_FLIGHT);
let mut source_ended = false;
while wave.len() < DIRECT_MULTIPART_PARTS_IN_FLIGHT {
let Some(bytes) = source.next_part().await? else {
source_ended = true;
break;
};
whole_object.update(&bytes);
size_bytes += bytes.len() as u64;
let part_number = next_part_number;
next_part_number += 1;
if let Some(landed) = landed.iter().find(|part| part.part_number == part_number) {
parts.push(landed.clone());
continue;
}
wave.push(PendingPart {
claim: UploadPartChecksumClaim {
part_number,
crc64nvme: StorageChecksum::crc64nvme(&bytes).value,
},
bytes,
});
}
if !wave.is_empty() {
let uploaded = self.upload_wave(namespace_id, upload_id, wave).await?;
if let Some(journal) = continuity.journal {
for part in &uploaded {
journal.part_completed(part);
}
}
parts.extend(uploaded);
}
if source_ended {
break;
}
}
parts.sort_by_key(|part| part.part_number);
Ok(UploadedObject {
size_bytes,
crc64nvme: whole_object.finish(),
parts,
})
}
async fn upload_wave(
&self,
namespace_id: &NamespaceId,
upload_id: &UploadId,
wave: Vec<PendingPart>,
) -> Result<Vec<CompletedUploadPart>> {
let claims = wave.iter().map(|part| part.claim.clone()).collect();
let signed = self
.sign_upload_parts(namespace_id, upload_id, claims)
.await?;
let mut in_flight = tokio::task::JoinSet::new();
for part in wave {
let access = signed_access(&signed.parts, part.claim.part_number)?;
let client = self.clone();
let namespace_id = namespace_id.clone();
let upload_id = upload_id.clone();
in_flight.spawn(async move {
client
.upload_one_part(&namespace_id, &upload_id, part, access)
.await
});
}
let mut uploaded = Vec::new();
let mut failure = None;
while let Some(joined) = in_flight.join_next().await {
match joined.map_err(|err| ClientError::Http(format!("part upload task failed: {err}")))
{
Ok(Ok(part)) => uploaded.push(part),
Ok(Err(error)) | Err(error) => failure = failure.or(Some(error)),
}
}
match failure {
Some(error) => Err(error),
None => Ok(uploaded),
}
}
async fn upload_one_part(
&self,
namespace_id: &NamespaceId,
upload_id: &UploadId,
part: PendingPart,
mut access: ObjectTransferAccess,
) -> Result<CompletedUploadPart> {
let part_number = part.claim.part_number;
for attempt in 1..=DIRECT_MULTIPART_PART_ATTEMPTS {
let result = self
.upload_part_via_presigned_url(
part_number,
&access,
part.claim.crc64nvme.clone(),
part.bytes.clone(),
)
.await;
match result {
Ok(uploaded) => return Ok(uploaded),
Err(error) if attempt == DIRECT_MULTIPART_PART_ATTEMPTS => return Err(error),
Err(_) => {
let signed = self
.sign_upload_parts(namespace_id, upload_id, vec![part.claim.clone()])
.await?;
access = signed_access(&signed.parts, part_number)?;
}
}
}
Err(ClientError::Http(format!(
"part {part_number} upload made no attempt"
)))
}
async fn stage_bytes_via_server(
&self,
namespace_id: &NamespaceId,
bytes: &[u8],
) -> Result<StagedContent> {
let upload = self
.begin_upload(namespace_id, &BeginUploadRequest::ServiceProxied {})
.await?;
let staged = self
.upload_content(namespace_id, &upload.upload_id, bytes)
.await?;
self.complete_staged(namespace_id, &upload.upload_id, staged)
.await
}
async fn stage_source_via_server(
&self,
namespace_id: &NamespaceId,
source: PayloadSource,
) -> Result<StagedContent> {
let upload = self
.begin_upload(namespace_id, &BeginUploadRequest::ServiceProxied {})
.await?;
let staged = self
.upload_streamed_content(namespace_id, &upload.upload_id, source)
.await;
let staged = match staged {
Ok(staged) => staged,
Err(error) => {
let _ = self.abort_upload(namespace_id, &upload.upload_id).await;
return Err(error);
}
};
self.complete_staged(namespace_id, &upload.upload_id, staged)
.await
}
async fn complete_staged(
&self,
namespace_id: &NamespaceId,
upload_id: &UploadId,
staged: UploadContentResponse,
) -> Result<StagedContent> {
let response = self
.complete_upload(
namespace_id,
upload_id,
&CompleteUploadRequest::for_content_ref(staged.content_ref),
)
.await?;
Ok(Self::staged_from_completion(response))
}
fn staged_from_completion(response: CompleteUploadResponse) -> StagedContent {
let validated_content_token =
response
.validated_content_token
.map(|token| ValidatedContentToken {
content_ref: response.content_ref.clone(),
token,
});
StagedContent {
content_ref: response.content_ref,
validated_content_token,
}
}
pub async fn put_file_bytes(
&self,
spec: &NamespacePath,
bytes: &[u8],
options: &PutFileOptions,
) -> Result<ApiCommitResponse> {
let staged = self
.stage_bytes_as_content_ref(spec.namespace(), bytes)
.await?;
self.commit_staged_file(spec, staged, options, UploadedContent::Bytes(bytes))
.await
}
pub async fn put_file_stream(
&self,
spec: &NamespacePath,
source: PayloadSource,
options: &PutFileOptions,
) -> Result<ApiCommitResponse> {
self.put_file_stream_continuing(spec, source, options, UploadContinuity::default())
.await
}
pub async fn put_file_stream_resumable(
&self,
spec: &NamespacePath,
source: PayloadSource,
options: &PutFileOptions,
journal: &dyn MultipartUploadJournal,
resume: Option<&MultipartUploadResume>,
) -> Result<ApiCommitResponse> {
self.put_file_stream_continuing(
spec,
source,
options,
UploadContinuity {
resume,
journal: Some(journal),
},
)
.await
}
async fn put_file_stream_continuing(
&self,
spec: &NamespacePath,
source: PayloadSource,
options: &PutFileOptions,
continuity: UploadContinuity<'_>,
) -> Result<ApiCommitResponse> {
let staged = self
.stage_source_as_content_ref(spec.namespace(), source, continuity)
.await?;
let uploaded = staged.content_ref.clone();
self.commit_staged_file(spec, staged, options, UploadedContent::Streamed(&uploaded))
.await
}
pub async fn commit_completed_upload(
&self,
spec: &NamespacePath,
content_ref: ContentRef,
validated_content_token: Option<String>,
options: &PutFileOptions,
) -> Result<ApiCommitResponse> {
let uploaded = content_ref.clone();
let staged = StagedContent {
validated_content_token: validated_content_token.map(|token| ValidatedContentToken {
content_ref: content_ref.clone(),
token,
}),
content_ref,
};
self.commit_staged_file(spec, staged, options, UploadedContent::Streamed(&uploaded))
.await
}
async fn commit_staged_file(
&self,
spec: &NamespacePath,
staged: StagedContent,
options: &PutFileOptions,
uploaded: UploadedContent<'_>,
) -> Result<ApiCommitResponse> {
let commit_id = options.commit_id.clone().unwrap_or_else(CommitId::generate);
let response = self
.commit(
spec.namespace(),
&CommitRequest {
commit_id: commit_id.clone(),
message: options.message.clone(),
content_tokens: staged.validated_content_token.into_iter().collect(),
operations: vec![FilesystemOperation::PutFile {
path: spec.absolute_path().clone(),
content_ref: staged.content_ref,
behavior: options.behavior,
expected_revision_no: options.expected_revision_no,
}],
},
)
.await;
match response {
Ok(response) => Ok(response),
Err(error) if error.code() == Some(ErrorCode::CommitIdReuseConflict) => {
self.reconcile_commit_id_reuse(spec, &commit_id, options, uploaded, error)
.await
}
Err(error) => Err(error),
}
}
async fn reconcile_commit_id_reuse(
&self,
spec: &NamespacePath,
commit_id: &CommitId,
options: &PutFileOptions,
uploaded: UploadedContent<'_>,
conflict: ClientError,
) -> Result<ApiCommitResponse> {
let namespace_id = spec.namespace();
let Some((committed_seq, committed_fingerprint)) = reported_commit_receipt(&conflict)
else {
return Err(conflict);
};
let Some(committed) = self
.read_committed_change(namespace_id, commit_id, committed_seq)
.await?
else {
return Err(conflict);
};
let Some(content_ref) = sole_committed_content_ref(&committed) else {
return Err(conflict);
};
let retried = loonfs_api::put_retry_fingerprint(
namespace_id,
spec.absolute_path(),
options.behavior,
options.expected_revision_no,
options.message.as_deref(),
content_ref,
);
if retried.ok().as_deref() != Some(committed_fingerprint.as_str()) {
return Err(conflict);
}
if !uploaded_matches_committed(&uploaded, content_ref) {
return Err(conflict);
}
Ok(ApiCommitResponse {
namespace_id: namespace_id.clone(),
commit_id: committed.commit_id,
committed_seq: committed.seq,
})
}
async fn read_committed_change(
&self,
namespace_id: &NamespaceId,
commit_id: &CommitId,
committed_seq: ChangeSeq,
) -> Result<Option<CommittedChange>> {
let after_seq = ChangeSeq(committed_seq.0.saturating_sub(1));
let page = match self.list_changes(namespace_id, after_seq, Some(1)).await {
Ok(page) => page,
Err(error) if error.code() == Some(ErrorCode::RebootstrapRequired) => {
return Ok(None);
}
Err(error) => return Err(error),
};
Ok(page
.changes
.into_iter()
.find(|change| change.seq == committed_seq && &change.commit_id == commit_id))
}
pub async fn create_directory(
&self,
spec: &NamespacePath,
options: &CreateDirectoryOptions,
) -> Result<ApiCommitResponse> {
let commit_id = options.commit_id.clone().unwrap_or_else(CommitId::generate);
let response = self
.commit(
spec.namespace(),
&CommitRequest::single(
commit_id,
options.message.clone(),
FilesystemOperation::CreateDirectory {
path: spec.absolute_path().clone(),
parents: options.parents,
},
),
)
.await?;
Ok(response)
}
pub async fn delete_path(
&self,
spec: &NamespacePath,
options: &DeleteOptions,
) -> Result<ApiCommitResponse> {
let commit_id = options.commit_id.clone().unwrap_or_else(CommitId::generate);
let response = self
.commit(
spec.namespace(),
&CommitRequest::single(
commit_id,
options.message.clone(),
FilesystemOperation::DeletePath {
path: spec.absolute_path().clone(),
behavior: options.behavior,
expected_inode_id: options.expected_inode_id,
},
),
)
.await?;
Ok(response)
}
pub async fn move_path(
&self,
from: &NamespacePath,
to: &NamespacePath,
options: &MoveOptions,
) -> Result<ApiCommitResponse> {
if from.namespace() != to.namespace() {
return Err(ClientError::InvalidNamespacePath(format!(
"cannot move across namespaces: {} -> {}",
from.namespace(),
to.namespace()
)));
}
let commit_id = options.commit_id.clone().unwrap_or_else(CommitId::generate);
let response = self
.commit(
from.namespace(),
&CommitRequest::single(
commit_id,
options.message.clone(),
FilesystemOperation::MovePath {
from_path: from.absolute_path().clone(),
to_path: to.absolute_path().clone(),
behavior: options.behavior,
},
),
)
.await?;
Ok(response)
}
pub async fn copy_path(
&self,
from: &NamespacePath,
to: &NamespacePath,
options: &CopyOptions,
) -> Result<ApiCommitResponse> {
if from.namespace() != to.namespace() {
return Err(ClientError::InvalidNamespacePath(format!(
"cannot copy across namespaces: {} -> {}",
from.namespace(),
to.namespace()
)));
}
let commit_id = options.commit_id.clone().unwrap_or_else(CommitId::generate);
let response = self
.commit(
from.namespace(),
&CommitRequest::single(
commit_id,
options.message.clone(),
FilesystemOperation::CopyPath {
from_path: from.absolute_path().clone(),
to_path: to.absolute_path().clone(),
behavior: options.behavior,
},
),
)
.await?;
Ok(response)
}
pub async fn undelete(
&self,
namespace: &NamespaceId,
inode_id: InodeId,
deleted_at_seq: ChangeSeq,
path: Option<&AbsolutePath>,
options: &UndeleteOptions,
) -> Result<ApiCommitResponse> {
let commit_id = options.commit_id.clone().unwrap_or_else(CommitId::generate);
let response = self
.commit(
namespace,
&CommitRequest::single(
commit_id,
options.message.clone(),
FilesystemOperation::Undelete {
inode_id,
deleted_at_seq,
path: path.cloned(),
},
),
)
.await?;
Ok(response)
}
pub async fn restore_file_revision(
&self,
spec: &NamespacePath,
source_revision_no: RevisionNo,
options: &RestoreRevisionOptions,
) -> Result<ApiCommitResponse> {
let commit_id = options.commit_id.clone().unwrap_or_else(CommitId::generate);
let response = self
.commit(
spec.namespace(),
&CommitRequest::single(
commit_id,
options.message.clone(),
FilesystemOperation::RestoreRevision {
path: spec.absolute_path().clone(),
source_revision_no,
},
),
)
.await?;
Ok(response)
}
}
impl NamespacePath {
pub fn parse(namespace: &str, absolute_path: &str) -> Result<Self> {
let namespace = NamespaceId::parse(namespace)
.map_err(|error| ClientError::InvalidNamespacePath(error.to_string()))?;
let absolute_path = AbsolutePath::parse(absolute_path)
.map_err(|error| ClientError::InvalidNamespacePath(error.to_string()))?;
Ok(Self {
namespace,
absolute_path,
})
}
pub fn new(namespace: NamespaceId, absolute_path: AbsolutePath) -> Self {
Self {
namespace,
absolute_path,
}
}
pub fn namespace(&self) -> &NamespaceId {
&self.namespace
}
pub fn absolute_path(&self) -> &AbsolutePath {
&self.absolute_path
}
}
fn append_optional_pagination_query(
url: &mut String,
has_query: bool,
limit: Option<u32>,
cursor: Option<&str>,
) {
let mut has_query = has_query;
if let Some(limit) = limit {
append_query_param(url, &mut has_query, "limit", &limit.to_string());
}
if let Some(cursor) = cursor {
append_query_param(url, &mut has_query, "cursor", cursor);
}
}
fn append_query_param(url: &mut String, has_query: &mut bool, name: &str, value: &str) {
url.push(if *has_query { '&' } else { '?' });
*has_query = true;
url.push_str(name);
url.push('=');
url.push_str(&urlencoding::encode(value));
}
#[cfg(test)]
mod download_tests;
#[cfg(test)]
mod streaming_tests;
#[cfg(test)]
mod tests {
use super::*;
use crate::transport::{transient_failure, MAX_TRANSIENT_ATTEMPTS};
use loonfs_api::{ContentId, ErrorCode, ErrorKind};
use std::fs;
use tempfile::tempdir;
fn test_content_ref(bytes: &[u8]) -> ContentRef {
ContentRef::blob_v1(ContentId::generate(), bytes)
}
fn direct_put_claim(bytes: &[u8]) -> DirectPutContentClaim {
let content_ref = test_content_ref(bytes);
DirectPutContentClaim {
size_bytes: content_ref.size_bytes,
sha256: content_ref.storage_checksum.value,
}
}
fn crc32c_content_ref(bytes: &[u8]) -> ContentRef {
ContentRef {
kind: loonfs_api::ContentRefKind::BlobV1,
content_id: ContentId::generate(),
size_bytes: bytes.len() as u64,
storage_checksum: StorageChecksum {
algorithm: ChecksumAlgorithm::Crc32c,
value: "0f5c0a1e".to_owned(),
},
whole_file_sha256: None,
}
}
#[test]
fn a_digest_this_client_cannot_compare_leaves_the_reuse_conflict_standing() {
let bytes = b"retried payload";
let committed = crc32c_content_ref(bytes);
let uploaded = UploadedContent::Bytes(bytes);
assert_eq!(
uploaded.matches(&committed.storage_checksum),
None,
"the fixture must reach the refusal, not a comparison"
);
assert!(!uploaded_matches_committed(&uploaded, &committed));
}
#[test]
fn a_whole_file_digest_over_the_same_bytes_proves_the_retry_did_this_work() {
let bytes = b"retried payload";
let committed = test_content_ref(bytes);
assert!(uploaded_matches_committed(
&UploadedContent::Bytes(bytes),
&committed
));
assert!(!uploaded_matches_committed(
&UploadedContent::Bytes(b"some other payload"),
&committed
));
}
#[test]
fn construction_validates_config_like_load_does() {
let error = super::Client::new(super::ClientConfig {
server_url: "ftp://example.com".to_owned(),
auth_token: None,
request_timeout_ms: None,
disable_transient_retry: false,
ca_cert_path: None,
})
.expect_err("ftp scheme must be rejected");
assert!(
matches!(
&error,
super::ClientError::ConfigValidation {
field: "server_url",
..
}
),
"unexpected error: {error:?}"
);
}
#[test]
fn an_unusable_ca_bundle_fails_construction_and_names_the_path() {
let dir = tempdir().expect("tempdir");
let missing = dir.path().join("absent.crt");
let garbage = dir.path().join("garbage.crt");
fs::write(&garbage, b"this is not a certificate\n").expect("write garbage");
for path in [missing, garbage] {
let display = path.display().to_string();
let error = super::Client::new(super::ClientConfig {
server_url: "https://example.com".to_owned(),
auth_token: None,
request_timeout_ms: None,
disable_transient_retry: false,
ca_cert_path: Some(display.clone()),
})
.expect_err("unusable ca bundle");
match &error {
super::ClientError::ConfigValidation {
field: "ca_cert_path",
reason,
} => assert!(
reason.contains(&display),
"the reason must name the path, got: {reason}"
),
other => unreachable!("unexpected error: {other:?}"),
}
}
}
#[test]
fn client_config_rejects_unknown_keys() {
let error = toml::from_str::<ClientConfig>(
"server_url = \"http://localhost:1\"\nauth_tokn = \"oops\"\n",
)
.expect_err("unknown key must fail decode");
assert!(error.to_string().contains("auth_tokn"), "{error}");
let config: ClientConfig =
toml::from_str("server_url = \"http://localhost:1\"\n").expect("minimal config");
assert!(config.auth_token.is_none());
}
#[test]
fn transient_failure_covers_transport_and_retryable_unavailability_only() {
let api = |code: &str| ClientError::Api {
status: 503,
code: code.to_owned(),
feature: None,
message: String::new(),
request_id: None,
details: None,
};
assert!(transient_failure(
true,
&ClientError::Http("reset".to_owned())
));
assert!(transient_failure(false, &api("server_busy")));
assert!(transient_failure(false, &api("commit_queue_full")));
assert!(transient_failure(false, &api("shutting_down")));
assert!(!transient_failure(false, &api("server_error")));
assert!(!transient_failure(false, &api("maintenance_required")));
assert!(!transient_failure(
false,
&ClientError::Http("http status 502 with a non-envelope body".to_owned())
));
}
#[tokio::test]
async fn transport_failures_resend_up_to_the_attempt_cap() {
let namespace_id = NamespaceId::parse("demo").expect("valid namespace id");
let transport = crate::transport::test_transport::failures(MAX_TRANSIENT_ATTEMPTS as usize);
let retrying = Client::new(ClientConfig {
server_url: "http://example.invalid".to_owned(),
auth_token: None,
request_timeout_ms: None,
disable_transient_retry: false,
ca_cert_path: None,
})
.expect("valid client config");
let error = retrying
.namespace_status(&namespace_id)
.await
.expect_err("dropped connections must fail");
assert!(matches!(error, ClientError::Http(_)), "{error:?}");
assert_eq!(transport.attempts(), MAX_TRANSIENT_ATTEMPTS as usize);
drop(transport);
let transport = crate::transport::test_transport::failures(1);
let single_shot = Client::new(ClientConfig {
server_url: "http://example.invalid".to_owned(),
auth_token: None,
request_timeout_ms: None,
disable_transient_retry: true,
ca_cert_path: None,
})
.expect("valid client config");
single_shot
.namespace_status(&namespace_id)
.await
.expect_err("dropped connection must fail without retry");
assert_eq!(transport.attempts(), 1);
}
fn retry_policy_client() -> Client {
Client::new(ClientConfig {
server_url: "http://example.invalid".to_owned(),
auth_token: None,
request_timeout_ms: None,
disable_transient_retry: false,
ca_cert_path: None,
})
.expect("valid client config")
}
fn single_attempt_probe() -> (crate::transport::test_transport::Guard, Client) {
(
crate::transport::test_transport::failure_then_success(b"{}".to_vec()),
retry_policy_client(),
)
}
fn assert_single_attempt<T>(
result: Result<T>,
transport: &crate::transport::test_transport::Guard,
) {
assert!(
matches!(result, Err(ClientError::Http(_))),
"expected the first transport failure to surface"
);
assert_eq!(transport.attempts(), 1);
}
#[tokio::test]
async fn retry_policy_lifecycle_mutations_are_single_attempt() {
let namespace_id = NamespaceId::parse("demo").expect("valid namespace id");
let fork_id = NamespaceId::parse("fork").expect("valid id");
let (transport, client) = single_attempt_probe();
assert_single_attempt(client.create_namespace(&namespace_id).await, &transport);
drop(transport);
let (transport, client) = single_attempt_probe();
assert_single_attempt(
client.fork_namespace(&namespace_id, &fork_id).await,
&transport,
);
drop(transport);
let (transport, client) = single_attempt_probe();
assert_single_attempt(
client
.delete_namespace(&namespace_id, Some(ChangeSeq(7)))
.await,
&transport,
);
}
#[tokio::test]
async fn retry_policy_commit_id_filesystem_mutation_retries() {
let namespace_id = NamespaceId::parse("demo").expect("valid namespace id");
let commit_id =
CommitId::parse("c_00000000000000000000000000000001").expect("valid commit id");
let response = ApiCommitResponse {
namespace_id: namespace_id.clone(),
commit_id: commit_id.clone(),
committed_seq: ChangeSeq(1),
};
let transport = crate::transport::test_transport::failure_then_success(
serde_json::to_vec(&response).expect("serialize response"),
);
let client = retry_policy_client();
let spec = NamespacePath::parse("demo", "/docs").expect("valid namespace path");
let actual = client
.create_directory(
&spec,
&CreateDirectoryOptions {
commit_id: Some(commit_id),
message: None,
..CreateDirectoryOptions::default()
},
)
.await
.expect("commit-id mutation should retry");
assert_eq!(actual, response);
assert_eq!(transport.attempts(), 2);
}
#[tokio::test]
async fn retry_policy_read_retries() {
let namespace_id = NamespaceId::parse("demo").expect("valid namespace id");
let response = NamespaceStatusResponse {
namespace_id: namespace_id.clone(),
head_seq: ChangeSeq(0),
current_manifest_id: None,
wal_tail_segments: 0,
retention_floor_seq: ChangeSeq(0),
};
let transport = crate::transport::test_transport::failure_then_success(
serde_json::to_vec(&response).expect("serialize response"),
);
let client = retry_policy_client();
let actual = client
.namespace_status(&namespace_id)
.await
.expect("read should retry");
assert_eq!(actual, response);
assert_eq!(transport.attempts(), 2);
}
#[tokio::test]
async fn retry_policy_upload_begins_are_single_attempt() {
let namespace_id = NamespaceId::parse("demo").expect("valid namespace id");
let (transport, client) = single_attempt_probe();
assert_single_attempt(
client
.begin_upload(&namespace_id, &BeginUploadRequest::ServiceProxied {})
.await,
&transport,
);
drop(transport);
let (transport, client) = single_attempt_probe();
assert_single_attempt(
client
.begin_direct_put(&namespace_id, direct_put_claim(b"direct"))
.await,
&transport,
);
}
#[tokio::test]
async fn retry_policy_presigned_upload_is_single_attempt() {
let transport = crate::transport::test_transport::failure_then_success(Vec::new());
let client = retry_policy_client();
let access = ObjectTransferAccess::PresignedUrl {
method: "PUT".to_owned(),
url: "http://example.invalid/upload".to_owned(),
headers: std::collections::BTreeMap::new(),
expires_at_ms: 1,
};
let result = client.upload_via_presigned_url(&access, b"direct").await;
assert!(matches!(result, Err(ClientError::Http(_))), "{result:?}");
assert_eq!(transport.attempts(), 1);
}
#[tokio::test]
async fn retry_policy_proxied_upload_content_retries() {
let namespace_id = NamespaceId::parse("demo").expect("valid namespace id");
let upload_id = loonfs_api::UploadId::parse("upl_00000000000000000000000000000001")
.expect("valid upload id");
let response = UploadContentResponse {
namespace_id: namespace_id.clone(),
upload_id: upload_id.clone(),
content_ref: test_content_ref(b"content"),
};
let transport = crate::transport::test_transport::failure_then_success(
serde_json::to_vec(&response).expect("serialize response"),
);
let client = retry_policy_client();
let actual = client
.upload_content(&namespace_id, &upload_id, b"content")
.await
.expect("identical content staging should retry");
assert_eq!(actual, response);
assert_eq!(transport.attempts(), 2);
}
#[tokio::test]
async fn retry_policy_upload_completion_retries() {
let namespace_id = NamespaceId::parse("demo").expect("valid namespace id");
let upload_id = loonfs_api::UploadId::parse("upl_00000000000000000000000000000001")
.expect("valid upload id");
let content_ref = test_content_ref(b"content");
let response = CompleteUploadResponse {
namespace_id: namespace_id.clone(),
upload_id: upload_id.clone(),
content_ref: content_ref.clone(),
validated_content_token: None,
};
let transport = crate::transport::test_transport::failure_then_success(
serde_json::to_vec(&response).expect("serialize response"),
);
let client = retry_policy_client();
let actual = client
.complete_upload(
&namespace_id,
&upload_id,
&CompleteUploadRequest::for_content_ref(content_ref),
)
.await
.expect("completed-session replay should retry");
assert_eq!(actual, response);
assert_eq!(transport.attempts(), 2);
}
#[test]
fn status_errors_keep_the_status_when_the_body_is_not_the_envelope() {
let error = crate::transport::map_status_error(502, b"<html>upstream error</html>");
let ClientError::Http(message) = error else {
unreachable!("expected Http error, got {error:?}");
};
assert!(message.contains("502"), "{message}");
assert!(message.contains("non-envelope body"), "{message}");
}
fn api_error(status: u16, code: &str) -> ClientError {
ClientError::Api {
status,
code: code.to_owned(),
feature: None,
message: "test".to_owned(),
request_id: None,
details: None,
}
}
#[test]
fn api_errors_with_known_codes_classify_through_the_registry() {
let error = api_error(409, "stale_revision");
assert_eq!(error.code(), Some(ErrorCode::StaleRevision));
assert_eq!(error.kind(), Some(ErrorKind::Conflict));
let error = api_error(409, "content_not_prepared");
assert_eq!(error.code(), Some(ErrorCode::ContentNotPrepared));
assert_eq!(error.kind(), Some(ErrorKind::Conflict));
let error = api_error(410, "namespace_deleted");
assert_eq!(error.code(), Some(ErrorCode::NamespaceDeleted));
assert_eq!(error.kind(), Some(ErrorKind::Gone));
let error = api_error(503, "commit_outcome_unknown");
assert_eq!(error.code(), Some(ErrorCode::CommitOutcomeUnknown));
assert_eq!(error.kind(), Some(ErrorKind::OutcomeUnknown));
let error = api_error(500, "index_corrupt");
assert_eq!(error.code(), Some(ErrorCode::IndexCorrupt));
assert_eq!(error.kind(), Some(ErrorKind::DataCorruption));
}
#[test]
fn api_errors_with_unknown_codes_fall_back_to_the_status_class() {
for (status, kind) in [
(400, ErrorKind::InvalidRequest),
(404, ErrorKind::InvalidRequest),
(500, ErrorKind::Internal),
(503, ErrorKind::Unavailable),
] {
let error = api_error(status, "code_from_a_newer_server");
assert_eq!(error.code(), None);
assert_eq!(error.kind(), Some(kind), "status {status}");
}
}
#[test]
fn non_api_errors_have_no_code_or_kind() {
let error = ClientError::Http("connection refused".to_owned());
assert_eq!(error.code(), None);
assert_eq!(error.kind(), None);
}
#[test]
fn load_rejects_invalid_server_url() {
let path = write_config(
r#"
server_url = "ftp://example.com"
auth_token = "dev-token"
"#,
);
let error = ClientConfig::load(&path).expect_err("invalid server url");
assert!(
matches!(error, ClientError::ConfigValidation { field, .. } if field == "server_url"),
"expected config validation error, got {error:?}"
);
}
#[test]
fn load_rejects_blank_auth_token() {
let path = write_config(
r#"
server_url = "http://127.0.0.1:9400"
auth_token = " "
"#,
);
let error = ClientConfig::load(&path).expect_err("blank auth token");
assert!(
matches!(error, ClientError::ConfigValidation { field, .. } if field == "auth_token"),
"expected config validation error, got {error:?}"
);
}
#[test]
fn load_preserves_missing_file_as_config_io() {
let temp_dir = tempdir().expect("tempdir");
let path = temp_dir.path().join("missing.toml");
let error = ClientConfig::load(&path).expect_err("missing config");
assert!(matches!(error, ClientError::ConfigIo(_)));
}
#[test]
fn load_preserves_decode_error() {
let path = write_config("server_url = [");
let error = ClientConfig::load(&path).expect_err("decode error");
assert!(matches!(error, ClientError::ConfigDecode(_)));
}
#[test]
fn namespace_path_parse_rejects_invalid_namespace_id() {
for namespace in ["bad/name", "Demo", "..", "demo?"] {
assert!(
matches!(
NamespacePath::parse(namespace, "/notes.txt"),
Err(ClientError::InvalidNamespacePath(_))
),
"expected invalid namespace path for id {namespace:?}"
);
}
}
#[test]
fn namespace_path_parse_rejects_invalid_paths() {
for path in ["notes.txt", "", "/docs/../a.txt", "/docs/./a.txt"] {
assert!(
matches!(
NamespacePath::parse("demo", path),
Err(ClientError::InvalidNamespacePath(_))
),
"expected invalid namespace path for path {path:?}"
);
}
}
fn write_config(contents: &str) -> std::path::PathBuf {
let temp_dir = tempdir().expect("tempdir");
let path = temp_dir.path().join("client.toml");
fs::write(&path, contents).expect("write config");
let _ = temp_dir.keep();
path
}
}