use crate::cache::{MetadataTableCache, WalTailProjectionCache};
use crate::checkpoint::{CheckpointFilesPage, CheckpointFilesPageCursor};
use crate::commit_engine::CommitCandidate;
use crate::context::MutationContext;
use crate::error::{CoreError, Result};
use crate::namespace::basis::MetadataBasis;
use crate::namespace::catalog::VerifiedNamespaceCatalogEntry;
use crate::namespace::{bootstrap, fork, BootstrapNamespaceError};
use crate::options::{BootstrapOptions, DeleteNamespaceOptions};
use crate::path::read::{
load_metadata_view, CurrentFileState, DirectDownloadTarget, LoadedMetadataView, ReadLoadContext,
};
use crate::protocol::CompletedUpload;
use crate::storage::content::FileContentStream;
use crate::storage::content_admission::{CompletedUploadReceipt, PreparedContent};
use crate::time::current_time_ms;
use loonfs_api::v0::{
AbortUploadResponse, BeginUploadRequest, BeginUploadResponse, ChangesResponse, CommitResponse,
CompleteUploadRequest, CompleteUploadResponse, DirectMultipartUploadOptions,
DirectPutContentClaim, UploadContentResponse, UploadPartChecksumClaim, UploadStatusResponse,
};
use loonfs_api::wire::control::{CheckpointOwner, HeadState, NamespaceState};
use loonfs_api::EffectiveLimit;
use loonfs_api::{
AdvanceRetentionResponse, AuthoritativeFileBytes, AuthoritativePathEntry, ChangeSeq,
CheckpointId, ContentRef, CreateCheckpointResponse, DeleteNamespaceResponse,
DirectoryPageCursor, FileRevision, FileRevisionsPageCursor, FlushWalResponse, InodeId,
ListCheckpointsResponse, NamespaceId, NamespaceSummary, Page, PageRequest,
ReleaseCheckpointResponse, RevisionNo, StorageChecksum, TrashEntry, TrashPageCursor, UploadId,
};
use loonfs_objectstore::{ByteStream, ObjectStore};
use std::num::NonZeroU64;
use std::sync::Arc;
use thiserror::Error;
#[derive(Debug, Clone)]
pub struct RuntimeReadContext {
pub head: HeadState,
pub head_etag: String,
pub basis: MetadataBasis,
pub table_cache: Arc<MetadataTableCache>,
pub tail_cache: Arc<WalTailProjectionCache>,
}
fn runtime_read_load_context(context: &RuntimeReadContext) -> ReadLoadContext<'_> {
ReadLoadContext::pinned_head(
&context.head,
Some(context.head_etag.as_str()),
&context.basis,
Some(&context.table_cache),
Some(&context.tail_cache),
)
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct DirectPutUploadTarget {
pub content_ref: ContentRef,
pub object_key: String,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct BeginDirectPutUploadTargetResponse {
pub namespace_id: NamespaceId,
pub upload_id: UploadId,
pub target: DirectPutUploadTarget,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct DirectMultipartUploadTarget {
pub object_key: String,
pub part_size_bytes: u64,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct BeginDirectMultipartUploadTargetResponse {
pub namespace_id: NamespaceId,
pub upload_id: UploadId,
pub target: DirectMultipartUploadTarget,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct MultipartPartTarget {
pub part_number: u32,
pub checksum: StorageChecksum,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct MultipartPartTargets {
pub object_key: String,
pub provider_upload_id: String,
pub parts: Vec<MultipartPartTarget>,
}
#[derive(Debug)]
struct EngineWriter {
writer_id: String,
}
#[derive(Debug)]
pub struct NamespaceEngine<S> {
store: S,
namespace_id: NamespaceId,
writer: Option<EngineWriter>,
}
impl<S: ObjectStore> NamespaceEngine<S> {
pub fn builder(store: S) -> NamespaceEngineBuilder<S> {
NamespaceEngineBuilder {
store,
namespace_id: None,
writer_id: None,
}
}
pub fn namespace_id(&self) -> &NamespaceId {
&self.namespace_id
}
pub fn writer_id(&self) -> Option<&str> {
self.writer.as_ref().map(|writer| writer.writer_id.as_str())
}
pub async fn bootstrap_namespace(
&self,
options: BootstrapOptions,
) -> std::result::Result<NamespaceSummary, BootstrapNamespaceError> {
bootstrap::bootstrap_namespace(
&self.store,
&self.namespace_id,
&self.mutation_context()?,
options.allow_existing,
)
.await
}
pub async fn fork_namespace(&self, target: &NamespaceId) -> Result<NamespaceSummary> {
fork::fork_namespace(
&self.store,
&self.namespace_id,
target,
&self.mutation_context()?,
)
.await
}
pub async fn delete_namespace(
&self,
options: DeleteNamespaceOptions,
) -> Result<DeleteNamespaceResponse> {
crate::commit_engine::delete_namespace(
&self.store,
&self.namespace_id,
options,
&self.mutation_context()?,
)
.await
}
pub async fn resolve_path(
&self,
path: impl AsRef<str>,
context: &RuntimeReadContext,
) -> Result<AuthoritativePathEntry> {
let view = self.load_read_view(context).await?;
view.resolve_path(path.as_ref()).await
}
pub async fn list_path_page(
&self,
path: impl AsRef<str>,
request: PageRequest<DirectoryPageCursor>,
context: &RuntimeReadContext,
) -> Result<Page<AuthoritativePathEntry, DirectoryPageCursor>> {
let view = self.load_read_view(context).await?;
view.list_path_page(path.as_ref(), request).await
}
pub async fn read_file(
&self,
path: impl AsRef<str>,
context: &RuntimeReadContext,
max_content_bytes: Option<u64>,
) -> Result<AuthoritativeFileBytes> {
let view = self.load_read_view(context).await?;
view.read_file_bytes(&self.store, path.as_ref(), max_content_bytes)
.await
}
pub async fn read_file_stream(
&self,
path: impl AsRef<str>,
context: &RuntimeReadContext,
chunk_bytes: NonZeroU64,
start_offset: u64,
) -> Result<FileContentStream<S>>
where
S: Clone,
{
let view = self.load_read_view(context).await?;
let (entry, content_ref) = view.resolve_file_content(path.as_ref()).await?;
if start_offset > content_ref.size_bytes {
return Err(CoreError::ResumeOffsetOutOfRange {
start_offset,
size_bytes: content_ref.size_bytes,
});
}
Ok(FileContentStream::open(
self.store.clone(),
view.content_store_id(),
entry,
content_ref,
chunk_bytes,
start_offset,
)
.await?)
}
pub async fn direct_download_target(
&self,
path: impl AsRef<str>,
revision_no: Option<RevisionNo>,
context: &RuntimeReadContext,
) -> Result<DirectDownloadTarget> {
let view = self.load_read_view(context).await?;
view.direct_download_target(path.as_ref(), revision_no)
.await
}
pub async fn list_file_revisions_page(
&self,
path: impl AsRef<str>,
request: PageRequest<FileRevisionsPageCursor>,
context: &RuntimeReadContext,
) -> Result<Page<FileRevision, FileRevisionsPageCursor>> {
let view = self.load_read_view(context).await?;
view.list_file_revisions_page(path.as_ref(), request).await
}
pub async fn list_trash_page(
&self,
request: PageRequest<TrashPageCursor>,
context: &RuntimeReadContext,
) -> Result<Page<TrashEntry, TrashPageCursor>> {
let view = self.load_read_view(context).await?;
view.list_trash_page(request).await
}
pub async fn list_checkpoint_files_page(
&self,
checkpoint_id: &CheckpointId,
request: PageRequest<CheckpointFilesPageCursor>,
context: &RuntimeReadContext,
) -> Result<CheckpointFilesPage> {
self.live_catalog(context)?;
crate::checkpoint::list_checkpoint_files_page(
&self.store,
Some(context.table_cache.as_ref()),
&self.namespace_id,
checkpoint_id,
request,
)
.await
}
pub async fn resolve_current_files(
&self,
inode_ids: &[InodeId],
context: &RuntimeReadContext,
) -> Result<Vec<CurrentFileState>> {
crate::path::read::ensure_resolve_batch_within_cap(inode_ids.len())?;
let view = self.load_read_view(context).await?;
crate::path::read::resolve_current_files(&view, inode_ids).await
}
pub async fn read_content_ref(
&self,
content_ref: &ContentRef,
max_bytes: u64,
context: &RuntimeReadContext,
) -> Result<Vec<u8>> {
let catalog = self.live_catalog(context)?;
crate::path::read::ensure_within_read_limit(content_ref.size_bytes, Some(max_bytes))?;
let read = crate::storage::content::read_durable_content_bytes(
&self.store,
catalog.content_store_id(),
content_ref,
)
.await?;
Ok(read.bytes)
}
fn live_catalog(&self, context: &RuntimeReadContext) -> Result<VerifiedNamespaceCatalogEntry> {
if context.head.namespace_id != self.namespace_id {
return Err(crate::error::CoreError::NamespaceCorrupt(format!(
"head namespace `{}` does not match requested namespace `{}`",
context.head.namespace_id, self.namespace_id
)));
}
if context.head.state == NamespaceState::Deleted {
return Err(crate::error::CoreError::NamespaceDeleted {
namespace_id: self.namespace_id.clone(),
});
}
Ok(VerifiedNamespaceCatalogEntry::from_head(&context.head))
}
pub async fn list_file_revisions_for_inode_page(
&self,
inode_id: InodeId,
request: PageRequest<FileRevisionsPageCursor>,
context: &RuntimeReadContext,
) -> Result<Page<FileRevision, FileRevisionsPageCursor>> {
let view = self.load_read_view(context).await?;
view.list_file_revisions_for_inode_page(inode_id, request)
.await
}
pub async fn read_file_revision(
&self,
path: impl AsRef<str>,
revision_no: RevisionNo,
context: &RuntimeReadContext,
max_content_bytes: Option<u64>,
) -> Result<AuthoritativeFileBytes> {
let view = self.load_read_view(context).await?;
view.read_file_revision_bytes(&self.store, path.as_ref(), revision_no, max_content_bytes)
.await
}
pub async fn read_file_revision_for_inode(
&self,
inode_id: InodeId,
revision_no: RevisionNo,
context: &RuntimeReadContext,
max_content_bytes: Option<u64>,
) -> Result<Vec<u8>> {
let view = self.load_read_view(context).await?;
view.read_file_revision_bytes_for_inode(
&self.store,
inode_id,
revision_no,
max_content_bytes,
)
.await
}
async fn load_read_view<'a>(
&'a self,
context: &'a RuntimeReadContext,
) -> Result<LoadedMetadataView<'a, S>> {
load_metadata_view(
&self.store,
&self.namespace_id,
runtime_read_load_context(context),
)
.await
}
pub async fn publish_namespace_commits_batch(
&self,
candidates: Vec<CommitCandidate>,
) -> Vec<Result<CommitResponse>> {
let context = match self.mutation_context() {
Ok(context) => context,
Err(error) => return candidates.iter().map(|_| Err(error.clone())).collect(),
};
crate::commit_engine::publish_namespace_commits_batch(
&self.store,
&self.namespace_id,
candidates,
&context,
)
.await
}
pub async fn list_changes_after(
&self,
after_seq: ChangeSeq,
limit: EffectiveLimit,
) -> Result<ChangesResponse> {
crate::protocol::list_changes_after(&self.store, &self.namespace_id, after_seq, limit).await
}
pub async fn begin_upload(&self, request: BeginUploadRequest) -> Result<BeginUploadResponse> {
crate::protocol::begin_upload(
&self.store,
&self.namespace_id,
request,
&self.mutation_context()?,
)
.await
}
pub async fn begin_direct_put_upload_target(
&self,
claim: DirectPutContentClaim,
) -> Result<BeginDirectPutUploadTargetResponse> {
crate::protocol::begin_direct_put_upload_target(
&self.store,
&self.namespace_id,
claim,
&self.mutation_context()?,
)
.await
}
pub async fn begin_direct_multipart_upload_target(
&self,
options: DirectMultipartUploadOptions,
) -> Result<BeginDirectMultipartUploadTargetResponse> {
crate::protocol::begin_direct_multipart_upload_target(
&self.store,
&self.namespace_id,
options,
&self.mutation_context()?,
)
.await
}
pub async fn direct_multipart_part_targets(
&self,
upload_id: &UploadId,
requested: &[UploadPartChecksumClaim],
) -> Result<MultipartPartTargets> {
crate::protocol::direct_multipart_part_targets(
&self.store,
&self.namespace_id,
upload_id,
requested,
)
.await
}
pub async fn upload_content(
&self,
upload_id: &UploadId,
bytes: &[u8],
) -> Result<UploadContentResponse> {
crate::protocol::upload_content(&self.store, &self.namespace_id, upload_id, bytes).await
}
pub async fn upload_streamed_content(
&self,
upload_id: &UploadId,
body: ByteStream,
) -> Result<UploadContentResponse> {
crate::protocol::upload_streamed_content(&self.store, &self.namespace_id, upload_id, body)
.await
}
pub async fn complete_upload(
&self,
upload_id: &UploadId,
request: &CompleteUploadRequest,
) -> Result<CompleteUploadResponse> {
Ok(self
.complete_upload_prepared(upload_id, request)
.await?
.response)
}
pub async fn complete_upload_prepared(
&self,
upload_id: &UploadId,
request: &CompleteUploadRequest,
) -> Result<CompletedUpload> {
let catalog = crate::namespace::catalog::load_namespace_catalog_entry(
&self.store,
&self.namespace_id,
)
.await?;
self.complete_upload_prepared_with_catalog(&catalog, upload_id, request)
.await
}
pub async fn complete_upload_prepared_with_catalog(
&self,
catalog: &VerifiedNamespaceCatalogEntry,
upload_id: &UploadId,
request: &CompleteUploadRequest,
) -> Result<CompletedUpload> {
let catalog = self.own_catalog(catalog)?;
crate::protocol::complete_upload(
&self.store,
&self.namespace_id,
catalog.content_store_id(),
upload_id,
request,
&self.mutation_context()?,
)
.await
}
pub async fn stage_owned_bytes(
&self,
catalog: &VerifiedNamespaceCatalogEntry,
bytes: &[u8],
) -> Result<PreparedContent> {
crate::protocol::stage_owned_bytes(
&self.store,
self.own_catalog(catalog)?,
bytes,
&self.mutation_context()?,
)
.await
}
pub async fn stage_owned_stream(
&self,
catalog: &VerifiedNamespaceCatalogEntry,
body: ByteStream,
) -> Result<PreparedContent> {
crate::protocol::stage_owned_stream(
&self.store,
self.own_catalog(catalog)?,
body,
&self.mutation_context()?,
)
.await
}
fn own_catalog<'c>(
&self,
catalog: &'c VerifiedNamespaceCatalogEntry,
) -> Result<&'c VerifiedNamespaceCatalogEntry> {
if catalog.namespace_id() != &self.namespace_id {
return Err(CoreError::Internal(format!(
"an operation on namespace `{}` was given namespace `{}`'s catalog",
self.namespace_id,
catalog.namespace_id()
)));
}
Ok(catalog)
}
pub async fn abort_upload(&self, upload_id: &UploadId) -> Result<AbortUploadResponse> {
let content_store_id = crate::namespace::catalog::load_namespace_content_store_id(
&self.store,
&self.namespace_id,
)
.await?;
crate::protocol::abort_upload(
&self.store,
&self.namespace_id,
&content_store_id,
upload_id,
&self.mutation_context()?,
)
.await
}
pub async fn read_upload_status(
&self,
upload_id: &UploadId,
) -> Result<(UploadStatusResponse, Option<CompletedUploadReceipt>)> {
let content_store_id = crate::namespace::catalog::load_namespace_content_store_id(
&self.store,
&self.namespace_id,
)
.await?;
crate::protocol::read_upload_status(
&self.store,
&self.namespace_id,
&content_store_id,
upload_id,
self.mutation_context()?.now_ms,
)
.await
}
pub async fn create_checkpoint(
&self,
name: String,
ttl_ms: Option<u64>,
) -> Result<CreateCheckpointResponse> {
let context = self.mutation_context()?;
let expires_at_ms = ttl_ms.map(|ttl_ms| context.now_ms.saturating_add(ttl_ms));
crate::checkpoint::create_checkpoint(
&self.store,
&self.namespace_id,
CheckpointOwner::User { name },
expires_at_ms,
&context,
)
.await
}
pub async fn list_checkpoints(&self) -> Result<ListCheckpointsResponse> {
crate::checkpoint::list_checkpoints(&self.store, &self.namespace_id).await
}
pub async fn release_checkpoint(
&self,
checkpoint_id: &CheckpointId,
) -> Result<ReleaseCheckpointResponse> {
crate::checkpoint::release_checkpoint(
&self.store,
&self.namespace_id,
checkpoint_id,
&self.mutation_context()?,
)
.await
}
pub async fn flush_wal(&self) -> Result<FlushWalResponse> {
crate::checkpoint::flush_wal(&self.store, &self.namespace_id, &self.mutation_context()?)
.await
}
pub async fn reorganize_metadata(&self) -> Result<crate::checkpoint::MetadataReorganizeReport> {
crate::checkpoint::reorganize_metadata_step(
&self.store,
&self.namespace_id,
&self.mutation_context()?,
crate::checkpoint::MetadataLsmPolicy::default(),
)
.await
}
pub async fn advance_retention_floor(&self) -> Result<AdvanceRetentionResponse> {
crate::checkpoint::advance_retention_floor(
&self.store,
&self.namespace_id,
&self.mutation_context()?,
)
.await
}
fn mutation_context(&self) -> Result<MutationContext> {
let writer = self.writer.as_ref().ok_or_else(|| {
crate::error::CoreError::Internal(
"engine built without writer identity cannot mutate".to_owned(),
)
})?;
Ok(MutationContext {
writer_id: writer.writer_id.clone(),
now_ms: current_time_ms()?,
})
}
}
#[derive(Debug)]
pub struct NamespaceEngineBuilder<S> {
store: S,
namespace_id: Option<NamespaceId>,
writer_id: Option<String>,
}
impl<S: ObjectStore> NamespaceEngineBuilder<S> {
pub fn namespace_id(mut self, namespace_id: NamespaceId) -> Self {
self.namespace_id = Some(namespace_id);
self
}
pub fn writer_id(mut self, writer_id: impl Into<String>) -> Self {
self.writer_id = Some(writer_id.into());
self
}
pub fn build(self) -> std::result::Result<NamespaceEngine<S>, NamespaceEngineBuildError> {
let namespace_id = self
.namespace_id
.ok_or(NamespaceEngineBuildError::MissingNamespace)?;
let writer_id = self
.writer_id
.ok_or(NamespaceEngineBuildError::MissingWriter)?;
if writer_id.trim().is_empty() {
return Err(NamespaceEngineBuildError::EmptyWriter);
}
Ok(NamespaceEngine {
store: self.store,
namespace_id,
writer: Some(EngineWriter { writer_id }),
})
}
pub fn build_reader(
self,
) -> std::result::Result<NamespaceEngine<S>, NamespaceEngineBuildError> {
let namespace_id = self
.namespace_id
.ok_or(NamespaceEngineBuildError::MissingNamespace)?;
Ok(NamespaceEngine {
store: self.store,
namespace_id,
writer: None,
})
}
}
#[derive(Debug, Error, Clone, PartialEq, Eq)]
pub enum NamespaceEngineBuildError {
#[error("namespace is required")]
MissingNamespace,
#[error("writer identity is required")]
MissingWriter,
#[error("writer identity must not be empty")]
EmptyWriter,
}
#[cfg(test)]
mod tests {
use super::*;
use loonfs_objectstore::local_fs_store::LocalFsStore;
use tempfile::tempdir;
#[test]
fn namespace_engine_builds_with_required_identity() {
let temp_dir = tempdir().expect("tempdir");
let store = LocalFsStore::new(temp_dir.path()).expect("store");
let namespace_id = NamespaceId::parse("demo").expect("valid namespace id");
let engine = NamespaceEngine::builder(store)
.namespace_id(namespace_id.clone())
.writer_id("writer-a")
.build()
.expect("engine builds");
assert_eq!(engine.namespace_id(), &namespace_id);
assert_eq!(engine.writer_id(), Some("writer-a"));
}
#[test]
fn reader_engine_builds_without_any_writer_identity() {
let temp_dir = tempdir().expect("tempdir");
let store = LocalFsStore::new(temp_dir.path()).expect("store");
let namespace_id = NamespaceId::parse("demo").expect("valid namespace id");
let engine = NamespaceEngine::builder(store)
.namespace_id(namespace_id.clone())
.build_reader()
.expect("reader engine builds without a writer");
assert_eq!(engine.namespace_id(), &namespace_id);
assert_eq!(engine.writer_id(), None);
}
#[tokio::test]
async fn reader_engine_still_serves_reads() {
let temp_dir = tempdir().expect("tempdir");
let namespace_id = NamespaceId::parse("demo").expect("valid namespace id");
NamespaceEngine::builder(LocalFsStore::new(temp_dir.path()).expect("store"))
.namespace_id(namespace_id.clone())
.writer_id("writer-a")
.build()
.expect("engine builds")
.bootstrap_namespace(BootstrapOptions::default())
.await
.expect("bootstrap namespace");
let reader = NamespaceEngine::builder(LocalFsStore::new(temp_dir.path()).expect("store"))
.namespace_id(namespace_id.clone())
.build_reader()
.expect("reader engine builds");
let changes = reader
.list_changes_after(
ChangeSeq(0),
loonfs_api::PaginationPolicy::default()
.resolve_limit(None)
.expect("default limit"),
)
.await
.expect("a reader-built engine serves reads");
assert_eq!(changes.namespace_id, namespace_id);
}
#[tokio::test]
async fn reader_engine_refuses_to_mutate() {
let temp_dir = tempdir().expect("tempdir");
let store = LocalFsStore::new(temp_dir.path()).expect("store");
let reader = NamespaceEngine::builder(store)
.namespace_id(NamespaceId::parse("demo").expect("valid namespace id"))
.build_reader()
.expect("reader engine builds");
let error = reader
.flush_wal()
.await
.expect_err("a reader-built engine must refuse mutations");
assert!(
error
.to_string()
.contains("engine built without writer identity cannot mutate"),
"unexpected error: {error}"
);
}
#[test]
fn namespace_engine_builder_rejects_missing_required_fields() {
let temp_dir = tempdir().expect("tempdir");
let store = LocalFsStore::new(temp_dir.path()).expect("store");
let err = NamespaceEngine::builder(store)
.build()
.expect_err("missing namespace");
assert_eq!(err, NamespaceEngineBuildError::MissingNamespace);
let temp_dir = tempdir().expect("tempdir");
let store = LocalFsStore::new(temp_dir.path()).expect("store");
let err = NamespaceEngine::builder(store)
.namespace_id(NamespaceId::parse("demo").expect("valid namespace id"))
.build()
.expect_err("missing writer");
assert_eq!(err, NamespaceEngineBuildError::MissingWriter);
}
}