1use crate::cache::{MetadataTableCache, WalTailProjectionCache};
5use crate::checkpoint::{CheckpointFilesPage, CheckpointFilesPageCursor};
6use crate::commit_engine::CommitCandidate;
7use crate::context::MutationContext;
8use crate::error::{CoreError, Result};
9use crate::namespace::basis::MetadataBasis;
10use crate::namespace::catalog::VerifiedNamespaceCatalogEntry;
11use crate::namespace::{bootstrap, fork, BootstrapNamespaceError};
12use crate::options::{BootstrapOptions, DeleteNamespaceOptions};
13use crate::path::read::{
14 load_metadata_view, CurrentFileState, DirectDownloadTarget, LoadedMetadataView, ReadLoadContext,
15};
16use crate::protocol::CompletedUpload;
17use crate::storage::content::FileContentStream;
18use crate::storage::content_admission::{CompletedUploadReceipt, PreparedContent};
19use crate::time::current_time_ms;
20use loonfs_api::v0::{
21 AbortUploadResponse, BeginUploadRequest, BeginUploadResponse, ChangesResponse, CommitResponse,
22 CompleteUploadRequest, CompleteUploadResponse, DirectMultipartUploadOptions,
23 DirectPutContentClaim, UploadContentResponse, UploadPartChecksumClaim, UploadStatusResponse,
24};
25use loonfs_api::wire::control::{CheckpointOwner, HeadState, NamespaceState};
26use loonfs_api::EffectiveLimit;
27use loonfs_api::{
28 AdvanceRetentionResponse, AuthoritativeFileBytes, AuthoritativePathEntry, ChangeSeq,
29 CheckpointId, ContentRef, CreateCheckpointResponse, DeleteNamespaceResponse,
30 DirectoryPageCursor, FileRevision, FileRevisionsPageCursor, FlushWalResponse, InodeId,
31 ListCheckpointsResponse, NamespaceId, NamespaceSummary, Page, PageRequest,
32 ReleaseCheckpointResponse, RevisionNo, StorageChecksum, TrashEntry, TrashPageCursor, UploadId,
33};
34use loonfs_objectstore::{ByteStream, ObjectStore};
35use std::num::NonZeroU64;
36use std::sync::Arc;
37use thiserror::Error;
38
39#[derive(Debug, Clone)]
50pub struct RuntimeReadContext {
51 pub head: HeadState,
52 pub head_etag: String,
53 pub basis: MetadataBasis,
57 pub table_cache: Arc<MetadataTableCache>,
58 pub tail_cache: Arc<WalTailProjectionCache>,
59}
60
61fn runtime_read_load_context(context: &RuntimeReadContext) -> ReadLoadContext<'_> {
62 ReadLoadContext::pinned_head(
63 &context.head,
64 Some(context.head_etag.as_str()),
65 &context.basis,
66 Some(&context.table_cache),
67 Some(&context.tail_cache),
68 )
69}
70
71#[derive(Debug, Clone, PartialEq, Eq)]
73pub struct DirectPutUploadTarget {
74 pub content_ref: ContentRef,
75 pub object_key: String,
76}
77
78#[derive(Debug, Clone, PartialEq, Eq)]
80pub struct BeginDirectPutUploadTargetResponse {
81 pub namespace_id: NamespaceId,
82 pub upload_id: UploadId,
83 pub target: DirectPutUploadTarget,
84}
85
86#[derive(Debug, Clone, PartialEq, Eq)]
92pub struct DirectMultipartUploadTarget {
93 pub object_key: String,
94 pub part_size_bytes: u64,
95}
96
97#[derive(Debug, Clone, PartialEq, Eq)]
99pub struct BeginDirectMultipartUploadTargetResponse {
100 pub namespace_id: NamespaceId,
101 pub upload_id: UploadId,
102 pub target: DirectMultipartUploadTarget,
103}
104
105#[derive(Debug, Clone, PartialEq, Eq)]
107pub struct MultipartPartTarget {
108 pub part_number: u32,
109 pub checksum: StorageChecksum,
110}
111
112#[derive(Debug, Clone, PartialEq, Eq)]
114pub struct MultipartPartTargets {
115 pub object_key: String,
116 pub provider_upload_id: String,
117 pub parts: Vec<MultipartPartTarget>,
118}
119
120#[derive(Debug)]
122struct EngineWriter {
123 writer_id: String,
124}
125
126#[derive(Debug)]
136pub struct NamespaceEngine<S> {
137 store: S,
138 namespace_id: NamespaceId,
139 writer: Option<EngineWriter>,
140}
141
142impl<S: ObjectStore> NamespaceEngine<S> {
143 pub fn builder(store: S) -> NamespaceEngineBuilder<S> {
148 NamespaceEngineBuilder {
149 store,
150 namespace_id: None,
151 writer_id: None,
152 }
153 }
154
155 pub fn namespace_id(&self) -> &NamespaceId {
157 &self.namespace_id
158 }
159
160 pub fn writer_id(&self) -> Option<&str> {
163 self.writer.as_ref().map(|writer| writer.writer_id.as_str())
164 }
165
166 pub async fn bootstrap_namespace(
170 &self,
171 options: BootstrapOptions,
172 ) -> std::result::Result<NamespaceSummary, BootstrapNamespaceError> {
173 bootstrap::bootstrap_namespace(
174 &self.store,
175 &self.namespace_id,
176 &self.mutation_context()?,
177 options.allow_existing,
178 )
179 .await
180 }
181
182 pub async fn fork_namespace(&self, target: &NamespaceId) -> Result<NamespaceSummary> {
186 fork::fork_namespace(
187 &self.store,
188 &self.namespace_id,
189 target,
190 &self.mutation_context()?,
191 )
192 .await
193 }
194
195 pub async fn delete_namespace(
199 &self,
200 options: DeleteNamespaceOptions,
201 ) -> Result<DeleteNamespaceResponse> {
202 crate::commit_engine::delete_namespace(
203 &self.store,
204 &self.namespace_id,
205 options,
206 &self.mutation_context()?,
207 )
208 .await
209 }
210
211 pub async fn resolve_path(
213 &self,
214 path: impl AsRef<str>,
215 context: &RuntimeReadContext,
216 ) -> Result<AuthoritativePathEntry> {
217 let view = self.load_read_view(context).await?;
218 view.resolve_path(path.as_ref()).await
219 }
220
221 pub async fn list_path_page(
223 &self,
224 path: impl AsRef<str>,
225 request: PageRequest<DirectoryPageCursor>,
226 context: &RuntimeReadContext,
227 ) -> Result<Page<AuthoritativePathEntry, DirectoryPageCursor>> {
228 let view = self.load_read_view(context).await?;
229 view.list_path_page(path.as_ref(), request).await
230 }
231
232 pub async fn read_file(
234 &self,
235 path: impl AsRef<str>,
236 context: &RuntimeReadContext,
237 max_content_bytes: Option<u64>,
238 ) -> Result<AuthoritativeFileBytes> {
239 let view = self.load_read_view(context).await?;
240 view.read_file_bytes(&self.store, path.as_ref(), max_content_bytes)
241 .await
242 }
243
244 pub async fn read_file_stream(
264 &self,
265 path: impl AsRef<str>,
266 context: &RuntimeReadContext,
267 chunk_bytes: NonZeroU64,
268 start_offset: u64,
269 ) -> Result<FileContentStream<S>>
270 where
271 S: Clone,
272 {
273 let view = self.load_read_view(context).await?;
274 let (entry, content_ref) = view.resolve_file_content(path.as_ref()).await?;
275 if start_offset > content_ref.size_bytes {
276 return Err(CoreError::ResumeOffsetOutOfRange {
277 start_offset,
278 size_bytes: content_ref.size_bytes,
279 });
280 }
281 Ok(FileContentStream::open(
282 self.store.clone(),
283 view.content_store_id(),
284 entry,
285 content_ref,
286 chunk_bytes,
287 start_offset,
288 )
289 .await?)
290 }
291
292 pub async fn direct_download_target(
300 &self,
301 path: impl AsRef<str>,
302 revision_no: Option<RevisionNo>,
303 context: &RuntimeReadContext,
304 ) -> Result<DirectDownloadTarget> {
305 let view = self.load_read_view(context).await?;
306 view.direct_download_target(path.as_ref(), revision_no)
307 .await
308 }
309
310 pub async fn list_file_revisions_page(
312 &self,
313 path: impl AsRef<str>,
314 request: PageRequest<FileRevisionsPageCursor>,
315 context: &RuntimeReadContext,
316 ) -> Result<Page<FileRevision, FileRevisionsPageCursor>> {
317 let view = self.load_read_view(context).await?;
318 view.list_file_revisions_page(path.as_ref(), request).await
319 }
320
321 pub async fn list_trash_page(
323 &self,
324 request: PageRequest<TrashPageCursor>,
325 context: &RuntimeReadContext,
326 ) -> Result<Page<TrashEntry, TrashPageCursor>> {
327 let view = self.load_read_view(context).await?;
328 view.list_trash_page(request).await
329 }
330
331 pub async fn list_checkpoint_files_page(
341 &self,
342 checkpoint_id: &CheckpointId,
343 request: PageRequest<CheckpointFilesPageCursor>,
344 context: &RuntimeReadContext,
345 ) -> Result<CheckpointFilesPage> {
346 self.live_catalog(context)?;
348 crate::checkpoint::list_checkpoint_files_page(
349 &self.store,
350 Some(context.table_cache.as_ref()),
351 &self.namespace_id,
352 checkpoint_id,
353 request,
354 )
355 .await
356 }
357
358 pub async fn resolve_current_files(
369 &self,
370 inode_ids: &[InodeId],
371 context: &RuntimeReadContext,
372 ) -> Result<Vec<CurrentFileState>> {
373 crate::path::read::ensure_resolve_batch_within_cap(inode_ids.len())?;
374 let view = self.load_read_view(context).await?;
375 crate::path::read::resolve_current_files(&view, inode_ids).await
376 }
377
378 pub async fn read_content_ref(
387 &self,
388 content_ref: &ContentRef,
389 max_bytes: u64,
390 context: &RuntimeReadContext,
391 ) -> Result<Vec<u8>> {
392 let catalog = self.live_catalog(context)?;
393 crate::path::read::ensure_within_read_limit(content_ref.size_bytes, Some(max_bytes))?;
394 let read = crate::storage::content::read_durable_content_bytes(
395 &self.store,
396 catalog.content_store_id(),
397 content_ref,
398 )
399 .await?;
400 Ok(read.bytes)
401 }
402
403 fn live_catalog(&self, context: &RuntimeReadContext) -> Result<VerifiedNamespaceCatalogEntry> {
409 if context.head.namespace_id != self.namespace_id {
410 return Err(crate::error::CoreError::NamespaceCorrupt(format!(
411 "head namespace `{}` does not match requested namespace `{}`",
412 context.head.namespace_id, self.namespace_id
413 )));
414 }
415 if context.head.state == NamespaceState::Deleted {
416 return Err(crate::error::CoreError::NamespaceDeleted {
417 namespace_id: self.namespace_id.clone(),
418 });
419 }
420 Ok(VerifiedNamespaceCatalogEntry::from_head(&context.head))
421 }
422
423 pub async fn list_file_revisions_for_inode_page(
425 &self,
426 inode_id: InodeId,
427 request: PageRequest<FileRevisionsPageCursor>,
428 context: &RuntimeReadContext,
429 ) -> Result<Page<FileRevision, FileRevisionsPageCursor>> {
430 let view = self.load_read_view(context).await?;
431 view.list_file_revisions_for_inode_page(inode_id, request)
432 .await
433 }
434
435 pub async fn read_file_revision(
438 &self,
439 path: impl AsRef<str>,
440 revision_no: RevisionNo,
441 context: &RuntimeReadContext,
442 max_content_bytes: Option<u64>,
443 ) -> Result<AuthoritativeFileBytes> {
444 let view = self.load_read_view(context).await?;
445 view.read_file_revision_bytes(&self.store, path.as_ref(), revision_no, max_content_bytes)
446 .await
447 }
448
449 pub async fn read_file_revision_for_inode(
451 &self,
452 inode_id: InodeId,
453 revision_no: RevisionNo,
454 context: &RuntimeReadContext,
455 max_content_bytes: Option<u64>,
456 ) -> Result<Vec<u8>> {
457 let view = self.load_read_view(context).await?;
458 view.read_file_revision_bytes_for_inode(
459 &self.store,
460 inode_id,
461 revision_no,
462 max_content_bytes,
463 )
464 .await
465 }
466
467 async fn load_read_view<'a>(
468 &'a self,
469 context: &'a RuntimeReadContext,
470 ) -> Result<LoadedMetadataView<'a, S>> {
471 load_metadata_view(
472 &self.store,
473 &self.namespace_id,
474 runtime_read_load_context(context),
475 )
476 .await
477 }
478
479 pub async fn publish_namespace_commits_batch(
482 &self,
483 candidates: Vec<CommitCandidate>,
484 ) -> Vec<Result<CommitResponse>> {
485 let context = match self.mutation_context() {
486 Ok(context) => context,
487 Err(error) => return candidates.iter().map(|_| Err(error.clone())).collect(),
488 };
489 crate::commit_engine::publish_namespace_commits_batch(
490 &self.store,
491 &self.namespace_id,
492 candidates,
493 &context,
494 )
495 .await
496 }
497
498 pub async fn list_changes_after(
500 &self,
501 after_seq: ChangeSeq,
502 limit: EffectiveLimit,
503 ) -> Result<ChangesResponse> {
504 crate::protocol::list_changes_after(&self.store, &self.namespace_id, after_seq, limit).await
505 }
506
507 pub async fn begin_upload(&self, request: BeginUploadRequest) -> Result<BeginUploadResponse> {
509 crate::protocol::begin_upload(
510 &self.store,
511 &self.namespace_id,
512 request,
513 &self.mutation_context()?,
514 )
515 .await
516 }
517
518 pub async fn begin_direct_put_upload_target(
521 &self,
522 claim: DirectPutContentClaim,
523 ) -> Result<BeginDirectPutUploadTargetResponse> {
524 crate::protocol::begin_direct_put_upload_target(
525 &self.store,
526 &self.namespace_id,
527 claim,
528 &self.mutation_context()?,
529 )
530 .await
531 }
532
533 pub async fn begin_direct_multipart_upload_target(
538 &self,
539 options: DirectMultipartUploadOptions,
540 ) -> Result<BeginDirectMultipartUploadTargetResponse> {
541 crate::protocol::begin_direct_multipart_upload_target(
542 &self.store,
543 &self.namespace_id,
544 options,
545 &self.mutation_context()?,
546 )
547 .await
548 }
549
550 pub async fn direct_multipart_part_targets(
553 &self,
554 upload_id: &UploadId,
555 requested: &[UploadPartChecksumClaim],
556 ) -> Result<MultipartPartTargets> {
557 crate::protocol::direct_multipart_part_targets(
558 &self.store,
559 &self.namespace_id,
560 upload_id,
561 requested,
562 )
563 .await
564 }
565
566 pub async fn upload_content(
568 &self,
569 upload_id: &UploadId,
570 bytes: &[u8],
571 ) -> Result<UploadContentResponse> {
572 crate::protocol::upload_content(&self.store, &self.namespace_id, upload_id, bytes).await
573 }
574
575 pub async fn upload_streamed_content(
578 &self,
579 upload_id: &UploadId,
580 body: ByteStream,
581 ) -> Result<UploadContentResponse> {
582 crate::protocol::upload_streamed_content(&self.store, &self.namespace_id, upload_id, body)
583 .await
584 }
585
586 pub async fn complete_upload(
588 &self,
589 upload_id: &UploadId,
590 request: &CompleteUploadRequest,
591 ) -> Result<CompleteUploadResponse> {
592 Ok(self
593 .complete_upload_prepared(upload_id, request)
594 .await?
595 .response)
596 }
597
598 pub async fn complete_upload_prepared(
603 &self,
604 upload_id: &UploadId,
605 request: &CompleteUploadRequest,
606 ) -> Result<CompletedUpload> {
607 let catalog = crate::namespace::catalog::load_namespace_catalog_entry(
608 &self.store,
609 &self.namespace_id,
610 )
611 .await?;
612 self.complete_upload_prepared_with_catalog(&catalog, upload_id, request)
613 .await
614 }
615
616 pub async fn complete_upload_prepared_with_catalog(
619 &self,
620 catalog: &VerifiedNamespaceCatalogEntry,
621 upload_id: &UploadId,
622 request: &CompleteUploadRequest,
623 ) -> Result<CompletedUpload> {
624 let catalog = self.own_catalog(catalog)?;
625 crate::protocol::complete_upload(
626 &self.store,
627 &self.namespace_id,
628 catalog.content_store_id(),
629 upload_id,
630 request,
631 &self.mutation_context()?,
632 )
633 .await
634 }
635
636 pub async fn stage_owned_bytes(
650 &self,
651 catalog: &VerifiedNamespaceCatalogEntry,
652 bytes: &[u8],
653 ) -> Result<PreparedContent> {
654 crate::protocol::stage_owned_bytes(
655 &self.store,
656 self.own_catalog(catalog)?,
657 bytes,
658 &self.mutation_context()?,
659 )
660 .await
661 }
662
663 pub async fn stage_owned_stream(
669 &self,
670 catalog: &VerifiedNamespaceCatalogEntry,
671 body: ByteStream,
672 ) -> Result<PreparedContent> {
673 crate::protocol::stage_owned_stream(
674 &self.store,
675 self.own_catalog(catalog)?,
676 body,
677 &self.mutation_context()?,
678 )
679 .await
680 }
681
682 fn own_catalog<'c>(
689 &self,
690 catalog: &'c VerifiedNamespaceCatalogEntry,
691 ) -> Result<&'c VerifiedNamespaceCatalogEntry> {
692 if catalog.namespace_id() != &self.namespace_id {
693 return Err(CoreError::Internal(format!(
694 "an operation on namespace `{}` was given namespace `{}`'s catalog",
695 self.namespace_id,
696 catalog.namespace_id()
697 )));
698 }
699 Ok(catalog)
700 }
701
702 pub async fn abort_upload(&self, upload_id: &UploadId) -> Result<AbortUploadResponse> {
707 let content_store_id = crate::namespace::catalog::load_namespace_content_store_id(
708 &self.store,
709 &self.namespace_id,
710 )
711 .await?;
712 crate::protocol::abort_upload(
713 &self.store,
714 &self.namespace_id,
715 &content_store_id,
716 upload_id,
717 &self.mutation_context()?,
718 )
719 .await
720 }
721
722 pub async fn read_upload_status(
725 &self,
726 upload_id: &UploadId,
727 ) -> Result<(UploadStatusResponse, Option<CompletedUploadReceipt>)> {
728 let content_store_id = crate::namespace::catalog::load_namespace_content_store_id(
729 &self.store,
730 &self.namespace_id,
731 )
732 .await?;
733 crate::protocol::read_upload_status(
734 &self.store,
735 &self.namespace_id,
736 &content_store_id,
737 upload_id,
738 self.mutation_context()?.now_ms,
739 )
740 .await
741 }
742
743 pub async fn create_checkpoint(
752 &self,
753 name: String,
754 ttl_ms: Option<u64>,
755 ) -> Result<CreateCheckpointResponse> {
756 let context = self.mutation_context()?;
757 let expires_at_ms = ttl_ms.map(|ttl_ms| context.now_ms.saturating_add(ttl_ms));
758 crate::checkpoint::create_checkpoint(
759 &self.store,
760 &self.namespace_id,
761 CheckpointOwner::User { name },
762 expires_at_ms,
763 &context,
764 )
765 .await
766 }
767
768 pub async fn list_checkpoints(&self) -> Result<ListCheckpointsResponse> {
774 crate::checkpoint::list_checkpoints(&self.store, &self.namespace_id).await
775 }
776
777 pub async fn release_checkpoint(
783 &self,
784 checkpoint_id: &CheckpointId,
785 ) -> Result<ReleaseCheckpointResponse> {
786 crate::checkpoint::release_checkpoint(
787 &self.store,
788 &self.namespace_id,
789 checkpoint_id,
790 &self.mutation_context()?,
791 )
792 .await
793 }
794
795 pub async fn flush_wal(&self) -> Result<FlushWalResponse> {
803 crate::checkpoint::flush_wal(&self.store, &self.namespace_id, &self.mutation_context()?)
804 .await
805 }
806
807 pub async fn reorganize_metadata(&self) -> Result<crate::checkpoint::MetadataReorganizeReport> {
814 crate::checkpoint::reorganize_metadata_step(
815 &self.store,
816 &self.namespace_id,
817 &self.mutation_context()?,
818 crate::checkpoint::MetadataLsmPolicy::default(),
819 )
820 .await
821 }
822
823 pub async fn advance_retention_floor(&self) -> Result<AdvanceRetentionResponse> {
825 crate::checkpoint::advance_retention_floor(
826 &self.store,
827 &self.namespace_id,
828 &self.mutation_context()?,
829 )
830 .await
831 }
832
833 fn mutation_context(&self) -> Result<MutationContext> {
840 let writer = self.writer.as_ref().ok_or_else(|| {
841 crate::error::CoreError::Internal(
842 "engine built without writer identity cannot mutate".to_owned(),
843 )
844 })?;
845 Ok(MutationContext {
846 writer_id: writer.writer_id.clone(),
847 now_ms: current_time_ms()?,
848 })
849 }
850}
851
852#[derive(Debug)]
857pub struct NamespaceEngineBuilder<S> {
858 store: S,
859 namespace_id: Option<NamespaceId>,
860 writer_id: Option<String>,
861}
862
863impl<S: ObjectStore> NamespaceEngineBuilder<S> {
864 pub fn namespace_id(mut self, namespace_id: NamespaceId) -> Self {
866 self.namespace_id = Some(namespace_id);
867 self
868 }
869
870 pub fn writer_id(mut self, writer_id: impl Into<String>) -> Self {
872 self.writer_id = Some(writer_id.into());
873 self
874 }
875
876 pub fn build(self) -> std::result::Result<NamespaceEngine<S>, NamespaceEngineBuildError> {
878 let namespace_id = self
879 .namespace_id
880 .ok_or(NamespaceEngineBuildError::MissingNamespace)?;
881 let writer_id = self
882 .writer_id
883 .ok_or(NamespaceEngineBuildError::MissingWriter)?;
884 if writer_id.trim().is_empty() {
885 return Err(NamespaceEngineBuildError::EmptyWriter);
886 }
887
888 Ok(NamespaceEngine {
889 store: self.store,
890 namespace_id,
891 writer: Some(EngineWriter { writer_id }),
892 })
893 }
894
895 pub fn build_reader(
900 self,
901 ) -> std::result::Result<NamespaceEngine<S>, NamespaceEngineBuildError> {
902 let namespace_id = self
903 .namespace_id
904 .ok_or(NamespaceEngineBuildError::MissingNamespace)?;
905 Ok(NamespaceEngine {
906 store: self.store,
907 namespace_id,
908 writer: None,
909 })
910 }
911}
912
913#[derive(Debug, Error, Clone, PartialEq, Eq)]
915pub enum NamespaceEngineBuildError {
916 #[error("namespace is required")]
918 MissingNamespace,
919 #[error("writer identity is required")]
921 MissingWriter,
922 #[error("writer identity must not be empty")]
924 EmptyWriter,
925}
926
927#[cfg(test)]
928mod tests {
929 use super::*;
930 use loonfs_objectstore::local_fs_store::LocalFsStore;
931 use tempfile::tempdir;
932
933 #[test]
934 fn namespace_engine_builds_with_required_identity() {
935 let temp_dir = tempdir().expect("tempdir");
936 let store = LocalFsStore::new(temp_dir.path()).expect("store");
937 let namespace_id = NamespaceId::parse("demo").expect("valid namespace id");
938
939 let engine = NamespaceEngine::builder(store)
940 .namespace_id(namespace_id.clone())
941 .writer_id("writer-a")
942 .build()
943 .expect("engine builds");
944
945 assert_eq!(engine.namespace_id(), &namespace_id);
946 assert_eq!(engine.writer_id(), Some("writer-a"));
947 }
948
949 #[test]
950 fn reader_engine_builds_without_any_writer_identity() {
951 let temp_dir = tempdir().expect("tempdir");
952 let store = LocalFsStore::new(temp_dir.path()).expect("store");
953 let namespace_id = NamespaceId::parse("demo").expect("valid namespace id");
954
955 let engine = NamespaceEngine::builder(store)
956 .namespace_id(namespace_id.clone())
957 .build_reader()
958 .expect("reader engine builds without a writer");
959
960 assert_eq!(engine.namespace_id(), &namespace_id);
961 assert_eq!(engine.writer_id(), None);
962 }
963
964 #[tokio::test]
965 async fn reader_engine_still_serves_reads() {
966 let temp_dir = tempdir().expect("tempdir");
967 let namespace_id = NamespaceId::parse("demo").expect("valid namespace id");
968 NamespaceEngine::builder(LocalFsStore::new(temp_dir.path()).expect("store"))
969 .namespace_id(namespace_id.clone())
970 .writer_id("writer-a")
971 .build()
972 .expect("engine builds")
973 .bootstrap_namespace(BootstrapOptions::default())
974 .await
975 .expect("bootstrap namespace");
976
977 let reader = NamespaceEngine::builder(LocalFsStore::new(temp_dir.path()).expect("store"))
978 .namespace_id(namespace_id.clone())
979 .build_reader()
980 .expect("reader engine builds");
981 let changes = reader
982 .list_changes_after(
983 ChangeSeq(0),
984 loonfs_api::PaginationPolicy::default()
985 .resolve_limit(None)
986 .expect("default limit"),
987 )
988 .await
989 .expect("a reader-built engine serves reads");
990 assert_eq!(changes.namespace_id, namespace_id);
991 }
992
993 #[tokio::test]
994 async fn reader_engine_refuses_to_mutate() {
995 let temp_dir = tempdir().expect("tempdir");
996 let store = LocalFsStore::new(temp_dir.path()).expect("store");
997 let reader = NamespaceEngine::builder(store)
998 .namespace_id(NamespaceId::parse("demo").expect("valid namespace id"))
999 .build_reader()
1000 .expect("reader engine builds");
1001
1002 let error = reader
1003 .flush_wal()
1004 .await
1005 .expect_err("a reader-built engine must refuse mutations");
1006 assert!(
1007 error
1008 .to_string()
1009 .contains("engine built without writer identity cannot mutate"),
1010 "unexpected error: {error}"
1011 );
1012 }
1013
1014 #[test]
1015 fn namespace_engine_builder_rejects_missing_required_fields() {
1016 let temp_dir = tempdir().expect("tempdir");
1017 let store = LocalFsStore::new(temp_dir.path()).expect("store");
1018 let err = NamespaceEngine::builder(store)
1019 .build()
1020 .expect_err("missing namespace");
1021 assert_eq!(err, NamespaceEngineBuildError::MissingNamespace);
1022
1023 let temp_dir = tempdir().expect("tempdir");
1024 let store = LocalFsStore::new(temp_dir.path()).expect("store");
1025 let err = NamespaceEngine::builder(store)
1026 .namespace_id(NamespaceId::parse("demo").expect("valid namespace id"))
1027 .build()
1028 .expect_err("missing writer");
1029 assert_eq!(err, NamespaceEngineBuildError::MissingWriter);
1030 }
1031}