Merge pull request #479 from EdouardVanbelle/feat/drive-impl

feat/drive impl
This commit is contained in:
Dionisio Pozo
2026-06-19 18:48:32 +02:00
committed by GitHub
102 changed files with 6362 additions and 1078 deletions
@@ -124,7 +124,7 @@ pub struct AuthApplicationService {
token_service: Arc<JwtTokenService>,
/// Dispatcher for user-lifecycle events. `None` only in tests that don't
/// exercise the lifecycle path; production DI always wires this.
/// HomeFolderLifecycleHook (registered on this dispatcher) owns the
/// PersonalDriveLifecycleHook (registered on this dispatcher) owns the
/// per-user folder provisioning that AuthApplicationService used to do
/// inline pre-PR 3.
user_lifecycle: Option<Arc<UserLifecycleService>>,
@@ -404,7 +404,7 @@ impl AuthApplicationService {
// Save user
let created_user = self.user_storage.create_user(user).await?;
// Lifecycle: HomeFolderLifecycleHook handles personal-folder
// Lifecycle: PersonalDriveLifecycleHook handles personal-folder
// creation (was inlined here pre-PR 3); audit log + future
// provisioning steps land here too.
if let Some(lc) = &self.user_lifecycle {
@@ -506,8 +506,8 @@ impl AuthApplicationService {
let created_user = self.user_storage.create_user(user).await?;
// Lifecycle: notify hooks. PR 3 moves home-folder creation into
// HomeFolderLifecycleHook fired here.
// Lifecycle: HomeFolderLifecycleHook provisions the admin's
// PersonalDriveLifecycleHook fired here.
// Lifecycle: PersonalDriveLifecycleHook provisions the admin's
// home folder. Audit logs the creation event.
if let Some(lc) = &self.user_lifecycle {
lc.dispatch_created(&created_user).await;
@@ -654,7 +654,7 @@ impl AuthApplicationService {
/// A second redemption attempt receives `Ok(false)` and is rejected
/// as `AccessDenied`.
/// 3. Load the user, verify they're active.
/// 4. Dispatch `on_user_login` (so HomeFolderLifecycleHook can
/// 4. Dispatch `on_user_login` (so PersonalDriveLifecycleHook can
/// safety-net any internal user whose first credential happens
/// to be a magic link — externals short-circuit by `is_external()`).
/// 5. Register login + persist + issue session in the same pipeline
@@ -1722,7 +1722,7 @@ impl AuthApplicationService {
.await?;
}
// Lifecycle: HomeFolderLifecycleHook handles the home-folder
// Lifecycle: PersonalDriveLifecycleHook handles the home-folder
// provisioning (idempotent + short-circuits on is_external).
// Audit logs the creation event.
if let Some(lc) = &self.user_lifecycle {
@@ -1785,7 +1785,7 @@ impl AuthApplicationService {
/// Runs the whole flow in a single transaction so the lifecycle
/// hooks (`SessionRevocationLifecycleHook` revoking sessions with
/// audit, `AuthzCacheLifecycleHook` invalidating the Moka cache,
/// `HomeFolderLifecycleHook` for future trash policy, …) can do
/// `PersonalDriveLifecycleHook` for future trash policy, …) can do
/// their work atomically with the user DELETE. If any hook returns
/// `Err`, the transaction rolls back and the user remains intact.
pub async fn delete_user_admin(&self, user_id: Uuid) -> Result<(), DomainError> {
@@ -2260,7 +2260,7 @@ impl AuthApplicationService {
// Lifecycle: created (audit + home-folder provisioning) +
// login (no register_login() for a fresh OIDC user means
// `last_login_at` is naturally None → first-login detection
// works). HomeFolderLifecycleHook creates the home folder.
// works). PersonalDriveLifecycleHook creates the home folder.
if let Some(lc) = &self.user_lifecycle {
lc.dispatch_created(&created_user).await;
lc.dispatch_login(&created_user).await;
@@ -2362,7 +2362,7 @@ impl AuthApplicationService {
// `create_personal_folder` was removed in PR 3 of the
// UserLifecycleHook migration — home-folder provisioning is now
// owned by `HomeFolderLifecycleHook` in folder_service.rs and runs
// owned by `PersonalDriveLifecycleHook` in folder_service.rs and runs
// via `dispatch_created` / `dispatch_login`.
}
@@ -623,6 +623,7 @@ impl DeltaUploadService {
Some(folder_id.clone()),
content_type,
blob,
caller_id,
)
.await
}
@@ -98,6 +98,7 @@ impl FileManagementService {
&self,
file_id: &str,
folder_id: Option<String>,
caller_id: Uuid,
) -> Result<FileDto, DomainError> {
info!(
"Moving file with ID: {} to folder: {:?}",
@@ -106,7 +107,7 @@ impl FileManagementService {
let moved_file = self
.file_repository
.move_file(file_id, folder_id)
.move_file(file_id, folder_id, caller_id)
.await
.map_err(|e| {
error!("Error moving file (ID: {}): {}", file_id, e);
@@ -128,6 +129,7 @@ impl FileManagementService {
file_id: &str,
target_folder_id: Option<String>,
new_name: Option<&str>,
caller_id: Uuid,
) -> Result<FileDto, DomainError> {
info!(
"Copying file with ID: {} to folder: {:?} as {:?}",
@@ -136,7 +138,7 @@ impl FileManagementService {
let copied_file = self
.file_repository
.copy_file(file_id, target_folder_id, new_name)
.copy_file(file_id, target_folder_id, new_name, caller_id)
.await
.map_err(|e| {
error!("Error copying file (ID: {}): {}", file_id, e);
@@ -157,7 +159,12 @@ impl FileManagementService {
Ok(dto)
}
async fn rename_file(&self, file_id: &str, new_name: &str) -> Result<FileDto, DomainError> {
async fn rename_file(
&self,
file_id: &str,
new_name: &str,
caller_id: Uuid,
) -> Result<FileDto, DomainError> {
if let Err(reason) = validate_storage_name(new_name) {
return Err(DomainError::validation_error(format!(
"Invalid file name '{new_name}': {reason}"
@@ -168,7 +175,7 @@ impl FileManagementService {
let renamed_file = self
.file_repository
.rename_file(file_id, new_name)
.rename_file(file_id, new_name, caller_id)
.await
.map_err(|e| {
error!("Error renaming file (ID: {}): {}", file_id, e);
@@ -253,7 +260,7 @@ impl FileManagementUseCase for FileManagementService {
.await?;
self.require_target_folder_perm(folder_id.as_deref(), Permission::Create, caller_id)
.await?;
self.move_file(file_id, folder_id).await
self.move_file(file_id, folder_id, caller_id).await
}
async fn copy_file_with_perms(
@@ -268,7 +275,7 @@ impl FileManagementUseCase for FileManagementService {
.await?;
self.require_target_folder_perm(target_folder_id.as_deref(), Permission::Create, caller_id)
.await?;
self.copy_file(file_id, target_folder_id, new_name.as_deref())
self.copy_file(file_id, target_folder_id, new_name.as_deref(), caller_id)
.await
}
@@ -280,7 +287,7 @@ impl FileManagementUseCase for FileManagementService {
) -> Result<FileDto, DomainError> {
self.require_file_perm(file_id, Permission::Update, caller_id)
.await?;
self.rename_file(file_id, new_name).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> {
@@ -286,13 +286,16 @@ impl FileRetrievalUseCase for FileRetrievalService {
}
// FIXME no authorisation at all
async fn get_file_by_path(&self, path: &str) -> Result<FileDto, DomainError> {
async fn get_file_by_path(&self, path: &str, drive_id: Uuid) -> Result<FileDto, DomainError> {
// Direct SQL lookup — O(folder_depth) queries instead of O(total_files)
// NOTE: This method does NOT perform any authorization check. Callers
// that surface its result to a user-driven request MUST resolve the
// file via get_file_owned afterwards, or call authz.require directly.
// (Tracked in the audit punch-list under "path-based lookups".)
if let Some(file) = self.file_read.find_file_by_path(path).await? {
// `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.
if let Some(file) = self.file_read.find_file_by_path(path, drive_id).await? {
return Ok(FileDto::from(file));
}
@@ -209,6 +209,7 @@ impl FileUploadService {
size: metadata.size,
is_new_blob: false,
},
caller_id,
)
.await?;
@@ -255,7 +256,7 @@ impl FileUploadService {
let file = file_read.get_file(file_id).await?;
let (new_hash, updated_at) = self
.file_write
.update_file_content_with_blob(file_id, &blob.hash, blob.size, None)
.update_file_content_with_blob(file_id, &blob.hash, blob.size, None, caller_id)
.await?;
// The file maps to a different blob now — stale cached content must
// never be served for the rest of its TTI window.
@@ -326,10 +327,18 @@ impl FileUploadUseCase for FileUploadService {
folder_id: Option<String>,
content_type: String,
blob: StoredBlob,
caller_id: Uuid,
) -> Result<FileDto, DomainError> {
let file = self
.file_write
.save_file_with_blob(name.clone(), folder_id, content_type, &blob.hash, blob.size)
.save_file_with_blob(
name.clone(),
folder_id,
content_type,
&blob.hash,
blob.size,
caller_id,
)
.await?;
let dto = FileDto::from(file);
info!(
@@ -348,18 +357,26 @@ impl FileUploadUseCase for FileUploadService {
async fn update_file_streaming(
&self,
path: &str,
drive_id: Uuid,
blob: StoredBlob,
content_type: &str,
modified_at: Option<i64>,
caller_id: Uuid,
) -> Result<FileDto, DomainError> {
// Try to find the existing file first
if let Some(file_read) = &self.file_read
&& let Some(file) = file_read.find_file_by_path(path).await?
&& let Some(file) = file_read.find_file_by_path(path, drive_id).await?
{
let file_id = file.id().to_string();
let (new_hash, updated_at) = self
.file_write
.update_file_content_with_blob(&file_id, &blob.hash, blob.size, modified_at)
.update_file_content_with_blob(
&file_id,
&blob.hash,
blob.size,
modified_at,
caller_id,
)
.await?;
// Invalidate content cache — file content has changed.
if let Some(cc) = &self.content_cache {
@@ -402,9 +419,15 @@ impl FileUploadUseCase for FileUploadService {
// get_parent_folder_id expects the full file path — it strips the
// last segment (filename) internally to find the parent folder.
// `drive_id` scopes the parent lookup to the same drive as the
// incoming write (post-D0 `storage.folders.path` repeats across
// drives).
let parent_id = if path_normalized.contains('/') {
if let Some(file_read) = &self.file_read {
file_read.get_parent_folder_id(path_normalized).await.ok()
file_read
.get_parent_folder_id(path_normalized, drive_id)
.await
.ok()
} else {
None
}
@@ -421,6 +444,7 @@ impl FileUploadUseCase for FileUploadService {
content_type.to_string(),
&blob.hash,
blob.size,
caller_id,
)
.await?;
let dto = FileDto::from(created);
+119 -119
View File
@@ -82,7 +82,11 @@ impl FolderService {
Ok(FolderDto::empty())
}
async fn get_folder_by_path(&self, _path: &str) -> Result<FolderDto, DomainError> {
async fn get_folder_by_path(
&self,
_path: &str,
_drive_id: Uuid,
) -> Result<FolderDto, DomainError> {
Ok(FolderDto::empty())
}
@@ -163,14 +167,6 @@ impl FolderService {
) -> Result<(), DomainError> {
Ok(())
}
async fn create_home_folder(
&self,
_user_id: Uuid,
_name: String,
) -> Result<FolderDto, DomainError> {
Ok(FolderDto::empty())
}
}
FolderServiceStub
@@ -229,31 +225,11 @@ impl FolderUseCase for FolderService {
let folder = self
.folder_storage
.create_folder(dto.name, dto.parent_id)
.create_folder(dto.name, dto.parent_id, caller_id)
.await?;
Ok(FolderDto::from(folder))
}
/// Creates a root-level home folder for a user during registration.
async fn create_home_folder(
&self,
user_id: Uuid,
name: String,
) -> Result<FolderDto, DomainError> {
let folder = self
.folder_storage
.create_home_folder(user_id, name)
.await
.map_err(|e| {
DomainError::internal_error(
"FolderStorage",
format!("Failed to create home folder: {}", e),
)
})?;
Ok(FolderDto::from(folder))
}
async fn list_subtree_folders(&self, folder_id: &str) -> Result<Vec<FolderDto>, DomainError> {
let folders = self.folder_storage.list_subtree_folders(folder_id).await?;
Ok(folders.into_iter().map(FolderDto::from).collect())
@@ -288,14 +264,17 @@ impl FolderUseCase for FolderService {
self.get_folder(id).await
}
/// Gets a folder by its path
async fn get_folder_by_path(&self, path: &str) -> Result<FolderDto, DomainError> {
// Convert the string path to StoragePath
/// Gets a folder by its path, scoped to a drive.
async fn get_folder_by_path(
&self,
path: &str,
drive_id: Uuid,
) -> Result<FolderDto, DomainError> {
let storage_path = StoragePath::from_string(path);
let folder = self
.folder_storage
.get_folder_by_path(&storage_path)
.get_folder_by_path(&storage_path, drive_id)
.await
.map_err(|e| {
DomainError::internal_error(
@@ -329,7 +308,7 @@ impl FolderUseCase for FolderService {
///
/// **Note (post PR 3):** the self-heal block that auto-created a
/// home folder when listing returned empty has been removed.
/// `HomeFolderLifecycleHook` (registered on `UserLifecycleService`)
/// `PersonalDriveLifecycleHook` (registered on `UserLifecycleService`)
/// now provisions the folder on `on_user_created` / `on_user_login`,
/// idempotently, so the listing path no longer needs to self-heal.
async fn list_folders_with_perms(
@@ -477,7 +456,7 @@ impl FolderUseCase for FolderService {
let folder = self
.folder_storage
.rename_folder(id, dto.name)
.rename_folder(id, dto.name, caller_id)
.await
.map_err(|e| {
DomainError::internal_error(
@@ -529,7 +508,7 @@ impl FolderUseCase for FolderService {
let parent_ref = dto.parent_id.as_deref();
let folder = self
.folder_storage
.move_folder(id, parent_ref)
.move_folder(id, parent_ref, caller_id)
.await
.map_err(|e| {
DomainError::internal_error(
@@ -614,62 +593,6 @@ impl FolderService {
Ok((rows, next_cursor))
}
/// Idempotently provision a home folder for a user.
///
/// Returns `Ok(true)` if a folder was newly created, `Ok(false)` if the
/// user already had at least one root folder.
///
/// **System-level operation** — bypasses authz because this runs on
/// the user's own behalf (during creation or login provisioning) at a
/// point where the caller may be the engine itself, not an HTTP user.
/// Callers must be inside trusted code paths (lifecycle hooks).
///
/// Used by [`HomeFolderLifecycleHook`] on `on_user_created` and
/// `on_user_login`. Replaces the old self-heal at the listing path
/// and the four eager `create_personal_folder` calls in
/// `AuthApplicationService` (removed in the same PR).
pub async fn ensure_home_folder(
&self,
user_id: Uuid,
username: Option<&str>,
) -> Result<bool, DomainError> {
let existing = self
.folder_storage
.list_folders_by_owner(None, user_id)
.await
.map_err(|e| {
DomainError::internal_error(
"FolderStorage",
format!("ensure_home_folder: list root folders: {}", e),
)
})?;
if !existing.is_empty() {
return Ok(false);
}
let folder_name = match username {
Some(u) => format!("My Folder - {}", u),
None => format!("My Folder - {}", user_id),
};
self.folder_storage
.create_home_folder(user_id, folder_name.clone())
.await
.map_err(|e| {
DomainError::internal_error(
"FolderStorage",
format!("ensure_home_folder: create: {}", e),
)
})?;
tracing::info!(
target: "user_lifecycle",
hook = "home_folder",
user_id = %user_id,
folder_name = %folder_name,
"Home folder provisioned"
);
Ok(true)
}
}
/// Build the next-page cursor from the last row of the current page.
@@ -725,7 +648,7 @@ fn build_folder_resource_cursor(
}
// ─────────────────────────────────────────────────────────────────────────────
// HomeFolderLifecycleHook
// PersonalDriveLifecycleHook
//
// Owns home-folder provisioning policy. Replaces:
// - the 4 eager `create_personal_folder` calls in AuthApplicationService
@@ -743,36 +666,115 @@ use async_trait::async_trait;
use crate::application::ports::user_lifecycle::{DeletionMode, LogoutReason, UserLifecycleHook};
use crate::domain::entities::user::User;
/// Lifecycle hook: provisions and (in PR 4) deprovisions a user's home folder.
pub struct HomeFolderLifecycleHook {
folder_service: Arc<FolderService>,
/// Lifecycle hook: provisions a user's default Personal drive at first
/// login (replaces the legacy `My Folder - <username>` wrapper as of D0).
///
/// Two writes happen on first provisioning:
/// 1. A row in `storage.drives` with `kind='personal'`,
/// `default_for_user=<uid>`, and the user's quota carried over from
/// `auth.users.storage_quota_bytes`.
/// 2. An Owner role grant in `storage.role_grants` so the user can
/// read/write/manage their own drive (the engine's owner short-
/// circuit applies to folders/files but not drives — see
/// `pg_acl_engine::check_inner` D0-6 rewrite).
///
/// Both writes are idempotent: `find_default_for_user` short-circuits
/// when the drive already exists; `set_role` is an UPSERT that no-ops
/// when the Owner row is already present.
pub struct PersonalDriveLifecycleHook {
drive_repo: Arc<dyn crate::domain::repositories::drive_repository::DriveRepository>,
// The `AuthorizationEngine` trait isn't `dyn`-compatible (native
// async-fn-in-trait methods are not object-safe), so we hold the
// concrete engine. This matches the convention already used by
// `AppState.authorization`. Only the idempotent-rerun path uses it
// now; the create path goes through the repo's atomic CTE which
// writes the role_grant inline.
authorization: Arc<crate::infrastructure::services::pg_acl_engine::PgAclEngine>,
}
impl HomeFolderLifecycleHook {
pub fn new(folder_service: Arc<FolderService>) -> Self {
Self { folder_service }
impl PersonalDriveLifecycleHook {
pub fn new(
drive_repo: Arc<dyn crate::domain::repositories::drive_repository::DriveRepository>,
authorization: Arc<crate::infrastructure::services::pg_acl_engine::PgAclEngine>,
) -> Self {
Self {
drive_repo,
authorization,
}
}
/// Idempotent provisioning shared by `on_user_created` and
/// `on_user_login`. External users are skipped per tip #2 in the
/// trait docstring.
/// trait docstring — they have no resources of their own, only
/// grants on other users' resources.
async fn provision_if_needed(&self, user: &User) -> Result<(), DomainError> {
use crate::domain::repositories::drive_repository::DriveRepositoryError;
use crate::domain::services::authorization::{Resource, Role, Subject};
if user.is_external() {
return Ok(());
}
// `ensure_home_folder` handles the "does the user already have a
// root folder?" check internally and is a no-op if so.
self.folder_service
.ensure_home_folder(user.id(), user.username())
// Idempotent shortcut: if the user already has a default drive,
// the atomic CTE already ran on a prior turn. The CTE writes
// the Owner role_grant inline, so there's nothing to repair —
// but we still re-emit the grant via `set_role` (UPSERT-safe)
// to cover the historical case where a pre-CTE provisioning
// path partially completed (drive created, grant missing).
match self.drive_repo.find_default_for_user(user.id()).await {
Ok(drive_with_name) => {
self.authorization
.set_role(
user.id(),
Subject::User(user.id()),
Role::Owner,
Resource::Drive(drive_with_name.drive.id),
None,
)
.await
.map(|_grant| ())?;
return Ok(());
}
Err(DriveRepositoryError::NotFound(_)) => { /* fall through to create */ }
Err(e) => {
return Err(DomainError::internal_error(
"PersonalDriveHook",
format!("find_default lookup: {e}"),
));
}
}
// One atomic CTE — drive row + root folder ("Personal",
// parent_id=NULL, drive_id pinned) + drives.root_folder_id
// wire-up + Owner role_grant. Single SQL statement, atomic
// against server crash mid-sequence (docs/plan/drive.md §3).
let drive_with_name = self
.drive_repo
.create_personal_drive_atomic(user.id(), Some(user.storage_quota_bytes()))
.await
.map(|_created| ())
.map_err(|e| {
DomainError::internal_error(
"PersonalDriveHook",
format!("create_personal_drive_atomic: {e}"),
)
})?;
tracing::info!(
target: "user_lifecycle",
hook = "personal_drive",
user_id = %user.id(),
drive_id = %drive_with_name.drive.id,
root_folder_id = %drive_with_name.drive.root_folder_id,
"Default personal drive + root folder + owner grant provisioned (atomic CTE)"
);
Ok(())
}
}
#[async_trait]
impl UserLifecycleHook for HomeFolderLifecycleHook {
impl UserLifecycleHook for PersonalDriveLifecycleHook {
fn name(&self) -> &'static str {
"home_folder"
"personal_drive"
}
async fn on_user_created(&self, user: &User) -> Result<(), DomainError> {
@@ -787,7 +789,7 @@ impl UserLifecycleHook for HomeFolderLifecycleHook {
}
async fn on_user_logout(&self, _user: &User, _reason: LogoutReason) -> Result<(), DomainError> {
// Folders don't react to logout. Explicit no-op per the
// Drives don't react to logout. Explicit no-op per the
// "no defaults" convention.
Ok(())
}
@@ -798,24 +800,22 @@ impl UserLifecycleHook for HomeFolderLifecycleHook {
mode: DeletionMode,
_tx: &mut sqlx::Transaction<'_, sqlx::Postgres>,
) -> Result<(), DomainError> {
// For both DeletionMode variants today the FK CASCADE on
// `storage.folders.user_id` (and downstream files/blobs)
// removes the home folder + contents when the user row goes.
// `storage.drives.default_for_user` has ON DELETE CASCADE
// referencing `auth.users(id)`, and `storage.folders.drive_id`
// / `storage.files.drive_id` both have ON DELETE CASCADE on
// `storage.drives(id)` (M3). So a user delete cascades:
// user → drive → folders → files in one transaction.
//
// The hook emits a per-mode tracing event so audit can tell
// AdminDelete (currently recoverable only via DB-level rollback
// before commit) from GdprPurge (no sweeper exists yet — the
// variant is reserved for a future PR that adds retention).
//
// The `tx` is provided per the trait contract but unused here:
// emitting a tracing event doesn't require DB access. Future
// policy (trash with retention) would write to `storage.trash`
// inside this same tx.
tracing::info!(
target: "user_lifecycle",
hook = "home_folder",
hook = "personal_drive",
user_id = %user.id(),
mode = ?mode,
"Home folder will be removed via FK CASCADE on user delete"
"Personal drive (and tree) will be removed via FK CASCADE on user delete"
);
Ok(())
}
@@ -101,7 +101,11 @@ impl FileReadPort for MockFileReadPort {
unimplemented!()
}
async fn get_parent_folder_id(&self, _path: &str) -> Result<String, DomainError> {
async fn get_parent_folder_id(
&self,
_path: &str,
_drive_id: Uuid,
) -> Result<String, DomainError> {
unimplemented!()
}
@@ -127,7 +131,11 @@ impl FileReadPort for MockFileReadPort {
Ok(0)
}
async fn get_folder_id_by_path(&self, _folder_path: &str) -> Result<String, DomainError> {
async fn get_folder_id_by_path(
&self,
_folder_path: &str,
_drive_id: Uuid,
) -> Result<String, DomainError> {
unimplemented!()
}
@@ -176,6 +184,7 @@ impl FileWritePort for MockFileWritePort {
_content_type: String,
_blob_hash: &str,
_size: u64,
_caller_id: Uuid,
) -> Result<File, DomainError> {
unimplemented!()
}
@@ -184,6 +193,7 @@ impl FileWritePort for MockFileWritePort {
&self,
file_id: &str,
_target_folder_id: Option<String>,
_caller_id: Uuid,
) -> Result<File, DomainError> {
let files = self.files.lock().unwrap();
files
@@ -192,7 +202,12 @@ impl FileWritePort for MockFileWritePort {
.ok_or_else(|| DomainError::not_found("File", file_id.to_string()))
}
async fn rename_file(&self, file_id: &str, _new_name: &str) -> Result<File, DomainError> {
async fn rename_file(
&self,
file_id: &str,
_new_name: &str,
_caller_id: Uuid,
) -> Result<File, DomainError> {
let files = self.files.lock().unwrap();
files
.get(file_id)
@@ -210,6 +225,7 @@ impl FileWritePort for MockFileWritePort {
_blob_hash: &str,
_size: u64,
_modified_at: Option<i64>,
_caller_id: Uuid,
) -> Result<(String, i64), DomainError> {
Ok((String::new(), 0))
}
@@ -220,6 +236,7 @@ impl FileWritePort for MockFileWritePort {
_folder_id: Option<String>,
_content_type: String,
_size: u64,
_caller_id: Uuid,
) -> Result<(File, PathBuf), DomainError> {
unimplemented!()
}
@@ -229,11 +246,12 @@ impl FileWritePort for MockFileWritePort {
_file_id: &str,
_target_folder_id: Option<String>,
_new_name: Option<&str>,
_caller_id: Uuid,
) -> Result<File, DomainError> {
unimplemented!()
}
async fn move_to_trash(&self, _file_id: &str) -> Result<(), DomainError> {
async fn move_to_trash(&self, _file_id: &str, _caller_id: Uuid) -> Result<(), DomainError> {
Ok(())
}
@@ -241,6 +259,7 @@ impl FileWritePort for MockFileWritePort {
&self,
_file_id: &str,
_original_path: &str,
_caller_id: Uuid,
) -> Result<(), DomainError> {
Ok(())
}
@@ -307,6 +307,23 @@ impl MagicLinkInviteService {
let (kind, resource_id) = match resource {
Resource::Folder(id) => (MagicLinkResourceKind::Folder, id),
Resource::File(id) => (MagicLinkResourceKind::File, id),
// Drive sharing — and therefore drive magic-link invitations —
// land in D2. The grant DTOs accept `Resource::Drive` from the
// wire today (see ResourceTypeDto) but no public API path
// actually grants on a drive in D0, so this arm is
// defensively unreachable. Treating it as an audit-logged
// no-op (grant is in place, mail suppressed) matches the
// ineligible-recipient branch above.
Resource::Drive(_) => {
tracing::info!(
target: "audit",
event = "magic_link.invitation_suppressed",
reason = "drive_resource_unsupported",
user_id = %recipient.id(),
"📭 magic-link invitation suppressed: drive resources aren't invitable until D2",
);
return Ok(());
}
};
// Invitation tokens are cross-device by design (recipient has
// no prior browser context with the server) — no challenge
@@ -329,6 +346,11 @@ impl MagicLinkInviteService {
let kind_key = match resource {
Resource::Folder(_) => "server.magic_link.email.kind_folder",
Resource::File(_) => "server.magic_link.email.kind_file",
// Unreachable — the early-return above exits before we get
// here for a Drive resource. The arm exists only to satisfy
// exhaustiveness; if you find this firing, the early-return
// was bypassed.
Resource::Drive(_) => "server.magic_link.email.kind_folder",
};
// PR C: render in the recipient's preferred locale (set by UI
// switcher, OIDC JIT claim, or inviter inheritance at row
@@ -731,6 +753,15 @@ impl From<ResourceKind> for MagicLinkResourceKind {
match kind {
ResourceKind::Folder => Self::Folder,
ResourceKind::File => Self::File,
// Drives aren't a magic-link invite target in D0. The
// grant DTO surface accepts drive resources, but the
// grant_handler doesn't issue magic-links for them
// (drive sharing lands in D2). Mapping Drive → Folder
// gives a non-panicking fallback that would still emit a
// valid token shape if the path were ever reached; the
// runtime branches above suppress drive invitations
// before reaching this conversion.
ResourceKind::Drive => Self::Folder,
}
}
}
@@ -3,6 +3,7 @@ use std::sync::{Arc, Mutex};
use std::time::{Duration, Instant};
use rand_core::RngCore;
use uuid::Uuid;
/// Maximum number of concurrent pending login flows to prevent memory exhaustion.
const MAX_PENDING_FLOWS: usize = 1000;
@@ -30,6 +31,13 @@ pub struct LoginResult {
struct PendingFlow {
created_at: Instant,
poll_token: String,
/// Set after the user authenticates on the login page **and** has more
/// than one root drive — the flow is paused until the user picks a
/// drive on the picker page. Consumed by `take_pending_user` when the
/// picker submission arrives, so the second step is single-use even
/// if the flow token leaks. `None` for single-drive accounts (legacy
/// path goes straight to `completed`).
pending_user_id: Option<Uuid>,
completed: Option<LoginResult>,
}
@@ -80,6 +88,7 @@ impl NextcloudLoginFlowService {
PendingFlow {
created_at: Instant::now(),
poll_token: poll_token.clone(),
pending_user_id: None,
completed: None,
},
);
@@ -101,6 +110,39 @@ impl NextcloudLoginFlowService {
state.flows.contains_key(flow_token)
}
/// Stash a verified user_id on the flow so a follow-up drive-pick
/// request can prove "this browser just authenticated" without
/// asking for the password again. Returns `false` if the flow
/// token is unknown or expired.
///
/// Only used on multi-drive accounts — single-drive logins go
/// straight to [`complete`](Self::complete).
pub fn mark_awaiting_drive(&self, flow_token: &str, user_id: Uuid) -> bool {
let mut state = self.state.lock().unwrap_or_else(|e| e.into_inner());
prune_expired(&mut state, self.ttl);
match state.flows.get_mut(flow_token) {
Some(pending) => {
pending.pending_user_id = Some(user_id);
true
}
None => false,
}
}
/// Consume the stashed user_id (single-use). Returns the user_id
/// when the flow is in "awaiting drive choice" state, or `None`
/// when the flow is unknown, expired, or was never marked. Single-
/// use semantics make this safe even if the flow token leaks: the
/// second drive-pick attempt finds nothing to consume.
pub fn take_pending_user(&self, flow_token: &str) -> Option<Uuid> {
let mut state = self.state.lock().unwrap_or_else(|e| e.into_inner());
prune_expired(&mut state, self.ttl);
state
.flows
.get_mut(flow_token)
.and_then(|pending| pending.pending_user_id.take())
}
pub fn complete(
&self,
flow_token: &str,
@@ -257,6 +299,35 @@ mod tests {
assert!(svc.poll(&info.poll_token).is_none());
}
#[test]
fn test_mark_awaiting_drive_then_take_pending_user() {
let svc = service();
let info = svc.initiate("https://cloud.example.com").unwrap();
let flow_token = info.login_url.rsplit('/').next().unwrap();
let uid = Uuid::new_v4();
assert!(svc.mark_awaiting_drive(flow_token, uid));
// First take consumes the slot.
assert_eq!(svc.take_pending_user(flow_token), Some(uid));
// Second take must return None (single-use).
assert_eq!(svc.take_pending_user(flow_token), None);
}
#[test]
fn test_mark_awaiting_drive_unknown_flow_returns_false() {
let svc = service();
assert!(!svc.mark_awaiting_drive("nonexistent", Uuid::new_v4()));
}
#[test]
fn test_take_pending_user_without_mark_returns_none() {
let svc = service();
let info = svc.initiate("https://cloud.example.com").unwrap();
let flow_token = info.login_url.rsplit('/').next().unwrap();
// Flow exists but mark_awaiting_drive was never called.
assert_eq!(svc.take_pending_user(flow_token), None);
}
#[test]
fn test_max_pending_flows_cap() {
let svc = NextcloudLoginFlowService::new(Duration::from_secs(600));
+1 -1
View File
@@ -199,7 +199,7 @@ impl PeopleService {
})
})
.collect();
out.sort_by(|a, b| b.face_count.cmp(&a.face_count));
out.sort_by_key(|p| std::cmp::Reverse(p.face_count));
Ok(out)
}
@@ -472,6 +472,11 @@ impl RecipientNotificationService {
let kind_key = match resource {
Resource::Folder(_) => "server.magic_link.email.kind_folder",
Resource::File(_) => "server.magic_link.email.kind_file",
// Drives don't generate share notifications in D0 — drive
// sharing lands in D2 and gets its own template key. Fall
// back to the folder label so any path that does reach
// here produces a readable, if generic, mail body.
Resource::Drive(_) => "server.magic_link.email.kind_folder",
};
let kind_label = self.i18n_or(kind_key, &locale, &[]).await;
// Short form for the subject, long form (with email) for the
+92 -3
View File
@@ -50,6 +50,20 @@ pub struct SearchService {
/// matches; hits are hydrated and re-filtered through SQL before use.
content_index: Option<Arc<dyn ContentIndexPort>>,
/// Optional authorization engine — needed to resolve the caller's
/// accessible drive set before querying the content index, and to
/// re-verify each Tantivy hit against `engine.check(Read, File(id))`
/// as a defense-in-depth measure (catches index staleness and
/// per-file grants that the drive-only Tantivy filter misses; see
/// `docs/plan/drive.md` §11). `None` short-circuits the content
/// index (the cheapest safe degradation).
authorization: Option<Arc<crate::infrastructure::services::pg_acl_engine::PgAclEngine>>,
/// Optional drive repository — used in tandem with the authorization
/// engine to resolve the caller's accessible drives for the Tantivy
/// filter. `None` short-circuits the content index.
drive_repo: Option<Arc<dyn crate::domain::repositories::drive_repository::DriveRepository>>,
/// Lock-free concurrent cache with automatic TTL and LRU eviction (moka).
/// Values are `Arc<SearchResultsDto>` so cache insert/hit is a single
/// atomic ref-count increment (~1 ns) instead of cloning thousands of Strings.
@@ -151,6 +165,8 @@ impl SearchService {
file_repository: Arc<FileBlobReadRepository>,
folder_repository: Arc<FolderDbRepository>,
content_index: Option<Arc<dyn ContentIndexPort>>,
authorization: Option<Arc<crate::infrastructure::services::pg_acl_engine::PgAclEngine>>,
drive_repo: Option<Arc<dyn crate::domain::repositories::drive_repository::DriveRepository>>,
cache_ttl: u64,
max_cache_size: usize,
) -> Self {
@@ -163,6 +179,8 @@ impl SearchService {
file_repository,
folder_repository,
content_index,
authorization,
drive_repo,
search_cache,
}
}
@@ -250,9 +268,18 @@ impl SearchService {
criteria: &SearchCriteriaDto,
user_id: Uuid,
) -> Vec<ContentHitDto> {
use crate::application::ports::authorization_ports::AuthorizationEngine;
use crate::domain::services::authorization::{Permission, Resource, Subject};
let Some(index) = &self.content_index else {
return Vec::new();
};
let Some(authz) = &self.authorization else {
return Vec::new();
};
let Some(drive_repo) = &self.drive_repo else {
return Vec::new();
};
if criteria.offset != 0 {
return Vec::new();
}
@@ -265,16 +292,78 @@ impl SearchService {
return Vec::new();
};
match index
.search_content(user_id, query, CONTENT_HITS_LIMIT)
// Resolve the caller's accessible drive set via the engine
// (handles group-mediated drive grants) + the repo lookup.
let caller = Subject::User(user_id);
let (subject_types, subject_ids) = match authz.expand_subject_for_listing(caller).await {
Ok(pair) => pair,
Err(e) => {
tracing::warn!("Content-index: subject expansion failed — degrading to empty: {e}");
return Vec::new();
}
};
let accessible_drives: Vec<Uuid> = match drive_repo
.list_for_subjects(&subject_types, &subject_ids)
.await
{
Ok(drives) => drives.into_iter().map(|d| d.drive.id).collect(),
Err(e) => {
tracing::warn!("Content-index: drive lookup failed — degrading to empty: {e}");
return Vec::new();
}
};
// Tantivy filter (Must drive_id ∈ accessible_drives) handles
// the cross-drive isolation. Empty drive list short-circuits
// inside `search_content`.
let hits = match index
.search_content(&accessible_drives, query, CONTENT_HITS_LIMIT)
.await
{
Ok(hits) => hits,
Err(e) => {
tracing::warn!("Content-index lookup failed — returning name-only results: {e}");
Vec::new()
return Vec::new();
}
};
// Defense in depth: re-verify each hit through the engine.
// Catches two cases the drive_id filter can't:
// * Index staleness — the file just moved drives and the
// worker hasn't caught up.
// * Per-file grants — ReBAC can grant a single file inside a
// drive the caller doesn't otherwise have. The Tantivy
// filter is drive-only; this re-check restores per-file
// resolution.
// Failures degrade conservatively (drop the hit, log it) —
// never leak.
let mut verified = Vec::with_capacity(hits.len());
for hit in hits {
let file_uuid = match Uuid::parse_str(&hit.file_id) {
Ok(u) => u,
Err(_) => {
tracing::warn!("Content-index hit had non-UUID file_id: {}", hit.file_id);
continue;
}
};
match authz
.check(caller, Permission::Read, Resource::File(file_uuid))
.await
{
Ok(true) => verified.push(hit),
Ok(false) => {
tracing::debug!(
target: "oxicloud::search",
file_id = %file_uuid,
"dropping content-index hit: ReBAC denies Read after Tantivy filter",
);
}
Err(e) => {
tracing::warn!("ReBAC re-check failed for {file_uuid}: {e}");
}
}
}
verified
}
/// Merge content-index hits into the name-search result page:
+21 -11
View File
@@ -838,11 +838,19 @@ mod tests {
unimplemented!()
}
async fn get_parent_folder_id(&self, _path: &str) -> Result<String, DomainError> {
async fn get_parent_folder_id(
&self,
_path: &str,
_drive_id: uuid::Uuid,
) -> Result<String, DomainError> {
unimplemented!()
}
async fn get_folder_id_by_path(&self, _folder_path: &str) -> Result<String, DomainError> {
async fn get_folder_id_by_path(
&self,
_folder_path: &str,
_drive_id: uuid::Uuid,
) -> Result<String, DomainError> {
unimplemented!()
}
@@ -898,6 +906,7 @@ mod tests {
&self,
_name: String,
_parent_id: Option<String>,
_caller_id: uuid::Uuid,
) -> Result<crate::domain::entities::folder::Folder, DomainError> {
unimplemented!()
}
@@ -925,6 +934,7 @@ mod tests {
async fn get_folder_by_path(
&self,
_storage_path: &crate::domain::services::path_service::StoragePath,
_drive_id: uuid::Uuid,
) -> Result<crate::domain::entities::folder::Folder, DomainError> {
unimplemented!()
}
@@ -971,6 +981,7 @@ mod tests {
&self,
_id: &str,
_new_name: String,
_caller_id: Uuid,
) -> Result<crate::domain::entities::folder::Folder, DomainError> {
unimplemented!()
}
@@ -979,6 +990,7 @@ mod tests {
&self,
_id: &str,
_new_parent_id: Option<&str>,
_caller_id: Uuid,
) -> Result<crate::domain::entities::folder::Folder, DomainError> {
unimplemented!()
}
@@ -990,6 +1002,7 @@ mod tests {
async fn folder_exists(
&self,
_storage_path: &crate::domain::services::path_service::StoragePath,
_drive_id: uuid::Uuid,
) -> Result<bool, DomainError> {
unimplemented!()
}
@@ -1001,7 +1014,11 @@ mod tests {
unimplemented!()
}
async fn move_to_trash(&self, _folder_id: &str) -> Result<(), DomainError> {
async fn move_to_trash(
&self,
_folder_id: &str,
_caller_id: Uuid,
) -> Result<(), DomainError> {
unimplemented!()
}
@@ -1009,6 +1026,7 @@ mod tests {
&self,
_folder_id: &str,
_original_path: &str,
_caller_id: Uuid,
) -> Result<(), DomainError> {
unimplemented!()
}
@@ -1016,14 +1034,6 @@ mod tests {
async fn delete_folder_permanently(&self, _folder_id: &str) -> Result<(), DomainError> {
unimplemented!()
}
async fn create_home_folder(
&self,
_user_id: Uuid,
_name: String,
) -> Result<crate::domain::entities::folder::Folder, DomainError> {
unimplemented!()
}
}
struct MockShareRepository {
+18 -6
View File
@@ -253,9 +253,10 @@ impl TrashUseCase for TrashService {
}
};
// Then physically move the file to trash
// Then physically move the file to trash.
// §14: caller_id stamps `updated_by` on the trashed row.
info!("Physically moving file to trash: {}", item_id);
match self.file_write_port.move_to_trash(item_id).await {
match self.file_write_port.move_to_trash(item_id, user_id).await {
Ok(_) => {
debug!("File physically moved to trash successfully: {}", item_id);
}
@@ -320,9 +321,10 @@ impl TrashUseCase for TrashService {
}
};
// Then physically move the folder to trash
// Then physically move the folder to trash.
// §14: caller_id stamps `updated_by` on every cascade-trashed row.
self.folder_storage_port
.move_to_trash(item_id)
.move_to_trash(item_id, user_id)
.await
.map_err(|e| {
DomainError::new(
@@ -391,7 +393,7 @@ impl TrashUseCase for TrashService {
);
match self
.file_write_port
.restore_from_trash(&file_id, &original_path)
.restore_from_trash(&file_id, &original_path, user_id)
.await
{
Ok(_) => {
@@ -431,7 +433,7 @@ impl TrashUseCase for TrashService {
);
match self
.folder_storage_port
.restore_from_trash(&folder_id, &original_path)
.restore_from_trash(&folder_id, &original_path, user_id)
.await
{
Ok(_) => {
@@ -821,12 +823,19 @@ fn row_to_item_dto(row: TrashResourceRow) -> TrashResourceItemDto {
path,
parent_id: row.parent_id.map(|u| u.to_string()),
owner_id: Some(row.owner_id.to_string()),
// Trash listing — drive_id is informational and the trash
// row doesn't currently SELECT it. Path-based lookups
// never enter this code path.
drive_id: uuid::Uuid::nil(),
created_at: row.resource_created_at.timestamp() as u64,
modified_at: row.modified_at.timestamp() as u64,
is_root: false,
icon_class: std::sync::Arc::from("fas fa-folder"),
icon_special_class: std::sync::Arc::from("folder-icon"),
category: std::sync::Arc::from("Folder"),
// §14 provenance not selected by the trash listing query.
created_by: None,
updated_by: None,
};
TrashResourceItemDto {
resource_type: ResourceTypeDto::Folder,
@@ -867,6 +876,9 @@ fn row_to_item_dto(row: TrashResourceRow) -> TrashResourceItemDto {
sort_date: None,
content_hash,
etag,
// §14 provenance not selected by the trash listing query.
created_by: None,
updated_by: None,
};
TrashResourceItemDto {
resource_type: ResourceTypeDto::File,
+33 -15
View File
@@ -139,7 +139,7 @@ where
)
})?;
self.file_write_port
.move_to_trash(item_id)
.move_to_trash(item_id, user_id)
.await
.map_err(|e| {
DomainError::new(
@@ -181,7 +181,7 @@ where
)
})?;
self.folder_storage_port
.move_to_trash(item_id)
.move_to_trash(item_id, user_id)
.await
.map_err(|e| {
DomainError::new(
@@ -215,7 +215,7 @@ where
let original_path = item.original_path().to_string();
let result = self
.file_write_port
.restore_from_trash(&file_id, &original_path)
.restore_from_trash(&file_id, &original_path, user_id)
.await;
if let Err(e) = result
&& !format!("{}", e).contains("not found")
@@ -232,7 +232,7 @@ where
let original_path = item.original_path().to_string();
let result = self
.folder_storage_port
.restore_from_trash(&folder_id, &original_path)
.restore_from_trash(&folder_id, &original_path, user_id)
.await;
if let Err(e) = result
&& !format!("{}", e).contains("not found")
@@ -500,13 +500,18 @@ impl FileReadPort for MockFileRepository {
unimplemented!()
}
async fn get_parent_folder_id(&self, _path: &str) -> std::result::Result<String, DomainError> {
async fn get_parent_folder_id(
&self,
_path: &str,
_drive_id: Uuid,
) -> std::result::Result<String, DomainError> {
unimplemented!()
}
async fn get_folder_id_by_path(
&self,
_folder_path: &str,
_drive_id: Uuid,
) -> std::result::Result<String, DomainError> {
unimplemented!()
}
@@ -561,6 +566,7 @@ impl FileWritePort for MockFileRepository {
_content_type: String,
_blob_hash: &str,
_size: u64,
_caller_id: Uuid,
) -> std::result::Result<File, DomainError> {
unimplemented!()
}
@@ -569,6 +575,7 @@ impl FileWritePort for MockFileRepository {
&self,
_file_id: &str,
_target_folder_id: Option<String>,
_caller_id: Uuid,
) -> std::result::Result<File, DomainError> {
unimplemented!()
}
@@ -577,6 +584,7 @@ impl FileWritePort for MockFileRepository {
&self,
_file_id: &str,
_new_name: &str,
_caller_id: Uuid,
) -> std::result::Result<File, DomainError> {
unimplemented!()
}
@@ -591,6 +599,7 @@ impl FileWritePort for MockFileRepository {
_blob_hash: &str,
_size: u64,
_modified_at: Option<i64>,
_caller_id: Uuid,
) -> std::result::Result<(String, i64), DomainError> {
Ok((String::new(), 0))
}
@@ -601,6 +610,7 @@ impl FileWritePort for MockFileRepository {
_folder_id: Option<String>,
_content_type: String,
_size: u64,
_caller_id: Uuid,
) -> std::result::Result<(File, PathBuf), DomainError> {
unimplemented!()
}
@@ -610,11 +620,16 @@ impl FileWritePort for MockFileRepository {
_file_id: &str,
_target_folder_id: Option<String>,
_new_name: Option<&str>,
_caller_id: Uuid,
) -> std::result::Result<File, DomainError> {
unimplemented!()
}
async fn move_to_trash(&self, id: &str) -> std::result::Result<(), DomainError> {
async fn move_to_trash(
&self,
id: &str,
_caller_id: Uuid,
) -> std::result::Result<(), DomainError> {
let mut files = self.files.lock().unwrap();
let mut trashed = self.trashed_files.lock().unwrap();
@@ -630,6 +645,7 @@ impl FileWritePort for MockFileRepository {
&self,
id: &str,
_original_path: &str,
_caller_id: Uuid,
) -> std::result::Result<(), DomainError> {
let mut files = self.files.lock().unwrap();
let mut trashed = self.trashed_files.lock().unwrap();
@@ -693,6 +709,7 @@ impl FolderRepository for MockFolderRepository {
&self,
_name: String,
_parent_id: Option<String>,
_caller_id: Uuid,
) -> std::result::Result<Folder, DomainError> {
unimplemented!()
}
@@ -709,6 +726,7 @@ impl FolderRepository for MockFolderRepository {
async fn get_folder_by_path(
&self,
_storage_path: &StoragePath,
_drive_id: Uuid,
) -> std::result::Result<Folder, DomainError> {
unimplemented!()
}
@@ -753,6 +771,7 @@ impl FolderRepository for MockFolderRepository {
&self,
_id: &str,
_new_name: String,
_caller_id: Uuid,
) -> std::result::Result<Folder, DomainError> {
unimplemented!()
}
@@ -761,6 +780,7 @@ impl FolderRepository for MockFolderRepository {
&self,
_id: &str,
_new_parent_id: Option<&str>,
_caller_id: Uuid,
) -> std::result::Result<Folder, DomainError> {
unimplemented!()
}
@@ -772,6 +792,7 @@ impl FolderRepository for MockFolderRepository {
async fn folder_exists(
&self,
_storage_path: &StoragePath,
_drive_id: Uuid,
) -> std::result::Result<bool, DomainError> {
Ok(false)
}
@@ -780,7 +801,11 @@ impl FolderRepository for MockFolderRepository {
Ok(StoragePath::from_string("/"))
}
async fn move_to_trash(&self, id: &str) -> std::result::Result<(), DomainError> {
async fn move_to_trash(
&self,
id: &str,
_caller_id: Uuid,
) -> std::result::Result<(), DomainError> {
let mut folders = self.folders.lock().unwrap();
let mut trashed = self.trashed_folders.lock().unwrap();
@@ -796,6 +821,7 @@ impl FolderRepository for MockFolderRepository {
&self,
id: &str,
_original_path: &str,
_caller_id: Uuid,
) -> std::result::Result<(), DomainError> {
let mut folders = self.folders.lock().unwrap();
let mut trashed = self.trashed_folders.lock().unwrap();
@@ -822,14 +848,6 @@ impl FolderRepository for MockFolderRepository {
))
}
}
async fn create_home_folder(
&self,
_user_id: Uuid,
_name: String,
) -> std::result::Result<Folder, DomainError> {
Ok(Folder::default())
}
}
#[cfg(integration_tests)]
@@ -132,7 +132,7 @@ impl UserLifecycleService {
//
// Always-on observer. Emits one structured `tracing::info!(target: "audit",
// ...)` line per event. The only hook registered in PR 1; subsequent PRs
// add HomeFolderLifecycleHook, AuthzCacheLifecycleHook, etc., each living
// add PersonalDriveLifecycleHook, AuthzCacheLifecycleHook, etc., each living
// next to the service it works for.
// ─────────────────────────────────────────────────────────────────────────────