1mod config;
16mod error;
17mod payload;
18mod transport;
19
20use bytes::Bytes;
21use futures::StreamExt as _;
22use loonfs_api::{
23 v0::{
24 AbortUploadResponse, BeginDownloadRequest, BeginDownloadResponse, BeginUploadRequest,
25 BeginUploadResponse, ChangesResponse, CommitResponse as ApiCommitResponse, CommittedChange,
26 CompleteUploadRequest, CompleteUploadResponse, CompletedUploadPart,
27 DirectMultipartContentClaim, DirectMultipartUploadOptions, DirectPutContentClaim,
28 DisableGrepIndexResponse, EnableGrepIndexResponse, FilesystemChange, GrepGcRequest,
29 GrepGcResponse, GrepIndexStatusResponse, ObjectTransferAccess, SignUploadPartsRequest,
30 SignUploadPartsResponse, SignedUploadPart, StoreProbeRequest, StoreProbeResponse,
31 UploadContentResponse, UploadPartChecksumClaim, UploadStatusResponse,
32 ValidatedContentToken,
33 },
34 AbsolutePath, AuthoritativePathEntry, CapabilityDocument, ChangeSeq, CheckpointId,
35 ChecksumAlgorithm, CommitId, CommitRequest, ContentRef, Crc64Nvme, CreateCheckpointRequest,
36 CreateCheckpointResponse, CreateNamespaceRequest, DeleteNamespaceResponse, ErrorCode,
37 FilesystemOperation, ForkNamespaceRequest, GrepRequest, GrepResponse, InodeId,
38 ListCheckpointsResponse, ListFileRevisionsResponse, ListPathEntriesResponse, ListTrashResponse,
39 MaintenanceStepRequest, MaintenanceStepResponse, NamespaceId, NamespaceStatusResponse,
40 NamespaceSummary, ReleaseCheckpointResponse, RevisionNo, Sha256, StorageChecksum, UploadId,
41 FEATURE_DOWNLOADS_DIRECT_GET, FEATURE_UPLOADS_DIRECT_MULTIPART,
42 LIMIT_DOWNLOAD_MAX_CONTENT_BYTES,
43};
44use payload::PartReader;
45use std::sync::{Arc, OnceLock};
46
47pub const STREAMING_PUT_MIN_BYTES: u64 = 8 * 1024 * 1024;
56
57const DIRECT_MULTIPART_PARTS_IN_FLIGHT: usize = 4;
60
61const DIRECT_MULTIPART_PART_ATTEMPTS: usize = 3;
65
66#[derive(Debug, Clone, Copy)]
75enum UploadedContent<'a> {
76 Bytes(&'a [u8]),
78 Streamed(&'a ContentRef),
81}
82
83impl UploadedContent<'_> {
84 fn matches(&self, expected: &StorageChecksum) -> Option<bool> {
90 match self {
91 Self::Bytes(bytes) => expected.matches(bytes),
92 Self::Streamed(content_ref) => {
93 let observed = digest_of(content_ref, expected.algorithm)?;
94 Some(observed == expected.value)
95 }
96 }
97 }
98}
99
100fn digest_of(content_ref: &ContentRef, algorithm: ChecksumAlgorithm) -> Option<&str> {
102 if content_ref.storage_checksum.algorithm == algorithm {
103 return Some(&content_ref.storage_checksum.value);
104 }
105 match algorithm {
106 ChecksumAlgorithm::Sha256 => content_ref.whole_file_sha256.as_deref(),
107 _ => None,
108 }
109}
110
111fn uploaded_matches_committed(uploaded: &UploadedContent<'_>, content_ref: &ContentRef) -> bool {
121 let evidence = match &content_ref.whole_file_sha256 {
122 Some(digest) => StorageChecksum {
123 algorithm: ChecksumAlgorithm::Sha256,
124 value: digest.clone(),
125 },
126 None => content_ref.storage_checksum.clone(),
127 };
128 uploaded.matches(&evidence) == Some(true)
129}
130
131fn reported_commit_receipt(error: &ClientError) -> Option<(ChangeSeq, String)> {
139 match error {
140 ClientError::Api { details, .. } => {
141 let details = details.as_ref()?;
142 Some((
143 details.committed_seq?,
144 details.committed_fingerprint.clone()?,
145 ))
146 }
147 _ => None,
148 }
149}
150
151fn sole_committed_content_ref(change: &CommittedChange) -> Option<&ContentRef> {
161 let mut content = change.events.iter().filter_map(|event| match event {
162 FilesystemChange::Created { content_ref, .. } => content_ref.as_ref(),
163 FilesystemChange::ContentChanged { content_ref, .. } => Some(content_ref),
164 _ => None,
165 });
166 let only = content.next()?;
167 content.next().is_none().then_some(only)
168}
169
170pub use config::ClientConfig;
171pub use error::ClientError;
172pub use payload::{PayloadSource, PayloadStream};
173use transport::{WireRequest, IO_INACTIVITY_TIMEOUT};
174pub use ClientError as Error;
175
176pub use loonfs_api::options::{
179 CopyOptions, CreateDirectoryOptions, DeleteOptions, MoveOptions, PutFileOptions,
180 RestoreRevisionOptions, UndeleteOptions,
181};
182
183pub type Result<T> = std::result::Result<T, ClientError>;
185
186#[derive(Debug, Clone)]
191pub struct Client {
192 base_url: String,
193 auth_token: Option<String>,
194 http: reqwest::Client,
195 transient_retry: bool,
198 capabilities: Arc<OnceLock<CapabilityDocument>>,
200}
201
202pub struct DirectDownloadStream {
209 body: payload::PayloadStream,
210 expected: ContentRef,
211 path: AbsolutePath,
212 sha256: Option<Sha256>,
213 size_bytes: u64,
214 resumed_from: u64,
217 prefix_folded: u64,
221 finished: bool,
222}
223
224impl DirectDownloadStream {
225 pub fn fold_resumed_prefix(&mut self, bytes: &[u8]) {
234 if let Some(sha256) = self.sha256.as_mut() {
235 sha256.update(bytes);
236 }
237 self.prefix_folded = self.prefix_folded.saturating_add(bytes.len() as u64);
238 }
239
240 pub async fn next_chunk(&mut self) -> Result<Option<Bytes>> {
243 if self.prefix_folded != self.resumed_from {
244 return Err(ClientError::Http(format!(
245 "a download of `{}` resumed at offset {} was given {} bytes of what it \
246 skipped; verification covers the whole object, so all of them are needed \
247 first",
248 self.path, self.resumed_from, self.prefix_folded
249 )));
250 }
251 if self.finished {
252 return Ok(None);
253 }
254 match self.body.next().await {
255 Some(Ok(chunk)) => {
256 self.size_bytes = self.size_bytes.saturating_add(chunk.len() as u64);
257 if self.size_bytes > self.expected.size_bytes {
258 self.finished = true;
259 return Err(ClientError::Http(format!(
260 "direct download of `{}` sent more than the {} bytes the grant named",
261 self.path, self.expected.size_bytes
262 )));
263 }
264 if let Some(sha256) = self.sha256.as_mut() {
265 sha256.update(&chunk);
266 }
267 Ok(Some(chunk))
268 }
269 Some(Err(error)) => {
270 self.finished = true;
271 Err(ClientError::Io(format!(
272 "read of `{}` failed: {error}",
273 self.path
274 )))
275 }
276 None => {
277 self.finished = true;
278 if self.size_bytes != self.expected.size_bytes {
279 return Err(ClientError::Http(format!(
280 "direct download of `{}` ended after {} bytes, not the {} the grant named",
281 self.path, self.size_bytes, self.expected.size_bytes
282 )));
283 }
284 if let (Some(sha256), Some(expected_sha256)) = (
285 self.sha256.take(),
286 self.expected.whole_file_sha256.as_deref(),
287 ) {
288 let observed = sha256.finish().value;
289 if observed != expected_sha256 {
290 return Err(ClientError::Http(format!(
291 "direct download of `{}` hashed to {observed}, not the \
292 {expected_sha256} the grant named",
293 self.path
294 )));
295 }
296 }
297 Ok(None)
298 }
299 }
300 }
301}
302
303#[derive(Debug, Clone, PartialEq, Eq)]
313pub struct MultipartUploadResume {
314 pub upload_id: UploadId,
315 pub part_size_bytes: u64,
319 pub parts: Vec<CompletedUploadPart>,
320}
321
322pub trait MultipartUploadJournal: Send + Sync {
330 fn began(&self, upload_id: &UploadId, part_size_bytes: u64);
332 fn part_completed(&self, part: &CompletedUploadPart);
334}
335
336#[derive(Clone, Copy, Default)]
339struct UploadContinuity<'a> {
340 resume: Option<&'a MultipartUploadResume>,
341 journal: Option<&'a dyn MultipartUploadJournal>,
342}
343
344#[derive(Debug, Clone, PartialEq, Eq)]
345struct StagedContent {
346 content_ref: ContentRef,
347 validated_content_token: Option<ValidatedContentToken>,
348}
349
350enum MultipartPayload<'a> {
357 Held(&'a [u8]),
359 Streamed(PayloadStream),
361}
362
363impl<'a> MultipartPayload<'a> {
364 fn into_parts_of(self, part_bytes: usize) -> MultipartParts<'a> {
366 match self {
367 Self::Held(bytes) => MultipartParts::Held {
368 bytes,
369 offset: 0,
370 part_bytes: part_bytes.max(1),
371 },
372 Self::Streamed(stream) => MultipartParts::Streamed(PartReader::new(stream, part_bytes)),
373 }
374 }
375}
376
377enum MultipartParts<'a> {
379 Held {
380 bytes: &'a [u8],
381 offset: usize,
382 part_bytes: usize,
383 },
384 Streamed(PartReader),
385}
386
387impl MultipartParts<'_> {
388 async fn next_part(&mut self) -> Result<Option<Bytes>> {
391 match self {
392 Self::Held {
393 bytes,
394 offset,
395 part_bytes,
396 } => {
397 if *offset >= bytes.len() {
398 return Ok(None);
399 }
400 let end = bytes.len().min(*offset + *part_bytes);
401 let part = Bytes::copy_from_slice(&bytes[*offset..end]);
402 *offset = end;
403 Ok(Some(part))
404 }
405 Self::Streamed(reader) => reader
406 .next_part()
407 .await
408 .map_err(|error| ClientError::Http(format!("reading the payload failed: {error}"))),
409 }
410 }
411}
412
413struct PendingPart {
416 claim: UploadPartChecksumClaim,
417 bytes: Bytes,
418}
419
420struct UploadedObject {
423 size_bytes: u64,
424 crc64nvme: StorageChecksum,
425 parts: Vec<CompletedUploadPart>,
426}
427
428fn signed_access(signed: &[SignedUploadPart], part_number: u32) -> Result<ObjectTransferAccess> {
430 signed
431 .iter()
432 .find(|part| part.part_number == part_number)
433 .map(|part| part.access.clone())
434 .ok_or_else(|| {
435 ClientError::Http(format!(
436 "server authorized no upload for part {part_number}"
437 ))
438 })
439}
440
441#[derive(Debug, Clone, PartialEq, Eq)]
447pub struct NamespacePath {
448 namespace: NamespaceId,
449 absolute_path: AbsolutePath,
450}
451
452impl Client {
453 pub fn new(config: ClientConfig) -> Result<Self> {
457 config.validate()?;
458 let mut builder = reqwest::Client::builder()
459 .read_timeout(IO_INACTIVITY_TIMEOUT)
462 .connect_timeout(IO_INACTIVITY_TIMEOUT);
463 if let Some(timeout_ms) = config.request_timeout_ms {
464 builder = builder.timeout(std::time::Duration::from_millis(timeout_ms));
465 }
466 for certificate in config.extra_root_certificates()? {
469 builder = builder.add_root_certificate(certificate);
470 }
471 Ok(Self {
472 base_url: config.server_url.trim().trim_end_matches('/').to_owned(),
473 auth_token: config.auth_token,
474 http: builder
475 .build()
476 .map_err(|err| ClientError::Http(err.to_string()))?,
477 transient_retry: !config.disable_transient_retry,
478 capabilities: Arc::new(OnceLock::new()),
479 })
480 }
481
482 pub async fn capabilities(&self) -> Result<CapabilityDocument> {
490 if let Some(document) = self.capabilities.get() {
491 return Ok(document.clone());
492 }
493 let url = format!("{}/v0/capabilities", self.base_url);
494 let mut document: CapabilityDocument =
495 self.request_json::<(), _>(self.get(&url), None).await?;
496 document.retain_well_formed();
497 let _ = self.capabilities.set(document);
500 Ok(self
501 .capabilities
502 .get()
503 .expect("capability cache was just filled")
504 .clone())
505 }
506
507 pub async fn create_namespace(&self, namespace_id: &NamespaceId) -> Result<NamespaceSummary> {
508 let url = format!("{}/v0/namespaces", self.base_url);
509 self.request_json_once::<_, NamespaceSummary>(
511 self.post(&url),
512 Some(&CreateNamespaceRequest {
513 namespace_id: namespace_id.clone(),
514 }),
515 )
516 .await
517 }
518
519 pub async fn namespace_status(
520 &self,
521 namespace_id: &NamespaceId,
522 ) -> Result<NamespaceStatusResponse> {
523 let url = format!("{}/v0/namespaces/{namespace_id}", self.base_url);
526 self.request_json::<(), NamespaceStatusResponse>(self.get(&url), None)
527 .await
528 }
529
530 pub async fn delete_namespace(
536 &self,
537 namespace_id: &NamespaceId,
538 expected_head_seq: Option<ChangeSeq>,
539 ) -> Result<DeleteNamespaceResponse> {
540 let mut url = format!("{}/v0/namespaces/{namespace_id}", self.base_url);
541 if let Some(expected) = expected_head_seq {
542 url.push_str(&format!("?expected_head_seq={}", expected.0));
543 }
544 self.request_json_once::<(), DeleteNamespaceResponse>(self.delete(&url), None)
546 .await
547 }
548
549 pub async fn fork_namespace(
550 &self,
551 source_namespace_id: &NamespaceId,
552 new_namespace_id: &NamespaceId,
553 ) -> Result<NamespaceSummary> {
554 let url = format!(
555 "{}/v0/namespaces/{source_namespace_id}/forks",
556 self.base_url
557 );
558 self.request_json_once::<_, NamespaceSummary>(
560 self.post(&url),
561 Some(&ForkNamespaceRequest {
562 new_namespace_id: new_namespace_id.clone(),
563 }),
564 )
565 .await
566 }
567
568 pub async fn list_path_entries_all(
576 &self,
577 spec: &NamespacePath,
578 ) -> Result<ListPathEntriesResponse> {
579 let mut entries = Vec::new();
580 let mut envelope = None;
581 let mut cursor = None;
582 loop {
583 let page = self
584 .list_path_entries_page(spec, None, cursor.as_deref())
585 .await?;
586 let envelope_ref = envelope.get_or_insert_with(|| ListPathEntriesResponse {
587 namespace_id: page.namespace_id.clone(),
588 absolute_path: page.absolute_path.clone(),
589 head_seq: page.head_seq,
590 entries: Vec::new(),
591 next_cursor: None,
592 });
593 envelope_ref.head_seq = envelope_ref.head_seq.max(page.head_seq);
594 entries.extend(page.entries);
595 cursor = page.next_cursor;
596 if cursor.is_none() {
597 envelope_ref.entries = entries;
600 return Ok(envelope.expect("first page initializes response envelope"));
601 }
602 }
603 }
604
605 pub async fn list_path_entries_page(
606 &self,
607 spec: &NamespacePath,
608 limit: Option<u32>,
609 cursor: Option<&str>,
610 ) -> Result<ListPathEntriesResponse> {
611 let mut url = format!(
612 "{}/v0/namespaces/{}/filesystem/list?path={}",
613 self.base_url,
614 spec.namespace().as_str(),
615 urlencoding::encode(spec.absolute_path().as_str())
616 );
617 let has_query = true;
618 append_optional_pagination_query(&mut url, has_query, limit, cursor);
619 self.request_json::<(), ListPathEntriesResponse>(self.get(&url), None)
620 .await
621 }
622
623 pub async fn stat_path(&self, spec: &NamespacePath) -> Result<AuthoritativePathEntry> {
624 let url = format!(
625 "{}/v0/namespaces/{}/filesystem/stat?path={}",
626 self.base_url,
627 spec.namespace().as_str(),
628 urlencoding::encode(spec.absolute_path().as_str())
629 );
630 self.request_json::<(), AuthoritativePathEntry>(self.get(&url), None)
631 .await
632 }
633
634 pub async fn get_file_bytes(&self, spec: &NamespacePath) -> Result<Vec<u8>> {
635 let url = format!(
636 "{}/v0/namespaces/{}/filesystem/content?path={}",
637 self.base_url,
638 spec.namespace().as_str(),
639 urlencoding::encode(spec.absolute_path().as_str())
640 );
641 self.request_bytes(&url).await
642 }
643
644 pub async fn get_file_revision_bytes(
645 &self,
646 spec: &NamespacePath,
647 revision_no: RevisionNo,
648 ) -> Result<Vec<u8>> {
649 let url = format!(
650 "{}/v0/namespaces/{}/filesystem/content?path={}&revision_no={}",
651 self.base_url,
652 spec.namespace().as_str(),
653 urlencoding::encode(spec.absolute_path().as_str()),
654 revision_no.0
655 );
656 self.request_bytes(&url).await
657 }
658
659 pub async fn offers_direct_download(&self, size_bytes: u64) -> bool {
672 let Ok(capabilities) = self.capabilities().await else {
673 return false;
674 };
675 capabilities.supports(FEATURE_DOWNLOADS_DIRECT_GET)
676 && capabilities
677 .limits
678 .get(LIMIT_DOWNLOAD_MAX_CONTENT_BYTES)
679 .is_some_and(|proxy_cap| size_bytes > *proxy_cap)
680 }
681
682 pub async fn begin_download(
685 &self,
686 spec: &NamespacePath,
687 revision_no: Option<RevisionNo>,
688 ) -> Result<BeginDownloadResponse> {
689 let url = format!(
690 "{}/v0/namespaces/{}/filesystem/downloads",
691 self.base_url,
692 spec.namespace().as_str()
693 );
694 let request = match revision_no {
695 Some(revision_no) => {
696 BeginDownloadRequest::for_revision(spec.absolute_path().clone(), revision_no)
697 }
698 None => BeginDownloadRequest::for_path(spec.absolute_path().clone()),
699 };
700 self.request_json::<_, BeginDownloadResponse>(self.post(&url), Some(&request))
703 .await
704 }
705
706 pub async fn open_direct_download(
709 &self,
710 download: &BeginDownloadResponse,
711 ) -> Result<DirectDownloadStream> {
712 self.open_direct_download_at(download, 0).await
713 }
714
715 pub async fn open_direct_download_at(
725 &self,
726 download: &BeginDownloadResponse,
727 start_offset: u64,
728 ) -> Result<DirectDownloadStream> {
729 let ObjectTransferAccess::PresignedUrl {
730 method,
731 url,
732 headers,
733 ..
734 } = &download.access;
735 if method != "GET" {
736 return Err(ClientError::Http(format!(
737 "unsupported presigned download method `{method}`"
738 )));
739 }
740 if start_offset > download.content_ref.size_bytes {
741 return Err(ClientError::Http(format!(
742 "cannot resume a download of `{}` at offset {start_offset} of {} bytes",
743 download.absolute_path, download.content_ref.size_bytes
744 )));
745 }
746 let mut request = WireRequest::presigned(reqwest::Method::GET, url);
747 for (name, value) in headers {
748 request = request.header(name, value);
749 }
750 if start_offset > 0 {
751 request = request.header("range", format!("bytes={start_offset}-"));
752 }
753 let body = self.call_for_response_stream(&request).await?;
754 Ok(DirectDownloadStream {
755 body,
756 expected: download.content_ref.clone(),
757 path: download.absolute_path.clone(),
758 sha256: download
759 .content_ref
760 .whole_file_sha256
761 .as_ref()
762 .map(|_| Sha256::new()),
763 size_bytes: start_offset,
766 resumed_from: start_offset,
767 prefix_folded: 0,
768 finished: false,
769 })
770 }
771
772 pub async fn download_via_presigned_url<W>(
791 &self,
792 download: &BeginDownloadResponse,
793 sink: &mut W,
794 ) -> Result<u64>
795 where
796 W: tokio::io::AsyncWrite + Unpin,
797 {
798 use tokio::io::AsyncWriteExt as _;
799 let path = &download.absolute_path;
800 let mut download = self.open_direct_download(download).await?;
801 let mut size_bytes = 0u64;
802 while let Some(chunk) = download.next_chunk().await? {
803 size_bytes += chunk.len() as u64;
804 sink.write_all(&chunk)
805 .await
806 .map_err(|err| ClientError::Io(format!("write of `{path}` failed: {err}")))?;
807 }
808 sink.flush()
809 .await
810 .map_err(|err| ClientError::Io(format!("write of `{path}` failed: {err}")))?;
811 Ok(size_bytes)
812 }
813
814 pub async fn list_file_revisions_page(
815 &self,
816 spec: &NamespacePath,
817 limit: Option<u32>,
818 cursor: Option<&str>,
819 ) -> Result<ListFileRevisionsResponse> {
820 let mut url = format!(
821 "{}/v0/namespaces/{}/filesystem/revisions?path={}",
822 self.base_url,
823 spec.namespace().as_str(),
824 urlencoding::encode(spec.absolute_path().as_str())
825 );
826 let has_query = true;
827 append_optional_pagination_query(&mut url, has_query, limit, cursor);
828 self.request_json::<(), ListFileRevisionsResponse>(self.get(&url), None)
829 .await
830 }
831
832 pub async fn list_trash_page(
833 &self,
834 namespace_id: &NamespaceId,
835 limit: Option<u32>,
836 cursor: Option<&str>,
837 ) -> Result<ListTrashResponse> {
838 let mut url = format!(
839 "{}/v0/namespaces/{}/filesystem/trash",
840 self.base_url,
841 namespace_id.as_str()
842 );
843 let has_query = false;
844 append_optional_pagination_query(&mut url, has_query, limit, cursor);
845 self.request_json::<(), ListTrashResponse>(self.get(&url), None)
846 .await
847 }
848
849 pub async fn health(&self) -> Result<()> {
850 let url = format!("{}/health", self.base_url);
851 self.call_with_transient_retry(&self.get(&url), None)
852 .await?;
853 Ok(())
854 }
855
856 pub async fn begin_upload(
857 &self,
858 namespace_id: &NamespaceId,
859 request: &BeginUploadRequest,
860 ) -> Result<BeginUploadResponse> {
861 let url = format!("{}/v0/namespaces/{namespace_id}/uploads", self.base_url);
862 self.request_json_once::<_, BeginUploadResponse>(self.post(&url), Some(request))
864 .await
865 }
866
867 pub async fn begin_direct_put(
874 &self,
875 namespace_id: &NamespaceId,
876 claim: DirectPutContentClaim,
877 ) -> Result<BeginUploadResponse> {
878 self.begin_upload(
879 namespace_id,
880 &BeginUploadRequest::DirectPut { content: claim },
881 )
882 .await
883 }
884
885 pub async fn begin_direct_multipart(
894 &self,
895 namespace_id: &NamespaceId,
896 options: DirectMultipartUploadOptions,
897 ) -> Result<BeginUploadResponse> {
898 self.begin_upload(
899 namespace_id,
900 &BeginUploadRequest::DirectMultipart {
901 multipart: Some(options),
902 },
903 )
904 .await
905 }
906
907 pub async fn sign_upload_parts(
913 &self,
914 namespace_id: &NamespaceId,
915 upload_id: &UploadId,
916 parts: Vec<UploadPartChecksumClaim>,
917 ) -> Result<SignUploadPartsResponse> {
918 let url = format!(
919 "{}/v0/namespaces/{namespace_id}/uploads/{upload_id}/parts",
920 self.base_url
921 );
922 self.request_json::<_, SignUploadPartsResponse>(
925 self.post(&url),
926 Some(&SignUploadPartsRequest { parts }),
927 )
928 .await
929 }
930
931 pub async fn upload_part_via_presigned_url(
937 &self,
938 part_number: u32,
939 access: &ObjectTransferAccess,
940 crc64nvme: String,
941 bytes: Bytes,
942 ) -> Result<CompletedUploadPart> {
943 let ObjectTransferAccess::PresignedUrl {
944 method,
945 url,
946 headers,
947 ..
948 } = access;
949 if method != "PUT" {
950 return Err(ClientError::Http(format!(
951 "unsupported presigned part method `{method}`"
952 )));
953 }
954 let mut request = WireRequest::presigned(reqwest::Method::PUT, url);
955 for (name, value) in headers {
956 request = request.header(name, value);
957 }
958 let response = self
962 .call_with_transient_retry_headers(&request, Some(&bytes))
963 .await?;
964 let etag = response
965 .get(http::header::ETAG)
966 .and_then(|value| value.to_str().ok())
967 .ok_or_else(|| {
968 ClientError::Http(format!("part {part_number} upload returned no etag"))
969 })?
970 .to_owned();
971 Ok(CompletedUploadPart {
972 part_number,
973 etag,
974 crc64nvme,
975 })
976 }
977
978 pub async fn upload_via_presigned_url(
979 &self,
980 access: &ObjectTransferAccess,
981 bytes: &[u8],
982 ) -> Result<()> {
983 let (method, url, headers) = match access {
984 ObjectTransferAccess::PresignedUrl {
985 method,
986 url,
987 headers,
988 ..
989 } => (method, url, headers),
990 };
991 if method != "PUT" {
992 return Err(ClientError::Http(format!(
993 "unsupported presigned upload method `{method}`"
994 )));
995 }
996 let mut request = WireRequest::presigned(reqwest::Method::PUT, url);
997 for (name, value) in headers {
998 request = request.header(name, value);
999 }
1000 self.call_once(&request, Some(&Bytes::copy_from_slice(bytes)))
1002 .await
1003 .map(|_| ())
1004 }
1005
1006 pub async fn upload_content(
1007 &self,
1008 namespace_id: &NamespaceId,
1009 upload_id: &UploadId,
1010 bytes: &[u8],
1011 ) -> Result<UploadContentResponse> {
1012 let request = self.upload_content_request(namespace_id, upload_id);
1013 let response = self
1016 .call_with_transient_retry(&request, Some(&Bytes::copy_from_slice(bytes)))
1017 .await?;
1018 serde_json::from_slice(&response).map_err(|err| ClientError::Json(err.to_string()))
1019 }
1020
1021 pub async fn upload_streamed_content(
1034 &self,
1035 namespace_id: &NamespaceId,
1036 upload_id: &UploadId,
1037 source: PayloadSource,
1038 ) -> Result<UploadContentResponse> {
1039 let request = self.upload_content_request(namespace_id, upload_id);
1040 let (stream, size_bytes) = source.into_stream();
1041 let response = self
1042 .call_streamed_once(&request, stream, size_bytes)
1043 .await?;
1044 serde_json::from_slice(&response).map_err(|err| ClientError::Json(err.to_string()))
1045 }
1046
1047 fn upload_content_request(
1048 &self,
1049 namespace_id: &NamespaceId,
1050 upload_id: &UploadId,
1051 ) -> WireRequest {
1052 let url = format!(
1053 "{}/v0/namespaces/{namespace_id}/uploads/{upload_id}/content",
1054 self.base_url
1055 );
1056 self.put(&url)
1057 .header("content-type", "application/octet-stream")
1058 }
1059
1060 pub async fn abort_upload(
1066 &self,
1067 namespace_id: &NamespaceId,
1068 upload_id: &UploadId,
1069 ) -> Result<AbortUploadResponse> {
1070 let url = format!(
1071 "{}/v0/namespaces/{namespace_id}/uploads/{upload_id}/abort",
1072 self.base_url
1073 );
1074 self.request_json::<(), AbortUploadResponse>(self.post(&url), None)
1076 .await
1077 }
1078
1079 pub async fn read_upload_status(
1087 &self,
1088 namespace_id: &NamespaceId,
1089 upload_id: &UploadId,
1090 ) -> Result<UploadStatusResponse> {
1091 let url = format!(
1092 "{}/v0/namespaces/{namespace_id}/uploads/{upload_id}",
1093 self.base_url
1094 );
1095 self.request_json::<(), UploadStatusResponse>(self.get(&url), None)
1096 .await
1097 }
1098
1099 pub async fn complete_upload(
1100 &self,
1101 namespace_id: &NamespaceId,
1102 upload_id: &UploadId,
1103 request: &CompleteUploadRequest,
1104 ) -> Result<CompleteUploadResponse> {
1105 let url = format!(
1106 "{}/v0/namespaces/{namespace_id}/uploads/{upload_id}/complete",
1107 self.base_url
1108 );
1109 self.request_json::<_, CompleteUploadResponse>(self.post(&url), Some(request))
1111 .await
1112 }
1113
1114 pub async fn list_changes(
1115 &self,
1116 namespace_id: &NamespaceId,
1117 after_seq: ChangeSeq,
1118 limit: Option<u32>,
1119 ) -> Result<ChangesResponse> {
1120 let mut url = format!(
1121 "{}/v0/namespaces/{namespace_id}/changes?after_seq={}",
1122 self.base_url, after_seq.0
1123 );
1124 if let Some(limit) = limit {
1125 url.push_str(&format!("&limit={limit}"));
1126 }
1127 self.request_json::<(), ChangesResponse>(self.get(&url), None)
1128 .await
1129 }
1130
1131 pub async fn create_checkpoint(
1136 &self,
1137 namespace_id: &NamespaceId,
1138 request: &CreateCheckpointRequest,
1139 ) -> Result<CreateCheckpointResponse> {
1140 let url = format!(
1141 "{}/v0/admin/namespaces/{namespace_id}/checkpoints",
1142 self.base_url
1143 );
1144 self.request_json(self.post(&url), Some(request)).await
1145 }
1146
1147 pub async fn list_checkpoints(
1155 &self,
1156 namespace_id: &NamespaceId,
1157 ) -> Result<ListCheckpointsResponse> {
1158 let url = format!(
1159 "{}/v0/admin/namespaces/{namespace_id}/checkpoints",
1160 self.base_url
1161 );
1162 self.request_json::<(), ListCheckpointsResponse>(self.get(&url), None)
1163 .await
1164 }
1165
1166 pub async fn release_checkpoint(
1169 &self,
1170 namespace_id: &NamespaceId,
1171 checkpoint_id: &CheckpointId,
1172 ) -> Result<ReleaseCheckpointResponse> {
1173 let url = format!(
1174 "{}/v0/admin/namespaces/{namespace_id}/checkpoints/{checkpoint_id}/release",
1175 self.base_url
1176 );
1177 self.request_json::<(), ReleaseCheckpointResponse>(self.post(&url), None)
1178 .await
1179 }
1180
1181 pub async fn maintenance_step(
1185 &self,
1186 namespace_id: &NamespaceId,
1187 request: &MaintenanceStepRequest,
1188 ) -> Result<MaintenanceStepResponse> {
1189 let url = format!(
1190 "{}/v0/admin/namespaces/{namespace_id}/maintenance/step",
1191 self.base_url
1192 );
1193 self.request_json(self.post(&url), Some(request)).await
1194 }
1195
1196 pub async fn probe_store(&self, request: &StoreProbeRequest) -> Result<StoreProbeResponse> {
1204 let url = format!("{}/v0/admin/store/probe", self.base_url);
1205 self.request_json(self.post(&url), Some(request)).await
1206 }
1207
1208 pub async fn grep(
1213 &self,
1214 namespace_id: &NamespaceId,
1215 request: &GrepRequest,
1216 ) -> Result<GrepResponse> {
1217 let url = format!("{}/v0/namespaces/{namespace_id}/query/grep", self.base_url);
1218 self.request_json(self.post(&url), Some(request)).await
1219 }
1220
1221 pub async fn grep_index_status(
1225 &self,
1226 namespace_id: &NamespaceId,
1227 ) -> Result<GrepIndexStatusResponse> {
1228 let url = format!(
1229 "{}/v0/admin/namespaces/{namespace_id}/grep/index",
1230 self.base_url
1231 );
1232 self.request_json::<(), GrepIndexStatusResponse>(self.get(&url), None)
1233 .await
1234 }
1235
1236 pub async fn enable_grep_index(
1239 &self,
1240 namespace_id: &NamespaceId,
1241 ) -> Result<EnableGrepIndexResponse> {
1242 let url = format!(
1243 "{}/v0/admin/namespaces/{namespace_id}/grep/index/enable",
1244 self.base_url
1245 );
1246 self.request_json::<(), EnableGrepIndexResponse>(self.post(&url), None)
1247 .await
1248 }
1249
1250 pub async fn disable_grep_index(
1253 &self,
1254 namespace_id: &NamespaceId,
1255 ) -> Result<DisableGrepIndexResponse> {
1256 let url = format!(
1257 "{}/v0/admin/namespaces/{namespace_id}/grep/index/disable",
1258 self.base_url
1259 );
1260 self.request_json::<(), DisableGrepIndexResponse>(self.post(&url), None)
1261 .await
1262 }
1263
1264 pub async fn gc_grep_index(
1269 &self,
1270 namespace_id: &NamespaceId,
1271 request: &GrepGcRequest,
1272 ) -> Result<GrepGcResponse> {
1273 let url = format!(
1274 "{}/v0/admin/namespaces/{namespace_id}/grep/index/gc",
1275 self.base_url
1276 );
1277 self.request_json(self.post(&url), Some(request)).await
1278 }
1279
1280 pub async fn commit(
1288 &self,
1289 namespace_id: &NamespaceId,
1290 request: &CommitRequest,
1291 ) -> Result<ApiCommitResponse> {
1292 let url = format!("{}/v0/namespaces/{namespace_id}/commits", self.base_url);
1293 self.request_json::<_, ApiCommitResponse>(self.post(&url), Some(request))
1295 .await
1296 }
1297
1298 async fn stage_bytes_as_content_ref(
1306 &self,
1307 namespace_id: &NamespaceId,
1308 bytes: &[u8],
1309 ) -> Result<StagedContent> {
1310 if bytes.len() as u64 >= STREAMING_PUT_MIN_BYTES && self.offers_direct_multipart().await {
1311 return self
1312 .stage_via_multipart(
1313 namespace_id,
1314 MultipartPayload::Held(bytes),
1315 UploadContinuity::default(),
1316 )
1317 .await;
1318 }
1319 self.stage_bytes_via_server(namespace_id, bytes).await
1320 }
1321
1322 async fn stage_source_as_content_ref(
1329 &self,
1330 namespace_id: &NamespaceId,
1331 source: PayloadSource,
1332 continuity: UploadContinuity<'_>,
1333 ) -> Result<StagedContent> {
1334 let known_small = source
1339 .size_bytes()
1340 .is_some_and(|size_bytes| size_bytes < STREAMING_PUT_MIN_BYTES);
1341 if !known_small && self.offers_direct_multipart().await {
1342 let (stream, _) = source.into_stream();
1343 return self
1344 .stage_via_multipart(namespace_id, MultipartPayload::Streamed(stream), continuity)
1345 .await;
1346 }
1347 self.stage_source_via_server(namespace_id, source).await
1350 }
1351
1352 async fn offers_direct_multipart(&self) -> bool {
1354 self.capabilities()
1355 .await
1356 .is_ok_and(|capabilities| capabilities.supports(FEATURE_UPLOADS_DIRECT_MULTIPART))
1357 }
1358
1359 async fn stage_via_multipart(
1372 &self,
1373 namespace_id: &NamespaceId,
1374 payload: MultipartPayload<'_>,
1375 continuity: UploadContinuity<'_>,
1376 ) -> Result<StagedContent> {
1377 let (upload_id, part_size_bytes) = match continuity.resume {
1381 Some(resume) => (resume.upload_id.clone(), resume.part_size_bytes),
1382 None => {
1383 let begin = self
1384 .begin_direct_multipart(namespace_id, DirectMultipartUploadOptions::default())
1385 .await?;
1386 let Some(multipart) = begin.direct_multipart else {
1387 return Err(ClientError::Http(
1388 "server accepted direct_multipart without part geometry".to_owned(),
1389 ));
1390 };
1391 if let Some(journal) = continuity.journal {
1392 journal.began(&begin.upload_id, multipart.part_size_bytes);
1393 }
1394 (begin.upload_id, multipart.part_size_bytes)
1395 }
1396 };
1397 let uploaded = self
1398 .upload_every_part(
1399 namespace_id,
1400 &upload_id,
1401 payload,
1402 part_size_bytes,
1403 continuity,
1404 )
1405 .await;
1406 let uploaded = match uploaded {
1407 Ok(uploaded) => uploaded,
1408 Err(error) => {
1409 let _ = self.abort_upload(namespace_id, &upload_id).await;
1414 return Err(error);
1415 }
1416 };
1417 if uploaded.parts.is_empty() {
1418 let _ = self.abort_upload(namespace_id, &upload_id).await;
1421 return self.stage_bytes_via_server(namespace_id, &[]).await;
1422 }
1423
1424 let response = self
1425 .complete_upload(
1426 namespace_id,
1427 &upload_id,
1428 &CompleteUploadRequest::for_multipart(
1429 DirectMultipartContentClaim {
1430 size_bytes: uploaded.size_bytes,
1431 crc64nvme: uploaded.crc64nvme.value,
1432 },
1433 uploaded.parts,
1434 ),
1435 )
1436 .await?;
1437 Ok(Self::staged_from_completion(response))
1438 }
1439
1440 async fn upload_every_part(
1454 &self,
1455 namespace_id: &NamespaceId,
1456 upload_id: &UploadId,
1457 payload: MultipartPayload<'_>,
1458 part_size_bytes: u64,
1459 continuity: UploadContinuity<'_>,
1460 ) -> Result<UploadedObject> {
1461 let part_size = usize::try_from(part_size_bytes)
1462 .map_err(|_| ClientError::Http("part size does not fit this platform".to_owned()))?;
1463 let landed = continuity
1464 .resume
1465 .map_or::<&[CompletedUploadPart], _>(&[], |resume| &resume.parts);
1466 let mut source = payload.into_parts_of(part_size);
1467 let mut whole_object = Crc64Nvme::new();
1468 let mut size_bytes = 0u64;
1469 let mut parts = Vec::new();
1470 let mut next_part_number = 1u32;
1471
1472 loop {
1473 let mut wave = Vec::with_capacity(DIRECT_MULTIPART_PARTS_IN_FLIGHT);
1474 let mut source_ended = false;
1475 while wave.len() < DIRECT_MULTIPART_PARTS_IN_FLIGHT {
1476 let Some(bytes) = source.next_part().await? else {
1477 source_ended = true;
1478 break;
1479 };
1480 whole_object.update(&bytes);
1481 size_bytes += bytes.len() as u64;
1482 let part_number = next_part_number;
1483 next_part_number += 1;
1484 if let Some(landed) = landed.iter().find(|part| part.part_number == part_number) {
1485 parts.push(landed.clone());
1486 continue;
1487 }
1488 wave.push(PendingPart {
1489 claim: UploadPartChecksumClaim {
1490 part_number,
1491 crc64nvme: StorageChecksum::crc64nvme(&bytes).value,
1492 },
1493 bytes,
1494 });
1495 }
1496 if !wave.is_empty() {
1497 let uploaded = self.upload_wave(namespace_id, upload_id, wave).await?;
1498 if let Some(journal) = continuity.journal {
1499 for part in &uploaded {
1500 journal.part_completed(part);
1501 }
1502 }
1503 parts.extend(uploaded);
1504 }
1505 if source_ended {
1506 break;
1507 }
1508 }
1509
1510 parts.sort_by_key(|part| part.part_number);
1511 Ok(UploadedObject {
1512 size_bytes,
1513 crc64nvme: whole_object.finish(),
1514 parts,
1515 })
1516 }
1517
1518 async fn upload_wave(
1520 &self,
1521 namespace_id: &NamespaceId,
1522 upload_id: &UploadId,
1523 wave: Vec<PendingPart>,
1524 ) -> Result<Vec<CompletedUploadPart>> {
1525 let claims = wave.iter().map(|part| part.claim.clone()).collect();
1526 let signed = self
1527 .sign_upload_parts(namespace_id, upload_id, claims)
1528 .await?;
1529 let mut in_flight = tokio::task::JoinSet::new();
1530 for part in wave {
1531 let access = signed_access(&signed.parts, part.claim.part_number)?;
1532 let client = self.clone();
1533 let namespace_id = namespace_id.clone();
1534 let upload_id = upload_id.clone();
1535 in_flight.spawn(async move {
1536 client
1537 .upload_one_part(&namespace_id, &upload_id, part, access)
1538 .await
1539 });
1540 }
1541 let mut uploaded = Vec::new();
1542 let mut failure = None;
1543 while let Some(joined) = in_flight.join_next().await {
1544 match joined.map_err(|err| ClientError::Http(format!("part upload task failed: {err}")))
1545 {
1546 Ok(Ok(part)) => uploaded.push(part),
1549 Ok(Err(error)) | Err(error) => failure = failure.or(Some(error)),
1550 }
1551 }
1552 match failure {
1553 Some(error) => Err(error),
1554 None => Ok(uploaded),
1555 }
1556 }
1557
1558 async fn upload_one_part(
1564 &self,
1565 namespace_id: &NamespaceId,
1566 upload_id: &UploadId,
1567 part: PendingPart,
1568 mut access: ObjectTransferAccess,
1569 ) -> Result<CompletedUploadPart> {
1570 let part_number = part.claim.part_number;
1571 for attempt in 1..=DIRECT_MULTIPART_PART_ATTEMPTS {
1572 let result = self
1573 .upload_part_via_presigned_url(
1574 part_number,
1575 &access,
1576 part.claim.crc64nvme.clone(),
1577 part.bytes.clone(),
1578 )
1579 .await;
1580 match result {
1581 Ok(uploaded) => return Ok(uploaded),
1582 Err(error) if attempt == DIRECT_MULTIPART_PART_ATTEMPTS => return Err(error),
1583 Err(_) => {
1584 let signed = self
1585 .sign_upload_parts(namespace_id, upload_id, vec![part.claim.clone()])
1586 .await?;
1587 access = signed_access(&signed.parts, part_number)?;
1588 }
1589 }
1590 }
1591 Err(ClientError::Http(format!(
1593 "part {part_number} upload made no attempt"
1594 )))
1595 }
1596
1597 async fn stage_bytes_via_server(
1598 &self,
1599 namespace_id: &NamespaceId,
1600 bytes: &[u8],
1601 ) -> Result<StagedContent> {
1602 let upload = self
1603 .begin_upload(namespace_id, &BeginUploadRequest::ServiceProxied {})
1604 .await?;
1605 let staged = self
1606 .upload_content(namespace_id, &upload.upload_id, bytes)
1607 .await?;
1608 self.complete_staged(namespace_id, &upload.upload_id, staged)
1609 .await
1610 }
1611
1612 async fn stage_source_via_server(
1618 &self,
1619 namespace_id: &NamespaceId,
1620 source: PayloadSource,
1621 ) -> Result<StagedContent> {
1622 let upload = self
1623 .begin_upload(namespace_id, &BeginUploadRequest::ServiceProxied {})
1624 .await?;
1625 let staged = self
1626 .upload_streamed_content(namespace_id, &upload.upload_id, source)
1627 .await;
1628 let staged = match staged {
1629 Ok(staged) => staged,
1630 Err(error) => {
1631 let _ = self.abort_upload(namespace_id, &upload.upload_id).await;
1632 return Err(error);
1633 }
1634 };
1635 self.complete_staged(namespace_id, &upload.upload_id, staged)
1636 .await
1637 }
1638
1639 async fn complete_staged(
1640 &self,
1641 namespace_id: &NamespaceId,
1642 upload_id: &UploadId,
1643 staged: UploadContentResponse,
1644 ) -> Result<StagedContent> {
1645 let response = self
1646 .complete_upload(
1647 namespace_id,
1648 upload_id,
1649 &CompleteUploadRequest::for_content_ref(staged.content_ref),
1650 )
1651 .await?;
1652 Ok(Self::staged_from_completion(response))
1653 }
1654
1655 fn staged_from_completion(response: CompleteUploadResponse) -> StagedContent {
1656 let validated_content_token =
1657 response
1658 .validated_content_token
1659 .map(|token| ValidatedContentToken {
1660 content_ref: response.content_ref.clone(),
1661 token,
1662 });
1663 StagedContent {
1664 content_ref: response.content_ref,
1665 validated_content_token,
1666 }
1667 }
1668
1669 pub async fn put_file_bytes(
1687 &self,
1688 spec: &NamespacePath,
1689 bytes: &[u8],
1690 options: &PutFileOptions,
1691 ) -> Result<ApiCommitResponse> {
1692 let staged = self
1693 .stage_bytes_as_content_ref(spec.namespace(), bytes)
1694 .await?;
1695 self.commit_staged_file(spec, staged, options, UploadedContent::Bytes(bytes))
1696 .await
1697 }
1698
1699 pub async fn put_file_stream(
1723 &self,
1724 spec: &NamespacePath,
1725 source: PayloadSource,
1726 options: &PutFileOptions,
1727 ) -> Result<ApiCommitResponse> {
1728 self.put_file_stream_continuing(spec, source, options, UploadContinuity::default())
1729 .await
1730 }
1731
1732 pub async fn put_file_stream_resumable(
1749 &self,
1750 spec: &NamespacePath,
1751 source: PayloadSource,
1752 options: &PutFileOptions,
1753 journal: &dyn MultipartUploadJournal,
1754 resume: Option<&MultipartUploadResume>,
1755 ) -> Result<ApiCommitResponse> {
1756 self.put_file_stream_continuing(
1757 spec,
1758 source,
1759 options,
1760 UploadContinuity {
1761 resume,
1762 journal: Some(journal),
1763 },
1764 )
1765 .await
1766 }
1767
1768 async fn put_file_stream_continuing(
1769 &self,
1770 spec: &NamespacePath,
1771 source: PayloadSource,
1772 options: &PutFileOptions,
1773 continuity: UploadContinuity<'_>,
1774 ) -> Result<ApiCommitResponse> {
1775 let staged = self
1776 .stage_source_as_content_ref(spec.namespace(), source, continuity)
1777 .await?;
1778 let uploaded = staged.content_ref.clone();
1782 self.commit_staged_file(spec, staged, options, UploadedContent::Streamed(&uploaded))
1783 .await
1784 }
1785
1786 pub async fn commit_completed_upload(
1794 &self,
1795 spec: &NamespacePath,
1796 content_ref: ContentRef,
1797 validated_content_token: Option<String>,
1798 options: &PutFileOptions,
1799 ) -> Result<ApiCommitResponse> {
1800 let uploaded = content_ref.clone();
1801 let staged = StagedContent {
1802 validated_content_token: validated_content_token.map(|token| ValidatedContentToken {
1803 content_ref: content_ref.clone(),
1804 token,
1805 }),
1806 content_ref,
1807 };
1808 self.commit_staged_file(spec, staged, options, UploadedContent::Streamed(&uploaded))
1809 .await
1810 }
1811
1812 async fn commit_staged_file(
1815 &self,
1816 spec: &NamespacePath,
1817 staged: StagedContent,
1818 options: &PutFileOptions,
1819 uploaded: UploadedContent<'_>,
1820 ) -> Result<ApiCommitResponse> {
1821 let commit_id = options.commit_id.clone().unwrap_or_else(CommitId::generate);
1822 let response = self
1823 .commit(
1824 spec.namespace(),
1825 &CommitRequest {
1826 commit_id: commit_id.clone(),
1827 message: options.message.clone(),
1828 content_tokens: staged.validated_content_token.into_iter().collect(),
1829 operations: vec![FilesystemOperation::PutFile {
1830 path: spec.absolute_path().clone(),
1831 content_ref: staged.content_ref,
1832 behavior: options.behavior,
1833 expected_revision_no: options.expected_revision_no,
1834 }],
1835 },
1836 )
1837 .await;
1838 match response {
1839 Ok(response) => Ok(response),
1840 Err(error) if error.code() == Some(ErrorCode::CommitIdReuseConflict) => {
1841 self.reconcile_commit_id_reuse(spec, &commit_id, options, uploaded, error)
1842 .await
1843 }
1844 Err(error) => Err(error),
1845 }
1846 }
1847
1848 async fn reconcile_commit_id_reuse(
1866 &self,
1867 spec: &NamespacePath,
1868 commit_id: &CommitId,
1869 options: &PutFileOptions,
1870 uploaded: UploadedContent<'_>,
1871 conflict: ClientError,
1872 ) -> Result<ApiCommitResponse> {
1873 let namespace_id = spec.namespace();
1874 let Some((committed_seq, committed_fingerprint)) = reported_commit_receipt(&conflict)
1875 else {
1876 return Err(conflict);
1877 };
1878 let Some(committed) = self
1879 .read_committed_change(namespace_id, commit_id, committed_seq)
1880 .await?
1881 else {
1882 return Err(conflict);
1883 };
1884 let Some(content_ref) = sole_committed_content_ref(&committed) else {
1885 return Err(conflict);
1886 };
1887 let retried = loonfs_api::put_retry_fingerprint(
1888 namespace_id,
1889 spec.absolute_path(),
1890 options.behavior,
1891 options.expected_revision_no,
1892 options.message.as_deref(),
1893 content_ref,
1894 );
1895 if retried.ok().as_deref() != Some(committed_fingerprint.as_str()) {
1896 return Err(conflict);
1897 }
1898 if !uploaded_matches_committed(&uploaded, content_ref) {
1899 return Err(conflict);
1900 }
1901 Ok(ApiCommitResponse {
1902 namespace_id: namespace_id.clone(),
1903 commit_id: committed.commit_id,
1904 committed_seq: committed.seq,
1905 })
1906 }
1907
1908 async fn read_committed_change(
1917 &self,
1918 namespace_id: &NamespaceId,
1919 commit_id: &CommitId,
1920 committed_seq: ChangeSeq,
1921 ) -> Result<Option<CommittedChange>> {
1922 let after_seq = ChangeSeq(committed_seq.0.saturating_sub(1));
1923 let page = match self.list_changes(namespace_id, after_seq, Some(1)).await {
1924 Ok(page) => page,
1925 Err(error) if error.code() == Some(ErrorCode::RebootstrapRequired) => {
1928 return Ok(None);
1929 }
1930 Err(error) => return Err(error),
1931 };
1932 Ok(page
1933 .changes
1934 .into_iter()
1935 .find(|change| change.seq == committed_seq && &change.commit_id == commit_id))
1936 }
1937
1938 pub async fn create_directory(
1939 &self,
1940 spec: &NamespacePath,
1941 options: &CreateDirectoryOptions,
1942 ) -> Result<ApiCommitResponse> {
1943 let commit_id = options.commit_id.clone().unwrap_or_else(CommitId::generate);
1944 let response = self
1945 .commit(
1946 spec.namespace(),
1947 &CommitRequest::single(
1948 commit_id,
1949 options.message.clone(),
1950 FilesystemOperation::CreateDirectory {
1951 path: spec.absolute_path().clone(),
1952 parents: options.parents,
1953 },
1954 ),
1955 )
1956 .await?;
1957 Ok(response)
1958 }
1959
1960 pub async fn delete_path(
1961 &self,
1962 spec: &NamespacePath,
1963 options: &DeleteOptions,
1964 ) -> Result<ApiCommitResponse> {
1965 let commit_id = options.commit_id.clone().unwrap_or_else(CommitId::generate);
1966 let response = self
1967 .commit(
1968 spec.namespace(),
1969 &CommitRequest::single(
1970 commit_id,
1971 options.message.clone(),
1972 FilesystemOperation::DeletePath {
1973 path: spec.absolute_path().clone(),
1974 behavior: options.behavior,
1975 expected_inode_id: options.expected_inode_id,
1976 },
1977 ),
1978 )
1979 .await?;
1980 Ok(response)
1981 }
1982
1983 pub async fn move_path(
1984 &self,
1985 from: &NamespacePath,
1986 to: &NamespacePath,
1987 options: &MoveOptions,
1988 ) -> Result<ApiCommitResponse> {
1989 if from.namespace() != to.namespace() {
1990 return Err(ClientError::InvalidNamespacePath(format!(
1991 "cannot move across namespaces: {} -> {}",
1992 from.namespace(),
1993 to.namespace()
1994 )));
1995 }
1996 let commit_id = options.commit_id.clone().unwrap_or_else(CommitId::generate);
1997 let response = self
1998 .commit(
1999 from.namespace(),
2000 &CommitRequest::single(
2001 commit_id,
2002 options.message.clone(),
2003 FilesystemOperation::MovePath {
2004 from_path: from.absolute_path().clone(),
2005 to_path: to.absolute_path().clone(),
2006 behavior: options.behavior,
2007 },
2008 ),
2009 )
2010 .await?;
2011 Ok(response)
2012 }
2013
2014 pub async fn copy_path(
2015 &self,
2016 from: &NamespacePath,
2017 to: &NamespacePath,
2018 options: &CopyOptions,
2019 ) -> Result<ApiCommitResponse> {
2020 if from.namespace() != to.namespace() {
2021 return Err(ClientError::InvalidNamespacePath(format!(
2022 "cannot copy across namespaces: {} -> {}",
2023 from.namespace(),
2024 to.namespace()
2025 )));
2026 }
2027 let commit_id = options.commit_id.clone().unwrap_or_else(CommitId::generate);
2028 let response = self
2029 .commit(
2030 from.namespace(),
2031 &CommitRequest::single(
2032 commit_id,
2033 options.message.clone(),
2034 FilesystemOperation::CopyPath {
2035 from_path: from.absolute_path().clone(),
2036 to_path: to.absolute_path().clone(),
2037 behavior: options.behavior,
2038 },
2039 ),
2040 )
2041 .await?;
2042 Ok(response)
2043 }
2044
2045 pub async fn undelete(
2049 &self,
2050 namespace: &NamespaceId,
2051 inode_id: InodeId,
2052 deleted_at_seq: ChangeSeq,
2053 path: Option<&AbsolutePath>,
2054 options: &UndeleteOptions,
2055 ) -> Result<ApiCommitResponse> {
2056 let commit_id = options.commit_id.clone().unwrap_or_else(CommitId::generate);
2059 let response = self
2060 .commit(
2061 namespace,
2062 &CommitRequest::single(
2063 commit_id,
2064 options.message.clone(),
2065 FilesystemOperation::Undelete {
2066 inode_id,
2067 deleted_at_seq,
2068 path: path.cloned(),
2069 },
2070 ),
2071 )
2072 .await?;
2073 Ok(response)
2074 }
2075
2076 pub async fn restore_file_revision(
2077 &self,
2078 spec: &NamespacePath,
2079 source_revision_no: RevisionNo,
2080 options: &RestoreRevisionOptions,
2081 ) -> Result<ApiCommitResponse> {
2082 let commit_id = options.commit_id.clone().unwrap_or_else(CommitId::generate);
2083 let response = self
2084 .commit(
2085 spec.namespace(),
2086 &CommitRequest::single(
2087 commit_id,
2088 options.message.clone(),
2089 FilesystemOperation::RestoreRevision {
2090 path: spec.absolute_path().clone(),
2091 source_revision_no,
2092 },
2093 ),
2094 )
2095 .await?;
2096 Ok(response)
2097 }
2098}
2099
2100impl NamespacePath {
2101 pub fn parse(namespace: &str, absolute_path: &str) -> Result<Self> {
2103 let namespace = NamespaceId::parse(namespace)
2104 .map_err(|error| ClientError::InvalidNamespacePath(error.to_string()))?;
2105 let absolute_path = AbsolutePath::parse(absolute_path)
2106 .map_err(|error| ClientError::InvalidNamespacePath(error.to_string()))?;
2107 Ok(Self {
2108 namespace,
2109 absolute_path,
2110 })
2111 }
2112
2113 pub fn new(namespace: NamespaceId, absolute_path: AbsolutePath) -> Self {
2115 Self {
2116 namespace,
2117 absolute_path,
2118 }
2119 }
2120
2121 pub fn namespace(&self) -> &NamespaceId {
2123 &self.namespace
2124 }
2125
2126 pub fn absolute_path(&self) -> &AbsolutePath {
2128 &self.absolute_path
2129 }
2130}
2131
2132fn append_optional_pagination_query(
2133 url: &mut String,
2134 has_query: bool,
2135 limit: Option<u32>,
2136 cursor: Option<&str>,
2137) {
2138 let mut has_query = has_query;
2139 if let Some(limit) = limit {
2140 append_query_param(url, &mut has_query, "limit", &limit.to_string());
2141 }
2142 if let Some(cursor) = cursor {
2143 append_query_param(url, &mut has_query, "cursor", cursor);
2144 }
2145}
2146
2147fn append_query_param(url: &mut String, has_query: &mut bool, name: &str, value: &str) {
2148 url.push(if *has_query { '&' } else { '?' });
2149 *has_query = true;
2150 url.push_str(name);
2151 url.push('=');
2152 url.push_str(&urlencoding::encode(value));
2153}
2154
2155#[cfg(test)]
2156mod download_tests;
2157
2158#[cfg(test)]
2159mod streaming_tests;
2160
2161#[cfg(test)]
2162mod tests {
2163 use super::*;
2164 use crate::transport::{transient_failure, MAX_TRANSIENT_ATTEMPTS};
2165 use loonfs_api::{ContentId, ErrorCode, ErrorKind};
2166 use std::fs;
2167 use tempfile::tempdir;
2168
2169 fn test_content_ref(bytes: &[u8]) -> ContentRef {
2170 ContentRef::blob_v1(ContentId::generate(), bytes)
2171 }
2172
2173 fn direct_put_claim(bytes: &[u8]) -> DirectPutContentClaim {
2174 let content_ref = test_content_ref(bytes);
2175 DirectPutContentClaim {
2176 size_bytes: content_ref.size_bytes,
2177 sha256: content_ref.storage_checksum.value,
2178 }
2179 }
2180
2181 fn crc32c_content_ref(bytes: &[u8]) -> ContentRef {
2186 ContentRef {
2187 kind: loonfs_api::ContentRefKind::BlobV1,
2188 content_id: ContentId::generate(),
2189 size_bytes: bytes.len() as u64,
2190 storage_checksum: StorageChecksum {
2191 algorithm: ChecksumAlgorithm::Crc32c,
2192 value: "0f5c0a1e".to_owned(),
2193 },
2194 whole_file_sha256: None,
2195 }
2196 }
2197
2198 #[test]
2199 fn a_digest_this_client_cannot_compare_leaves_the_reuse_conflict_standing() {
2200 let bytes = b"retried payload";
2201 let committed = crc32c_content_ref(bytes);
2202 let uploaded = UploadedContent::Bytes(bytes);
2203
2204 assert_eq!(
2205 uploaded.matches(&committed.storage_checksum),
2206 None,
2207 "the fixture must reach the refusal, not a comparison"
2208 );
2209 assert!(!uploaded_matches_committed(&uploaded, &committed));
2210 }
2211
2212 #[test]
2213 fn a_whole_file_digest_over_the_same_bytes_proves_the_retry_did_this_work() {
2214 let bytes = b"retried payload";
2215 let committed = test_content_ref(bytes);
2216
2217 assert!(uploaded_matches_committed(
2218 &UploadedContent::Bytes(bytes),
2219 &committed
2220 ));
2221 assert!(!uploaded_matches_committed(
2222 &UploadedContent::Bytes(b"some other payload"),
2223 &committed
2224 ));
2225 }
2226
2227 #[test]
2230 fn construction_validates_config_like_load_does() {
2231 let error = super::Client::new(super::ClientConfig {
2232 server_url: "ftp://example.com".to_owned(),
2233 auth_token: None,
2234 request_timeout_ms: None,
2235 disable_transient_retry: false,
2236 ca_cert_path: None,
2237 })
2238 .expect_err("ftp scheme must be rejected");
2239 assert!(
2240 matches!(
2241 &error,
2242 super::ClientError::ConfigValidation {
2243 field: "server_url",
2244 ..
2245 }
2246 ),
2247 "unexpected error: {error:?}"
2248 );
2249 }
2250
2251 #[test]
2255 fn an_unusable_ca_bundle_fails_construction_and_names_the_path() {
2256 let dir = tempdir().expect("tempdir");
2257 let missing = dir.path().join("absent.crt");
2258 let garbage = dir.path().join("garbage.crt");
2259 fs::write(&garbage, b"this is not a certificate\n").expect("write garbage");
2260
2261 for path in [missing, garbage] {
2262 let display = path.display().to_string();
2263 let error = super::Client::new(super::ClientConfig {
2264 server_url: "https://example.com".to_owned(),
2265 auth_token: None,
2266 request_timeout_ms: None,
2267 disable_transient_retry: false,
2268 ca_cert_path: Some(display.clone()),
2269 })
2270 .expect_err("unusable ca bundle");
2271 match &error {
2272 super::ClientError::ConfigValidation {
2273 field: "ca_cert_path",
2274 reason,
2275 } => assert!(
2276 reason.contains(&display),
2277 "the reason must name the path, got: {reason}"
2278 ),
2279 other => unreachable!("unexpected error: {other:?}"),
2280 }
2281 }
2282 }
2283
2284 #[test]
2288 fn client_config_rejects_unknown_keys() {
2289 let error = toml::from_str::<ClientConfig>(
2290 "server_url = \"http://localhost:1\"\nauth_tokn = \"oops\"\n",
2291 )
2292 .expect_err("unknown key must fail decode");
2293 assert!(error.to_string().contains("auth_tokn"), "{error}");
2294
2295 let config: ClientConfig =
2296 toml::from_str("server_url = \"http://localhost:1\"\n").expect("minimal config");
2297 assert!(config.auth_token.is_none());
2298 }
2299 #[test]
2304 fn transient_failure_covers_transport_and_retryable_unavailability_only() {
2305 let api = |code: &str| ClientError::Api {
2306 status: 503,
2307 code: code.to_owned(),
2308 feature: None,
2309 message: String::new(),
2310 request_id: None,
2311 details: None,
2312 };
2313 assert!(transient_failure(
2314 true,
2315 &ClientError::Http("reset".to_owned())
2316 ));
2317 assert!(transient_failure(false, &api("server_busy")));
2318 assert!(transient_failure(false, &api("commit_queue_full")));
2319 assert!(transient_failure(false, &api("shutting_down")));
2320 assert!(!transient_failure(false, &api("server_error")));
2321 assert!(!transient_failure(false, &api("maintenance_required")));
2322 assert!(!transient_failure(
2323 false,
2324 &ClientError::Http("http status 502 with a non-envelope body".to_owned())
2325 ));
2326 }
2327
2328 #[tokio::test]
2331 async fn transport_failures_resend_up_to_the_attempt_cap() {
2332 let namespace_id = NamespaceId::parse("demo").expect("valid namespace id");
2333 let transport = crate::transport::test_transport::failures(MAX_TRANSIENT_ATTEMPTS as usize);
2334 let retrying = Client::new(ClientConfig {
2335 server_url: "http://example.invalid".to_owned(),
2336 auth_token: None,
2337 request_timeout_ms: None,
2338 disable_transient_retry: false,
2339 ca_cert_path: None,
2340 })
2341 .expect("valid client config");
2342 let error = retrying
2343 .namespace_status(&namespace_id)
2344 .await
2345 .expect_err("dropped connections must fail");
2346 assert!(matches!(error, ClientError::Http(_)), "{error:?}");
2347 assert_eq!(transport.attempts(), MAX_TRANSIENT_ATTEMPTS as usize);
2348 drop(transport);
2349
2350 let transport = crate::transport::test_transport::failures(1);
2351 let single_shot = Client::new(ClientConfig {
2352 server_url: "http://example.invalid".to_owned(),
2353 auth_token: None,
2354 request_timeout_ms: None,
2355 disable_transient_retry: true,
2356 ca_cert_path: None,
2357 })
2358 .expect("valid client config");
2359 single_shot
2360 .namespace_status(&namespace_id)
2361 .await
2362 .expect_err("dropped connection must fail without retry");
2363 assert_eq!(transport.attempts(), 1);
2364 }
2365
2366 fn retry_policy_client() -> Client {
2367 Client::new(ClientConfig {
2368 server_url: "http://example.invalid".to_owned(),
2369 auth_token: None,
2370 request_timeout_ms: None,
2371 disable_transient_retry: false,
2372 ca_cert_path: None,
2373 })
2374 .expect("valid client config")
2375 }
2376
2377 fn single_attempt_probe() -> (crate::transport::test_transport::Guard, Client) {
2381 (
2382 crate::transport::test_transport::failure_then_success(b"{}".to_vec()),
2383 retry_policy_client(),
2384 )
2385 }
2386
2387 fn assert_single_attempt<T>(
2388 result: Result<T>,
2389 transport: &crate::transport::test_transport::Guard,
2390 ) {
2391 assert!(
2392 matches!(result, Err(ClientError::Http(_))),
2393 "expected the first transport failure to surface"
2394 );
2395 assert_eq!(transport.attempts(), 1);
2396 }
2397
2398 #[tokio::test]
2399 async fn retry_policy_lifecycle_mutations_are_single_attempt() {
2400 let namespace_id = NamespaceId::parse("demo").expect("valid namespace id");
2401 let fork_id = NamespaceId::parse("fork").expect("valid id");
2402
2403 let (transport, client) = single_attempt_probe();
2404 assert_single_attempt(client.create_namespace(&namespace_id).await, &transport);
2405 drop(transport);
2406
2407 let (transport, client) = single_attempt_probe();
2408 assert_single_attempt(
2409 client.fork_namespace(&namespace_id, &fork_id).await,
2410 &transport,
2411 );
2412 drop(transport);
2413
2414 let (transport, client) = single_attempt_probe();
2415 assert_single_attempt(
2416 client
2417 .delete_namespace(&namespace_id, Some(ChangeSeq(7)))
2418 .await,
2419 &transport,
2420 );
2421 }
2422
2423 #[tokio::test]
2424 async fn retry_policy_commit_id_filesystem_mutation_retries() {
2425 let namespace_id = NamespaceId::parse("demo").expect("valid namespace id");
2426 let commit_id =
2427 CommitId::parse("c_00000000000000000000000000000001").expect("valid commit id");
2428 let response = ApiCommitResponse {
2429 namespace_id: namespace_id.clone(),
2430 commit_id: commit_id.clone(),
2431 committed_seq: ChangeSeq(1),
2432 };
2433 let transport = crate::transport::test_transport::failure_then_success(
2434 serde_json::to_vec(&response).expect("serialize response"),
2435 );
2436 let client = retry_policy_client();
2437 let spec = NamespacePath::parse("demo", "/docs").expect("valid namespace path");
2438
2439 let actual = client
2440 .create_directory(
2441 &spec,
2442 &CreateDirectoryOptions {
2443 commit_id: Some(commit_id),
2444 message: None,
2445 ..CreateDirectoryOptions::default()
2446 },
2447 )
2448 .await
2449 .expect("commit-id mutation should retry");
2450 assert_eq!(actual, response);
2451 assert_eq!(transport.attempts(), 2);
2452 }
2453
2454 #[tokio::test]
2455 async fn retry_policy_read_retries() {
2456 let namespace_id = NamespaceId::parse("demo").expect("valid namespace id");
2457 let response = NamespaceStatusResponse {
2458 namespace_id: namespace_id.clone(),
2459 head_seq: ChangeSeq(0),
2460 current_manifest_id: None,
2461 wal_tail_segments: 0,
2462 retention_floor_seq: ChangeSeq(0),
2463 };
2464 let transport = crate::transport::test_transport::failure_then_success(
2465 serde_json::to_vec(&response).expect("serialize response"),
2466 );
2467 let client = retry_policy_client();
2468
2469 let actual = client
2470 .namespace_status(&namespace_id)
2471 .await
2472 .expect("read should retry");
2473 assert_eq!(actual, response);
2474 assert_eq!(transport.attempts(), 2);
2475 }
2476
2477 #[tokio::test]
2478 async fn retry_policy_upload_begins_are_single_attempt() {
2479 let namespace_id = NamespaceId::parse("demo").expect("valid namespace id");
2480
2481 let (transport, client) = single_attempt_probe();
2482 assert_single_attempt(
2483 client
2484 .begin_upload(&namespace_id, &BeginUploadRequest::ServiceProxied {})
2485 .await,
2486 &transport,
2487 );
2488 drop(transport);
2489
2490 let (transport, client) = single_attempt_probe();
2491 assert_single_attempt(
2492 client
2493 .begin_direct_put(&namespace_id, direct_put_claim(b"direct"))
2494 .await,
2495 &transport,
2496 );
2497 }
2498
2499 #[tokio::test]
2500 async fn retry_policy_presigned_upload_is_single_attempt() {
2501 let transport = crate::transport::test_transport::failure_then_success(Vec::new());
2502 let client = retry_policy_client();
2503 let access = ObjectTransferAccess::PresignedUrl {
2504 method: "PUT".to_owned(),
2505 url: "http://example.invalid/upload".to_owned(),
2506 headers: std::collections::BTreeMap::new(),
2507 expires_at_ms: 1,
2508 };
2509
2510 let result = client.upload_via_presigned_url(&access, b"direct").await;
2511
2512 assert!(matches!(result, Err(ClientError::Http(_))), "{result:?}");
2513 assert_eq!(transport.attempts(), 1);
2514 }
2515
2516 #[tokio::test]
2517 async fn retry_policy_proxied_upload_content_retries() {
2518 let namespace_id = NamespaceId::parse("demo").expect("valid namespace id");
2519 let upload_id = loonfs_api::UploadId::parse("upl_00000000000000000000000000000001")
2520 .expect("valid upload id");
2521 let response = UploadContentResponse {
2522 namespace_id: namespace_id.clone(),
2523 upload_id: upload_id.clone(),
2524 content_ref: test_content_ref(b"content"),
2525 };
2526 let transport = crate::transport::test_transport::failure_then_success(
2527 serde_json::to_vec(&response).expect("serialize response"),
2528 );
2529 let client = retry_policy_client();
2530
2531 let actual = client
2532 .upload_content(&namespace_id, &upload_id, b"content")
2533 .await
2534 .expect("identical content staging should retry");
2535 assert_eq!(actual, response);
2536 assert_eq!(transport.attempts(), 2);
2537 }
2538
2539 #[tokio::test]
2540 async fn retry_policy_upload_completion_retries() {
2541 let namespace_id = NamespaceId::parse("demo").expect("valid namespace id");
2542 let upload_id = loonfs_api::UploadId::parse("upl_00000000000000000000000000000001")
2543 .expect("valid upload id");
2544 let content_ref = test_content_ref(b"content");
2545 let response = CompleteUploadResponse {
2546 namespace_id: namespace_id.clone(),
2547 upload_id: upload_id.clone(),
2548 content_ref: content_ref.clone(),
2549 validated_content_token: None,
2550 };
2551 let transport = crate::transport::test_transport::failure_then_success(
2552 serde_json::to_vec(&response).expect("serialize response"),
2553 );
2554 let client = retry_policy_client();
2555
2556 let actual = client
2557 .complete_upload(
2558 &namespace_id,
2559 &upload_id,
2560 &CompleteUploadRequest::for_content_ref(content_ref),
2561 )
2562 .await
2563 .expect("completed-session replay should retry");
2564 assert_eq!(actual, response);
2565 assert_eq!(transport.attempts(), 2);
2566 }
2567
2568 #[test]
2572 fn status_errors_keep_the_status_when_the_body_is_not_the_envelope() {
2573 let error = crate::transport::map_status_error(502, b"<html>upstream error</html>");
2574
2575 let ClientError::Http(message) = error else {
2576 unreachable!("expected Http error, got {error:?}");
2577 };
2578 assert!(message.contains("502"), "{message}");
2579 assert!(message.contains("non-envelope body"), "{message}");
2580 }
2581
2582 fn api_error(status: u16, code: &str) -> ClientError {
2583 ClientError::Api {
2584 status,
2585 code: code.to_owned(),
2586 feature: None,
2587 message: "test".to_owned(),
2588 request_id: None,
2589 details: None,
2590 }
2591 }
2592
2593 #[test]
2594 fn api_errors_with_known_codes_classify_through_the_registry() {
2595 let error = api_error(409, "stale_revision");
2596 assert_eq!(error.code(), Some(ErrorCode::StaleRevision));
2597 assert_eq!(error.kind(), Some(ErrorKind::Conflict));
2598
2599 let error = api_error(409, "content_not_prepared");
2600 assert_eq!(error.code(), Some(ErrorCode::ContentNotPrepared));
2601 assert_eq!(error.kind(), Some(ErrorKind::Conflict));
2602
2603 let error = api_error(410, "namespace_deleted");
2604 assert_eq!(error.code(), Some(ErrorCode::NamespaceDeleted));
2605 assert_eq!(error.kind(), Some(ErrorKind::Gone));
2606
2607 let error = api_error(503, "commit_outcome_unknown");
2608 assert_eq!(error.code(), Some(ErrorCode::CommitOutcomeUnknown));
2609 assert_eq!(error.kind(), Some(ErrorKind::OutcomeUnknown));
2610
2611 let error = api_error(500, "index_corrupt");
2612 assert_eq!(error.code(), Some(ErrorCode::IndexCorrupt));
2613 assert_eq!(error.kind(), Some(ErrorKind::DataCorruption));
2614 }
2615
2616 #[test]
2617 fn api_errors_with_unknown_codes_fall_back_to_the_status_class() {
2618 for (status, kind) in [
2619 (400, ErrorKind::InvalidRequest),
2620 (404, ErrorKind::InvalidRequest),
2621 (500, ErrorKind::Internal),
2622 (503, ErrorKind::Unavailable),
2623 ] {
2624 let error = api_error(status, "code_from_a_newer_server");
2625 assert_eq!(error.code(), None);
2626 assert_eq!(error.kind(), Some(kind), "status {status}");
2627 }
2628 }
2629
2630 #[test]
2631 fn non_api_errors_have_no_code_or_kind() {
2632 let error = ClientError::Http("connection refused".to_owned());
2633 assert_eq!(error.code(), None);
2634 assert_eq!(error.kind(), None);
2635 }
2636
2637 #[test]
2638 fn load_rejects_invalid_server_url() {
2639 let path = write_config(
2640 r#"
2641server_url = "ftp://example.com"
2642auth_token = "dev-token"
2643"#,
2644 );
2645
2646 let error = ClientConfig::load(&path).expect_err("invalid server url");
2647
2648 assert!(
2649 matches!(error, ClientError::ConfigValidation { field, .. } if field == "server_url"),
2650 "expected config validation error, got {error:?}"
2651 );
2652 }
2653
2654 #[test]
2655 fn load_rejects_blank_auth_token() {
2656 let path = write_config(
2657 r#"
2658server_url = "http://127.0.0.1:9400"
2659auth_token = " "
2660"#,
2661 );
2662
2663 let error = ClientConfig::load(&path).expect_err("blank auth token");
2664
2665 assert!(
2666 matches!(error, ClientError::ConfigValidation { field, .. } if field == "auth_token"),
2667 "expected config validation error, got {error:?}"
2668 );
2669 }
2670
2671 #[test]
2672 fn load_preserves_missing_file_as_config_io() {
2673 let temp_dir = tempdir().expect("tempdir");
2674 let path = temp_dir.path().join("missing.toml");
2675
2676 let error = ClientConfig::load(&path).expect_err("missing config");
2677
2678 assert!(matches!(error, ClientError::ConfigIo(_)));
2679 }
2680
2681 #[test]
2682 fn load_preserves_decode_error() {
2683 let path = write_config("server_url = [");
2684
2685 let error = ClientConfig::load(&path).expect_err("decode error");
2686
2687 assert!(matches!(error, ClientError::ConfigDecode(_)));
2688 }
2689
2690 #[test]
2691 fn namespace_path_parse_rejects_invalid_namespace_id() {
2692 for namespace in ["bad/name", "Demo", "..", "demo?"] {
2693 assert!(
2694 matches!(
2695 NamespacePath::parse(namespace, "/notes.txt"),
2696 Err(ClientError::InvalidNamespacePath(_))
2697 ),
2698 "expected invalid namespace path for id {namespace:?}"
2699 );
2700 }
2701 }
2702
2703 #[test]
2707 fn namespace_path_parse_rejects_invalid_paths() {
2708 for path in ["notes.txt", "", "/docs/../a.txt", "/docs/./a.txt"] {
2709 assert!(
2710 matches!(
2711 NamespacePath::parse("demo", path),
2712 Err(ClientError::InvalidNamespacePath(_))
2713 ),
2714 "expected invalid namespace path for path {path:?}"
2715 );
2716 }
2717 }
2718
2719 fn write_config(contents: &str) -> std::path::PathBuf {
2720 let temp_dir = tempdir().expect("tempdir");
2721 let path = temp_dir.path().join("client.toml");
2722 fs::write(&path, contents).expect("write config");
2723 let _ = temp_dir.keep();
2724 path
2725 }
2726}