Merge pull request #521 from BCNelson/feat/external-file-mounts

feat(mounts): external file mounts — pluggable provider, read-write, WebDAV/NextCloud, admin UI
This commit is contained in:
Dionisio Pozo
2026-07-22 06:52:43 +02:00
committed by GitHub
40 changed files with 5920 additions and 76 deletions
@@ -0,0 +1,228 @@
//! 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<'a> =
Pin<Box<dyn Stream<Item = Result<Bytes, std::io::Error>> + Send + 'a>>;
/// 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,
}
/// The persistable columns of an `external_mounts` row (admin create).
#[derive(Debug, Clone)]
pub struct NewExternalMount {
/// The mount-root folder UUID this mount attaches to.
pub mount_folder_id: Uuid,
/// Provider kind.
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,
}
/// 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>;
/// Insert a new mount row. Default errors — only the PG repo implements it
/// (test doubles need `list_all` only).
async fn create(&self, _mount: &NewExternalMount) -> Result<(), DomainError> {
Err(DomainError::operation_not_supported(
"ExternalMount",
"create is not supported by this repository",
))
}
/// Delete a mount row by its mount-root folder id. Returns `true` when a
/// row was removed. Default errors (see `create`).
async fn delete(&self, _mount_folder_id: Uuid) -> Result<bool, DomainError> {
Err(DomainError::operation_not_supported(
"ExternalMount",
"delete is not supported by this repository",
))
}
}
/// 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>;
}
+1
View File
@@ -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;
@@ -101,10 +101,16 @@ mod tests {
None,
authz.clone(),
));
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,
Arc::new(FileLifecycleService::new()),
mount_router,
));
let _batch_service = BatchOperationService::new(
@@ -0,0 +1,192 @@
//! 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 { .. }
)
}
/// Path-based lookup for the protocol surfaces (WebDAV / NextCloud): does
/// `internal_path` descend into a mount within `drive_id`? Returns the mount
/// config plus the remainder relpath (empty when the path IS the mount root).
pub fn find_path(
&self,
drive_id: uuid::Uuid,
internal_path: &str,
) -> Option<(Arc<MountConfig>, String)> {
self.registry.find_mount_for_path(drive_id, internal_path)
}
}
#[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
));
}
}
@@ -0,0 +1,55 @@
//! Streams an upload straight to an external mount's provider, bypassing the
//! content-addressable store entirely (no BLAKE3 / dedup).
//!
//! The REST upload handler detects a mount destination BEFORE ingesting into the
//! CAS and routes here. Authorization stays in this service (the mount-root
//! `Create` grant); the handler only classifies and supplies the body stream.
use std::sync::Arc;
use crate::application::dtos::file_dto::FileDto;
use crate::application::ports::authorization_ports::AuthorizationEngine;
use crate::application::ports::external_mount_ports::MountByteStream;
use crate::application::services::mount_dto::{audit_mount_write, mount_file_dto, mount_parent_id};
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::services::pg_acl_engine::PgAclEngine;
use uuid::Uuid;
/// Writes uploaded bytes to a mount provider with authorization + auditing.
pub struct ExternalUploadService {
authz: Arc<PgAclEngine>,
}
impl ExternalUploadService {
/// Construct over the ReBAC engine.
pub fn new(authz: Arc<PgAclEngine>) -> Self {
Self { authz }
}
/// Authorize (`Create` on the mount root) then stream `body` to the provider
/// as `name` under `parent_node`. Returns the synthesized `FileDto`.
pub async fn write_file(
&self,
cfg: &MountConfig,
parent_node: &NodeId,
name: &str,
body: MountByteStream<'_>,
caller_id: Uuid,
) -> Result<FileDto, DomainError> {
self.authz
.require(
Subject::User(caller_id),
Permission::Create,
Resource::Folder(cfg.mount_id),
)
.await?;
let stat = cfg.provider.write_stream(parent_node, name, body).await?;
audit_mount_write("upload", cfg, caller_id, stat.node_id.as_str());
let parent = mount_parent_id(cfg, stat.node_id.as_str());
Ok(mount_file_dto(cfg, &parent, &stat))
}
}
@@ -7,9 +7,13 @@ use crate::application::ports::file_ports::FileManagementUseCase;
use crate::application::ports::resource_access_hook::ResourceAccessHook;
use crate::application::ports::storage_ports::{CopyFolderTreeResult, FileWritePort};
use crate::application::ports::trash_ports::TrashUseCase;
use crate::application::services::external_mount_router::{MountRouter, ResolvedId};
use crate::application::services::mount_dto::{audit_mount_write, mount_file_dto, mount_parent_id};
use crate::application::services::mount_registry::MountConfig;
use crate::application::services::trash_service::TrashService;
use crate::common::errors::DomainError;
use crate::domain::services::authorization::{Permission, Resource, Subject};
use crate::domain::services::external_mount_id::NodeId;
use crate::domain::services::path_service::validate_storage_name;
use crate::infrastructure::repositories::pg::file_blob_read_repository::FileBlobReadRepository;
use crate::infrastructure::repositories::pg::file_blob_write_repository::FileBlobWriteRepository;
@@ -32,6 +36,9 @@ pub struct FileManagementService {
authz: Arc<PgAclEngine>,
/// Lifecycle hook dispatcher — fired on file created (copy) and deleted.
file_lifecycle_hook: Option<Arc<dyn FileLifecycleHook>>,
/// External-mount classifier. `None` in stub/test construction → all ids
/// are treated as native.
mount_router: Option<Arc<MountRouter>>,
/// Read/write access hook — fired so Recent reflects "this is the file
/// I just copied / renamed / moved", same way the read paths surface
/// downloads. Distinct from the lifecycle hook because lifecycle hooks
@@ -72,6 +79,7 @@ impl FileManagementService {
content_cache,
authz,
file_lifecycle_hook: None,
mount_router: None,
resource_access_hook: None,
drive_repo: None,
storage_usage: None,
@@ -84,6 +92,59 @@ impl FileManagementService {
self
}
/// Injects the external-mount classifier so file mutations can branch
/// `ext:` ids to the provider.
pub fn with_mount_router(mut self, router: Arc<MountRouter>) -> Self {
self.mount_router = Some(router);
self
}
/// Classify an id via the mount router (if configured). Returns `Regular`
/// when no router is wired.
fn classify(&self, id: &str) -> ResolvedId {
match &self.mount_router {
Some(r) => r.classify(id),
None => ResolvedId::Regular,
}
}
/// Authorize a mutation inside a mount (gates on the mount-root folder).
async fn require_mount_perm(
&self,
cfg: &MountConfig,
perm: Permission,
caller_id: Uuid,
) -> Result<(), DomainError> {
self.authz
.require(
Subject::User(caller_id),
perm,
Resource::Folder(cfg.mount_id),
)
.await
}
/// Resolve a move destination within the same mount as `cfg`. Errors when
/// the destination is absent, native, or in a different mount.
fn mount_dest_node(
&self,
cfg: &MountConfig,
folder_id: Option<&str>,
) -> Result<NodeId, DomainError> {
let Some(folder_id) = folder_id else {
return Err(cross_boundary_move_err());
};
match self.classify(folder_id) {
ResolvedId::MountRoot { cfg: dest } if dest.mount_id == cfg.mount_id => {
Ok(NodeId::default())
}
ResolvedId::MountChild { cfg: dest, node_id } if dest.mount_id == cfg.mount_id => {
Ok(node_id)
}
_ => Err(cross_boundary_move_err()),
}
}
/// Registers the read/write access hook (Recent list recorder).
pub fn with_resource_access_hook(mut self, hook: Arc<dyn ResourceAccessHook>) -> Self {
self.resource_access_hook = Some(hook);
@@ -316,6 +377,29 @@ impl FileManagementUseCase for FileManagementService {
caller_id: Uuid,
folder_id: Option<String>,
) -> Result<FileDto, DomainError> {
// External mount: moves stay within one mount; cross-backend is forbidden.
match self.classify(file_id) {
ResolvedId::Regular => {
if let Some(dst) = folder_id.as_deref()
&& !matches!(self.classify(dst), ResolvedId::Regular)
{
return Err(cross_boundary_move_err());
}
}
ResolvedId::MountRoot { .. } => return Err(DomainError::not_found("File", file_id)),
ResolvedId::MountChild { cfg, node_id } => {
let dest = self.mount_dest_node(&cfg, folder_id.as_deref())?;
self.require_mount_perm(&cfg, Permission::Update, caller_id)
.await?;
self.require_mount_perm(&cfg, Permission::Create, caller_id)
.await?;
let stat = cfg.provider.move_within(&node_id, &dest).await?;
audit_mount_write("move", &cfg, caller_id, stat.node_id.as_str());
let parent = mount_parent_id(&cfg, stat.node_id.as_str());
return Ok(mount_file_dto(&cfg, &parent, &stat));
}
}
// Move = Update on the file + Create on the target folder (if any).
self.require_file_perm(file_id, Permission::Update, caller_id)
.await?;
@@ -460,12 +544,32 @@ impl FileManagementUseCase for FileManagementService {
caller_id: Uuid,
new_name: &str,
) -> Result<FileDto, DomainError> {
if let ResolvedId::MountChild { cfg, node_id } = self.classify(file_id) {
if let Err(reason) = validate_storage_name(new_name) {
return Err(DomainError::validation_error(format!(
"Invalid file name '{new_name}': {reason}"
)));
}
self.require_mount_perm(&cfg, Permission::Update, caller_id)
.await?;
let stat = cfg.provider.rename(&node_id, new_name).await?;
audit_mount_write("rename", &cfg, caller_id, stat.node_id.as_str());
let parent = mount_parent_id(&cfg, stat.node_id.as_str());
return Ok(mount_file_dto(&cfg, &parent, &stat));
}
self.require_file_perm(file_id, Permission::Update, caller_id)
.await?;
self.rename_file(file_id, new_name, caller_id).await
}
async fn delete_file_with_perms(&self, id: &str, caller_id: Uuid) -> Result<(), DomainError> {
if let ResolvedId::MountChild { cfg, node_id } = self.classify(id) {
self.require_mount_perm(&cfg, Permission::Delete, caller_id)
.await?;
cfg.provider.delete(&node_id).await?;
audit_mount_write("delete", &cfg, caller_id, node_id.as_str());
return Ok(());
}
self.require_file_perm(id, Permission::Delete, caller_id)
.await?;
self.delete_file(id).await
@@ -482,6 +586,15 @@ impl FileManagementUseCase for FileManagementService {
id: &str,
caller_id: Uuid,
) -> Result<bool, DomainError> {
// External mount: permanent provider delete (mounts have no trash).
if let ResolvedId::MountChild { cfg, node_id } = self.classify(id) {
self.require_mount_perm(&cfg, Permission::Delete, caller_id)
.await?;
cfg.provider.delete(&node_id).await?;
audit_mount_write("delete", &cfg, caller_id, node_id.as_str());
return Ok(false); // permanently deleted (no trash)
}
self.require_file_perm(id, Permission::Delete, caller_id)
.await?;
// Step 1: Try trash (soft delete — file row stays, blob stays referenced)
@@ -552,3 +665,12 @@ impl FileManagementUseCase for FileManagementService {
.await
}
}
/// Error for a move/copy that would cross a storage backend boundary
/// (mount ↔ native, or between two different mounts). Forbidden in v1.
fn cross_boundary_move_err() -> DomainError {
DomainError::operation_not_supported(
"File",
"moving between external mounts and regular storage is not supported",
)
}
@@ -5,13 +5,17 @@ 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, RangeContent,
};
use crate::application::ports::resource_access_hook::ResourceAccessHook;
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::{
@@ -36,6 +40,9 @@ pub struct FileRetrievalService {
content_cache: Option<Arc<FileContentCache>>,
transcode: Option<Arc<ImageTranscodeService>>,
authz: Option<Arc<PgAclEngine>>,
/// External-mount classifier for path-based resolution (WebDAV/NextCloud).
/// `None` in the simple/test constructor → no mount support.
mount_router: Option<Arc<crate::application::services::external_mount_router::MountRouter>>,
/// Optional read-event observer. Currently fans out to the Recent-list
/// recorder; future observers (audit trail, "last seen by", …) attach
/// to the same hook so service code only knows the trait, not the impl.
@@ -53,6 +60,7 @@ impl FileRetrievalService {
content_cache: None,
transcode: None,
authz: None,
mount_router: None,
resource_access_hook: None,
}
}
@@ -70,10 +78,21 @@ impl FileRetrievalService {
content_cache: Some(content_cache),
transcode: Some(transcode),
authz: Some(authz),
mount_router: None,
resource_access_hook: None,
}
}
/// Injects the external-mount classifier so path-based lookups
/// (`get_file_by_path`) can resolve mount paths to the provider.
pub fn with_mount_router(
mut self,
router: Arc<crate::application::services::external_mount_router::MountRouter>,
) -> Self {
self.mount_router = Some(router);
self
}
/// Builder: attach a [`ResourceAccessHook`] that fires after every
/// authorised `_with_perms` read. Without it the service is silent —
/// existing behaviour for stub / test paths.
@@ -82,6 +101,24 @@ impl FileRetrievalService {
self
}
/// 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),
mount_router: None,
resource_access_hook: None,
}
}
/// Fire the access hook if registered. Called from every `_with_perms`
/// read after the authZ + lookup has succeeded (never on failure
/// paths — denied reads must not surface in Recent).
@@ -176,6 +213,66 @@ 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
}
/// If `id` is an `ext:` mount FILE id, return the mount config + node id.
/// `None` for native ids, mount roots, or when no router is wired.
fn mount_file_node(
&self,
id: &str,
) -> Option<(
Arc<crate::application::services::mount_registry::MountConfig>,
NodeId,
)> {
use crate::application::services::external_mount_router::ResolvedId;
match self.mount_router.as_ref()?.classify(id) {
ResolvedId::MountChild { cfg, node_id } => Some((cfg, node_id)),
_ => None,
}
}
/// Try to transcode image content to WebP and return transcoded variant.
async fn try_transcode(
&self,
@@ -451,6 +548,26 @@ impl FileRetrievalUseCase for FileRetrievalService {
// `drive_id` scope axis prevents cross-drive resolution — without
// it, `find_file_by_path` would return a non-deterministic row
// when the same path exists in multiple drives.
// External mount: a path descending past a mount root resolves on the
// provider (stat). The mount root itself has no file at its path.
if let Some(router) = &self.mount_router
&& let Some((cfg, remainder)) = router.find_path(drive_id, path)
&& !remainder.is_empty()
{
let node = cfg.provider.resolve_path(&remainder);
let stat = cfg.provider.stat(&node).await?;
if stat.is_dir {
return Err(DomainError::not_found("File", path));
}
let parent = crate::application::services::mount_dto::mount_parent_id(
&cfg,
stat.node_id.as_str(),
);
return Ok(crate::application::services::mount_dto::mount_file_dto(
&cfg, &parent, &stat,
));
}
if let Some(file) = self.file_read.find_file_by_path(path, drive_id).await? {
return Ok(FileDto::from(file));
}
@@ -488,6 +605,11 @@ impl FileRetrievalUseCase for FileRetrievalService {
&self,
id: &str,
) -> Result<Box<dyn Stream<Item = Result<Bytes, std::io::Error>> + Send>, DomainError> {
if let Some((cfg, node)) = self.mount_file_node(id) {
let s = cfg.provider.open_read_stream(&node, None).await?;
// `Pin<Box<dyn Stream>>` is itself a `Stream`, so re-box it.
return Ok(Box::new(s));
}
self.file_read.get_file_stream(id).await
}
@@ -548,6 +670,13 @@ impl FileRetrievalUseCase for FileRetrievalService {
start: u64,
end: Option<u64>,
) -> Result<Box<dyn Stream<Item = Result<Bytes, std::io::Error>> + Send>, DomainError> {
if let Some((cfg, node)) = self.mount_file_node(id) {
// The native range convention is exclusive-end; the provider wants
// an inclusive end.
let range = Some((start, end.map(|e| e.saturating_sub(1))));
let s = cfg.provider.open_read_stream(&node, range).await?;
return Ok(Box::new(s));
}
self.file_read.get_file_range_stream(id, start, end).await
}
@@ -596,6 +725,58 @@ impl FileRetrievalUseCase for FileRetrievalService {
after_name: Option<&str>,
limit: i64,
) -> Result<Vec<FileDto>, DomainError> {
// External mount: list files from the provider (WebDAV/NextCloud
// PROPFIND Depth:1 file loop). Authz collapses on the mount root.
// Keyset pagination by name mirrors `paginate_mount_entries` —
// provider order isn't guaranteed, so sort before slicing on
// `after_name`.
if let Some(fid) = folder_id
&& let Some(router) = &self.mount_router
{
use crate::application::services::external_mount_router::ResolvedId;
let resolved = match router.classify(fid) {
ResolvedId::Regular => None,
ResolvedId::MountRoot { cfg } => Some((cfg, NodeId::default())),
ResolvedId::MountChild { cfg, node_id } => Some((cfg, node_id)),
};
if let Some((cfg, node)) = resolved {
if let Some(authz) = &self.authz {
authz
.require(
Subject::User(owner_id),
Permission::Read,
Resource::Folder(cfg.mount_id),
)
.await?;
}
let mut entries: Vec<_> = cfg
.provider
.list_dir(&node)
.await?
.into_iter()
.filter(|e| !e.is_dir)
.collect();
entries.sort_by_key(|e| e.name.to_lowercase());
let start = match after_name {
Some(name) => entries
.iter()
.position(|e| name.eq_ignore_ascii_case(&e.name))
.map(|i| i + 1)
.unwrap_or(0),
None => 0,
};
let files: Vec<FileDto> = entries
.into_iter()
.skip(start)
.take(limit.max(0) as usize)
.map(|e| {
crate::application::services::mount_dto::mount_entry_file_dto(&cfg, fid, &e)
})
.collect();
return Ok(files);
}
}
// Post-D0: every file lives in a folder — `storage.files.folder_id`
// is NOT NULL. `folder_id = None` means the caller is asking for
// "root-level files", which by design return an empty set: the
File diff suppressed because it is too large Load Diff
+4
View File
@@ -9,6 +9,8 @@ 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 external_upload_service;
pub mod favorites_service;
pub mod file_lifecycle_service;
pub mod file_management_service;
@@ -18,6 +20,8 @@ pub mod file_use_case_factory;
pub mod folder_service;
pub mod i18n_application_service;
pub mod magic_link_invite_service;
pub mod mount_dto;
pub mod mount_registry;
pub mod music_service;
pub mod nextcloud_file_id_service;
pub mod nextcloud_login_flow_service;
+145
View File
@@ -0,0 +1,145 @@
//! Builders that synthesize `FolderDto` / `FileDto` from a provider [`MountStat`].
//!
//! Mount entries have no `storage.folders`/`storage.files` row, so the normal
//! `FolderDto::from(Folder)` path doesn't apply. These helpers produce the same
//! DTO shape from a provider stat plus the mount config, with a synthetic `ext:`
//! id and a virtual etag. Shared by the folder/file services and the handlers.
use std::sync::Arc;
use uuid::Uuid;
use crate::application::dtos::display_helpers::{
category_for, format_file_size, icon_class_for, icon_special_class_for,
};
use crate::application::dtos::file_dto::FileDto;
use crate::application::dtos::folder_dto::FolderDto;
use crate::application::ports::external_mount_ports::{MountEntry, MountStat};
use crate::application::services::mount_registry::MountConfig;
use crate::domain::services::external_mount_id::{
encode_child_id, virtual_file_etag, virtual_folder_etag,
};
/// Final path segment of a node id (the display name).
fn node_name(node_id: &str) -> &str {
node_id.rsplit('/').next().unwrap_or(node_id)
}
/// Emit the structured audit line for a mount mutation (per AGENTS.md). Every
/// write op (upload / mkdir / rename / delete / move) calls this.
pub fn audit_mount_write(action: &str, cfg: &MountConfig, caller_id: Uuid, node_id: &str) {
tracing::info!(
target: "audit",
event = "external_mount.write",
action,
mount_id = %cfg.mount_id,
caller_id = %caller_id,
node_id = %node_id,
reason = "external_mount_op",
"👮🏻‍♂️ external mount mutation",
);
}
/// The id-string of a mount entry's parent: the parent's `ext:` id, or the
/// mount-root folder UUID when the entry is a direct child of the root.
pub fn mount_parent_id(cfg: &MountConfig, node_id: &str) -> String {
match node_id.rsplit_once('/') {
Some((parent, _)) => encode_child_id(cfg.mount_id, parent),
None => cfg.mount_id.to_string(),
}
}
/// Build a `FolderDto` for a mount directory from its stat. `parent_id` is the
/// id-string of the containing directory (mount-root UUID or an `ext:` id).
pub fn mount_folder_dto(cfg: &MountConfig, parent_id: &str, stat: &MountStat) -> FolderDto {
FolderDto {
etag: virtual_folder_etag(stat.modified_at),
id: encode_child_id(cfg.mount_id, stat.node_id.clone()),
name: node_name(stat.node_id.as_str()).to_owned(),
path: String::new(),
parent_id: Some(parent_id.to_owned()),
drive_id: cfg.drive_id,
created_at: stat.created_at,
modified_at: stat.modified_at,
is_root: false,
icon_class: Arc::from("fas fa-folder"),
icon_special_class: Arc::from("folder-icon"),
category: Arc::from("Folder"),
created_by: Some(cfg.owner_id),
updated_by: Some(cfg.owner_id),
}
}
/// Build a `FolderDto` from a directory listing entry. `parent_id` is the
/// id-string of the directory being listed.
pub fn mount_entry_folder_dto(cfg: &MountConfig, parent_id: &str, entry: &MountEntry) -> FolderDto {
FolderDto {
etag: virtual_folder_etag(entry.modified_at),
id: encode_child_id(cfg.mount_id, entry.node_id.clone()),
name: entry.name.clone(),
path: String::new(),
parent_id: Some(parent_id.to_owned()),
drive_id: cfg.drive_id,
created_at: entry.created_at,
modified_at: entry.modified_at,
is_root: false,
icon_class: Arc::from("fas fa-folder"),
icon_special_class: Arc::from("folder-icon"),
category: Arc::from("Folder"),
created_by: Some(cfg.owner_id),
updated_by: Some(cfg.owner_id),
}
}
/// Build a `FileDto` from a directory listing entry (mime sniffed from name).
pub fn mount_entry_file_dto(cfg: &MountConfig, parent_id: &str, entry: &MountEntry) -> FileDto {
let name = entry.name.as_str();
let mime = mime_guess::from_path(name)
.first_or_octet_stream()
.to_string();
FileDto {
id: encode_child_id(cfg.mount_id, entry.node_id.clone()),
name: name.to_owned(),
path: String::new(),
size: entry.size,
mime_type: Arc::from(mime.as_str()),
folder_id: Some(parent_id.to_owned()),
created_at: entry.created_at,
modified_at: entry.modified_at,
icon_class: Arc::from(icon_class_for(name, &mime)),
icon_special_class: Arc::from(icon_special_class_for(name, &mime)),
category: Arc::from(category_for(name, &mime)),
size_formatted: format_file_size(entry.size),
sort_date: None,
content_hash: String::new(),
etag: virtual_file_etag(entry.size, entry.modified_at),
created_by: Some(cfg.owner_id),
updated_by: Some(cfg.owner_id),
}
}
/// Build a `FileDto` for a mount file from its stat. Virtual files have no blob
/// hash (`content_hash` empty) and a size+mtime etag.
pub fn mount_file_dto(cfg: &MountConfig, parent_id: &str, stat: &MountStat) -> FileDto {
let name = node_name(stat.node_id.as_str());
let mime = stat.mime_type.as_str();
FileDto {
id: encode_child_id(cfg.mount_id, stat.node_id.clone()),
name: name.to_owned(),
path: String::new(),
size: stat.size,
mime_type: Arc::from(mime),
folder_id: Some(parent_id.to_owned()),
created_at: stat.created_at,
modified_at: stat.modified_at,
icon_class: Arc::from(icon_class_for(name, mime)),
icon_special_class: Arc::from(icon_special_class_for(name, mime)),
category: Arc::from(category_for(name, mime)),
size_formatted: format_file_size(stat.size),
sort_date: None,
content_hash: String::new(),
etag: virtual_file_etag(stat.size, stat.modified_at),
created_by: Some(cfg.owner_id),
updated_by: Some(cfg.owner_id),
}
}
+387
View File
@@ -0,0 +1,387 @@
//! 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();
// Normalize away a leading slash: materialized folder paths arrive both
// as `Personal/Media` (raw `folders.path`) and `/Personal/Media`
// (FolderDto / WebDAV internal paths). The index keys are stored without
// a leading slash (see `reload`).
let internal_path = internal_path.trim_start_matches('/');
// 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;
}
};
// Store the path key without a leading slash so lookups normalize
// consistently (see `find_mount_for_path`).
by_path.insert(
(
rec.drive_id,
rec.mount_path.trim_start_matches('/').to_string(),
),
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");
}
}
+14
View File
@@ -1072,6 +1072,12 @@ pub struct FeaturesConfig {
/// thumbnail through the same WebP pipeline as photos; otherwise videos have
/// no thumbnail. Env: `OXICLOUD_ENABLE_VIDEO_THUMBNAILS`.
pub enable_video_thumbnails: bool,
/// Expose admin-configured external filesystem mounts (raw host fs, …) as
/// folders inside a user's drive. Contents are read live from the backend
/// and are a deliberately limited, separate storage type (no dedup/sharing/
/// trash/search). OFF by default — opt-in per deployment.
/// Env: `OXICLOUD_ENABLE_EXTERNAL_MOUNTS`.
pub enable_external_mounts: bool,
/// Expose `/api/admin/internal/*` test-only endpoints that trigger
/// background sweeps on demand (storage-usage reconciliation, blob
/// GC). Intended for Hurl / integration tests that need to wait
@@ -1156,6 +1162,7 @@ impl Default for FeaturesConfig {
enable_faces: false, // People/faces (biometric) — opt-in, off by default
expose_system_users: true, // Expose OxiCloud users as address book by default
enable_video_thumbnails: true, // Video thumbs via ffmpeg (if detected)
enable_external_mounts: false, // External mounts — opt-in, off by default
// Test-only sweep triggers — strictly opt-in. Production
// deployments do NOT need this; the periodic ticker handles
// reconciliation transparently.
@@ -1904,6 +1911,13 @@ impl AppConfig {
config.features.enable_faces = val;
}
if let Ok(enable_external_mounts) =
env::var("OXICLOUD_ENABLE_EXTERNAL_MOUNTS").map(|v| v.parse::<bool>())
&& let Ok(val) = enable_external_mounts
{
config.features.enable_external_mounts = val;
}
// Faces (People) ONNX runtime + models — operator-provided at runtime.
if let Ok(v) = env::var("OXICLOUD_FACES_ORT_DYLIB").or_else(|_| env::var("ORT_DYLIB_PATH"))
&& !v.is_empty()
+46 -1
View File
@@ -522,6 +522,7 @@ impl AppServiceFactory {
plugin_dispatch: Option<
Arc<dyn crate::application::ports::plugin_ports::PluginDispatchPort>,
>,
mount_router: Arc<crate::application::services::external_mount_router::MountRouter>,
resource_access_hook: Option<
Arc<dyn crate::application::ports::resource_access_hook::ResourceAccessHook>,
>,
@@ -535,6 +536,7 @@ impl AppServiceFactory {
// `delete_folder_with_perms` fans out to the same handlers
// (thumbnails, metadata, …) as a single-file delete.
core.file_lifecycle.clone(),
mount_router.clone(),
)
// D5 cross-drive move gate reads policies via the same
// drive repo every other policy uses. Wired here so
@@ -558,7 +560,8 @@ impl AppServiceFactory {
core.file_content_cache.clone(),
core.image_transcode_service.clone(),
authz.clone(),
);
)
.with_mount_router(mount_router.clone());
if let Some(hook) = resource_access_hook.clone() {
svc = svc.with_resource_access_hook(hook);
}
@@ -623,6 +626,7 @@ impl AppServiceFactory {
authz.clone(),
)
.with_file_lifecycle_hook(file_lifecycle.clone())
.with_mount_router(mount_router.clone())
// D5 cross-drive move gate reads policies via the same
// drive repo every other policy uses. Wired here so
// `move_file_with_perms` can enforce `forbid_cross_drive_move`
@@ -637,6 +641,13 @@ impl AppServiceFactory {
svc
});
// Streams uploads to external mount providers (bypasses the CAS).
let external_upload_service = Arc::new(
crate::application::services::external_upload_service::ExternalUploadService::new(
authz.clone(),
),
);
let file_use_case_factory = Arc::new(AppFileUseCaseFactory::new(
repos.file_read_repository.clone(),
repos.file_write_repository.clone(),
@@ -675,6 +686,7 @@ impl AppServiceFactory {
delta_upload_service,
file_retrieval_service,
file_management_service,
external_upload_service,
file_use_case_factory,
i18n_service,
trash_service, // Already set via parameter
@@ -1278,6 +1290,27 @@ impl AppServiceFactory {
let (plugin_dispatch, plugin_management) = self.create_plugin_ports();
// 4. Application services (with trash + authz already wired)
// External mount registry + router. Built before the application
// services because `FolderService` holds the router to branch listing
// onto the provider. The router is always present (an empty registry is
// a cheap no-op); when the feature is enabled we load the configured
// mounts and build their providers up front. The registry is
// interior-mutable (arc-swap), so reloading here is visible to every
// holder of the shared router.
let mount_registry =
Arc::new(crate::application::services::mount_registry::MountRegistry::empty());
if self.config.features.enable_external_mounts {
let repo = crate::infrastructure::repositories::pg::ExternalMountPgRepository::new(
pool.clone(),
);
let factory =
crate::infrastructure::services::mount_provider_factory::DefaultMountProviderFactory::new();
mount_registry.reload(&repo, &factory).await;
}
let mount_router = Arc::new(
crate::application::services::external_mount_router::MountRouter::new(mount_registry),
);
let mut apps = self.create_application_services(
&core,
&repos,
@@ -1287,6 +1320,7 @@ impl AppServiceFactory {
&storage_usage,
content_index.as_ref().map(|(idx, _)| idx.clone()),
plugin_dispatch.clone(),
mount_router.clone(),
Some(resource_access_hook.clone()),
);
@@ -1641,6 +1675,7 @@ impl AppServiceFactory {
locale_registry: self.locale_registry.clone(),
db_pool: Some(pool.clone()),
maintenance_pool: Some(maintenance_pool),
mount_router,
auth_service: auth_services,
nextcloud: nextcloud_services,
admin_settings_service: None,
@@ -2068,6 +2103,9 @@ pub struct ApplicationServices {
Arc<crate::application::services::delta_upload_service::DeltaUploadService>,
pub file_retrieval_service: Arc<FileRetrievalService>,
pub file_management_service: Arc<FileManagementService>,
/// Streams uploads straight to an external mount provider (bypasses the CAS).
pub external_upload_service:
Arc<crate::application::services::external_upload_service::ExternalUploadService>,
pub file_use_case_factory: Arc<dyn FileUseCaseFactory>,
pub i18n_service: Arc<I18nApplicationService>,
pub trash_service: Option<Arc<TrashService>>,
@@ -2111,6 +2149,13 @@ pub struct AppState {
pub db_pool: Option<Arc<PgPool>>,
/// Isolated pool for background / batch operations.
pub maintenance_pool: Option<Arc<PgPool>>,
/// External-mount classifier + registry. Always present; an empty registry
/// (feature disabled or no mounts configured) makes `classify` a cheap no-op
/// that routes every id to native handling. Handlers consult this before
/// parsing an id as a UUID, then call the matching service-layer mount
/// method (which still owns the authorization check).
pub mount_router:
Arc<crate::application::services::external_mount_router::MountRouter>,
pub auth_service: Option<AuthServices>,
pub nextcloud: Option<NextcloudServices>,
pub admin_settings_service: Option<Arc<AdminSettingsService>>,
+258
View File
@@ -0,0 +1,258 @@
//! External mount identifiers (domain value objects).
//!
//! Files and folders *below* an external mount root have no database row. They
//! are addressed by a synthetic id that wraps a provider-owned node identity:
//!
//! ```text
//! ext:<mount_id>:<base64url(node_id)>
//! ```
//!
//! * `<mount_id>` is the mount-root folder's UUID (simple form, no hyphens) — the
//! only part the system interprets, used to find the provider in the registry.
//! * `<node_id>` is **assigned and owned by the provider** and is **opaque** to the
//! rest of the system. Most providers use the entry's path (`local_fs`, `sftp`,
//! `webdav`); a provider with a stronger stable handle (inode, object id, href)
//! may use that. It is base64url-encoded (no padding) so it survives URLs and
//! WebDAV hrefs and never collides with the `:` / `/` separators.
//!
//! The system NEVER parses or validates a `node_id` — it only splits the envelope
//! and hands the decoded bytes back to the provider verbatim.
use std::fmt;
use std::str::FromStr;
use base64::{Engine as _, engine::general_purpose::URL_SAFE_NO_PAD};
use uuid::Uuid;
/// The `ext:` scheme prefix that marks a synthetic external-mount id.
pub const EXTERNAL_ID_PREFIX: &str = "ext:";
/// A provider-owned, opaque node identity for one entry inside a mount.
///
/// The system treats this as opaque bytes. For path-based providers it is the
/// POSIX path relative to the mount root (no leading `/`).
#[derive(Debug, Clone, PartialEq, Eq, Hash, Default)]
pub struct NodeId(pub String);
impl NodeId {
/// Borrow the inner string.
pub fn as_str(&self) -> &str {
&self.0
}
/// Consume into the inner string.
pub fn into_string(self) -> String {
self.0
}
}
impl From<String> for NodeId {
fn from(s: String) -> Self {
NodeId(s)
}
}
impl From<&str> for NodeId {
fn from(s: &str) -> Self {
NodeId(s.to_owned())
}
}
impl fmt::Display for NodeId {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.write_str(&self.0)
}
}
/// The decoded parts of a synthetic external-mount child id.
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct MountChildId {
/// The mount-root folder UUID (envelope; system-interpreted).
pub mount_id: Uuid,
/// The provider-owned node identity (opaque payload).
pub node_id: NodeId,
}
impl MountChildId {
/// Build a child id from a mount root and a provider node id.
pub fn new(mount_id: Uuid, node_id: impl Into<NodeId>) -> Self {
Self {
mount_id,
node_id: node_id.into(),
}
}
}
/// Renders as `ext:<mount_id>:<base64url(node_id)>`.
impl fmt::Display for MountChildId {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
write!(
f,
"{EXTERNAL_ID_PREFIX}{}:{}",
self.mount_id.simple(),
URL_SAFE_NO_PAD.encode(self.node_id.0.as_bytes())
)
}
}
/// Error returned when a string is not a well-formed external-mount child id.
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ParseMountChildIdError;
impl fmt::Display for ParseMountChildIdError {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.write_str("not a valid external mount id (expected ext:<mount_id>:<token>)")
}
}
impl std::error::Error for ParseMountChildIdError {}
impl FromStr for MountChildId {
type Err = ParseMountChildIdError;
fn from_str(s: &str) -> Result<Self, Self::Err> {
// ext:<mount_id>:<token> — split into exactly three logical parts.
let rest = s
.strip_prefix(EXTERNAL_ID_PREFIX)
.ok_or(ParseMountChildIdError)?;
let (mount_part, token) = rest.split_once(':').ok_or(ParseMountChildIdError)?;
let mount_id = Uuid::parse_str(mount_part).map_err(|_| ParseMountChildIdError)?;
let bytes = URL_SAFE_NO_PAD
.decode(token)
.map_err(|_| ParseMountChildIdError)?;
let node = String::from_utf8(bytes).map_err(|_| ParseMountChildIdError)?;
Ok(MountChildId {
mount_id,
node_id: NodeId(node),
})
}
}
/// Cheap check: does this id use the external-mount scheme?
///
/// Used at the top of service methods to route `ext:` ids away from
/// `Uuid::parse_str` and the PostgreSQL repositories.
pub fn is_external_id(id: &str) -> bool {
id.starts_with(EXTERNAL_ID_PREFIX)
}
/// Encode a child id string from its parts.
pub fn encode_child_id(mount_id: Uuid, node_id: impl Into<NodeId>) -> String {
MountChildId::new(mount_id, node_id).to_string()
}
/// Parse a child id string into its parts, or `None` when not `ext:`-shaped.
pub fn parse_child_id(id: &str) -> Option<MountChildId> {
id.parse().ok()
}
/// ETag for a virtual file (no blob hash): `ext-{size:x}-{modified_at}`.
///
/// Combines size and mtime so it changes on any content edit without reading
/// the file. Mirrors the native `{blob_hash[..16]}-{modified_at}` shape closely
/// enough for conditional requests.
pub fn virtual_file_etag(size: u64, modified_at: u64) -> String {
format!("ext-{size:x}-{modified_at}")
}
/// ETag for a virtual folder: `ext-{modified_at}` (directory mtime).
pub fn virtual_folder_etag(modified_at: u64) -> String {
format!("ext-{modified_at}")
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn round_trips_simple_path() {
let mount = Uuid::new_v4();
let id = encode_child_id(mount, "docs/report.txt");
assert!(is_external_id(&id));
let parsed = parse_child_id(&id).expect("parse");
assert_eq!(parsed.mount_id, mount);
assert_eq!(parsed.node_id.as_str(), "docs/report.txt");
}
#[test]
fn round_trips_via_fromstr_display() {
let mount = Uuid::new_v4();
let child = MountChildId::new(mount, "a/b/c.bin");
let rendered = child.to_string();
let parsed: MountChildId = rendered.parse().expect("parse");
assert_eq!(child, parsed);
}
#[test]
fn token_survives_separators_and_unicode() {
// node ids containing ':' '/' and non-ascii must survive the envelope.
let mount = Uuid::new_v4();
let tricky = "weird: name/with:colons/café.txt";
let id = encode_child_id(mount, tricky);
// The encoded form must not be ambiguous to the splitter: exactly two
// colons (the `ext:` scheme and the `<mount_id>:` separator); the
// base64url token never contains `:` or `/`.
assert_eq!(id.matches(':').count(), 2, "only scheme + mount separators");
let parsed = parse_child_id(&id).expect("parse");
assert_eq!(parsed.node_id.as_str(), tricky);
}
#[test]
fn rejects_non_external_ids() {
assert!(parse_child_id("not-an-ext-id").is_none());
assert!(parse_child_id(&Uuid::new_v4().to_string()).is_none());
assert!(parse_child_id("ext:not-a-uuid:dG9rZW4").is_none());
assert!(!is_external_id(&Uuid::new_v4().to_string()));
}
#[test]
fn rejects_malformed_envelopes() {
// Missing the second colon (no token separator).
assert!(parse_child_id("ext:abc").is_none());
// `ext:` with a valid uuid but a non-base64url token.
let u = Uuid::new_v4().simple().to_string();
assert!(parse_child_id(&format!("ext:{u}:!!!not-base64!!!")).is_none());
// Empty string / bare scheme.
assert!(parse_child_id("").is_none());
assert!(parse_child_id("ext:").is_none());
// is_external_id is a pure prefix check.
assert!(is_external_id("ext:anything"));
assert!(!is_external_id("EXT:upper"));
}
#[test]
fn round_trips_empty_node_id() {
// The mount root's own node id is empty; it must still round-trip
// (e.g. if ever encoded), encoding to a token-less-but-present form.
let mount = Uuid::new_v4();
let id = encode_child_id(mount, "");
let parsed = parse_child_id(&id).expect("parse empty node");
assert_eq!(parsed.mount_id, mount);
assert_eq!(parsed.node_id.as_str(), "");
}
#[test]
fn round_trips_long_and_nested_paths() {
let mount = Uuid::new_v4();
let deep = "a/".repeat(64) + "leaf.bin";
let id = encode_child_id(mount, deep.clone());
assert_eq!(parse_child_id(&id).unwrap().node_id.as_str(), deep);
}
#[test]
fn node_id_conversions() {
assert_eq!(NodeId::from("x").as_str(), "x");
assert_eq!(NodeId::from(String::from("y")).into_string(), "y");
assert_eq!(NodeId::default().as_str(), "");
}
#[test]
fn etags_change_with_inputs() {
assert_ne!(virtual_file_etag(10, 100), virtual_file_etag(11, 100));
assert_ne!(virtual_file_etag(10, 100), virtual_file_etag(10, 101));
assert_ne!(virtual_folder_etag(100), virtual_folder_etag(101));
// Format is stable and documented.
assert_eq!(virtual_file_etag(255, 16), "ext-ff-16");
assert_eq!(virtual_folder_etag(42), "ext-42");
}
}
+1
View File
@@ -1,5 +1,6 @@
pub mod authorization;
pub mod email_normalize;
pub mod external_mount_id;
pub mod i18n_service;
pub mod path_service;
@@ -0,0 +1,214 @@
//! PostgreSQL persistence for external mount configuration.
//!
//! `list_all` (re)builds the in-memory registry; `create`/`delete` back the
//! admin CRUD endpoints.
use std::sync::Arc;
use async_trait::async_trait;
use sqlx::{PgPool, Row};
use uuid::Uuid;
use crate::application::ports::external_mount_ports::{
ExternalMountRecord, ExternalMountRepositoryPort, NewExternalMount,
};
use crate::domain::errors::DomainError;
/// PostgreSQL implementation of [`ExternalMountRepositoryPort`].
pub struct ExternalMountPgRepository {
pool: Arc<PgPool>,
}
impl ExternalMountPgRepository {
/// Construct over a connection pool.
pub fn new(pool: Arc<PgPool>) -> Self {
Self { pool }
}
}
#[async_trait]
impl ExternalMountRepositoryPort for ExternalMountPgRepository {
async fn list_all(&self) -> Result<Vec<ExternalMountRecord>, DomainError> {
// Join each mount to its (non-trashed) root folder to pick up the
// drive scope and the materialized path needed for path resolution.
let rows = sqlx::query(
r#"
SELECT
m.mount_folder_id AS mount_folder_id,
m.kind AS kind,
m.config AS config,
m.name AS name,
m.owner_id AS owner_id,
m.read_only AS read_only,
f.drive_id AS drive_id,
f.path AS mount_path
FROM storage.external_mounts m
JOIN storage.folders f ON f.id = m.mount_folder_id
WHERE NOT f.is_trashed
"#,
)
.fetch_all(self.pool.as_ref())
.await
.map_err(|e| DomainError::database_error(format!("failed to list external mounts: {e}")))?;
let mut out = Vec::with_capacity(rows.len());
for row in rows {
let drive_id: Option<Uuid> = row
.try_get("drive_id")
.map_err(|e| DomainError::database_error(format!("external mount row: {e}")))?;
let Some(drive_id) = drive_id else {
// A mount root without a drive shouldn't exist post-D0; skip safely.
tracing::warn!(
target: "oxicloud::external_mounts",
"skipping external mount with NULL drive_id"
);
continue;
};
out.push(ExternalMountRecord {
mount_folder_id: row
.try_get("mount_folder_id")
.map_err(|e| DomainError::database_error(format!("external mount row: {e}")))?,
kind: row
.try_get("kind")
.map_err(|e| DomainError::database_error(format!("external mount row: {e}")))?,
config: row
.try_get("config")
.map_err(|e| DomainError::database_error(format!("external mount row: {e}")))?,
name: row
.try_get("name")
.map_err(|e| DomainError::database_error(format!("external mount row: {e}")))?,
owner_id: row
.try_get("owner_id")
.map_err(|e| DomainError::database_error(format!("external mount row: {e}")))?,
read_only: row
.try_get("read_only")
.map_err(|e| DomainError::database_error(format!("external mount row: {e}")))?,
drive_id,
mount_path: row
.try_get("mount_path")
.map_err(|e| DomainError::database_error(format!("external mount row: {e}")))?,
});
}
Ok(out)
}
async fn create(&self, mount: &NewExternalMount) -> Result<(), DomainError> {
sqlx::query(
"INSERT INTO storage.external_mounts
(mount_folder_id, kind, config, name, owner_id, read_only)
VALUES ($1, $2, $3, $4, $5, $6)",
)
.bind(mount.mount_folder_id)
.bind(&mount.kind)
.bind(&mount.config)
.bind(&mount.name)
.bind(mount.owner_id)
.bind(mount.read_only)
.execute(self.pool.as_ref())
.await
.map_err(|e| {
DomainError::database_error(format!("failed to create external mount: {e}"))
})?;
Ok(())
}
async fn delete(&self, mount_folder_id: Uuid) -> Result<bool, DomainError> {
let res = sqlx::query("DELETE FROM storage.external_mounts WHERE mount_folder_id = $1")
.bind(mount_folder_id)
.execute(self.pool.as_ref())
.await
.map_err(|e| {
DomainError::database_error(format!("failed to delete external mount: {e}"))
})?;
Ok(res.rows_affected() > 0)
}
}
// Gated on `test` too: the module uses the `testcontainers` dev-dependency,
// which is only linked into test targets — a plain `--cfg integration_tests`
// lib build (e.g. clippy's lib pass) must not try to compile it.
#[cfg(all(test, integration_tests))]
mod integration_tests {
use super::*;
use crate::mount_it_support::{fresh_db, insert_mount, provision_folder};
#[tokio::test]
async fn list_all_returns_mount_joined_with_folder() {
let (_c, pool) = fresh_db().await;
let p = provision_folder(&pool, "mountowner", "Media").await;
insert_mount(&pool, &p, "/srv/media").await;
let repo = ExternalMountPgRepository::new(pool.clone());
let mounts = repo.list_all().await.expect("list_all");
assert_eq!(mounts.len(), 1);
let m = &mounts[0];
assert_eq!(m.mount_folder_id, p.mount_folder_id);
assert_eq!(m.kind, "local_fs");
assert_eq!(m.owner_id, p.owner_id);
assert_eq!(m.drive_id, p.drive_id);
assert!(!m.read_only);
// The joined folder path (drive-scoped materialized path) contains the
// mount folder's name.
assert!(
m.mount_path.contains("Media"),
"mount_path was {:?}",
m.mount_path
);
assert_eq!(m.config["path"], "/srv/media");
}
#[tokio::test]
async fn list_all_skips_trashed_mount_folder() {
let (_c, pool) = fresh_db().await;
let p = provision_folder(&pool, "mountowner", "Media").await;
insert_mount(&pool, &p, "/srv/media").await;
// Soft-delete the mount-root folder; the join filters NOT is_trashed.
sqlx::query("UPDATE storage.folders SET is_trashed = true WHERE id = $1")
.bind(p.mount_folder_id)
.execute(pool.as_ref())
.await
.unwrap();
let repo = ExternalMountPgRepository::new(pool.clone());
let mounts = repo.list_all().await.expect("list_all");
assert!(mounts.is_empty());
}
#[tokio::test]
async fn list_all_empty_when_no_mounts() {
let (_c, pool) = fresh_db().await;
let repo = ExternalMountPgRepository::new(pool.clone());
assert!(repo.list_all().await.expect("list_all").is_empty());
}
#[tokio::test]
async fn create_and_delete_round_trip() {
use crate::application::ports::external_mount_ports::NewExternalMount;
let (_c, pool) = fresh_db().await;
let p = provision_folder(&pool, "owner", "Media").await;
let repo = ExternalMountPgRepository::new(pool.clone());
repo.create(&NewExternalMount {
mount_folder_id: p.mount_folder_id,
kind: "local_fs".to_string(),
config: serde_json::json!({ "path": "/srv/x" }),
name: "Media".to_string(),
owner_id: p.owner_id,
read_only: true,
})
.await
.expect("create");
let mounts = repo.list_all().await.expect("list");
assert_eq!(mounts.len(), 1);
assert!(mounts[0].read_only);
assert_eq!(mounts[0].mount_folder_id, p.mount_folder_id);
assert!(repo.delete(p.mount_folder_id).await.expect("delete"));
assert!(repo.list_all().await.expect("list").is_empty());
// Deleting a non-existent mount returns false.
assert!(!repo.delete(p.mount_folder_id).await.expect("delete again"));
}
}
@@ -7,6 +7,7 @@ mod contact_persistence_dto;
mod contact_pg_repository;
mod device_code_pg_repository;
mod drive_pg_repository;
mod external_mount_repository;
mod face_pg_repository;
mod favorites_pg_repository;
pub mod file_metadata_repository;
@@ -36,6 +37,7 @@ pub use contact_persistence_dto::*;
pub use contact_pg_repository::ContactPgRepository;
pub use device_code_pg_repository::DeviceCodePgRepository;
pub use drive_pg_repository::DrivePgRepository;
pub use external_mount_repository::ExternalMountPgRepository;
pub use face_pg_repository::FacePgRepository;
pub use favorites_pg_repository::FavoritesPgRepository;
pub use file_blob_read_repository::FileBlobReadRepository;
File diff suppressed because it is too large Load Diff
+2
View File
@@ -16,11 +16,13 @@ pub mod grant_cleanup_service;
pub mod image_transcode_service;
pub mod jwt_service;
pub mod local_blob_backend;
pub mod local_fs_mount_provider;
pub mod login_lockout_service;
pub mod media_metadata_service;
pub mod migration_blob_backend;
pub mod migration_job;
pub mod mock_email_sender;
pub mod mount_provider_factory;
pub mod nextcloud_chunked_upload_service;
pub mod noop_face_analyzer;
pub mod oidc_service;
@@ -0,0 +1,122 @@
//! The single place external mount provider kinds are registered.
//!
//! Adding a new backend (`sftp`, `webdav`, …) is: implement
//! [`ExternalMountProvider`](crate::application::ports::external_mount_ports::ExternalMountProvider)
//! and add one arm to [`DefaultMountProviderFactory::build`]. Nothing else in the
//! router / listing / authz / path-resolution layers changes.
use std::sync::Arc;
use async_trait::async_trait;
use crate::application::ports::external_mount_ports::{
ExternalMountProvider, MountProviderFactory,
};
use crate::domain::errors::DomainError;
use crate::infrastructure::services::local_fs_mount_provider::LocalFsMountProvider;
/// Default factory: knows the built-in provider kinds.
#[derive(Default)]
pub struct DefaultMountProviderFactory;
impl DefaultMountProviderFactory {
/// Construct the factory.
pub fn new() -> Self {
Self
}
}
#[async_trait]
impl MountProviderFactory for DefaultMountProviderFactory {
async fn build(
&self,
kind: &str,
config: &serde_json::Value,
) -> Result<Arc<dyn ExternalMountProvider>, DomainError> {
match kind {
"local_fs" => {
let path = config.get("path").and_then(|v| v.as_str()).ok_or_else(|| {
DomainError::validation_error(
"local_fs mount config requires a string \"path\"",
)
})?;
let read_only = config
.get("read_only")
.and_then(|v| v.as_bool())
.unwrap_or(false);
let provider = LocalFsMountProvider::new(path, read_only)?;
Ok(Arc::new(provider))
}
other => Err(DomainError::operation_not_supported(
"ExternalMount",
format!("unknown mount provider kind: {other}"),
)),
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::domain::errors::{DomainError, ErrorKind};
/// `Arc<dyn ExternalMountProvider>` isn't `Debug`, so `unwrap_err` won't
/// compile — extract the error by matching instead.
fn expect_err(r: Result<Arc<dyn ExternalMountProvider>, DomainError>) -> DomainError {
match r {
Ok(_) => panic!("expected an error"),
Err(e) => e,
}
}
#[tokio::test]
async fn builds_local_fs_provider_from_valid_config() {
let dir = tempfile::tempdir().unwrap();
let factory = DefaultMountProviderFactory::new();
let cfg = serde_json::json!({ "path": dir.path().to_str().unwrap() });
let provider = factory.build("local_fs", &cfg).await.expect("builds");
assert_eq!(provider.kind(), "local_fs");
}
#[tokio::test]
async fn local_fs_honours_read_only_flag() {
let dir = tempfile::tempdir().unwrap();
let factory = DefaultMountProviderFactory::new();
let cfg = serde_json::json!({ "path": dir.path().to_str().unwrap(), "read_only": true });
let provider = factory.build("local_fs", &cfg).await.unwrap();
assert!(provider.capabilities().read_only);
}
#[tokio::test]
async fn unknown_kind_is_unsupported() {
let factory = DefaultMountProviderFactory::new();
let err = expect_err(factory.build("sftp", &serde_json::json!({})).await);
assert_eq!(err.kind, ErrorKind::UnsupportedOperation);
}
#[tokio::test]
async fn local_fs_missing_path_is_validation_error() {
let factory = DefaultMountProviderFactory::new();
let err = expect_err(
factory
.build("local_fs", &serde_json::json!({ "read_only": true }))
.await,
);
assert_eq!(err.kind, ErrorKind::InvalidInput);
}
#[tokio::test]
async fn local_fs_nonexistent_path_errors() {
let factory = DefaultMountProviderFactory::new();
let err = expect_err(
factory
.build(
"local_fs",
&serde_json::json!({ "path": "/no/such/dir/xyz123" }),
)
.await,
);
// Propagated from LocalFsMountProvider::new (canonicalize failure).
assert_eq!(err.kind, ErrorKind::InternalError);
}
}
@@ -0,0 +1,214 @@
//! Admin CRUD for external file mounts (`/api/admin/external-mounts`).
//!
//! Creating a mount: validate the backend config, create a mount-root folder
//! under the admin's drive, insert the `external_mounts` row, then hot-reload
//! the in-memory registry. Deleting: remove the row + the folder and reload.
//! Every endpoint is admin-gated.
use std::sync::Arc;
use axum::{
Json,
extract::{Path, State},
http::{HeaderMap, StatusCode},
response::IntoResponse,
};
use serde::{Deserialize, Serialize};
use uuid::Uuid;
use crate::application::ports::external_mount_ports::{
ExternalMountRecord, ExternalMountRepositoryPort, MountProviderFactory, NewExternalMount,
};
use crate::common::di::AppState;
use crate::domain::repositories::drive_repository::DriveRepository;
use crate::domain::repositories::folder_repository::FolderRepository;
use crate::infrastructure::repositories::pg::ExternalMountPgRepository;
use crate::infrastructure::services::mount_provider_factory::DefaultMountProviderFactory;
use crate::interfaces::errors::AppError;
use crate::interfaces::middleware::admin::require_admin;
/// JSON view of a configured mount.
#[derive(Debug, Serialize)]
pub struct ExternalMountResponse {
pub mount_folder_id: String,
pub name: String,
pub kind: String,
pub owner_id: String,
pub read_only: bool,
pub drive_id: String,
pub mount_path: String,
pub config: serde_json::Value,
}
impl From<ExternalMountRecord> for ExternalMountResponse {
fn from(r: ExternalMountRecord) -> Self {
Self {
mount_folder_id: r.mount_folder_id.to_string(),
name: r.name,
kind: r.kind,
owner_id: r.owner_id.to_string(),
read_only: r.read_only,
drive_id: r.drive_id.to_string(),
mount_path: r.mount_path,
config: r.config,
}
}
}
/// Request body for creating a mount.
#[derive(Debug, Deserialize)]
pub struct CreateExternalMountRequest {
/// Display name (also the mount-root folder name).
pub name: String,
/// Absolute host path for the `local_fs` provider.
pub host_path: String,
/// Provider kind. Defaults to `local_fs`.
#[serde(default = "default_kind")]
pub kind: String,
/// When true, the mount refuses all mutations.
#[serde(default)]
pub read_only: bool,
}
fn default_kind() -> String {
"local_fs".to_string()
}
fn pool(state: &AppState) -> Result<Arc<sqlx::PgPool>, AppError> {
state
.db_pool
.clone()
.ok_or_else(|| AppError::internal_error("Database not available"))
}
/// `GET /api/admin/external-mounts` — list all configured mounts.
pub async fn list_external_mounts(
State(state): State<Arc<AppState>>,
headers: HeaderMap,
) -> Result<impl IntoResponse, AppError> {
require_admin(&state, &headers).await?;
let repo = ExternalMountPgRepository::new(pool(&state)?);
let mounts = repo
.list_all()
.await
.map_err(|e| AppError::internal_error(format!("list external mounts: {e}")))?;
let out: Vec<ExternalMountResponse> = mounts.into_iter().map(Into::into).collect();
Ok(Json(out))
}
/// `POST /api/admin/external-mounts` — create a mount in the admin's drive.
pub async fn create_external_mount(
State(state): State<Arc<AppState>>,
headers: HeaderMap,
Json(req): Json<CreateExternalMountRequest>,
) -> Result<impl IntoResponse, AppError> {
let (admin_id, _role) = require_admin(&state, &headers).await?;
if req.name.trim().is_empty() {
return Err(AppError::bad_request("Mount name must not be empty"));
}
// Build the provider config and validate it up front (path exists, etc.).
let config = serde_json::json!({ "path": req.host_path, "read_only": req.read_only });
let factory = DefaultMountProviderFactory::new();
factory
.build(&req.kind, &config)
.await
.map_err(|e| AppError::bad_request(format!("invalid mount configuration: {e}")))?;
// Create the mount-root folder under the admin's default drive root.
let drive = state
.drive_repo
.find_default_for_user(admin_id)
.await
.map_err(|e| AppError::internal_error(format!("find default drive: {e}")))?;
let root_folder_id = drive.drive.root_folder_id.to_string();
let folder = state
.repositories
.folder_repository
.create_folder(req.name.clone(), Some(root_folder_id), admin_id)
.await
.map_err(AppError::from)?;
let mount_folder_id =
Uuid::parse_str(folder.id()).map_err(|_| AppError::internal_error("bad folder id"))?;
let repo = ExternalMountPgRepository::new(pool(&state)?);
repo.create(&NewExternalMount {
mount_folder_id,
kind: req.kind.clone(),
config,
name: req.name.clone(),
owner_id: admin_id,
read_only: req.read_only,
})
.await
.map_err(|e| AppError::internal_error(format!("create mount row: {e}")))?;
// Hot-reload so the new mount is live immediately.
state.mount_router.registry().reload(&repo, &factory).await;
tracing::info!(
target: "audit",
event = "external_mount.config",
action = "create",
mount_id = %mount_folder_id,
caller_id = %admin_id,
kind = %req.kind,
reason = "external_mount_admin",
"👮🏻‍♂️ external mount created",
);
// Return the freshly created mount.
let created = repo
.list_all()
.await
.map_err(|e| AppError::internal_error(format!("reload mounts: {e}")))?
.into_iter()
.find(|m| m.mount_folder_id == mount_folder_id)
.map(ExternalMountResponse::from)
.ok_or_else(|| AppError::internal_error("created mount not found"))?;
Ok((StatusCode::CREATED, Json(created)))
}
/// `DELETE /api/admin/external-mounts/{id}` — remove a mount (and its root
/// folder). The host filesystem content is untouched.
pub async fn delete_external_mount(
State(state): State<Arc<AppState>>,
headers: HeaderMap,
Path(id): Path<Uuid>,
) -> Result<impl IntoResponse, AppError> {
let (admin_id, _role) = require_admin(&state, &headers).await?;
let repo = ExternalMountPgRepository::new(pool(&state)?);
let removed = repo
.delete(id)
.await
.map_err(|e| AppError::internal_error(format!("delete mount row: {e}")))?;
if !removed {
return Err(AppError::not_found("Mount not found"));
}
// Remove the mount-root folder row (host content is left intact).
state
.repositories
.folder_repository
.delete_folder(&id.to_string())
.await
.map_err(AppError::from)?;
let factory = DefaultMountProviderFactory::new();
state.mount_router.registry().reload(&repo, &factory).await;
tracing::info!(
target: "audit",
event = "external_mount.config",
action = "delete",
mount_id = %id,
caller_id = %admin_id,
reason = "external_mount_admin",
"👮🏻‍♂️ external mount deleted",
);
Ok(StatusCode::NO_CONTENT)
}
@@ -52,7 +52,17 @@ struct AdminUsersPageResponse {
/// Admin API routes — all require admin role.
pub fn admin_routes() -> Router<Arc<AppState>> {
use super::admin_external_mounts as ext_mounts;
Router::new()
// External file mounts
.route(
"/external-mounts",
get(ext_mounts::list_external_mounts).post(ext_mounts::create_external_mount),
)
.route(
"/external-mounts/{id}",
delete(ext_mounts::delete_external_mount),
)
// OIDC settings
.route("/settings/oidc", get(get_oidc_settings))
.route("/settings/oidc", put(save_oidc_settings))
+268
View File
@@ -11,13 +11,18 @@ use serde::Deserialize;
use std::collections::HashMap;
use utoipa::ToSchema;
use crate::application::ports::external_mount_ports::MountStat;
use crate::application::ports::file_ports::{
FileManagementUseCase, FileRetrievalUseCase, FileUploadUseCase, RangeContent,
};
use crate::application::ports::storage_ports::{FileReadPort, StorageUsagePort};
use crate::application::ports::thumbnail_ports::ThumbnailPort;
use crate::application::ports::{file_ports::OptimizedFileContent, folder_ports::FolderUseCase};
use crate::application::services::external_mount_router::ResolvedId;
use crate::application::services::mount_registry::MountConfig;
use crate::common::di::AppState;
use crate::domain::errors::DomainError;
use crate::domain::services::external_mount_id::{NodeId, virtual_file_etag};
use crate::interfaces::errors::AppError;
use crate::interfaces::middleware::auth::AuthUser;
use crate::interfaces::range_requests::not_modified_response;
@@ -251,6 +256,35 @@ impl FileHandler {
}
}
// ── External mount destination? Stream to the provider ──
// Detected BEFORE the CAS ingest so the bytes never touch
// BLAKE3/dedup. Authorization happens inside the service.
if let Some(ref fid) = folder_id {
let (mount_cfg, parent_node) = match state.mount_router.classify(fid) {
ResolvedId::MountRoot { cfg } => (Some(cfg), NodeId::default()),
ResolvedId::MountChild { cfg, node_id } => (Some(cfg), node_id),
ResolvedId::Regular => (None, NodeId::default()),
};
if let Some(cfg) = mount_cfg {
use futures::StreamExt;
let body: crate::application::ports::external_mount_ports::MountByteStream<
'_,
> = Box::pin(
upload_ingest::multipart_field_stream(field)
.map(|r| r.map_err(|e| std::io::Error::other(e.to_string()))),
);
return match state
.applications
.external_upload_service
.write_file(&cfg, &parent_node, &filename, body, auth_user.id)
.await
{
Ok(file) => Ok((file, String::new())),
Err(err) => Err(Self::domain_error_response(err)),
};
}
}
// ── Stream the field into the CDC chunk store ────────
// Chunking (FastCDC) + hashing (BLAKE3) + dedup checks +
// MIME sniffing all happen while the bytes arrive; chunks
@@ -683,6 +717,22 @@ impl FileHandler {
Query(params): Query<HashMap<String, String>>,
headers: &HeaderMap,
) -> impl IntoResponse + use<> {
// External mount: download a file living on the provider's backend.
// (A mount-root UUID is a folder and is not downloadable — it falls
// through and 404s as a non-file.)
if let ResolvedId::MountChild { cfg, node_id } = state.mount_router.classify(&id) {
return Self::download_mount_file(
&state,
&cfg,
&node_id,
&id,
auth_user.id,
&params,
headers,
)
.await;
}
let retrieval = &state.applications.file_retrieval_service;
// ── Get file metadata (ownership-scoped) ────────────────────────
@@ -838,6 +888,134 @@ impl FileHandler {
}
}
/// Download a file living inside an external mount: stat via the provider
/// (authorized against the mount root), then serve metadata / 304 / Range /
/// full stream straight from the backend. No blob cache, dedup, or WebP
/// transcode — mount content is served as-is.
#[allow(clippy::too_many_arguments)]
pub(super) async fn download_mount_file(
state: &AppState,
cfg: &MountConfig,
node_id: &NodeId,
id: &str,
caller_id: uuid::Uuid,
params: &HashMap<String, String>,
headers: &HeaderMap,
) -> axum::response::Response {
let retrieval = &state.applications.file_retrieval_service;
let stat: MountStat = match retrieval
.stat_mount_file_with_perms(cfg, node_id, caller_id)
.await
{
Ok(s) => s,
Err(err) => return AppError::from(err).into_response(),
};
if stat.is_dir {
// Directories are not downloadable through this endpoint.
return AppError::from(DomainError::not_found("File", id)).into_response();
}
let name = node_id
.as_str()
.rsplit('/')
.next()
.unwrap_or_else(|| node_id.as_str());
// ── Metadata-only request ────────────────────────────────────
if params
.get("metadata")
.is_some_and(|v| v == "true" || v == "1")
{
return (
StatusCode::OK,
Json(serde_json::json!({
"id": id,
"name": name,
"size": stat.size,
"mime_type": stat.mime_type,
"modified_at": stat.modified_at,
})),
)
.into_response();
}
let etag = format!("\"{}\"", virtual_file_etag(stat.size, stat.modified_at));
if let Some(resp) = not_modified_response(headers, &etag) {
return resp.into_response();
}
// ── Range Requests ───────────────────────────────────────────
let range_header = headers.get(header::RANGE).and_then(|v| v.to_str().ok());
match plan_mount_range(stat.size, range_header) {
MountRangePlan::Full => {}
MountRangePlan::NotSatisfiable => {
return Response::builder()
.status(StatusCode::RANGE_NOT_SATISFIABLE)
.header(header::CONTENT_RANGE, format!("bytes */{}", stat.size))
.body(Body::empty())
.unwrap()
.into_response();
}
MountRangePlan::Range { start, end } => {
let range_length = end - start + 1;
let disposition = Self::content_disposition(name, &stat.mime_type, params);
match retrieval
.open_mount_file_with_perms(cfg, node_id, caller_id, Some((start, Some(end))))
.await
{
Ok(stream) => {
return Response::builder()
.status(StatusCode::PARTIAL_CONTENT)
.header(header::CONTENT_TYPE, &stat.mime_type)
.header(header::CONTENT_DISPOSITION, &disposition)
.header(header::CONTENT_LENGTH, range_length)
.header(
header::CONTENT_RANGE,
format!("bytes {}-{}/{}", start, end, stat.size),
)
.header(header::ACCEPT_RANGES, "bytes")
.header(header::ETAG, &etag)
.header(
header::CACHE_CONTROL,
"private, max-age=3600, must-revalidate",
)
.body(Body::from_stream(stream))
.unwrap()
.into_response();
}
Err(err) => {
tracing::error!("Error creating mount range stream: {}", err);
// fall through to full download
}
}
}
}
// ── Normal download ──────────────────────────────────────────
let disposition = Self::content_disposition(name, &stat.mime_type, params);
match retrieval
.open_mount_file_with_perms(cfg, node_id, caller_id, None)
.await
{
Ok(stream) => Response::builder()
.status(StatusCode::OK)
.header(header::CONTENT_TYPE, &stat.mime_type)
.header(header::CONTENT_DISPOSITION, &disposition)
.header(header::CONTENT_LENGTH, stat.size)
.header(header::ETAG, &etag)
.header(
header::CACHE_CONTROL,
"private, max-age=3600, must-revalidate",
)
.header(header::ACCEPT_RANGES, "bytes")
.body(Body::from_stream(stream))
.unwrap()
.into_response(),
Err(err) => AppError::from(err).into_response(),
}
}
// ═══════════════════════════════════════════════════════════════════════
// LIST
// ═══════════════════════════════════════════════════════════════════════
@@ -1472,3 +1650,93 @@ pub async fn move_file_simple(
) -> impl IntoResponse {
FileHandler::move_file_simple_impl(state, auth_user, path, json).await
}
/// The download decision for a mount file given a `Range` header — the gnarly
/// parse-and-validate logic, extracted so it is unit-testable without I/O.
#[derive(Debug, PartialEq, Eq)]
pub(super) enum MountRangePlan {
/// Serve the whole file (no/invalid range header).
Full,
/// Serve `start..=end` (inclusive) as 206 Partial Content.
Range { start: u64, end: u64 },
/// The requested range is unsatisfiable for this size → 416.
NotSatisfiable,
}
/// Decide how to serve a mount file for a given size + optional `Range` header.
/// A missing or unparseable range → `Full`; a valid range → `Range`; an
/// out-of-bounds range → `NotSatisfiable`.
pub(super) fn plan_mount_range(size: u64, range_header: Option<&str>) -> MountRangePlan {
let Some(rh) = range_header else {
return MountRangePlan::Full;
};
let Ok(ranges) = parse_range_header(rh) else {
// Malformed range header: ignore it and serve the whole file (RFC 7233).
return MountRangePlan::Full;
};
match ranges.validate(size) {
Ok(valid) => match valid.first() {
Some(r) => MountRangePlan::Range {
start: *r.start(),
end: *r.end(),
},
None => MountRangePlan::Full,
},
Err(_) => MountRangePlan::NotSatisfiable,
}
}
#[cfg(test)]
mod mount_range_tests {
use super::{MountRangePlan, plan_mount_range};
#[test]
fn no_range_header_is_full() {
assert_eq!(plan_mount_range(100, None), MountRangePlan::Full);
}
#[test]
fn malformed_range_falls_back_to_full() {
assert_eq!(
plan_mount_range(100, Some("not-a-range")),
MountRangePlan::Full
);
assert_eq!(
plan_mount_range(100, Some("bytes=abc")),
MountRangePlan::Full
);
}
#[test]
fn valid_range_is_parsed_inclusive() {
assert_eq!(
plan_mount_range(100, Some("bytes=10-19")),
MountRangePlan::Range { start: 10, end: 19 }
);
}
#[test]
fn open_ended_range_extends_to_eof() {
assert_eq!(
plan_mount_range(100, Some("bytes=90-")),
MountRangePlan::Range { start: 90, end: 99 }
);
}
#[test]
fn suffix_range_counts_from_end() {
// last 10 bytes of a 100-byte file => 90..=99
assert_eq!(
plan_mount_range(100, Some("bytes=-10")),
MountRangePlan::Range { start: 90, end: 99 }
);
}
#[test]
fn out_of_bounds_range_is_not_satisfiable() {
assert_eq!(
plan_mount_range(100, Some("bytes=200-300")),
MountRangePlan::NotSatisfiable
);
}
}
+217 -1
View File
@@ -7,7 +7,8 @@ use axum::{
use std::sync::Arc;
use crate::application::dtos::display_helpers::{
classify_display, format_file_size, intern_display, intern_mime,
category_for, classify_display, format_file_size, icon_class_for, icon_special_class_for,
intern_display, intern_mime,
};
use crate::application::dtos::file_dto::FileDto;
use crate::application::dtos::folder_dto::{
@@ -15,11 +16,17 @@ use crate::application::dtos::folder_dto::{
ListResourcesOptions, MoveFolderDto, RenameFolderDto,
};
use crate::application::dtos::grant_dto::{ResourceContentDto, ResourceTypeDto};
use crate::application::ports::external_mount_ports::MountEntry;
use crate::application::ports::folder_ports::FolderUseCase;
use crate::application::ports::trash_ports::TrashUseCase;
use crate::application::services::external_mount_router::ResolvedId;
use crate::application::services::folder_service::FolderService;
use crate::application::services::mount_registry::MountConfig;
use crate::common::di::AppState as GlobalAppState;
use crate::domain::entities::file::File;
use crate::domain::services::external_mount_id::{
NodeId, encode_child_id, virtual_file_etag, virtual_folder_etag,
};
use crate::interfaces::errors::AppError;
use crate::interfaces::middleware::auth::AuthUser;
@@ -190,6 +197,22 @@ impl FolderHandler {
Path(id): Path<String>,
) -> impl IntoResponse {
let user_id = auth_user.id;
// External mounts have no trash — a permanent provider delete is the
// only option. Route `ext:` ids straight to the mount-aware service
// delete, skipping the (always-failing) trash attempt.
if state.mount_router.is_mount_id(&id) {
return match state
.applications
.folder_service
.delete_folder_with_perms(&id, user_id)
.await
{
Ok(_) => StatusCode::NO_CONTENT.into_response(),
Err(err) => AppError::from(err).into_response(),
};
}
// Check if trash service is available
// FIXME: permissions !!
if let Some(trash_service) = &state.trash_service {
@@ -481,6 +504,28 @@ pub async fn list_folder_resources(
reverse: q.reverse,
};
// External mount branch: a mount-root UUID or an `ext:` id lists live from
// the provider instead of the PostgreSQL UNION. The parent of each entry is
// the requested id itself.
match service.mount_router().classify(&id) {
ResolvedId::MountRoot { cfg } => {
return list_mount_dir_response(
&service,
&cfg,
&NodeId::default(),
&id,
auth_user.id,
opts,
)
.await;
}
ResolvedId::MountChild { cfg, node_id } => {
return list_mount_dir_response(&service, &cfg, &node_id, &id, auth_user.id, opts)
.await;
}
ResolvedId::Regular => {}
}
match service
.list_resources_paged_with_perms(&id, auth_user.id, opts)
.await
@@ -585,3 +630,174 @@ pub async fn list_folder_resources(
Err(e) => AppError::from(e).into_response(),
}
}
/// List one directory inside an external mount and render the standard
/// `/resources` envelope, mapping each live provider entry to a
/// `FolderResourceItemDto` with a synthetic `ext:` id. `parent_id` is the
/// requested id (the directory being listed), which becomes each entry's parent.
async fn list_mount_dir_response(
service: &FolderService,
cfg: &MountConfig,
node_id: &NodeId,
parent_id: &str,
caller_id: uuid::Uuid,
opts: ListResourcesOptions<'_>,
) -> axum::response::Response {
match service
.list_mount_dir_with_perms(cfg, node_id, caller_id, opts)
.await
{
Ok((entries, next_cursor)) => {
let items: Vec<FolderResourceItemDto> = entries
.into_iter()
.map(|entry| mount_entry_to_item(cfg, parent_id, entry))
.collect();
(
StatusCode::OK,
Json(FolderResourcesDto::with_cursor(items, next_cursor)),
)
.into_response()
}
Err(e) => AppError::from(e).into_response(),
}
}
/// Map a live mount entry to a `/resources` item with a synthetic `ext:` id and
/// virtual (size+mtime / mtime) etag. Mount entries have no blob hash.
fn mount_entry_to_item(
cfg: &MountConfig,
parent_id: &str,
entry: MountEntry,
) -> FolderResourceItemDto {
let id = encode_child_id(cfg.mount_id, entry.node_id.clone());
if entry.is_dir {
let dto = FolderDto {
etag: virtual_folder_etag(entry.modified_at),
id,
name: entry.name.clone(),
path: String::new(),
parent_id: Some(parent_id.to_owned()),
drive_id: cfg.drive_id,
created_at: entry.created_at,
modified_at: entry.modified_at,
is_root: false,
icon_class: Arc::from("fas fa-folder"),
icon_special_class: Arc::from("folder-icon"),
category: Arc::from("Folder"),
created_by: Some(cfg.owner_id),
updated_by: Some(cfg.owner_id),
};
FolderResourceItemDto {
resource_type: ResourceTypeDto::Folder,
resource: ResourceContentDto::Folder(dto),
}
} else {
let mime = mime_guess::from_path(&entry.name)
.first_or_octet_stream()
.to_string();
let dto = FileDto {
id,
name: entry.name.clone(),
path: String::new(),
size: entry.size,
mime_type: Arc::from(mime.as_str()),
folder_id: Some(parent_id.to_owned()),
created_at: entry.created_at,
modified_at: entry.modified_at,
icon_class: Arc::from(icon_class_for(&entry.name, &mime)),
icon_special_class: Arc::from(icon_special_class_for(&entry.name, &mime)),
category: Arc::from(category_for(&entry.name, &mime)),
size_formatted: format_file_size(entry.size),
sort_date: None,
content_hash: String::new(),
etag: virtual_file_etag(entry.size, entry.modified_at),
created_by: Some(cfg.owner_id),
updated_by: Some(cfg.owner_id),
};
FolderResourceItemDto {
resource_type: ResourceTypeDto::File,
resource: ResourceContentDto::File(dto),
}
}
}
#[cfg(test)]
mod mount_mapping_tests {
use super::*;
use crate::application::services::mount_registry::MountConfig;
use crate::infrastructure::services::local_fs_mount_provider::LocalFsMountProvider;
use uuid::Uuid;
fn config() -> MountConfig {
let dir = tempfile::tempdir().unwrap();
// Leak the tempdir so the path stays valid for the provider's lifetime;
// the provider is never exercised here (mapping is pure metadata).
let path = dir.keep();
MountConfig {
mount_id: Uuid::new_v4(),
kind: "local_fs".to_string(),
name: "Media".to_string(),
owner_id: Uuid::new_v4(),
drive_id: Uuid::new_v4(),
read_only: false,
mount_path: "Personal/Media".to_string(),
provider: Arc::new(LocalFsMountProvider::new(&path, false).unwrap()),
}
}
fn mount_entry(name: &str, node_id: &str, is_dir: bool, size: u64, mtime: u64) -> MountEntry {
MountEntry {
name: name.to_string(),
node_id: NodeId(node_id.to_string()),
is_dir,
size,
modified_at: mtime,
created_at: mtime,
}
}
#[test]
fn maps_folder_entry_to_item() {
let cfg = config();
let parent = cfg.mount_id.to_string();
let item = mount_entry_to_item(&cfg, &parent, mount_entry("docs", "docs", true, 0, 1234));
assert!(matches!(item.resource_type, ResourceTypeDto::Folder));
let ResourceContentDto::Folder(dto) = item.resource else {
panic!("expected folder");
};
assert_eq!(dto.name, "docs");
// id is the synthetic ext: envelope for (mount_id, node_id).
assert_eq!(dto.id, encode_child_id(cfg.mount_id, "docs"));
assert_eq!(dto.parent_id.as_deref(), Some(parent.as_str()));
assert_eq!(dto.etag, virtual_folder_etag(1234));
assert_eq!(dto.drive_id, cfg.drive_id);
assert!(!dto.is_root);
// Hierarchy is intentionally cleared on this listing.
assert_eq!(dto.path, "");
}
#[test]
fn maps_file_entry_to_item_with_virtual_etag_and_no_hash() {
let cfg = config();
let parent = encode_child_id(cfg.mount_id, "docs");
let item = mount_entry_to_item(
&cfg,
&parent,
mount_entry("report.json", "docs/report.json", false, 42, 999),
);
assert!(matches!(item.resource_type, ResourceTypeDto::File));
let ResourceContentDto::File(dto) = item.resource else {
panic!("expected file");
};
assert_eq!(dto.id, encode_child_id(cfg.mount_id, "docs/report.json"));
assert_eq!(dto.folder_id.as_deref(), Some(parent.as_str()));
assert_eq!(dto.size, 42);
assert_eq!(dto.etag, virtual_file_etag(42, 999));
// Virtual files have no blob hash.
assert_eq!(dto.content_hash, "");
// Mime is sniffed from the name.
assert_eq!(&*dto.mime_type, "application/json");
}
}
+1
View File
@@ -1,3 +1,4 @@
pub mod admin_external_mounts;
pub mod admin_handler;
pub mod app_password_handler;
pub mod auth_handler;
+6
View File
@@ -12,6 +12,12 @@ pub mod interfaces;
#[cfg(integration_tests)]
pub mod integration_test_support;
// Shared testcontainers-backed harness for external-mount integration tests.
// Gated on `test` too because it links the `testcontainers` dev-dependency,
// which is only available to test targets (not the plain lib build).
#[cfg(all(test, integration_tests))]
mod mount_it_support;
// Phase 0 perf-benchmark support: deterministic image corpus generation/loading
// shared by `benches/thumbnails.rs` and `examples/bench_thumbnails_mem.rs`.
// Gated behind the `bench` feature so it adds nothing to normal builds.
+118
View File
@@ -0,0 +1,118 @@
//! Shared testcontainers harness for external-mount integration tests.
//!
//! Compiled only under `#[cfg(all(test, integration_tests))]` (the
//! `testcontainers` dev-dependency is linked into test targets only). Each
//! `fresh_db()` spins up an ephemeral Postgres, applies every migration, and
//! returns a pool — so the DB-backed tests are self-contained and need only
//! docker, not the external `spawn-db.sh` compose harness.
use std::sync::Arc;
use sqlx::PgPool;
use sqlx::postgres::PgPoolOptions;
use testcontainers_modules::postgres::Postgres;
use testcontainers_modules::testcontainers::runners::AsyncRunner;
use testcontainers_modules::testcontainers::{ContainerAsync, ImageExt};
use uuid::Uuid;
use crate::domain::repositories::drive_repository::DriveRepository;
use crate::domain::repositories::folder_repository::FolderRepository;
use crate::infrastructure::repositories::pg::{DrivePgRepository, FolderDbRepository};
/// Postgres image tag — must match production (PG13+): the schema uses
/// `CREATE OR REPLACE TRIGGER` (PG14+) and the `pg_trgm` / `ltree` contrib
/// extensions. The `testcontainers` module default (`11-alpine`) is too old.
pub const PG_IMAGE_TAG: &str = "17-alpine";
/// A provisioned mount: a user with a personal drive, a child folder serving as
/// the mount root, and (when created via [`provision_mount`]) an
/// `external_mounts` row.
pub struct Provisioned {
pub owner_id: Uuid,
pub drive_id: Uuid,
pub mount_folder_id: Uuid,
}
/// Bring up an ephemeral Postgres, apply every migration, return a pool. Keep
/// the returned container handle alive for the test's duration.
pub async fn fresh_db() -> (ContainerAsync<Postgres>, Arc<PgPool>) {
let container = Postgres::default()
.with_tag(PG_IMAGE_TAG)
.start()
.await
.expect("start postgres testcontainer (is docker running?)");
let port = container
.get_host_port_ipv4(5432)
.await
.expect("container port");
let url = format!("postgres://postgres:postgres@127.0.0.1:{port}/postgres");
let pool = PgPoolOptions::new()
.max_connections(4)
.connect(&url)
.await
.expect("connect to ephemeral postgres");
sqlx::migrate!().run(&pool).await.expect("apply migrations");
(container, Arc::new(pool))
}
/// Insert a minimal `user` role account, returning its id.
pub async fn make_user(pool: &PgPool, name: &str) -> Uuid {
sqlx::query_scalar::<_, Uuid>(
"INSERT INTO auth.users (username, email, password_hash, role)
VALUES ($1, $2, 'x', 'user') RETURNING id",
)
.bind(name)
.bind(format!("{name}@example.test"))
.fetch_one(pool)
.await
.expect("insert user")
}
/// Provision user → personal drive → a child folder named `folder_name` that
/// will act as the mount root. Does NOT insert an `external_mounts` row.
pub async fn provision_folder(
pool: &Arc<PgPool>,
user_name: &str,
folder_name: &str,
) -> Provisioned {
let owner_id = make_user(pool, user_name).await;
let drive_repo = DrivePgRepository::new(pool.clone());
let drive = drive_repo
.create_personal_drive_atomic(owner_id, None)
.await
.expect("create personal drive");
let root_folder_id = drive.drive.root_folder_id;
let folder_repo = FolderDbRepository::new(pool.clone());
let folder = folder_repo
.create_folder(
folder_name.to_string(),
Some(root_folder_id.to_string()),
owner_id,
)
.await
.expect("create mount-root folder");
let mount_folder_id = Uuid::parse_str(folder.id()).expect("uuid");
Provisioned {
owner_id,
drive_id: drive.drive.id,
mount_folder_id,
}
}
/// Insert an `external_mounts` row for `mount_folder_id` with a `local_fs`
/// provider pointed at `host_path`.
pub async fn insert_mount(pool: &PgPool, p: &Provisioned, host_path: &str) {
sqlx::query(
"INSERT INTO storage.external_mounts
(mount_folder_id, kind, config, name, owner_id, read_only)
VALUES ($1, 'local_fs', $2, 'Media', $3, false)",
)
.bind(p.mount_folder_id)
.bind(serde_json::json!({ "path": host_path }))
.bind(p.owner_id)
.execute(pool)
.await
.expect("insert external mount");
}