feat(mounts): external file mounts P1 — pluggable provider + read-only REST
Adds the foundation for external file mounts: admin-configured backends (raw host filesystem in v1; sftp/webdav/… as future provider kinds) surfaced as a folder inside a user's drive. Mount contents are virtual/live-passthrough — read straight from the backend, never stored in storage.files — and are a deliberately separate, limited storage type (no dedup/sharing/trash/search). The feature is dark by default (OXICLOUD_ENABLE_EXTERNAL_MOUNTS=false). P1 scope (this PR): data model, the pluggable provider abstraction, and the read-only REST surface (mount listing + download). Read-write (P2), WebDAV/NextCloud path resolution (P3), and the admin UI (P4) follow. Core model - Mount root = a real storage.folders row; authorization for everything inside collapses onto that folder UUID (ltree-ancestry grant cascade). - Children are virtual, addressed by ext:<mount_id>:<base64url(node_id)> where node_id is provider-owned and opaque to the rest of the system. - A lock-free (arc-swap) MountRegistry maps mount-root UUID -> provider; a thin MountRouter::classify() is the single cheap hook handlers call before parsing an id as a UUID. With no mounts configured it always returns Regular, so existing code paths are unchanged. Added - migrations/20260805000000_external_mounts.sql (storage.external_mounts, kind + config JSONB) - domain/services/external_mount_id (id envelope + virtual etags) - application/ports/external_mount_ports (ExternalMountProvider, MountProviderFactory, repo port) - infrastructure local_fs_mount_provider (tokio::fs, symlink-escape-safe) + factory - application MountRegistry + MountRouter, pg ExternalMountRepository - DI wiring (AppState.mount_router), FeaturesConfig.enable_external_mounts - listing branch (FolderService::list_mount_dir_with_perms + folder_handler) and download branch (FileRetrievalService stat/open mount methods + file_handler) Authorization stays in the service layer (authz.require(Resource::Folder(mount_id))); handlers only classify. Cross-backend operations are out of scope for P1. Tests: 529 unit tests + 5 testcontainers integration tests (real Postgres 17), including end-to-end authorization (owner allowed, stranger denied). Line coverage of the new modules is 84–100% (cargo-llvm-cov). Known gap: file_handler::download_mount_file (HTTP glue) needs a full-app test (P4).
This commit is contained in:
@@ -0,0 +1,192 @@
|
||||
//! External mount provider port — the pluggable backend abstraction for mounts.
|
||||
//!
|
||||
//! An [`ExternalMountProvider`] exposes a filesystem-style I/O surface for one
|
||||
//! mount's backend (raw host fs in v1; SFTP/WebDAV/… later). It is the *lowest
|
||||
//! common denominator* of browse + CRUD: deliberately small, so new backend
|
||||
//! `kind`s drop in without touching the router, listing, authz, or path
|
||||
//! resolution. Rich native features (sharing, trash, search, …) are NOT part of
|
||||
//! this trait — they compose above it.
|
||||
//!
|
||||
//! Each provider instance is **bound to one mount's root location at
|
||||
//! construction**, so methods take only a provider-owned [`NodeId`], never a host
|
||||
//! `Path` (an SFTP/WebDAV provider has no local path). The `NodeId` is opaque to
|
||||
//! the rest of the system (see [`crate::domain::services::external_mount_id`]).
|
||||
//!
|
||||
//! The trait returns boxed futures via `#[async_trait]` and takes a boxed write
|
||||
//! stream, so it is dyn-compatible (`Arc<dyn ExternalMountProvider>`) and a
|
||||
//! single mount registry can hold providers of different kinds.
|
||||
|
||||
use std::sync::Arc;
|
||||
|
||||
use async_trait::async_trait;
|
||||
use bytes::Bytes;
|
||||
use futures::Stream;
|
||||
use std::pin::Pin;
|
||||
use uuid::Uuid;
|
||||
|
||||
use crate::application::ports::blob_storage_ports::BlobStream;
|
||||
use crate::domain::errors::DomainError;
|
||||
use crate::domain::services::external_mount_id::NodeId;
|
||||
|
||||
/// A byte stream handed to [`ExternalMountProvider::write_stream`].
|
||||
///
|
||||
/// Boxed (not generic) so the trait stays object-safe. Callers map their body's
|
||||
/// error type to `std::io::Error` before constructing it.
|
||||
pub type MountByteStream = Pin<Box<dyn Stream<Item = Result<Bytes, std::io::Error>> + Send>>;
|
||||
|
||||
/// One entry returned by [`ExternalMountProvider::list_dir`].
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct MountEntry {
|
||||
/// Final path segment (display name).
|
||||
pub name: String,
|
||||
/// Provider-assigned, opaque identity for this entry.
|
||||
pub node_id: NodeId,
|
||||
/// Whether the entry is a directory.
|
||||
pub is_dir: bool,
|
||||
/// Size in bytes (0 for directories).
|
||||
pub size: u64,
|
||||
/// Last-modified time, unix seconds.
|
||||
pub modified_at: u64,
|
||||
/// Creation time, unix seconds (falls back to `modified_at` when unavailable).
|
||||
pub created_at: u64,
|
||||
}
|
||||
|
||||
/// Metadata for a single entry returned by [`ExternalMountProvider::stat`] and
|
||||
/// by the mutating ops (so the caller learns the new entry's `node_id`).
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct MountStat {
|
||||
/// Provider-assigned, opaque identity for this entry.
|
||||
pub node_id: NodeId,
|
||||
/// Whether the entry is a directory.
|
||||
pub is_dir: bool,
|
||||
/// Size in bytes (0 for directories).
|
||||
pub size: u64,
|
||||
/// Last-modified time, unix seconds.
|
||||
pub modified_at: u64,
|
||||
/// Creation time, unix seconds.
|
||||
pub created_at: u64,
|
||||
/// MIME type (sniffed from extension for files; `"directory"` for dirs).
|
||||
pub mime_type: String,
|
||||
}
|
||||
|
||||
/// Static capability flags a provider advertises.
|
||||
#[derive(Debug, Clone, Copy)]
|
||||
pub struct MountCaps {
|
||||
/// Provider can serve byte ranges (HTTP Range / partial reads).
|
||||
pub supports_range: bool,
|
||||
/// Provider refuses all mutations.
|
||||
pub read_only: bool,
|
||||
/// `node_id`s are stable across renames/moves (e.g. inode / object id).
|
||||
/// `false` for path-based providers — relevant to the future sharing path.
|
||||
pub stable_ids: bool,
|
||||
}
|
||||
|
||||
/// Pluggable I/O surface for one mount's backend, bound to its root location.
|
||||
///
|
||||
/// Implementations: `LocalFsMountProvider` (v1). All node ids are
|
||||
/// provider-owned and opaque; the system never parses them.
|
||||
#[async_trait]
|
||||
pub trait ExternalMountProvider: Send + Sync + 'static {
|
||||
/// Provider kind identifier (matches the `kind` column / factory arm).
|
||||
fn kind(&self) -> &'static str;
|
||||
|
||||
/// Static capabilities.
|
||||
fn capabilities(&self) -> MountCaps;
|
||||
|
||||
/// Map an internal path (relative to the mount root) to a `node_id`.
|
||||
///
|
||||
/// For path-based providers this is identity (the default). Providers whose
|
||||
/// identity is not a path override this. Does not assert existence — use
|
||||
/// [`stat`](Self::stat) for that.
|
||||
fn resolve_path(&self, path: &str) -> NodeId {
|
||||
NodeId(path.to_string())
|
||||
}
|
||||
|
||||
/// List the directory identified by `node_id` (root = the provider's bound
|
||||
/// location, addressed via `resolve_path("")`).
|
||||
async fn list_dir(&self, node_id: &NodeId) -> Result<Vec<MountEntry>, DomainError>;
|
||||
|
||||
/// Stat a single entry.
|
||||
async fn stat(&self, node_id: &NodeId) -> Result<MountStat, DomainError>;
|
||||
|
||||
/// Open a (optionally ranged) read stream over a file's bytes.
|
||||
///
|
||||
/// `range` is `(start, end_inclusive_opt)`; `None` reads the whole file.
|
||||
async fn open_read_stream(
|
||||
&self,
|
||||
node_id: &NodeId,
|
||||
range: Option<(u64, Option<u64>)>,
|
||||
) -> Result<BlobStream, DomainError>;
|
||||
|
||||
/// Create a child directory `name` under `parent`. Returns the new dir's stat.
|
||||
async fn create_dir(&self, parent: &NodeId, name: &str) -> Result<MountStat, DomainError>;
|
||||
|
||||
/// Stream-write a child file `name` under `parent`. Returns the new file's stat.
|
||||
async fn write_stream(
|
||||
&self,
|
||||
parent: &NodeId,
|
||||
name: &str,
|
||||
body: MountByteStream,
|
||||
) -> Result<MountStat, DomainError>;
|
||||
|
||||
/// Rename an entry in place (same parent). Returns the renamed entry's stat.
|
||||
async fn rename(&self, node_id: &NodeId, new_name: &str) -> Result<MountStat, DomainError>;
|
||||
|
||||
/// Delete an entry (recursively for directories). Permanent — no trash.
|
||||
async fn delete(&self, node_id: &NodeId) -> Result<(), DomainError>;
|
||||
|
||||
/// Move an entry into `dest_parent`, keeping its name. Returns the new stat.
|
||||
async fn move_within(
|
||||
&self,
|
||||
node_id: &NodeId,
|
||||
dest_parent: &NodeId,
|
||||
) -> Result<MountStat, DomainError>;
|
||||
}
|
||||
|
||||
/// A persisted external mount joined with its mount-root folder.
|
||||
///
|
||||
/// Returned by [`ExternalMountRepositoryPort::list_all`] to (re)build the
|
||||
/// in-memory registry.
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct ExternalMountRecord {
|
||||
/// The mount-root folder UUID (also the mount's identity in the registry).
|
||||
pub mount_folder_id: Uuid,
|
||||
/// Provider kind (factory discriminator).
|
||||
pub kind: String,
|
||||
/// Provider-specific connection config.
|
||||
pub config: serde_json::Value,
|
||||
/// Display name.
|
||||
pub name: String,
|
||||
/// Owner of the mount configuration.
|
||||
pub owner_id: Uuid,
|
||||
/// Whether the mount is read-only.
|
||||
pub read_only: bool,
|
||||
/// Drive the mount-root folder belongs to (for path resolution).
|
||||
pub drive_id: Uuid,
|
||||
/// Materialized internal path of the mount-root folder (for path resolution).
|
||||
pub mount_path: String,
|
||||
}
|
||||
|
||||
/// Persistence port for external mount configuration.
|
||||
#[async_trait]
|
||||
pub trait ExternalMountRepositoryPort: Send + Sync {
|
||||
/// Load every (non-trashed) mount joined with its folder, for registry build.
|
||||
async fn list_all(&self) -> Result<Vec<ExternalMountRecord>, DomainError>;
|
||||
}
|
||||
|
||||
/// Builds [`ExternalMountProvider`]s from a `kind` + `config` pair.
|
||||
///
|
||||
/// The single extension point for new backends: adding a provider is
|
||||
/// implementing the trait plus one arm here.
|
||||
#[async_trait]
|
||||
pub trait MountProviderFactory: Send + Sync {
|
||||
/// Construct a provider for `kind`, parsing its `config` JSON.
|
||||
///
|
||||
/// Errors with `UnsupportedOperation` for an unknown kind, or
|
||||
/// `validation_error` for malformed config.
|
||||
async fn build(
|
||||
&self,
|
||||
kind: &str,
|
||||
config: &serde_json::Value,
|
||||
) -> Result<Arc<dyn ExternalMountProvider>, DomainError>;
|
||||
}
|
||||
@@ -10,6 +10,7 @@ pub mod compression_ports;
|
||||
pub mod content_index_ports;
|
||||
pub mod dedup_ports;
|
||||
pub mod email_sender;
|
||||
pub mod external_mount_ports;
|
||||
pub mod face_ports;
|
||||
pub mod favorites_ports;
|
||||
pub mod file_lifecycle;
|
||||
|
||||
@@ -100,7 +100,12 @@ mod tests {
|
||||
None,
|
||||
authz.clone(),
|
||||
));
|
||||
let folder_service = Arc::new(FolderService::new(folder_repo, authz));
|
||||
let mount_router = Arc::new(
|
||||
crate::application::services::external_mount_router::MountRouter::new(Arc::new(
|
||||
crate::application::services::mount_registry::MountRegistry::empty(),
|
||||
)),
|
||||
);
|
||||
let folder_service = Arc::new(FolderService::new(folder_repo, authz, mount_router));
|
||||
|
||||
let _batch_service = BatchOperationService::new(
|
||||
file_retrieval,
|
||||
|
||||
@@ -0,0 +1,181 @@
|
||||
//! Classifies file/folder ids into native vs. external-mount handling.
|
||||
//!
|
||||
//! This is the single, cheap hook the service layer calls before any
|
||||
//! `Uuid::parse_str`, so synthetic `ext:` ids and mount-root UUIDs branch to the
|
||||
//! provider while everything else flows to the PostgreSQL repositories unchanged.
|
||||
|
||||
use std::sync::Arc;
|
||||
|
||||
use uuid::Uuid;
|
||||
|
||||
use crate::application::services::mount_registry::{MountConfig, MountRegistry};
|
||||
use crate::domain::services::external_mount_id::{NodeId, is_external_id, parse_child_id};
|
||||
|
||||
/// The result of classifying an id.
|
||||
pub enum ResolvedId {
|
||||
/// Plain native resource (UUID not registered as a mount root, or an
|
||||
/// unrecognized id). Handle exactly as today.
|
||||
Regular,
|
||||
/// A real UUID that IS a mount root. Listing/metadata branch to the provider;
|
||||
/// the row itself still exists natively.
|
||||
MountRoot { cfg: Arc<MountConfig> },
|
||||
/// A synthetic id addressing an entry inside a mount.
|
||||
MountChild {
|
||||
cfg: Arc<MountConfig>,
|
||||
node_id: NodeId,
|
||||
},
|
||||
}
|
||||
|
||||
/// Thin, cloneable classifier over the mount registry.
|
||||
#[derive(Clone)]
|
||||
pub struct MountRouter {
|
||||
registry: Arc<MountRegistry>,
|
||||
}
|
||||
|
||||
impl MountRouter {
|
||||
/// Construct from the shared registry.
|
||||
pub fn new(registry: Arc<MountRegistry>) -> Self {
|
||||
Self { registry }
|
||||
}
|
||||
|
||||
/// Borrow the underlying registry (for path-based resolution / admin reload).
|
||||
pub fn registry(&self) -> &Arc<MountRegistry> {
|
||||
&self.registry
|
||||
}
|
||||
|
||||
/// Fast path: are there no mounts at all? Lets callers skip classification.
|
||||
pub fn is_empty(&self) -> bool {
|
||||
self.registry.is_empty()
|
||||
}
|
||||
|
||||
/// Classify an id. Never parses a provider `node_id` — only the envelope.
|
||||
pub fn classify(&self, id: &str) -> ResolvedId {
|
||||
if is_external_id(id) {
|
||||
if let Some(child) = parse_child_id(id)
|
||||
&& let Some(cfg) = self.registry.get(&child.mount_id)
|
||||
{
|
||||
return ResolvedId::MountChild {
|
||||
cfg,
|
||||
node_id: child.node_id,
|
||||
};
|
||||
}
|
||||
// Malformed or dangling `ext:` id — fall through to Regular so it
|
||||
// surfaces a clean NotFound downstream rather than hitting the repos.
|
||||
return ResolvedId::Regular;
|
||||
}
|
||||
if let Ok(uuid) = Uuid::parse_str(id)
|
||||
&& let Some(cfg) = self.registry.get(&uuid)
|
||||
{
|
||||
return ResolvedId::MountRoot { cfg };
|
||||
}
|
||||
ResolvedId::Regular
|
||||
}
|
||||
|
||||
/// True when `id` addresses anything inside a mount (root or child).
|
||||
pub fn is_mount_id(&self, id: &str) -> bool {
|
||||
matches!(
|
||||
self.classify(id),
|
||||
ResolvedId::MountRoot { .. } | ResolvedId::MountChild { .. }
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use crate::application::ports::external_mount_ports::{
|
||||
ExternalMountRecord, ExternalMountRepositoryPort,
|
||||
};
|
||||
use crate::domain::errors::DomainError;
|
||||
use crate::domain::services::external_mount_id::encode_child_id;
|
||||
use crate::infrastructure::services::mount_provider_factory::DefaultMountProviderFactory;
|
||||
use async_trait::async_trait;
|
||||
use tempfile::TempDir;
|
||||
|
||||
struct FakeRepo(Vec<ExternalMountRecord>);
|
||||
#[async_trait]
|
||||
impl ExternalMountRepositoryPort for FakeRepo {
|
||||
async fn list_all(&self) -> Result<Vec<ExternalMountRecord>, DomainError> {
|
||||
Ok(self.0.clone())
|
||||
}
|
||||
}
|
||||
|
||||
async fn router_with_mount(mount_id: Uuid, dir: &TempDir) -> MountRouter {
|
||||
let repo = FakeRepo(vec![ExternalMountRecord {
|
||||
mount_folder_id: mount_id,
|
||||
kind: "local_fs".to_string(),
|
||||
config: serde_json::json!({ "path": dir.path().to_str().unwrap() }),
|
||||
name: "M".to_string(),
|
||||
owner_id: Uuid::new_v4(),
|
||||
read_only: false,
|
||||
drive_id: Uuid::new_v4(),
|
||||
mount_path: "Personal/M".to_string(),
|
||||
}]);
|
||||
let reg = Arc::new(MountRegistry::empty());
|
||||
reg.reload(&repo, &DefaultMountProviderFactory::new()).await;
|
||||
MountRouter::new(reg)
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn empty_registry_classifies_everything_regular() {
|
||||
let router = MountRouter::new(Arc::new(MountRegistry::empty()));
|
||||
assert!(router.is_empty());
|
||||
assert!(matches!(
|
||||
router.classify(&Uuid::new_v4().to_string()),
|
||||
ResolvedId::Regular
|
||||
));
|
||||
assert!(matches!(
|
||||
router.classify("ext:deadbeef:dG9rZW4"),
|
||||
ResolvedId::Regular
|
||||
));
|
||||
assert!(matches!(router.classify("garbage"), ResolvedId::Regular));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn classifies_mount_root_uuid() {
|
||||
let dir = TempDir::new().unwrap();
|
||||
let mount_id = Uuid::new_v4();
|
||||
let router = router_with_mount(mount_id, &dir).await;
|
||||
|
||||
match router.classify(&mount_id.to_string()) {
|
||||
ResolvedId::MountRoot { cfg } => assert_eq!(cfg.mount_id, mount_id),
|
||||
_ => panic!("expected MountRoot"),
|
||||
}
|
||||
assert!(router.is_mount_id(&mount_id.to_string()));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn classifies_ext_child_id() {
|
||||
let dir = TempDir::new().unwrap();
|
||||
let mount_id = Uuid::new_v4();
|
||||
let router = router_with_mount(mount_id, &dir).await;
|
||||
|
||||
let child = encode_child_id(mount_id, "docs/a.txt");
|
||||
match router.classify(&child) {
|
||||
ResolvedId::MountChild { cfg, node_id } => {
|
||||
assert_eq!(cfg.mount_id, mount_id);
|
||||
assert_eq!(node_id.as_str(), "docs/a.txt");
|
||||
}
|
||||
_ => panic!("expected MountChild"),
|
||||
}
|
||||
assert!(router.is_mount_id(&child));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn ext_id_for_unregistered_mount_is_regular() {
|
||||
let dir = TempDir::new().unwrap();
|
||||
let mount_id = Uuid::new_v4();
|
||||
let router = router_with_mount(mount_id, &dir).await;
|
||||
|
||||
// A well-formed ext: id but for a DIFFERENT (unknown) mount → Regular,
|
||||
// so it 404s downstream rather than hitting the repos.
|
||||
let dangling = encode_child_id(Uuid::new_v4(), "x");
|
||||
assert!(matches!(router.classify(&dangling), ResolvedId::Regular));
|
||||
|
||||
// A plain (non-mount) UUID is also Regular.
|
||||
assert!(matches!(
|
||||
router.classify(&Uuid::new_v4().to_string()),
|
||||
ResolvedId::Regular
|
||||
));
|
||||
}
|
||||
}
|
||||
@@ -5,10 +5,14 @@ use std::sync::Arc;
|
||||
|
||||
use crate::application::dtos::file_dto::FileDto;
|
||||
use crate::application::ports::authorization_ports::AuthorizationEngine;
|
||||
use crate::application::ports::blob_storage_ports::BlobStream;
|
||||
use crate::application::ports::external_mount_ports::MountStat;
|
||||
use crate::application::ports::file_ports::{FileRetrievalUseCase, OptimizedFileContent};
|
||||
use crate::application::ports::storage_ports::FileReadPort;
|
||||
use crate::application::services::mount_registry::MountConfig;
|
||||
use crate::common::errors::DomainError;
|
||||
use crate::domain::services::authorization::{Permission, Resource, Subject};
|
||||
use crate::domain::services::external_mount_id::NodeId;
|
||||
use crate::infrastructure::repositories::pg::file_blob_read_repository::FileBlobReadRepository;
|
||||
use crate::infrastructure::services::file_content_cache::FileContentCache;
|
||||
use crate::infrastructure::services::image_transcode_service::{
|
||||
@@ -64,6 +68,22 @@ impl FileRetrievalService {
|
||||
}
|
||||
}
|
||||
|
||||
/// Test-only constructor: authorization engine without the cache/transcode
|
||||
/// tiers. The external-mount read methods only consult `authz` + the
|
||||
/// provider, so this is sufficient to exercise their authorization.
|
||||
#[cfg(all(test, integration_tests))]
|
||||
pub(crate) fn new_with_authz_for_test(
|
||||
file_read: Arc<FileBlobReadRepository>,
|
||||
authz: Arc<PgAclEngine>,
|
||||
) -> Self {
|
||||
Self {
|
||||
file_read,
|
||||
content_cache: None,
|
||||
transcode: None,
|
||||
authz: Some(authz),
|
||||
}
|
||||
}
|
||||
|
||||
// ── private helpers ──────────────────────────────────────────
|
||||
|
||||
/// Read a file's full content through the streaming API into a single
|
||||
@@ -122,6 +142,50 @@ impl FileRetrievalService {
|
||||
.await
|
||||
}
|
||||
|
||||
/// Authorize then `stat` a file inside an external mount. Authorization
|
||||
/// collapses onto the mount-root folder (a `Read` grant there covers
|
||||
/// everything in the mount).
|
||||
pub async fn stat_mount_file_with_perms(
|
||||
&self,
|
||||
cfg: &MountConfig,
|
||||
node_id: &NodeId,
|
||||
caller_id: Uuid,
|
||||
) -> Result<MountStat, DomainError> {
|
||||
let authz = self.authz.as_ref().ok_or_else(|| {
|
||||
DomainError::internal_error("FileRetrieval", "Authorization engine unavailable")
|
||||
})?;
|
||||
authz
|
||||
.require(
|
||||
Subject::User(caller_id),
|
||||
Permission::Read,
|
||||
Resource::Folder(cfg.mount_id),
|
||||
)
|
||||
.await?;
|
||||
cfg.provider.stat(node_id).await
|
||||
}
|
||||
|
||||
/// Authorize then open a (optionally ranged) read stream over a mount file.
|
||||
/// `range` is `(start, end_inclusive_opt)`.
|
||||
pub async fn open_mount_file_with_perms(
|
||||
&self,
|
||||
cfg: &MountConfig,
|
||||
node_id: &NodeId,
|
||||
caller_id: Uuid,
|
||||
range: Option<(u64, Option<u64>)>,
|
||||
) -> Result<BlobStream, DomainError> {
|
||||
let authz = self.authz.as_ref().ok_or_else(|| {
|
||||
DomainError::internal_error("FileRetrieval", "Authorization engine unavailable")
|
||||
})?;
|
||||
authz
|
||||
.require(
|
||||
Subject::User(caller_id),
|
||||
Permission::Read,
|
||||
Resource::Folder(cfg.mount_id),
|
||||
)
|
||||
.await?;
|
||||
cfg.provider.open_read_stream(node_id, range).await
|
||||
}
|
||||
|
||||
/// Try to transcode image content to WebP and return transcoded variant.
|
||||
async fn try_transcode(
|
||||
&self,
|
||||
|
||||
@@ -4,10 +4,14 @@ use crate::application::dtos::folder_dto::{
|
||||
MoveFolderDto, RenameFolderDto,
|
||||
};
|
||||
use crate::application::ports::authorization_ports::AuthorizationEngine;
|
||||
use crate::application::ports::external_mount_ports::MountEntry;
|
||||
use crate::application::ports::folder_ports::FolderUseCase;
|
||||
use crate::application::services::external_mount_router::MountRouter;
|
||||
use crate::application::services::mount_registry::MountConfig;
|
||||
use crate::common::errors::{DomainError, ErrorKind};
|
||||
use crate::domain::repositories::folder_repository::FolderRepository;
|
||||
use crate::domain::services::authorization::{Permission, Resource, Subject};
|
||||
use crate::domain::services::authorization::{Permission, Resource, ResourceKind, Subject};
|
||||
use crate::domain::services::external_mount_id::NodeId;
|
||||
use crate::domain::services::path_service::{StoragePath, validate_storage_name};
|
||||
use crate::infrastructure::repositories::pg::folder_db_repository::FolderDbRepository;
|
||||
use crate::infrastructure::services::pg_acl_engine::PgAclEngine;
|
||||
@@ -18,17 +22,31 @@ use uuid::Uuid;
|
||||
pub struct FolderService {
|
||||
folder_storage: Arc<FolderDbRepository>,
|
||||
authz: Arc<PgAclEngine>,
|
||||
/// External-mount classifier. Lets folder operations branch a mount-root or
|
||||
/// `ext:` id onto the provider instead of the PostgreSQL repositories.
|
||||
mount_router: Arc<MountRouter>,
|
||||
}
|
||||
|
||||
impl FolderService {
|
||||
/// Creates a new folder service
|
||||
pub fn new(folder_storage: Arc<FolderDbRepository>, authz: Arc<PgAclEngine>) -> Self {
|
||||
pub fn new(
|
||||
folder_storage: Arc<FolderDbRepository>,
|
||||
authz: Arc<PgAclEngine>,
|
||||
mount_router: Arc<MountRouter>,
|
||||
) -> Self {
|
||||
Self {
|
||||
folder_storage,
|
||||
authz,
|
||||
mount_router,
|
||||
}
|
||||
}
|
||||
|
||||
/// Borrow the external-mount classifier (handlers branch on this before
|
||||
/// treating an id as a native UUID).
|
||||
pub fn mount_router(&self) -> &MountRouter {
|
||||
&self.mount_router
|
||||
}
|
||||
|
||||
/// Batch counterpart of `get_folder`: resolve many folder ids in ONE
|
||||
/// query instead of one per id. Like `get_folder` it performs no
|
||||
/// per-folder authorization — both current callers (ACL grant listing,
|
||||
@@ -615,6 +633,125 @@ impl FolderService {
|
||||
|
||||
Ok((rows, next_cursor))
|
||||
}
|
||||
|
||||
/// List one directory inside an external mount (the mount root when
|
||||
/// `node_id` is empty, or a nested virtual folder otherwise).
|
||||
///
|
||||
/// Authorization collapses onto the mount-root folder: a caller who may
|
||||
/// `Read` the mount root may browse everything inside it. The provider
|
||||
/// reads the live backend; entries are sorted in memory and paginated with
|
||||
/// a name-keyset cursor (directories are bounded, see provider cap).
|
||||
///
|
||||
/// Returns the page of raw [`MountEntry`]s plus an encoded next cursor; the
|
||||
/// handler maps each entry to a `FolderResourceItemDto` with a synthetic
|
||||
/// `ext:` id.
|
||||
pub async fn list_mount_dir_with_perms(
|
||||
&self,
|
||||
cfg: &MountConfig,
|
||||
node_id: &NodeId,
|
||||
caller_id: Uuid,
|
||||
opts: ListResourcesOptions<'_>,
|
||||
) -> Result<(Vec<MountEntry>, Option<String>), DomainError> {
|
||||
// AuthZ — everything in the mount is gated by the mount-root folder.
|
||||
self.authz
|
||||
.require(
|
||||
Subject::User(caller_id),
|
||||
Permission::Read,
|
||||
Resource::Folder(cfg.mount_id),
|
||||
)
|
||||
.await?;
|
||||
|
||||
let entries = cfg.provider.list_dir(node_id).await?;
|
||||
let cursor_name = opts.cursor.as_ref().and_then(|c| c.sort_str.as_deref());
|
||||
Ok(paginate_mount_entries(
|
||||
entries,
|
||||
opts.kinds,
|
||||
opts.order_by,
|
||||
opts.reverse,
|
||||
opts.limit,
|
||||
cursor_name,
|
||||
))
|
||||
}
|
||||
}
|
||||
|
||||
/// Filter, sort, and page a directory's worth of mount entries, returning the
|
||||
/// page plus an encoded next cursor. Pure (no I/O / authz) so it can be tested
|
||||
/// exhaustively.
|
||||
///
|
||||
/// The cursor is a **name keyset**: names are unique within a directory, so the
|
||||
/// last emitted name is a stable resume key under any sort dimension. Resume is
|
||||
/// best-effort — if the cursor's entry was deleted out-of-band the page restarts
|
||||
/// from the top (documented; avoids an infinite loop).
|
||||
fn paginate_mount_entries(
|
||||
mut entries: Vec<MountEntry>,
|
||||
kinds: Option<&[ResourceKind]>,
|
||||
order_by: &str,
|
||||
reverse: bool,
|
||||
limit: usize,
|
||||
cursor_name: Option<&str>,
|
||||
) -> (Vec<MountEntry>, Option<String>) {
|
||||
if let Some(kinds) = kinds {
|
||||
let want_files = kinds.contains(&ResourceKind::File);
|
||||
let want_folders = kinds.contains(&ResourceKind::Folder);
|
||||
entries.retain(|e| if e.is_dir { want_folders } else { want_files });
|
||||
}
|
||||
|
||||
sort_mount_entries(&mut entries, order_by, reverse);
|
||||
|
||||
let start = match cursor_name {
|
||||
Some(name) => entries
|
||||
.iter()
|
||||
.position(|e| name.eq_ignore_ascii_case(&e.name))
|
||||
.map(|i| i + 1)
|
||||
.unwrap_or(0),
|
||||
None => 0,
|
||||
};
|
||||
|
||||
let has_more = entries.len() > start + limit;
|
||||
let page: Vec<MountEntry> = entries.into_iter().skip(start).take(limit).collect();
|
||||
|
||||
let next_cursor = if has_more {
|
||||
page.last().map(|last| {
|
||||
FolderResourceCursor {
|
||||
order_by: order_by.to_owned(),
|
||||
resource_id: Uuid::nil(),
|
||||
sort_str: Some(last.name.clone()),
|
||||
sort_int: None,
|
||||
sort_ts: None,
|
||||
reverse,
|
||||
}
|
||||
.encode()
|
||||
})
|
||||
} else {
|
||||
None
|
||||
};
|
||||
|
||||
(page, next_cursor)
|
||||
}
|
||||
|
||||
/// Sort mount entries in place. Folders sort before files for the `name`/`type`
|
||||
/// dimensions; otherwise by the requested key with name as the tie-breaker.
|
||||
/// `reverse` flips the final order.
|
||||
fn sort_mount_entries(entries: &mut [MountEntry], order_by: &str, reverse: bool) {
|
||||
use std::cmp::Ordering;
|
||||
let name_key = |e: &MountEntry| e.name.to_lowercase();
|
||||
entries.sort_by(|a, b| {
|
||||
let primary = match order_by {
|
||||
"modified_at" => a.modified_at.cmp(&b.modified_at),
|
||||
"created_at" => a.created_at.cmp(&b.created_at),
|
||||
"size" => a.size.cmp(&b.size),
|
||||
// "name" / "type" / anything else: folders first, then by name.
|
||||
_ => b.is_dir.cmp(&a.is_dir),
|
||||
};
|
||||
let ord = primary.then_with(|| name_key(a).cmp(&name_key(b)));
|
||||
if ord == Ordering::Equal {
|
||||
Ordering::Equal
|
||||
} else if reverse {
|
||||
ord.reverse()
|
||||
} else {
|
||||
ord
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
/// Build the next-page cursor from the last row of the current page.
|
||||
@@ -842,3 +979,341 @@ impl UserLifecycleHook for PersonalDriveLifecycleHook {
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod mount_listing_tests {
|
||||
use super::{paginate_mount_entries, sort_mount_entries};
|
||||
use crate::application::dtos::cursor::PageCursor;
|
||||
use crate::application::dtos::folder_dto::FolderResourceCursor;
|
||||
use crate::application::ports::external_mount_ports::MountEntry;
|
||||
use crate::domain::services::authorization::ResourceKind;
|
||||
use crate::domain::services::external_mount_id::NodeId;
|
||||
|
||||
fn entry(name: &str, is_dir: bool, size: u64, modified: u64) -> MountEntry {
|
||||
MountEntry {
|
||||
name: name.to_string(),
|
||||
node_id: NodeId(name.to_string()),
|
||||
is_dir,
|
||||
size,
|
||||
modified_at: modified,
|
||||
created_at: modified,
|
||||
}
|
||||
}
|
||||
|
||||
fn names(entries: &[MountEntry]) -> Vec<String> {
|
||||
entries.iter().map(|e| e.name.clone()).collect()
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn sorts_folders_first_then_name_case_insensitive() {
|
||||
let mut e = vec![
|
||||
entry("Banana.txt", false, 1, 1),
|
||||
entry("apple", true, 0, 1),
|
||||
entry("Cherry", true, 0, 1),
|
||||
entry("almond.txt", false, 1, 1),
|
||||
];
|
||||
sort_mount_entries(&mut e, "name", false);
|
||||
assert_eq!(names(&e), ["apple", "Cherry", "almond.txt", "Banana.txt"]);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn reverse_flips_order() {
|
||||
let mut e = vec![
|
||||
entry("a", false, 1, 1),
|
||||
entry("b", false, 1, 1),
|
||||
entry("d", true, 0, 1),
|
||||
];
|
||||
sort_mount_entries(&mut e, "name", true);
|
||||
// folders-first then name, reversed.
|
||||
assert_eq!(names(&e), ["b", "a", "d"]);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn sorts_by_size_modified_created() {
|
||||
let mut by_size = vec![
|
||||
entry("big", false, 100, 1),
|
||||
entry("small", false, 1, 1),
|
||||
entry("mid", false, 50, 1),
|
||||
];
|
||||
sort_mount_entries(&mut by_size, "size", false);
|
||||
assert_eq!(names(&by_size), ["small", "mid", "big"]);
|
||||
|
||||
let mut by_mtime = vec![
|
||||
entry("new", false, 1, 300),
|
||||
entry("old", false, 1, 100),
|
||||
entry("mid", false, 1, 200),
|
||||
];
|
||||
sort_mount_entries(&mut by_mtime, "modified_at", false);
|
||||
assert_eq!(names(&by_mtime), ["old", "mid", "new"]);
|
||||
|
||||
let mut by_ctime = vec![entry("z", false, 1, 9), entry("a", false, 1, 5)];
|
||||
sort_mount_entries(&mut by_ctime, "created_at", false);
|
||||
assert_eq!(names(&by_ctime), ["a", "z"]);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn filters_by_kind() {
|
||||
let make = || vec![entry("dir", true, 0, 1), entry("file.txt", false, 1, 1)];
|
||||
let (files_only, _) =
|
||||
paginate_mount_entries(make(), Some(&[ResourceKind::File]), "name", false, 50, None);
|
||||
assert_eq!(names(&files_only), ["file.txt"]);
|
||||
|
||||
let (folders_only, _) = paginate_mount_entries(
|
||||
make(),
|
||||
Some(&[ResourceKind::Folder]),
|
||||
"name",
|
||||
false,
|
||||
50,
|
||||
None,
|
||||
);
|
||||
assert_eq!(names(&folders_only), ["dir"]);
|
||||
|
||||
let (both, _) = paginate_mount_entries(
|
||||
make(),
|
||||
Some(&[ResourceKind::File, ResourceKind::Folder]),
|
||||
"name",
|
||||
false,
|
||||
50,
|
||||
None,
|
||||
);
|
||||
assert_eq!(both.len(), 2);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn paginates_with_name_keyset_cursor() {
|
||||
let all = || {
|
||||
vec![
|
||||
entry("a", false, 1, 1),
|
||||
entry("b", false, 1, 1),
|
||||
entry("c", false, 1, 1),
|
||||
entry("d", false, 1, 1),
|
||||
entry("e", false, 1, 1),
|
||||
]
|
||||
};
|
||||
|
||||
// Page 1: limit 2 → [a, b], cursor present.
|
||||
let (p1, c1) = paginate_mount_entries(all(), None, "name", false, 2, None);
|
||||
assert_eq!(names(&p1), ["a", "b"]);
|
||||
let c1 = c1.expect("cursor after first page");
|
||||
let decoded = FolderResourceCursor::decode(&c1).expect("decodes");
|
||||
assert_eq!(decoded.sort_str.as_deref(), Some("b"));
|
||||
assert!(!decoded.reverse);
|
||||
|
||||
// Page 2: resume after "b" → [c, d], cursor present.
|
||||
let (p2, c2) = paginate_mount_entries(all(), None, "name", false, 2, Some("b"));
|
||||
assert_eq!(names(&p2), ["c", "d"]);
|
||||
assert!(c2.is_some());
|
||||
|
||||
// Page 3: resume after "d" → [e], no further cursor.
|
||||
let (p3, c3) = paginate_mount_entries(all(), None, "name", false, 2, Some("d"));
|
||||
assert_eq!(names(&p3), ["e"]);
|
||||
assert!(c3.is_none());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn no_cursor_when_page_is_last() {
|
||||
let e = vec![entry("a", false, 1, 1), entry("b", false, 1, 1)];
|
||||
let (page, cursor) = paginate_mount_entries(e, None, "name", false, 50, None);
|
||||
assert_eq!(page.len(), 2);
|
||||
assert!(cursor.is_none());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn deleted_cursor_entry_restarts_best_effort() {
|
||||
// Cursor names "zzz" which is not present → start from the top.
|
||||
let e = vec![entry("a", false, 1, 1), entry("b", false, 1, 1)];
|
||||
let (page, _) = paginate_mount_entries(e, None, "name", false, 50, Some("zzz"));
|
||||
assert_eq!(names(&page), ["a", "b"]);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn empty_directory_yields_empty_page() {
|
||||
let (page, cursor) = paginate_mount_entries(vec![], None, "name", false, 50, None);
|
||||
assert!(page.is_empty());
|
||||
assert!(cursor.is_none());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn limit_larger_than_len_returns_all_without_cursor() {
|
||||
let e = vec![entry("a", false, 1, 1), entry("b", false, 1, 1)];
|
||||
let (page, cursor) = paginate_mount_entries(e, None, "name", false, 100, None);
|
||||
assert_eq!(page.len(), 2);
|
||||
assert!(cursor.is_none());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn kind_filter_excluding_all_yields_empty() {
|
||||
let e = vec![entry("only_dir", true, 0, 1)];
|
||||
let (page, cursor) =
|
||||
paginate_mount_entries(e, Some(&[ResourceKind::File]), "name", false, 50, None);
|
||||
assert!(page.is_empty());
|
||||
assert!(cursor.is_none());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn cursor_preserves_reverse_flag() {
|
||||
let e = vec![
|
||||
entry("a", false, 1, 1),
|
||||
entry("b", false, 1, 1),
|
||||
entry("c", false, 1, 1),
|
||||
];
|
||||
let (_p, c) = paginate_mount_entries(e, None, "name", true, 1, None);
|
||||
let decoded = FolderResourceCursor::decode(&c.unwrap()).unwrap();
|
||||
assert!(decoded.reverse);
|
||||
assert_eq!(decoded.order_by, "name");
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(all(test, integration_tests))]
|
||||
mod mount_authz_integration {
|
||||
use super::*;
|
||||
use crate::application::dtos::folder_dto::ListResourcesOptions;
|
||||
use crate::application::services::external_mount_router::{MountRouter, ResolvedId};
|
||||
use crate::application::services::file_retrieval_service::FileRetrievalService;
|
||||
use crate::application::services::mount_registry::MountRegistry;
|
||||
use crate::domain::services::external_mount_id::{NodeId, encode_child_id};
|
||||
use crate::infrastructure::repositories::pg::{
|
||||
ExternalMountPgRepository, FileBlobReadRepository, SubjectGroupPgRepository,
|
||||
};
|
||||
use crate::infrastructure::services::mount_provider_factory::DefaultMountProviderFactory;
|
||||
use crate::mount_it_support::{fresh_db, insert_mount, make_user, provision_folder};
|
||||
use std::sync::Arc;
|
||||
|
||||
fn opts<'a>() -> ListResourcesOptions<'a> {
|
||||
ListResourcesOptions {
|
||||
limit: 50,
|
||||
cursor: None,
|
||||
order_by: "name",
|
||||
kinds: None,
|
||||
reverse: false,
|
||||
}
|
||||
}
|
||||
|
||||
/// Build a real PgAclEngine over the live pool. The folder-ancestry cascade
|
||||
/// uses the engine's own pool; the file repo is a stub (not exercised by
|
||||
/// folder checks).
|
||||
fn acl(pool: &Arc<sqlx::PgPool>) -> Arc<PgAclEngine> {
|
||||
Arc::new(PgAclEngine::new(
|
||||
pool.clone(),
|
||||
Arc::new(FolderDbRepository::new(pool.clone())),
|
||||
Arc::new(FileBlobReadRepository::new_stub()),
|
||||
Arc::new(SubjectGroupPgRepository::new(pool.clone())),
|
||||
))
|
||||
}
|
||||
|
||||
/// Full read path: owner can list a mount's live contents; a stranger with
|
||||
/// no grant is denied. Exercises the REAL authorization cascade
|
||||
/// (`authz.require(Resource::Folder(mount_id))`) over ltree ancestry.
|
||||
#[tokio::test]
|
||||
async fn owner_lists_mount_contents_stranger_denied() {
|
||||
let (_c, pool) = fresh_db().await;
|
||||
|
||||
// Real host directory the mount points at.
|
||||
let host = tempfile::tempdir().unwrap();
|
||||
std::fs::write(host.path().join("a.txt"), b"hello").unwrap();
|
||||
std::fs::create_dir(host.path().join("sub")).unwrap();
|
||||
|
||||
let p = provision_folder(&pool, "owner", "Media").await;
|
||||
insert_mount(&pool, &p, host.path().to_str().unwrap()).await;
|
||||
|
||||
// Build the registry from the DB (also exercises reload + provider build).
|
||||
let registry = Arc::new(MountRegistry::empty());
|
||||
registry
|
||||
.reload(
|
||||
&ExternalMountPgRepository::new(pool.clone()),
|
||||
&DefaultMountProviderFactory::new(),
|
||||
)
|
||||
.await;
|
||||
let router = Arc::new(MountRouter::new(registry.clone()));
|
||||
let folder_service = FolderService::new(
|
||||
Arc::new(FolderDbRepository::new(pool.clone())),
|
||||
acl(&pool),
|
||||
router.clone(),
|
||||
);
|
||||
|
||||
let cfg = registry.get(&p.mount_folder_id).expect("mount registered");
|
||||
|
||||
// The mount root UUID classifies as a MountRoot.
|
||||
assert!(matches!(
|
||||
router.classify(&p.mount_folder_id.to_string()),
|
||||
ResolvedId::MountRoot { .. }
|
||||
));
|
||||
|
||||
// Owner lists the live directory contents.
|
||||
let (entries, _cursor) = folder_service
|
||||
.list_mount_dir_with_perms(&cfg, &NodeId::default(), p.owner_id, opts())
|
||||
.await
|
||||
.expect("owner may list");
|
||||
let mut names: Vec<_> = entries.iter().map(|e| e.name.clone()).collect();
|
||||
names.sort();
|
||||
assert_eq!(names, ["a.txt", "sub"]);
|
||||
|
||||
// A stranger with no grant on the mount-root folder is denied
|
||||
// (NotFound — anti-enumeration).
|
||||
let stranger = make_user(&pool, "stranger").await;
|
||||
let err = folder_service
|
||||
.list_mount_dir_with_perms(&cfg, &NodeId::default(), stranger, opts())
|
||||
.await
|
||||
.expect_err("stranger must be denied");
|
||||
assert_eq!(err.kind, crate::domain::errors::ErrorKind::NotFound);
|
||||
}
|
||||
|
||||
/// Download path authz: owner can stat/open a mount file; stranger denied.
|
||||
#[tokio::test]
|
||||
async fn owner_reads_mount_file_stranger_denied() {
|
||||
let (_c, pool) = fresh_db().await;
|
||||
|
||||
let host = tempfile::tempdir().unwrap();
|
||||
std::fs::write(host.path().join("doc.txt"), b"payload").unwrap();
|
||||
|
||||
let p = provision_folder(&pool, "owner", "Media").await;
|
||||
insert_mount(&pool, &p, host.path().to_str().unwrap()).await;
|
||||
|
||||
let registry = Arc::new(MountRegistry::empty());
|
||||
registry
|
||||
.reload(
|
||||
&ExternalMountPgRepository::new(pool.clone()),
|
||||
&DefaultMountProviderFactory::new(),
|
||||
)
|
||||
.await;
|
||||
let cfg = registry.get(&p.mount_folder_id).expect("registered");
|
||||
|
||||
let retrieval = FileRetrievalService::new_with_authz_for_test(
|
||||
Arc::new(FileBlobReadRepository::new_stub()),
|
||||
acl(&pool),
|
||||
);
|
||||
|
||||
let node = NodeId::from("doc.txt");
|
||||
|
||||
// Owner: stat succeeds with the real size.
|
||||
let stat = retrieval
|
||||
.stat_mount_file_with_perms(&cfg, &node, p.owner_id)
|
||||
.await
|
||||
.expect("owner may stat");
|
||||
assert_eq!(stat.size, 7);
|
||||
assert!(!stat.is_dir);
|
||||
|
||||
// Owner: open succeeds (smoke — stream is consumed elsewhere).
|
||||
assert!(
|
||||
retrieval
|
||||
.open_mount_file_with_perms(&cfg, &node, p.owner_id, None)
|
||||
.await
|
||||
.is_ok()
|
||||
);
|
||||
|
||||
// The synthetic id for this file round-trips through the router.
|
||||
let ext_id = encode_child_id(p.mount_folder_id, "doc.txt");
|
||||
assert!(matches!(
|
||||
MountRouter::new(registry.clone()).classify(&ext_id),
|
||||
ResolvedId::MountChild { .. }
|
||||
));
|
||||
|
||||
// Stranger: denied.
|
||||
let stranger = make_user(&pool, "stranger").await;
|
||||
let err = retrieval
|
||||
.stat_mount_file_with_perms(&cfg, &node, stranger)
|
||||
.await
|
||||
.expect_err("stranger denied");
|
||||
assert_eq!(err.kind, crate::domain::errors::ErrorKind::NotFound);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -9,6 +9,7 @@ pub mod delta_upload_service;
|
||||
pub mod device_auth_service;
|
||||
pub mod drive_management_service;
|
||||
pub mod external_identity_service;
|
||||
pub mod external_mount_router;
|
||||
pub mod favorites_service;
|
||||
pub mod file_lifecycle_service;
|
||||
pub mod file_management_service;
|
||||
@@ -18,6 +19,7 @@ pub mod file_use_case_factory;
|
||||
pub mod folder_service;
|
||||
pub mod i18n_application_service;
|
||||
pub mod magic_link_invite_service;
|
||||
pub mod mount_registry;
|
||||
pub mod music_service;
|
||||
pub mod nextcloud_file_id_service;
|
||||
pub mod nextcloud_login_flow_service;
|
||||
|
||||
@@ -0,0 +1,374 @@
|
||||
//! In-memory registry of configured external mounts.
|
||||
//!
|
||||
//! Holds, per mount-root folder UUID, the constructed provider plus the metadata
|
||||
//! the service layer needs to authorize and synthesize DTOs. Reads are lock-free
|
||||
//! (`arc-swap`) so the hot path ("is this UUID a mount root?") never blocks; the
|
||||
//! whole index is rebuilt on admin mutation via [`MountRegistry::reload`].
|
||||
|
||||
use std::collections::HashMap;
|
||||
use std::sync::Arc;
|
||||
|
||||
use arc_swap::ArcSwap;
|
||||
use uuid::Uuid;
|
||||
|
||||
use crate::application::ports::external_mount_ports::{
|
||||
ExternalMountProvider, ExternalMountRepositoryPort, MountProviderFactory,
|
||||
};
|
||||
|
||||
/// One configured mount, with its live provider.
|
||||
pub struct MountConfig {
|
||||
/// Mount-root folder UUID — the mount's identity and the authz resource.
|
||||
pub mount_id: Uuid,
|
||||
/// Provider kind.
|
||||
pub kind: String,
|
||||
/// Display name.
|
||||
pub name: String,
|
||||
/// Owner of the mount configuration.
|
||||
pub owner_id: Uuid,
|
||||
/// Drive the mount root belongs to.
|
||||
pub drive_id: Uuid,
|
||||
/// Whether the mount refuses mutations.
|
||||
pub read_only: bool,
|
||||
/// Materialized internal path of the mount root (e.g. `"Personal/Media"`),
|
||||
/// used by path-based resolution for WebDAV / NextCloud.
|
||||
pub mount_path: String,
|
||||
/// The bound provider for this mount's backend.
|
||||
pub provider: Arc<dyn ExternalMountProvider>,
|
||||
}
|
||||
|
||||
/// Immutable snapshot swapped atomically on reload.
|
||||
#[derive(Default)]
|
||||
struct MountIndex {
|
||||
/// mount-root folder UUID → config.
|
||||
by_folder: HashMap<Uuid, Arc<MountConfig>>,
|
||||
/// (drive_id, mount_path) → mount-root UUID, for path resolution (P3).
|
||||
by_path: HashMap<(Uuid, String), Uuid>,
|
||||
}
|
||||
|
||||
/// Lock-free registry of mounts.
|
||||
pub struct MountRegistry {
|
||||
inner: ArcSwap<MountIndex>,
|
||||
}
|
||||
|
||||
impl Default for MountRegistry {
|
||||
fn default() -> Self {
|
||||
Self::empty()
|
||||
}
|
||||
}
|
||||
|
||||
impl MountRegistry {
|
||||
/// An empty registry (no mounts).
|
||||
pub fn empty() -> Self {
|
||||
Self {
|
||||
inner: ArcSwap::from_pointee(MountIndex::default()),
|
||||
}
|
||||
}
|
||||
|
||||
/// Look up a mount by its root folder UUID.
|
||||
pub fn get(&self, mount_id: &Uuid) -> Option<Arc<MountConfig>> {
|
||||
self.inner.load().by_folder.get(mount_id).cloned()
|
||||
}
|
||||
|
||||
/// Is this UUID the root of a configured mount?
|
||||
pub fn is_mount_root(&self, id: &Uuid) -> bool {
|
||||
self.inner.load().by_folder.contains_key(id)
|
||||
}
|
||||
|
||||
/// True when no mounts are configured (lets callers skip work entirely).
|
||||
pub fn is_empty(&self) -> bool {
|
||||
self.inner.load().by_folder.is_empty()
|
||||
}
|
||||
|
||||
/// Find the mount whose root path is a segment-aligned prefix of
|
||||
/// `internal_path` within `drive_id`. Returns the config plus the remainder
|
||||
/// path relative to the mount root (`""` when the path IS the mount root).
|
||||
///
|
||||
/// Used by path-based resolution (WebDAV / NextCloud) in P3.
|
||||
pub fn find_mount_for_path(
|
||||
&self,
|
||||
drive_id: Uuid,
|
||||
internal_path: &str,
|
||||
) -> Option<(Arc<MountConfig>, String)> {
|
||||
let index = self.inner.load();
|
||||
// Walk ancestor paths from the full path up to the root, longest first,
|
||||
// so the deepest matching mount wins.
|
||||
let mut candidate = internal_path;
|
||||
loop {
|
||||
if let Some(mount_id) = index.by_path.get(&(drive_id, candidate.to_string()))
|
||||
&& let Some(cfg) = index.by_folder.get(mount_id)
|
||||
{
|
||||
let remainder = internal_path
|
||||
.strip_prefix(candidate)
|
||||
.map(|r| r.trim_start_matches('/').to_string())
|
||||
.unwrap_or_default();
|
||||
return Some((cfg.clone(), remainder));
|
||||
}
|
||||
match candidate.rsplit_once('/') {
|
||||
Some((parent, _)) => candidate = parent,
|
||||
None => return None,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Rebuild the registry from persisted records, constructing each provider
|
||||
/// via the factory. A mount whose provider fails to build is skipped (logged)
|
||||
/// rather than failing the whole reload.
|
||||
pub async fn reload(
|
||||
&self,
|
||||
repo: &dyn ExternalMountRepositoryPort,
|
||||
factory: &dyn MountProviderFactory,
|
||||
) {
|
||||
let records = match repo.list_all().await {
|
||||
Ok(r) => r,
|
||||
Err(e) => {
|
||||
tracing::error!(
|
||||
target: "oxicloud::external_mounts",
|
||||
"failed to load external mounts: {e}"
|
||||
);
|
||||
return;
|
||||
}
|
||||
};
|
||||
|
||||
let mut by_folder = HashMap::with_capacity(records.len());
|
||||
let mut by_path = HashMap::with_capacity(records.len());
|
||||
for rec in records {
|
||||
let provider = match factory.build(&rec.kind, &rec.config).await {
|
||||
Ok(p) => p,
|
||||
Err(e) => {
|
||||
tracing::error!(
|
||||
target: "oxicloud::external_mounts",
|
||||
mount_id = %rec.mount_folder_id,
|
||||
kind = %rec.kind,
|
||||
"skipping mount: provider build failed: {e}"
|
||||
);
|
||||
continue;
|
||||
}
|
||||
};
|
||||
by_path.insert((rec.drive_id, rec.mount_path.clone()), rec.mount_folder_id);
|
||||
by_folder.insert(
|
||||
rec.mount_folder_id,
|
||||
Arc::new(MountConfig {
|
||||
mount_id: rec.mount_folder_id,
|
||||
kind: rec.kind,
|
||||
name: rec.name,
|
||||
owner_id: rec.owner_id,
|
||||
drive_id: rec.drive_id,
|
||||
read_only: rec.read_only,
|
||||
mount_path: rec.mount_path,
|
||||
provider,
|
||||
}),
|
||||
);
|
||||
}
|
||||
|
||||
let count = by_folder.len();
|
||||
self.inner
|
||||
.store(Arc::new(MountIndex { by_folder, by_path }));
|
||||
tracing::info!(
|
||||
target: "oxicloud::external_mounts",
|
||||
count, "external mount registry loaded"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use crate::application::ports::external_mount_ports::{
|
||||
ExternalMountRecord, ExternalMountRepositoryPort,
|
||||
};
|
||||
use crate::domain::errors::DomainError;
|
||||
use crate::infrastructure::services::mount_provider_factory::DefaultMountProviderFactory;
|
||||
use async_trait::async_trait;
|
||||
use std::path::Path;
|
||||
use tempfile::TempDir;
|
||||
|
||||
struct FakeRepo {
|
||||
records: Vec<ExternalMountRecord>,
|
||||
}
|
||||
|
||||
#[async_trait]
|
||||
impl ExternalMountRepositoryPort for FakeRepo {
|
||||
async fn list_all(&self) -> Result<Vec<ExternalMountRecord>, DomainError> {
|
||||
Ok(self.records.clone())
|
||||
}
|
||||
}
|
||||
|
||||
fn record(
|
||||
mount_id: Uuid,
|
||||
drive_id: Uuid,
|
||||
mount_path: &str,
|
||||
path: &Path,
|
||||
) -> ExternalMountRecord {
|
||||
ExternalMountRecord {
|
||||
mount_folder_id: mount_id,
|
||||
kind: "local_fs".to_string(),
|
||||
config: serde_json::json!({ "path": path.to_str().unwrap() }),
|
||||
name: "Test Mount".to_string(),
|
||||
owner_id: Uuid::new_v4(),
|
||||
read_only: false,
|
||||
drive_id,
|
||||
mount_path: mount_path.to_string(),
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn empty_registry_is_inert() {
|
||||
let r = MountRegistry::empty();
|
||||
assert!(r.is_empty());
|
||||
assert!(!r.is_mount_root(&Uuid::new_v4()));
|
||||
assert!(r.get(&Uuid::new_v4()).is_none());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn reload_populates_from_records() {
|
||||
let dir = TempDir::new().unwrap();
|
||||
let mount_id = Uuid::new_v4();
|
||||
let drive_id = Uuid::new_v4();
|
||||
let repo = FakeRepo {
|
||||
records: vec![record(mount_id, drive_id, "Personal/Media", dir.path())],
|
||||
};
|
||||
let factory = DefaultMountProviderFactory::new();
|
||||
|
||||
let reg = MountRegistry::empty();
|
||||
reg.reload(&repo, &factory).await;
|
||||
|
||||
assert!(!reg.is_empty());
|
||||
assert!(reg.is_mount_root(&mount_id));
|
||||
let cfg = reg.get(&mount_id).expect("present");
|
||||
assert_eq!(cfg.kind, "local_fs");
|
||||
assert_eq!(cfg.drive_id, drive_id);
|
||||
assert_eq!(cfg.mount_path, "Personal/Media");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn reload_skips_mount_whose_provider_fails_to_build() {
|
||||
let mount_id = Uuid::new_v4();
|
||||
// Point at a path that doesn't exist → LocalFsMountProvider::new errors.
|
||||
let repo = FakeRepo {
|
||||
records: vec![ExternalMountRecord {
|
||||
mount_folder_id: mount_id,
|
||||
kind: "local_fs".to_string(),
|
||||
config: serde_json::json!({ "path": "/nonexistent/path/xyz-123" }),
|
||||
name: "Bad".to_string(),
|
||||
owner_id: Uuid::new_v4(),
|
||||
read_only: false,
|
||||
drive_id: Uuid::new_v4(),
|
||||
mount_path: "Personal/Bad".to_string(),
|
||||
}],
|
||||
};
|
||||
let factory = DefaultMountProviderFactory::new();
|
||||
let reg = MountRegistry::empty();
|
||||
reg.reload(&repo, &factory).await;
|
||||
// The bad mount is skipped, not fatal.
|
||||
assert!(reg.is_empty());
|
||||
assert!(!reg.is_mount_root(&mount_id));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn find_mount_for_path_matches_prefix_and_remainder() {
|
||||
let dir = TempDir::new().unwrap();
|
||||
let mount_id = Uuid::new_v4();
|
||||
let drive_id = Uuid::new_v4();
|
||||
let repo = FakeRepo {
|
||||
records: vec![record(mount_id, drive_id, "Personal/Media", dir.path())],
|
||||
};
|
||||
let reg = MountRegistry::empty();
|
||||
reg.reload(&repo, &DefaultMountProviderFactory::new()).await;
|
||||
|
||||
// Exact match → empty remainder.
|
||||
let (cfg, rem) = reg
|
||||
.find_mount_for_path(drive_id, "Personal/Media")
|
||||
.expect("exact match");
|
||||
assert_eq!(cfg.mount_id, mount_id);
|
||||
assert_eq!(rem, "");
|
||||
|
||||
// Nested path → remainder is the suffix.
|
||||
let (_cfg, rem) = reg
|
||||
.find_mount_for_path(drive_id, "Personal/Media/docs/a.txt")
|
||||
.expect("nested match");
|
||||
assert_eq!(rem, "docs/a.txt");
|
||||
|
||||
// Non-matching path within the drive → None.
|
||||
assert!(
|
||||
reg.find_mount_for_path(drive_id, "Personal/Other")
|
||||
.is_none()
|
||||
);
|
||||
|
||||
// Same path but a DIFFERENT drive → None (drive-scoped).
|
||||
assert!(
|
||||
reg.find_mount_for_path(Uuid::new_v4(), "Personal/Media")
|
||||
.is_none()
|
||||
);
|
||||
|
||||
// A sibling that merely shares a name prefix must NOT match
|
||||
// (segment-aligned only).
|
||||
assert!(
|
||||
reg.find_mount_for_path(drive_id, "Personal/MediaLibrary")
|
||||
.is_none()
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn reload_replaces_previous_state() {
|
||||
let dir = TempDir::new().unwrap();
|
||||
let drive = Uuid::new_v4();
|
||||
let first = Uuid::new_v4();
|
||||
let reg = MountRegistry::empty();
|
||||
|
||||
reg.reload(
|
||||
&FakeRepo {
|
||||
records: vec![record(first, drive, "Personal/A", dir.path())],
|
||||
},
|
||||
&DefaultMountProviderFactory::new(),
|
||||
)
|
||||
.await;
|
||||
assert!(reg.is_mount_root(&first));
|
||||
|
||||
// A second reload with a different mount set replaces the first entirely.
|
||||
let second = Uuid::new_v4();
|
||||
reg.reload(
|
||||
&FakeRepo {
|
||||
records: vec![record(second, drive, "Personal/B", dir.path())],
|
||||
},
|
||||
&DefaultMountProviderFactory::new(),
|
||||
)
|
||||
.await;
|
||||
assert!(reg.is_mount_root(&second));
|
||||
assert!(
|
||||
!reg.is_mount_root(&first),
|
||||
"stale mount must be gone after reload"
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn find_mount_for_path_deepest_wins() {
|
||||
let outer = TempDir::new().unwrap();
|
||||
let inner = TempDir::new().unwrap();
|
||||
let drive_id = Uuid::new_v4();
|
||||
let outer_id = Uuid::new_v4();
|
||||
let inner_id = Uuid::new_v4();
|
||||
let repo = FakeRepo {
|
||||
records: vec![
|
||||
record(outer_id, drive_id, "Personal", outer.path()),
|
||||
record(inner_id, drive_id, "Personal/Media", inner.path()),
|
||||
],
|
||||
};
|
||||
let reg = MountRegistry::empty();
|
||||
reg.reload(&repo, &DefaultMountProviderFactory::new()).await;
|
||||
|
||||
// A path under the deeper mount resolves to the deeper mount.
|
||||
let (cfg, rem) = reg
|
||||
.find_mount_for_path(drive_id, "Personal/Media/x")
|
||||
.expect("match");
|
||||
assert_eq!(cfg.mount_id, inner_id);
|
||||
assert_eq!(rem, "x");
|
||||
|
||||
// A path under the shallower mount (but not the deeper one) resolves
|
||||
// to the shallower mount.
|
||||
let (cfg, rem) = reg
|
||||
.find_mount_for_path(drive_id, "Personal/Other/y")
|
||||
.expect("match");
|
||||
assert_eq!(cfg.mount_id, outer_id);
|
||||
assert_eq!(rem, "Other/y");
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user