from __future__ import annotations import json import os import tempfile from dataclasses import replace from datetime import datetime from io import BytesIO from typing import Any, Literal from urllib.parse import quote, urljoin from fastapi import APIRouter, Depends, File as FastAPIFile, Form, HTTPException, Query, UploadFile, status from fastapi.responses import FileResponse, StreamingResponse from starlette.background import BackgroundTask from sqlalchemy import func, or_ from sqlalchemy.orm import Session from govoplan_core.auth import ApiPrincipal, has_scope, require_any_scope, require_scope from govoplan_core.api.v1.schemas import DeltaDeletedItem from govoplan_core.audit.logging import audit_from_principal from govoplan_core.core.campaigns import CAPABILITY_CAMPAIGNS_ACCESS, CampaignAccessProvider from govoplan_core.core.change_sequence import ( ChangeSequenceEntry, decode_sequence_watermark, encode_sequence_watermark, max_sequence_id, record_change, sequence_entries_since, sequence_watermark_is_expired, ) from govoplan_core.core.pagination import ( KeysetCursorError, decode_keyset_cursor, encode_keyset_cursor, keyset_query_fingerprint, ) from govoplan_files.backend.change_tracking import ( FILES_ASSETS_COLLECTION, FILES_CONNECTOR_CREDENTIALS_COLLECTION, FILES_CONNECTOR_POLICIES_COLLECTION, FILES_CONNECTOR_PROFILES_COLLECTION, FILES_CONNECTOR_SPACES_COLLECTION, FILES_FOLDERS_COLLECTION, FILES_MODULE_ID, ) from govoplan_files.backend.schemas import ( ArchiveRequest, BulkFileShareRequest, BulkFileShareResponse, BulkDeleteRequest, BulkDeleteResponse, ConflictResolutionRequest, FileAssetResponse, FileConnectorBrowseItem, FileConnectorBrowseResponse, FileConnectorCredentialCreateRequest, FileConnectorCredentialResponse, FileConnectorCredentialsResponse, FileConnectorCredentialUpdateRequest, FileConnectorDiscoveryRequest, FileConnectorDiscoveryResponse, FileConnectorSettingsDeltaResponse, FileConnectorSpaceCreateRequest, FileConnectorSpaceResponse, FileConnectorSpacesResponse, FileConnectorSpaceUpdateRequest, FileConnectorPolicyEvaluateRequest, FileConnectorPolicyEvaluateResponse, FileConnectorPolicyResponse, FileConnectorPolicyUpdateRequest, FileConnectorProfileCreateRequest, FileConnectorProfileResponse, FileConnectorProfilesResponse, FileConnectorProfileUpdateRequest, FileConnectorProviderResponse, FileConnectorProvidersResponse, FileConnectorImportRequest, FileConnectorSyncResponse, FileDeltaResponse, FileFolderCreateRequest, FileFolderDeleteRequest, FileFolderDeleteResponse, FileFolderResponse, FileFoldersResponse, FileListResponse, FileShareRequest, FileShareResponse, FileSpaceResponse, FileSpacesResponse, FileUploadResponse, PatternMatchResponse, PatternResolveRequest, PatternResolveResponse, RenamePreviewItem, RenameRequest, RenameResponse, TransferRequest, TransferResponse, _conflict_resolutions, ) from govoplan_files.backend.db.models import CampaignAttachmentUse, FileAsset, FileBlob, FileConnectorSpace, FileFolder, FileShare, FileVersion from govoplan_core.db.session import get_session from govoplan_files.backend.runtime import get_registry, settings from govoplan_files.backend.storage.paths import UnsafeFilePathError, filename_from_path, normalize_logical_path from govoplan_files.backend.storage.access import ensure_group_access, group_refs_for_ids, user_group_ids from govoplan_files.backend.storage.archives import create_zip_file, extract_zip_upload from govoplan_files.backend.storage.common import FileStorageError from govoplan_files.backend.storage.connector_credential_store import ( connector_credential_from_row, create_connector_credential_row, deactivate_connector_credential_row, get_connector_credential_row, list_database_connector_credentials, update_connector_credential_row, ) from govoplan_files.backend.storage.connector_browse import ( ConnectorBrowseError, ConnectorBrowseUnsupported, browse_connector_profile, normalize_connector_browse_path, ) from govoplan_files.backend.storage.connector_imports import ConnectorImportError, ConnectorImportUnsupported, read_connector_file from govoplan_files.backend.storage.connector_profile_store import ( connector_profile_from_row, create_connector_profile_row, deactivate_connector_profile_row, get_connector_profile_row, list_database_connector_profiles, update_connector_profile_row, ) from govoplan_files.backend.storage.connector_profiles import ConnectorProfile, connector_profiles_from_settings from govoplan_files.backend.storage.connector_providers import connector_provider_descriptors from govoplan_files.backend.storage.connector_spaces import ( connector_space_owner_id, create_connector_space, get_connector_space_for_user, list_connector_spaces_for_user, soft_delete_connector_space, update_connector_space, ) from govoplan_files.backend.storage.connector_policy import ( ConnectorAccessRequest, ConnectorPolicyDenied, connector_policy_decision, connector_policy_sources_for_fields, connector_policy_sources_from_payload, ensure_connector_policy_allows, ) from govoplan_files.backend.storage.connector_policy_store import ( connector_policy_response, effective_connector_policy_sources, parent_connector_policy, set_connector_policy, validate_connector_policy_allowed_by_parent, ) from govoplan_files.backend.storage.files import ( asset_is_audit_relevant, create_file_asset, current_version_and_blob, get_asset_for_user, list_assets_for_user, list_assets_for_user_window, read_asset_bytes, share_file, share_files, soft_delete_assets, sync_file_asset_from_source, ) from govoplan_files.backend.storage.folders import create_folder, list_folders_for_user, list_folders_for_user_window, soft_delete_folder from govoplan_files.backend.storage.provenance import source_metadata, source_provenance_from_metadata, source_revision_from_metadata from govoplan_files.backend.storage.search import resolve_patterns from govoplan_files.backend.storage.transfers import rename_selection, transfer_selection router = APIRouter(prefix="/files", tags=["files"]) FILES_CONNECTOR_SETTINGS_COLLECTIONS = ( FILES_CONNECTOR_PROFILES_COLLECTION, FILES_CONNECTOR_CREDENTIALS_COLLECTION, FILES_CONNECTOR_POLICIES_COLLECTION, FILES_CONNECTOR_SPACES_COLLECTION, ) FILES_CONNECTOR_PROFILE_RESOURCE = "file_connector_profile" FILES_CONNECTOR_CREDENTIAL_RESOURCE = "file_connector_credential" FILES_CONNECTOR_POLICY_RESOURCE = "file_connector_policy" FILES_CONNECTOR_SPACE_RESOURCE = "file_connector_space" def _campaign_access_provider() -> CampaignAccessProvider: registry = get_registry() if registry is None or not hasattr(registry, "has_capability") or not registry.has_capability(CAPABILITY_CAMPAIGNS_ACCESS): raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail="Campaign module is not installed") capability = registry.require_capability(CAPABILITY_CAMPAIGNS_ACCESS) if not isinstance(capability, CampaignAccessProvider): raise HTTPException(status_code=status.HTTP_500_INTERNAL_SERVER_ERROR, detail="Campaign access capability is invalid") return capability def _is_admin(principal: ApiPrincipal) -> bool: return has_scope(principal, "files:file:admin") def _can_read_disabled_connector_profiles(principal: ApiPrincipal) -> bool: return any( has_scope(principal, scope) for scope in ( "files:file:admin", "system:settings:read", "system:settings:write", "admin:settings:read", "admin:settings:write", ) ) def _record_connector_settings_change( session: Session, *, collection: str, resource_type: str, resource_id: str | None, operation: str, principal: ApiPrincipal, tenant_id: str | None, payload: dict[str, object] | None = None, ) -> None: if not resource_id: return record_change( session, module_id=FILES_MODULE_ID, collection=collection, resource_type=resource_type, resource_id=resource_id, operation=operation, tenant_id=tenant_id, actor_type="user", actor_id=principal.user.id, payload=payload or {}, ) def _file_connector_settings_query(session: Session, *, tenant_id: str, since_sequence: int): return session.query(ChangeSequenceEntry).filter( ChangeSequenceEntry.id > since_sequence, ChangeSequenceEntry.module_id == FILES_MODULE_ID, ChangeSequenceEntry.collection.in_(FILES_CONNECTOR_SETTINGS_COLLECTIONS), or_(ChangeSequenceEntry.tenant_id == tenant_id, ChangeSequenceEntry.tenant_id.is_(None)), ) def _file_connector_settings_watermark(session: Session, *, tenant_id: str) -> str: sequence_id = _file_connector_settings_query(session, tenant_id=tenant_id, since_sequence=0).with_entities(func.max(ChangeSequenceEntry.id)).scalar() return encode_sequence_watermark(int(sequence_id or 0)) def _file_connector_settings_entries(session: Session, *, tenant_id: str, since: str, limit: int): try: since_sequence = decode_sequence_watermark(since) except ValueError as exc: raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail=str(exc)) from exc expired_for_tenant = sequence_watermark_is_expired( session, since=since_sequence, tenant_id=tenant_id, module_id=FILES_MODULE_ID, collections=FILES_CONNECTOR_SETTINGS_COLLECTIONS, ) expired_for_system = sequence_watermark_is_expired( session, since=since_sequence, tenant_id=None, module_id=FILES_MODULE_ID, collections=FILES_CONNECTOR_SETTINGS_COLLECTIONS, ) if expired_for_tenant or expired_for_system: return None, False entries_plus_one = _file_connector_settings_query( session, tenant_id=tenant_id, since_sequence=since_sequence, ).order_by(ChangeSequenceEntry.id.asc()).limit(limit + 1).all() has_more = len(entries_plus_one) > limit return entries_plus_one[:limit], has_more def _file_connector_settings_response_watermark(session: Session, *, tenant_id: str, entries, has_more: bool) -> str: return encode_sequence_watermark(entries[-1].id) if has_more and entries else _file_connector_settings_watermark(session, tenant_id=tenant_id) def _connector_deleted_item(entry) -> DeltaDeletedItem: return DeltaDeletedItem( id=entry.resource_id, resource_type=entry.resource_type, revision=encode_sequence_watermark(entry.id), deleted_at=entry.created_at if entry.operation == "deleted" else None, ) def _visible_connector_credentials( session: Session, principal: ApiPrincipal, *, provider: str | None = None, include_disabled: bool = False, ): provider_norm = provider.strip().casefold() if provider else None credentials = list_database_connector_credentials( session, tenant_id=principal.tenant_id, include_disabled=include_disabled and _can_read_disabled_connector_profiles(principal), ) return [ credential for credential in credentials if provider_norm is None or credential.provider in {None, provider_norm} ] def _visible_connector_spaces( session: Session, principal: ApiPrincipal, *, owner_type: Literal["user", "group"] | None = None, owner_id: str | None = None, include_inactive: bool = False, ) -> list[FileConnectorSpace]: if not (has_scope(principal, "files:file:read") or has_scope(principal, "files:file:admin")): return [] return list_connector_spaces_for_user( session, tenant_id=principal.tenant_id, user_id=principal.user.id, owner_type=owner_type, owner_id=owner_id, include_inactive=include_inactive and _is_admin(principal), is_admin=_is_admin(principal), ) def _file_connector_policy_resource_id(scope_type: str, scope_id: str | None) -> str: return f"{scope_type.strip().casefold()}:{scope_id or ''}" async def _read_limited_upload(upload: UploadFile, *, max_bytes: int) -> bytes: data = await upload.read(max_bytes + 1) if len(data) > max_bytes: raise HTTPException( status_code=status.HTTP_413_REQUEST_ENTITY_TOO_LARGE, detail=f"Upload exceeds limit of {max_bytes} bytes", ) return data async def _spool_limited_upload_to_temp(upload: UploadFile, *, max_bytes: int, suffix: str = ".upload") -> str: tmp = tempfile.NamedTemporaryFile(prefix="govoplan-upload-", suffix=suffix, delete=False) total = 0 try: while True: chunk = await upload.read(1024 * 1024) if not chunk: break total += len(chunk) if total > max_bytes: raise HTTPException( status_code=status.HTTP_413_REQUEST_ENTITY_TOO_LARGE, detail=f"Upload exceeds limit of {max_bytes} bytes", ) tmp.write(chunk) except Exception: tmp.close() _cleanup_temp_file(tmp.name) raise tmp.close() return tmp.name def _cleanup_temp_file(path: str) -> None: try: os.unlink(path) except FileNotFoundError: pass def _attachment_disposition(filename: str) -> str: safe = filename.replace("\\", "_").replace("/", "_").replace("\r", "_").replace("\n", "_").strip() or "download" ascii_name = "".join( char if 32 <= ord(char) < 127 and char not in {'"', "\\", ";"} else "_" for char in safe ).strip() or "download" encoded = quote(safe, safe="") return f'attachment; filename="{ascii_name}"; filename*=UTF-8' + "''" + encoded def _ensure_campaign_file_access(session: Session, principal: ApiPrincipal, campaign_id: str | None) -> None: if not campaign_id: return if not has_scope(principal, "campaigns:campaign:read"): raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail="Missing scope: campaign:read") provider = _campaign_access_provider() if not provider.campaign_exists(session, tenant_id=principal.tenant_id, campaign_id=campaign_id): raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Campaign not found") group_ids = set(user_group_ids(session, tenant_id=principal.tenant_id, user_id=principal.user.id)) if provider.can_read_campaign( session, tenant_id=principal.tenant_id, campaign_id=campaign_id, user_id=principal.user.id, group_ids=group_ids, tenant_admin=has_scope(principal, "tenant:*"), ): return raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail="No access to this campaign") def _http_error(exc: Exception, *, not_found: bool = False) -> HTTPException: code = status.HTTP_404_NOT_FOUND if not_found else status.HTTP_400_BAD_REQUEST return HTTPException(status_code=code, detail=str(exc)) def _connector_policy_error(exc: ConnectorPolicyDenied) -> HTTPException: return HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail=exc.decision.to_dict()) def _owner_id(asset: FileAsset) -> str: return asset.owner_user_id if asset.owner_type == "user" else asset.owner_group_id # type: ignore[return-value] def _source_metadata_from_form(source_provenance_json: str | None, source_revision: str | None) -> dict[str, object] | None: source_provenance = json.loads(source_provenance_json) if source_provenance_json else None return source_metadata(source_provenance=source_provenance, source_revision=source_revision) def _connector_policy_sources_from_form(connector_policy_json: str | None): if not connector_policy_json: return [] return connector_policy_sources_from_payload(json.loads(connector_policy_json)) def _enforce_connector_policy(source_provenance_json: str | None, connector_policy_json: str | None, *, operation: str) -> None: if not source_provenance_json: return sources = _connector_policy_sources_from_form(connector_policy_json) if not sources: return provenance = json.loads(source_provenance_json) request = ConnectorAccessRequest.from_provenance(provenance if isinstance(provenance, dict) else {}, operation=operation) ensure_connector_policy_allows(request, sources) def _connector_profile_visible( session: Session, principal: ApiPrincipal, profile: ConnectorProfile, *, campaign_id: str | None, group_ids: set[str], ) -> bool: if profile.scope_type == "system": return True if profile.scope_type == "tenant": return profile.scope_id == principal.tenant_id if profile.scope_type == "user": return profile.scope_id == principal.user.id if profile.scope_type == "group": return bool(profile.scope_id and profile.scope_id in group_ids) if profile.scope_type == "campaign": if not profile.scope_id: return False if campaign_id and profile.scope_id != campaign_id: return False try: _ensure_campaign_file_access(session, principal, profile.scope_id) except HTTPException: return False return True return False def _visible_connector_profiles( session: Session, principal: ApiPrincipal, *, provider: str | None = None, campaign_id: str | None = None, include_disabled: bool = False, include_admin_scopes: bool = False, include_effective_policy: bool = True, ) -> list[ConnectorProfile]: provider_norm = provider.strip().casefold() if provider else None group_ids = set(user_group_ids(session, tenant_id=principal.tenant_id, user_id=principal.user.id, include_admin_groups=_is_admin(principal))) database_profiles = list_database_connector_profiles( session, tenant_id=principal.tenant_id, include_disabled=include_disabled, ) env_profiles = connector_profiles_from_settings(settings) profiles_by_id: dict[str, ConnectorProfile] = {} for profile in [*database_profiles, *env_profiles]: if profile.id not in profiles_by_id: profiles_by_id[profile.id] = profile profiles = list(profiles_by_id.values()) visible: list[ConnectorProfile] = [] for profile in profiles: if not (profile.enabled or include_disabled): continue if provider_norm is not None and profile.provider != provider_norm: continue admin_scope_visible = include_admin_scopes and profile.scope_type in {"user", "group", "campaign"} if not admin_scope_visible and not _connector_profile_visible(session, principal, profile, campaign_id=campaign_id, group_ids=group_ids): continue visible.append(_with_effective_connector_policy(session, principal, profile) if include_effective_policy else profile) return visible def _webdav_discovery_candidates(payload: FileConnectorDiscoveryRequest) -> list[str]: endpoint_url = str(payload.endpoint_url or "").strip() if not endpoint_url: return [] direct_url = endpoint_url if endpoint_url.endswith("/") else endpoint_url + "/" candidates = [direct_url] if payload.provider == "nextcloud": root_url = _nextcloud_root_url(direct_url) username = (payload.credentials.username or "").strip() if username: candidates.append(urljoin(root_url, f"remote.php/dav/files/{quote(username, safe='')}/")) candidates.append(urljoin(root_url, "remote.php/webdav/")) seen: set[str] = set() unique: list[str] = [] for candidate in candidates: normalized = candidate.rstrip("/") if not normalized or normalized in seen: continue seen.add(normalized) unique.append(candidate) return unique def _nextcloud_root_url(endpoint_url: str) -> str: marker = "/remote.php/" if marker in endpoint_url: return endpoint_url.split(marker, 1)[0].rstrip("/") + "/" return endpoint_url.rstrip("/") + "/" def _same_endpoint(left: str, right: str) -> bool: return left.strip().rstrip("/") == right.strip().rstrip("/") def _discovery_profile_from_payload( payload: FileConnectorDiscoveryRequest, endpoint_url: str, *, principal: ApiPrincipal, ) -> ConnectorProfile: credentials = payload.credentials return ConnectorProfile( id=f"discovery-{principal.tenant_id}", label="Connector discovery", provider=payload.provider, endpoint_url=endpoint_url, base_path=payload.base_path, scope_type="tenant", scope_id=principal.tenant_id, credential_mode=payload.credential_mode, username=credentials.username, password_env=credentials.password_env, token_env=credentials.token_env, secret_ref=credentials.secret_ref, password_value=credentials.password, token_value=credentials.token, capabilities=("browse",), metadata=payload.metadata, source_kind="discovery", ) def _visible_connector_profile( session: Session, principal: ApiPrincipal, profile_id: str, *, campaign_id: str | None = None, include_disabled: bool = False, include_effective_policy: bool = True, ) -> ConnectorProfile: for profile in _visible_connector_profiles( session, principal, campaign_id=campaign_id, include_disabled=include_disabled, include_admin_scopes=False, include_effective_policy=include_effective_policy, ): if profile.id == profile_id: return profile raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Connector profile not found") def _with_effective_connector_policy(session: Session, principal: ApiPrincipal, profile: ConnectorProfile) -> ConnectorProfile: sources = effective_connector_policy_sources( session, tenant_id=principal.tenant_id, scope_type=profile.scope_type, scope_id=profile.scope_id, ) if not sources: return profile return replace(profile, policy_sources=tuple([*sources, *profile.policy_sources])) def _require_connector_profile_write(principal: ApiPrincipal, scope_type: str) -> None: clean_scope = scope_type.strip().casefold() if clean_scope == "system": if has_scope(principal, "system:settings:write") or has_scope(principal, "files:file:admin"): return raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail="Missing scope: system:settings:write") if has_scope(principal, "files:file:admin") or has_scope(principal, "admin:settings:write") or has_scope(principal, "system:settings:write"): return raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail="Missing scope: files:file:admin") def _require_connector_credential_write(principal: ApiPrincipal, scope_type: str) -> None: _require_connector_profile_write(principal, scope_type) def _require_connector_policy_read(principal: ApiPrincipal, scope_type: str) -> None: clean_scope = scope_type.strip().casefold() if clean_scope == "system": if any(has_scope(principal, scope) for scope in ("files:file:admin", "system:settings:read", "system:settings:write")): return raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail="Missing scope: system:settings:read") if any(has_scope(principal, scope) for scope in ("files:file:admin", "admin:settings:read", "admin:settings:write", "system:settings:read", "system:settings:write")): return raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail="Missing scope: files:file:admin") def _connector_policy_configure_sources( session: Session, principal: ApiPrincipal, *, scope_type: str, scope_id: str | None, fields: tuple[str, ...], ): sources = effective_connector_policy_sources(session, tenant_id=principal.tenant_id, scope_type=scope_type, scope_id=scope_id) return connector_policy_sources_for_fields(sources, fields) def _ensure_connector_configuration_allowed( session: Session, principal: ApiPrincipal, *, connector_id: str | None, credential_id: str | None, provider: str | None, endpoint_url: str | None, base_path: str | None, scope_type: str, scope_id: str | None, operation: str, ) -> None: sources = _connector_policy_configure_sources( session, principal, scope_type=scope_type, scope_id=scope_id, fields=("connectors", "credentials", "providers", "external_paths", "external_urls"), ) if not sources: return ensure_connector_policy_allows( ConnectorAccessRequest( connector_id=connector_id, credential_id=credential_id, provider=provider, external_path=base_path, external_url=endpoint_url, operation=operation, ), sources, ) def _ensure_connector_credential_configuration_allowed( session: Session, principal: ApiPrincipal, *, credential_id: str, provider: str | None, scope_type: str, scope_id: str | None, operation: str, ) -> None: sources = _connector_policy_configure_sources( session, principal, scope_type=scope_type, scope_id=scope_id, fields=("credentials", "providers"), ) if not sources: return ensure_connector_policy_allows( ConnectorAccessRequest( credential_id=credential_id, provider=provider, operation=operation, ), sources, ) def _ensure_connector_local_policy_allowed( session: Session, principal: ApiPrincipal, *, scope_type: str, scope_id: str | None, policy: dict[str, Any] | None, ) -> None: if policy is None: return validate_connector_policy_allowed_by_parent( parent_connector_policy(session, tenant_id=principal.tenant_id, scope_type=scope_type, scope_id=scope_id), policy, ) def _connector_credential_response(row) -> FileConnectorCredentialResponse: return FileConnectorCredentialResponse(**connector_credential_from_row(row).to_response()) def _credential_row_for_profile( session: Session, principal: ApiPrincipal, *, credential_profile_id: str | None, provider: str, include_disabled: bool = False, ): credential_id = (credential_profile_id or "").strip() if not credential_id: return None row = get_connector_credential_row( session, tenant_id=principal.tenant_id, credential_id=credential_id, include_disabled=include_disabled, ) if row.provider and row.provider != provider: raise FileStorageError(f"Connector credential {credential_id} is limited to {row.provider} connections") return row def _download_connector_payload( profile: ConnectorProfile, payload: FileConnectorImportRequest, *, operation: str, ) -> tuple[str, Any, dict[str, Any]]: source_path = normalize_connector_browse_path(payload.path) decision = connector_policy_decision( ConnectorAccessRequest( connector_id=profile.id, credential_id=profile.credential_profile_id, provider=profile.provider, external_id=f"{payload.library_id}:{source_path}", external_path=source_path, external_url=profile.endpoint_url, operation=operation, ), profile.policy_sources, ) if not decision.allowed: raise ConnectorPolicyDenied(decision) downloaded = read_connector_file( profile, library_id=payload.library_id, path=source_path, max_bytes=settings.file_upload_max_bytes, ) provenance_metadata = { "profile_id": profile.id, "library_id": payload.library_id, "library_path": source_path, **downloaded.metadata, **payload.metadata, } metadata = source_metadata( source_provenance={ "source_type": "connector", "connector_id": profile.id, "provider": profile.provider, "external_id": downloaded.external_id or f"{payload.library_id}:{source_path}", "external_path": source_path, "external_url": downloaded.external_url, "metadata": provenance_metadata, }, source_revision=payload.source_revision or downloaded.revision, ) return source_path, downloaded, metadata or {} def _connector_audit_details(asset: FileAsset, version: FileVersion, blob: FileBlob, *, operation: str) -> dict[str, object] | None: metadata = asset.metadata_ or {} provenance = source_provenance_from_metadata(metadata) if not provenance: return None return { "operation": operation, "file_id": asset.id, "display_path": asset.display_path, "version_id": version.id, "blob_id": blob.id, "checksum_sha256": blob.checksum_sha256, "size_bytes": blob.size_bytes, "source_revision": source_revision_from_metadata(metadata), "source_provenance": provenance, } def _audit_connector_event( session: Session, principal: ApiPrincipal, *, action: str, asset: FileAsset, version: FileVersion, blob: FileBlob, operation: str, extra_details: dict[str, object] | None = None, commit: bool = False, ) -> bool: details = _connector_audit_details(asset, version, blob, operation=operation) if details is None: return False if extra_details: details.update(extra_details) audit_from_principal( session, principal, action=action, object_type="file", object_id=asset.id, details=details, commit=commit, ) return True def _audit_connector_imports(session: Session, principal: ApiPrincipal, assets: list[FileAsset]) -> None: for asset in assets: version, blob = current_version_and_blob(session, asset) _audit_connector_event( session, principal, action="files.connector.imported", asset=asset, version=version, blob=blob, operation="import", ) def _audit_connector_sync(session: Session, principal: ApiPrincipal, asset: FileAsset, *, sync_action: str, previous_version_id: str | None) -> None: version, blob = current_version_and_blob(session, asset) _audit_connector_event( session, principal, action="files.connector.synced", asset=asset, version=version, blob=blob, operation="sync", extra_details={ "sync_action": sync_action, "previous_version_id": previous_version_id, }, ) def _audit_connector_access(session: Session, principal: ApiPrincipal, assets: list[FileAsset], *, operation: str) -> None: recorded = False for asset in assets: version, blob = current_version_and_blob(session, asset) recorded = _audit_connector_event( session, principal, action="files.connector.accessed", asset=asset, version=version, blob=blob, operation=operation, ) or recorded if recorded: session.commit() def _asset_response(session: Session, asset: FileAsset, *, include_shares: bool = False) -> FileAssetResponse: version, blob = current_version_and_blob(session, asset) metadata = asset.metadata_ or {} shares: list[FileShareResponse] = [] if include_shares: rows = session.query(FileShare).filter(FileShare.file_asset_id == asset.id).order_by(FileShare.created_at.desc()).all() shares = [ FileShareResponse( id=row.id, target_type=row.target_type, target_id=row.target_id, permission=row.permission, created_at=row.created_at.isoformat(), revoked_at=row.revoked_at.isoformat() if row.revoked_at else None, ) for row in rows ] return FileAssetResponse( id=asset.id, tenant_id=asset.tenant_id, owner_type=asset.owner_type, owner_id=_owner_id(asset), display_path=asset.display_path, filename=asset.filename, description=asset.description, size_bytes=blob.size_bytes, content_type=blob.content_type, checksum_sha256=blob.checksum_sha256, version_id=version.id, created_at=asset.created_at.isoformat(), updated_at=asset.updated_at.isoformat(), deleted_at=asset.deleted_at.isoformat() if asset.deleted_at else None, audit_relevant=asset_is_audit_relevant(session, asset), metadata=metadata, source_provenance=source_provenance_from_metadata(metadata), source_revision=source_revision_from_metadata(metadata), shares=shares, ) def _asset_list_response(session: Session, assets: list[FileAsset], *, include_shares: bool = False) -> list[FileAssetResponse]: if not assets: return [] asset_ids = [asset.id for asset in assets] version_ids = [asset.current_version_id for asset in assets if asset.current_version_id] if len(version_ids) != len(assets): raise FileStorageError("File has no current version") version_blob_rows: list[tuple[FileVersion, FileBlob]] = [] for chunk in _chunks(version_ids): version_blob_rows.extend( session.query(FileVersion, FileBlob) .join(FileBlob, FileBlob.id == FileVersion.blob_id) .filter(FileVersion.id.in_(chunk)) .all() ) versions_by_id = {version.id: version for version, _blob in version_blob_rows} blobs_by_version_id = {version.id: blob for version, blob in version_blob_rows} shares_by_asset_id: dict[str, list[FileShareResponse]] = {asset_id: [] for asset_id in asset_ids} if include_shares: for chunk in _chunks(asset_ids): share_rows = ( session.query(FileShare) .filter(FileShare.file_asset_id.in_(chunk)) .order_by(FileShare.created_at.desc()) .all() ) for row in share_rows: shares_by_asset_id.setdefault(row.file_asset_id, []).append( FileShareResponse( id=row.id, target_type=row.target_type, target_id=row.target_id, permission=row.permission, created_at=row.created_at.isoformat(), revoked_at=row.revoked_at.isoformat() if row.revoked_at else None, ) ) sent_asset_ids: set[str] = set() for chunk in _chunks(asset_ids): sent_rows = ( session.query(CampaignAttachmentUse.file_asset_id) .filter(CampaignAttachmentUse.file_asset_id.in_(chunk), CampaignAttachmentUse.use_stage == "sent") .distinct() .all() ) sent_asset_ids.update(row[0] for row in sent_rows) responses: list[FileAssetResponse] = [] for asset in assets: version = versions_by_id.get(asset.current_version_id or "") blob = blobs_by_version_id.get(asset.current_version_id or "") if not version or not blob: raise FileStorageError("File version not found") metadata = asset.metadata_ or {} responses.append( FileAssetResponse( id=asset.id, tenant_id=asset.tenant_id, owner_type=asset.owner_type, owner_id=_owner_id(asset), display_path=asset.display_path, filename=asset.filename, description=asset.description, size_bytes=blob.size_bytes, content_type=blob.content_type, checksum_sha256=blob.checksum_sha256, version_id=version.id, created_at=asset.created_at.isoformat(), updated_at=asset.updated_at.isoformat(), deleted_at=asset.deleted_at.isoformat() if asset.deleted_at else None, audit_relevant=asset.id in sent_asset_ids, metadata=metadata, source_provenance=source_provenance_from_metadata(metadata), source_revision=source_revision_from_metadata(metadata), shares=shares_by_asset_id.get(asset.id, []), ) ) return responses def _chunks(values: list[str], size: int = 900): for index in range(0, len(values), size): yield values[index:index + size] def _folder_owner_id(folder: FileFolder) -> str: return folder.owner_user_id if folder.owner_type == "user" else folder.owner_group_id # type: ignore[return-value] def _folder_response(folder: FileFolder) -> FileFolderResponse: return FileFolderResponse( id=folder.id, tenant_id=folder.tenant_id, owner_type=folder.owner_type, owner_id=_folder_owner_id(folder), path=folder.path, created_at=folder.created_at.isoformat(), updated_at=folder.updated_at.isoformat(), deleted_at=folder.deleted_at.isoformat() if folder.deleted_at else None, ) def _connector_space_response(space: FileConnectorSpace) -> FileConnectorSpaceResponse: return FileConnectorSpaceResponse( id=space.id, tenant_id=space.tenant_id, owner_type=space.owner_type, # type: ignore[arg-type] owner_id=connector_space_owner_id(space), label=space.label, connector_profile_id=space.connector_profile_id, provider=space.provider, library_id=space.library_id, remote_path=space.remote_path, sync_mode=space.sync_mode, read_only=space.read_only, is_active=space.is_active, created_at=space.created_at.isoformat(), updated_at=space.updated_at.isoformat(), deleted_at=space.deleted_at.isoformat() if space.deleted_at else None, metadata=space.metadata_ or {}, ) def _connector_space_file_space_response(space: FileConnectorSpace) -> FileSpaceResponse: owner_id = connector_space_owner_id(space) return FileSpaceResponse( id=f"connector:{space.id}", label=space.label, owner_type=space.owner_type, # type: ignore[arg-type] owner_id=owner_id, description=f"Linked {space.provider} folder.", space_type="connector", connector_space_id=space.id, connector_profile_id=space.connector_profile_id, provider=space.provider, library_id=space.library_id, remote_path=space.remote_path, sync_mode=space.sync_mode, read_only=space.read_only, ) def _connector_space_policy_decision( profile: ConnectorProfile, *, library_id: str | None, remote_path: str | None, operation: str, ): browse_path = normalize_connector_browse_path(remote_path) external_id = f"{library_id}:{browse_path}" if library_id else browse_path or profile.id return connector_policy_decision( ConnectorAccessRequest( connector_id=profile.id, credential_id=profile.credential_profile_id, provider=profile.provider, external_id=external_id, external_path=browse_path, external_url=profile.endpoint_url, operation=operation, ), profile.policy_sources, ) def _ensure_list_owner_access(session: Session, principal: ApiPrincipal, owner_type: str | None, owner_id: str | None) -> None: if not owner_type: return if owner_type == "user" and owner_id and owner_id != principal.user.id and not _is_admin(principal): raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail="No access to this user file space") if owner_type == "group" and owner_id: try: ensure_group_access( session, tenant_id=principal.tenant_id, group_id=owner_id, user_id=principal.user.id, is_admin=_is_admin(principal), ) except FileStorageError as exc: raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail=str(exc)) from exc _FILES_DELTA_COLLECTIONS = (FILES_ASSETS_COLLECTION, FILES_FOLDERS_COLLECTION) FILES_LIST_CURSOR_SCOPE = "files.list.v1" FOLDERS_LIST_CURSOR_SCOPE = "files.folders.list.v1" DEFAULT_FILE_LIST_PAGE_SIZE = 500 def _cursor_http_error(exc: Exception) -> HTTPException: return HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail=str(exc)) def _cursor_page_size(scope: str, cursor: str | None, explicit_page_size: int | None) -> int | None: if explicit_page_size is not None: return explicit_page_size if not cursor: return None try: values = decode_keyset_cursor(scope, cursor) except KeysetCursorError as exc: raise _cursor_http_error(exc) from exc raw_page_size = values.get("page_size") if not isinstance(raw_page_size, int) or raw_page_size < 1 or raw_page_size > 1000: raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail="Invalid pagination cursor") return raw_page_size def _files_list_fingerprint( principal: ApiPrincipal, *, owner_type: Literal["user", "group"] | None, owner_id: str | None, campaign_id: str | None, path_prefix: str | None, page_size: int, ) -> str: return keyset_query_fingerprint( FILES_LIST_CURSOR_SCOPE, { "tenant_id": principal.tenant_id, "actor": "admin" if _is_admin(principal) else principal.user.id, "owner_type": owner_type, "owner_id": owner_id, "campaign_id": campaign_id, "path_prefix": path_prefix or "", "sort": "display_path.asc,updated_at.desc,id.asc", "page_size": page_size, }, ) def _folders_list_fingerprint( principal: ApiPrincipal, *, owner_type: Literal["user", "group"], owner_id: str, page_size: int, ) -> str: return keyset_query_fingerprint( FOLDERS_LIST_CURSOR_SCOPE, { "tenant_id": principal.tenant_id, "actor": "admin" if _is_admin(principal) else principal.user.id, "owner_type": owner_type, "owner_id": owner_id, "sort": "path.asc,id.asc", "page_size": page_size, }, ) def _file_cursor_values(cursor: str | None, *, fingerprint: str) -> tuple[str | None, datetime | None, str | None]: try: values = decode_keyset_cursor(FILES_LIST_CURSOR_SCOPE, cursor, fingerprint=fingerprint) except KeysetCursorError as exc: raise _cursor_http_error(exc) from exc if values is None: return None, None, None display_path = values.get("display_path") updated_at = values.get("updated_at") asset_id = values.get("id") if not isinstance(display_path, str) or not isinstance(updated_at, str) or not isinstance(asset_id, str): raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail="Invalid pagination cursor") try: parsed_updated_at = datetime.fromisoformat(updated_at) except ValueError as exc: raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail="Invalid pagination cursor") from exc return display_path, parsed_updated_at, asset_id def _folder_cursor_values(cursor: str | None, *, fingerprint: str) -> tuple[str | None, str | None]: try: values = decode_keyset_cursor(FOLDERS_LIST_CURSOR_SCOPE, cursor, fingerprint=fingerprint) except KeysetCursorError as exc: raise _cursor_http_error(exc) from exc if values is None: return None, None folder_path = values.get("path") folder_id = values.get("id") if not isinstance(folder_path, str) or not isinstance(folder_id, str): raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail="Invalid pagination cursor") return folder_path, folder_id def _next_file_list_cursor( principal: ApiPrincipal, assets: list[FileAsset], *, owner_type: Literal["user", "group"] | None, owner_id: str | None, campaign_id: str | None, path_prefix: str | None, page_size: int, has_more: bool, ) -> str | None: if not has_more or not assets: return None last = assets[-1] return encode_keyset_cursor( FILES_LIST_CURSOR_SCOPE, fingerprint=_files_list_fingerprint( principal, owner_type=owner_type, owner_id=owner_id, campaign_id=campaign_id, path_prefix=path_prefix, page_size=page_size, ), values={"display_path": last.display_path, "updated_at": last.updated_at, "id": last.id, "page_size": page_size}, ) def _next_folder_list_cursor( principal: ApiPrincipal, folders: list[FileFolder], *, owner_type: Literal["user", "group"], owner_id: str, page_size: int, has_more: bool, ) -> str | None: if not has_more or not folders: return None last = folders[-1] return encode_keyset_cursor( FOLDERS_LIST_CURSOR_SCOPE, fingerprint=_folders_list_fingerprint(principal, owner_type=owner_type, owner_id=owner_id, page_size=page_size), values={"path": last.path, "id": last.id, "page_size": page_size}, ) def _files_delta_watermark(session: Session, tenant_id: str) -> str: return encode_sequence_watermark( max_sequence_id( session, tenant_id=tenant_id, module_id=FILES_MODULE_ID, collections=_FILES_DELTA_COLLECTIONS, ) ) def _full_file_delta_response( session: Session, *, principal: ApiPrincipal, owner_type: Literal["user", "group"] | None, owner_id: str | None, campaign_id: str | None, path_prefix: str | None, ) -> FileDeltaResponse: assets = list_assets_for_user( session, tenant_id=principal.tenant_id, user_id=principal.user.id, owner_type=owner_type, owner_id=owner_id, campaign_id=campaign_id, path_prefix=path_prefix, is_admin=_is_admin(principal), ) folders = _visible_folders_for_delta( session, principal=principal, owner_type=owner_type, owner_id=owner_id, path_prefix=path_prefix, ) return FileDeltaResponse( files=_asset_list_response(session, assets, include_shares=True), folders=[_folder_response(folder) for folder in folders], deleted=[], watermark=_files_delta_watermark(session, principal.tenant_id), has_more=False, full=True, ) def _visible_folders_for_delta( session: Session, *, principal: ApiPrincipal, owner_type: Literal["user", "group"] | None, owner_id: str | None, path_prefix: str | None, ) -> list[FileFolder]: if not owner_type or not owner_id: return [] folders = list_folders_for_user( session, tenant_id=principal.tenant_id, user_id=principal.user.id, owner_type=owner_type, owner_id=owner_id, is_admin=_is_admin(principal), ) if not path_prefix: return folders normalized = path_prefix.strip().strip("/") if not normalized: return folders return [folder for folder in folders if folder.path == normalized or folder.path.startswith(f"{normalized}/")] def _entry_path_matches(entry_payload: dict[str, object], path_prefix: str | None) -> bool: if not path_prefix: return True normalized = path_prefix.strip().strip("/") if not normalized: return True for key in ("path", "previous_path"): value = entry_payload.get(key) if isinstance(value, str) and (value == normalized or value.startswith(f"{normalized}/")): return True return False def _entry_owner_matches( entry_payload: dict[str, object], owner_type: Literal["user", "group"] | None, owner_id: str | None, ) -> bool: if not owner_type or not owner_id: return True return entry_payload.get("owner_type") == owner_type and entry_payload.get("owner_id") == owner_id def _entry_campaign_matches(session: Session, entry, campaign_id: str | None) -> bool: if not campaign_id: return True payload = entry.payload or {} if payload.get("share_target_type") == "campaign" and payload.get("share_target_id") == campaign_id: return True if entry.resource_type != "file": return False return ( session.query(FileShare) .filter( FileShare.tenant_id == entry.tenant_id, FileShare.file_asset_id == entry.resource_id, FileShare.target_type == "campaign", FileShare.target_id == campaign_id, FileShare.revoked_at.is_(None), ) .first() is not None ) def _principal_group_ids_for_delta(session: Session, principal: ApiPrincipal) -> set[str]: cache = session.info.setdefault("files_delta_group_ids", {}) key = (principal.tenant_id, principal.user.id) if key not in cache: cache[key] = set(user_group_ids(session, tenant_id=principal.tenant_id, user_id=principal.user.id)) return cache[key] def _entry_subject_matches_principal(session: Session, principal: ApiPrincipal, entry_payload: dict[str, object]) -> bool: if _is_admin(principal): return True owner_type = entry_payload.get("owner_type") owner_id = entry_payload.get("owner_id") if owner_type == "user" and owner_id == principal.user.id: return True if owner_type == "group" and isinstance(owner_id, str): if owner_id in _principal_group_ids_for_delta(session, principal): return True share_target_type = entry_payload.get("share_target_type") share_target_id = entry_payload.get("share_target_id") if share_target_type == "user" and share_target_id == principal.user.id: return True if share_target_type == "tenant" and share_target_id == principal.tenant_id: return True if share_target_type == "group" and isinstance(share_target_id, str): if share_target_id in _principal_group_ids_for_delta(session, principal): return True return False def _entry_matches_delta_scope( session: Session, principal: ApiPrincipal, entry, *, owner_type: Literal["user", "group"] | None, owner_id: str | None, campaign_id: str | None, path_prefix: str | None, ) -> bool: payload = entry.payload or {} if not owner_type and not campaign_id and not _entry_subject_matches_principal(session, principal, payload): return False return ( _entry_owner_matches(payload, owner_type, owner_id) and _entry_path_matches(payload, path_prefix) and _entry_campaign_matches(session, entry, campaign_id) ) def _files_delta_response( session: Session, *, principal: ApiPrincipal, owner_type: Literal["user", "group"] | None, owner_id: str | None, campaign_id: str | None, path_prefix: str | None, since: str, limit: int, ) -> FileDeltaResponse: try: since_sequence = decode_sequence_watermark(since) except ValueError as exc: raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail=str(exc)) from exc if sequence_watermark_is_expired( session, since=since_sequence, tenant_id=principal.tenant_id, module_id=FILES_MODULE_ID, collections=_FILES_DELTA_COLLECTIONS, ): return _full_file_delta_response( session, principal=principal, owner_type=owner_type, owner_id=owner_id, campaign_id=campaign_id, path_prefix=path_prefix, ) entries_plus_one = sequence_entries_since( session, since=since_sequence, tenant_id=principal.tenant_id, module_id=FILES_MODULE_ID, collections=_FILES_DELTA_COLLECTIONS, limit=limit + 1, ) has_more = len(entries_plus_one) > limit entries = entries_plus_one[:limit] changed_file_ids = list(dict.fromkeys(entry.resource_id for entry in entries if entry.resource_type == "file")) changed_folder_ids = list(dict.fromkeys(entry.resource_id for entry in entries if entry.resource_type == "folder")) visible_assets = { asset.id: asset for asset in list_assets_for_user( session, tenant_id=principal.tenant_id, user_id=principal.user.id, owner_type=owner_type, owner_id=owner_id, campaign_id=campaign_id, path_prefix=path_prefix, is_admin=_is_admin(principal), ) if asset.id in changed_file_ids } visible_folders = { folder.id: folder for folder in _visible_folders_for_delta( session, principal=principal, owner_type=owner_type, owner_id=owner_id, path_prefix=path_prefix, ) if folder.id in changed_folder_ids } deleted: dict[tuple[str, str], DeltaDeletedItem] = {} for entry in entries: is_visible = ( (entry.resource_type == "file" and entry.resource_id in visible_assets) or (entry.resource_type == "folder" and entry.resource_id in visible_folders) ) if is_visible: continue if not _entry_matches_delta_scope( session, principal, entry, owner_type=owner_type, owner_id=owner_id, campaign_id=campaign_id, path_prefix=path_prefix, ): continue deleted[(entry.resource_type, entry.resource_id)] = DeltaDeletedItem( id=entry.resource_id, resource_type=entry.resource_type, revision=encode_sequence_watermark(entry.id), deleted_at=entry.created_at if entry.operation == "deleted" else None, ) watermark = encode_sequence_watermark(entries[-1].id) if has_more and entries else _files_delta_watermark(session, principal.tenant_id) return FileDeltaResponse( files=_asset_list_response(session, list(visible_assets.values()), include_shares=True), folders=[_folder_response(folder) for folder in visible_folders.values()], deleted=list(deleted.values()), watermark=watermark, has_more=has_more, full=False, ) @router.get("/spaces", response_model=FileSpacesResponse) def list_file_spaces( session: Session = Depends(get_session), principal: ApiPrincipal = Depends(require_scope("files:file:read")), ): spaces = [ FileSpaceResponse( id=f"user:{principal.user.id}", label="My files", owner_type="user", owner_id=principal.user.id, description="Files owned by your user account.", ) ] group_ids = user_group_ids(session, tenant_id=principal.tenant_id, user_id=principal.user.id, include_admin_groups=_is_admin(principal)) if group_ids: groups = group_refs_for_ids(tenant_id=principal.tenant_id, group_ids=group_ids) spaces.extend( FileSpaceResponse( id=f"group:{group.id}", label=f"{group.name} files", owner_type="group", owner_id=group.id, description="Files owned by this group.", ) for group in groups ) connector_spaces = list_connector_spaces_for_user( session, tenant_id=principal.tenant_id, user_id=principal.user.id, is_admin=_is_admin(principal), ) spaces.extend(_connector_space_file_space_response(space) for space in connector_spaces) return FileSpacesResponse(spaces=spaces) @router.get("/connector-spaces", response_model=FileConnectorSpacesResponse) def list_file_connector_spaces( owner_type: Literal["user", "group"] | None = None, owner_id: str | None = None, include_inactive: bool = False, session: Session = Depends(get_session), principal: ApiPrincipal = Depends(require_scope("files:file:read")), ): try: spaces = list_connector_spaces_for_user( session, tenant_id=principal.tenant_id, user_id=principal.user.id, owner_type=owner_type, owner_id=owner_id, include_inactive=include_inactive and _is_admin(principal), is_admin=_is_admin(principal), ) return FileConnectorSpacesResponse(spaces=[_connector_space_response(space) for space in spaces]) except FileStorageError as exc: raise _http_error(exc) from exc @router.post("/connector-spaces", response_model=FileConnectorSpaceResponse, status_code=status.HTTP_201_CREATED) def create_file_connector_space( payload: FileConnectorSpaceCreateRequest, session: Session = Depends(get_session), principal: ApiPrincipal = Depends(require_scope("files:file:organize")), ): if payload.owner_type == "group" and not payload.owner_id: raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail="owner_id is required for group connector spaces") target_owner = payload.owner_id or principal.user.id try: profile = _visible_connector_profile(session, principal, payload.connector_profile_id) decision = _connector_space_policy_decision( profile, library_id=payload.library_id, remote_path=payload.remote_path, operation="link", ) if not decision.allowed: raise ConnectorPolicyDenied(decision) space = create_connector_space( session, tenant_id=principal.tenant_id, owner_type=payload.owner_type, owner_id=target_owner, user_id=principal.user.id, label=payload.label, profile=profile, library_id=payload.library_id, remote_path=payload.remote_path, sync_mode=payload.sync_mode, metadata=payload.metadata, is_admin=_is_admin(principal), ) _record_connector_settings_change( session, collection=FILES_CONNECTOR_SPACES_COLLECTION, resource_type=FILES_CONNECTOR_SPACE_RESOURCE, resource_id=space.id, operation="created", principal=principal, tenant_id=space.tenant_id, payload={"owner_type": space.owner_type, "owner_id": connector_space_owner_id(space)}, ) session.commit() return _connector_space_response(space) except ConnectorPolicyDenied as exc: session.rollback() raise _connector_policy_error(exc) from exc except (FileStorageError, ValueError, json.JSONDecodeError) as exc: session.rollback() raise _http_error(exc) from exc @router.patch("/connector-spaces/{space_id}", response_model=FileConnectorSpaceResponse) def update_file_connector_space( space_id: str, payload: FileConnectorSpaceUpdateRequest, session: Session = Depends(get_session), principal: ApiPrincipal = Depends(require_scope("files:file:organize")), ): try: space = get_connector_space_for_user( session, tenant_id=principal.tenant_id, user_id=principal.user.id, space_id=space_id, include_inactive=_is_admin(principal), is_admin=_is_admin(principal), ) except FileStorageError as exc: raise _http_error(exc, not_found=True) from exc try: if payload.library_id is not None or payload.remote_path is not None: profile = _visible_connector_profile(session, principal, space.connector_profile_id) decision = _connector_space_policy_decision( profile, library_id=payload.library_id if payload.library_id is not None else space.library_id, remote_path=payload.remote_path if payload.remote_path is not None else space.remote_path, operation="link", ) if not decision.allowed: raise ConnectorPolicyDenied(decision) update_connector_space( session, space, user_id=principal.user.id, label=payload.label, library_id=payload.library_id, remote_path=payload.remote_path, sync_mode=payload.sync_mode, is_active=payload.is_active, metadata=payload.metadata, is_admin=_is_admin(principal), ) _record_connector_settings_change( session, collection=FILES_CONNECTOR_SPACES_COLLECTION, resource_type=FILES_CONNECTOR_SPACE_RESOURCE, resource_id=space.id, operation="updated", principal=principal, tenant_id=space.tenant_id, payload={"owner_type": space.owner_type, "owner_id": connector_space_owner_id(space)}, ) session.commit() return _connector_space_response(space) except ConnectorPolicyDenied as exc: session.rollback() raise _connector_policy_error(exc) from exc except (FileStorageError, ValueError, json.JSONDecodeError) as exc: session.rollback() raise _http_error(exc) from exc @router.delete("/connector-spaces/{space_id}", response_model=FileConnectorSpaceResponse) def delete_file_connector_space( space_id: str, session: Session = Depends(get_session), principal: ApiPrincipal = Depends(require_scope("files:file:organize")), ): try: space = get_connector_space_for_user( session, tenant_id=principal.tenant_id, user_id=principal.user.id, space_id=space_id, include_inactive=True, is_admin=_is_admin(principal), ) soft_delete_connector_space(session, space, user_id=principal.user.id, is_admin=_is_admin(principal)) _record_connector_settings_change( session, collection=FILES_CONNECTOR_SPACES_COLLECTION, resource_type=FILES_CONNECTOR_SPACE_RESOURCE, resource_id=space.id, operation="deleted", principal=principal, tenant_id=space.tenant_id, payload={"owner_type": space.owner_type, "owner_id": connector_space_owner_id(space)}, ) session.commit() return _connector_space_response(space) except FileStorageError as exc: session.rollback() raise _http_error(exc, not_found=True) from exc @router.get("/folders", response_model=FileFoldersResponse) def list_file_folders( owner_type: Literal["user", "group"], owner_id: str, page_size: int | None = Query(default=None, ge=1, le=1000), cursor: str | None = None, session: Session = Depends(get_session), principal: ApiPrincipal = Depends(require_scope("files:file:read")), ): try: watermark = _files_delta_watermark(session, principal.tenant_id) effective_page_size = _cursor_page_size(FOLDERS_LIST_CURSOR_SCOPE, cursor, page_size) if effective_page_size is None: folders = list_folders_for_user( session, tenant_id=principal.tenant_id, user_id=principal.user.id, owner_type=owner_type, owner_id=owner_id, is_admin=_is_admin(principal), ) return FileFoldersResponse(folders=[_folder_response(folder) for folder in folders], watermark=watermark) fingerprint = _folders_list_fingerprint(principal, owner_type=owner_type, owner_id=owner_id, page_size=effective_page_size) after_path, after_id = _folder_cursor_values(cursor, fingerprint=fingerprint) folders, has_more = list_folders_for_user_window( session, tenant_id=principal.tenant_id, user_id=principal.user.id, owner_type=owner_type, owner_id=owner_id, is_admin=_is_admin(principal), page_size=effective_page_size, after_path=after_path, after_id=after_id, ) return FileFoldersResponse( folders=[_folder_response(folder) for folder in folders], cursor=cursor, next_cursor=_next_folder_list_cursor( principal, folders, owner_type=owner_type, owner_id=owner_id, page_size=effective_page_size, has_more=has_more, ), watermark=watermark, ) except FileStorageError as exc: raise _http_error(exc) from exc @router.post("/folders", response_model=FileFolderResponse) def create_file_folder( payload: FileFolderCreateRequest, session: Session = Depends(get_session), principal: ApiPrincipal = Depends(require_scope("files:file:organize")), ): try: folder = create_folder( session, tenant_id=principal.tenant_id, owner_type=payload.owner_type, owner_id=payload.owner_id, user_id=principal.user.id, path=payload.path, is_admin=_is_admin(principal), ) session.commit() return _folder_response(folder) except (FileStorageError, UnsafeFilePathError, ValueError) as exc: session.rollback() raise _http_error(exc) from exc @router.post("/folders/delete", response_model=FileFolderDeleteResponse) def delete_file_folder( payload: FileFolderDeleteRequest, session: Session = Depends(get_session), principal: ApiPrincipal = Depends(require_scope("files:file:delete")), ): try: deleted_folders, deleted_files = soft_delete_folder( session, tenant_id=principal.tenant_id, owner_type=payload.owner_type, owner_id=payload.owner_id, user_id=principal.user.id, path=payload.path, recursive=payload.recursive, is_admin=_is_admin(principal), ) session.commit() return FileFolderDeleteResponse(deleted_folders=deleted_folders, deleted_files=deleted_files) except (FileStorageError, UnsafeFilePathError, ValueError) as exc: session.rollback() raise _http_error(exc) from exc @router.get("/delta", response_model=FileDeltaResponse) def files_delta( owner_type: Literal["user", "group"] | None = None, owner_id: str | None = None, campaign_id: str | None = None, path_prefix: str | None = None, since: str | None = None, limit: int = Query(default=500, ge=1, le=1000), session: Session = Depends(get_session), principal: ApiPrincipal = Depends(require_scope("files:file:read")), ): _ensure_list_owner_access(session, principal, owner_type, owner_id) _ensure_campaign_file_access(session, principal, campaign_id) if since is None: return _full_file_delta_response( session, principal=principal, owner_type=owner_type, owner_id=owner_id, campaign_id=campaign_id, path_prefix=path_prefix, ) return _files_delta_response( session, principal=principal, owner_type=owner_type, owner_id=owner_id, campaign_id=campaign_id, path_prefix=path_prefix, since=since, limit=limit, ) @router.get("", response_model=FileListResponse) def list_files( owner_type: Literal["user", "group"] | None = None, owner_id: str | None = None, campaign_id: str | None = None, path_prefix: str | None = None, page_size: int | None = Query(default=None, ge=1, le=1000), cursor: str | None = None, session: Session = Depends(get_session), principal: ApiPrincipal = Depends(require_scope("files:file:read")), ): _ensure_list_owner_access(session, principal, owner_type, owner_id) _ensure_campaign_file_access(session, principal, campaign_id) watermark = _files_delta_watermark(session, principal.tenant_id) effective_page_size = _cursor_page_size(FILES_LIST_CURSOR_SCOPE, cursor, page_size) if effective_page_size is not None: fingerprint = _files_list_fingerprint( principal, owner_type=owner_type, owner_id=owner_id, campaign_id=campaign_id, path_prefix=path_prefix, page_size=effective_page_size, ) after_display_path, after_updated_at, after_id = _file_cursor_values(cursor, fingerprint=fingerprint) assets, has_more = list_assets_for_user_window( session, tenant_id=principal.tenant_id, user_id=principal.user.id, owner_type=owner_type, owner_id=owner_id, campaign_id=campaign_id, path_prefix=path_prefix, is_admin=_is_admin(principal), page_size=effective_page_size, after_display_path=after_display_path, after_updated_at=after_updated_at, after_id=after_id, ) return FileListResponse( files=_asset_list_response(session, assets, include_shares=True), cursor=cursor, next_cursor=_next_file_list_cursor( principal, assets, owner_type=owner_type, owner_id=owner_id, campaign_id=campaign_id, path_prefix=path_prefix, page_size=effective_page_size, has_more=has_more, ), watermark=watermark, ) assets = list_assets_for_user( session, tenant_id=principal.tenant_id, user_id=principal.user.id, owner_type=owner_type, owner_id=owner_id, campaign_id=campaign_id, path_prefix=path_prefix, is_admin=_is_admin(principal), ) return FileListResponse(files=_asset_list_response(session, assets, include_shares=True), watermark=watermark) @router.post("/upload", response_model=FileUploadResponse) async def upload_files( files: list[UploadFile] = FastAPIFile(...), owner_type: Literal["user", "group"] = Form(default="user"), owner_id: str | None = Form(default=None), path: str = Form(default=""), campaign_id: str | None = Form(default=None), unpack_zip: bool = Form(default=False), conflict_strategy: Literal["reject", "overwrite", "rename"] = Form(default="reject"), conflict_resolutions_json: str | None = Form(default=None), source_provenance_json: str | None = Form(default=None), source_revision: str | None = Form(default=None), connector_policy_json: str | None = Form(default=None), session: Session = Depends(get_session), principal: ApiPrincipal = Depends(require_scope("files:file:upload")), ): target_owner = owner_id or principal.user.id uploaded_assets: list[FileAsset] = [] try: raw_resolutions = json.loads(conflict_resolutions_json) if conflict_resolutions_json else [] upload_resolutions = _conflict_resolutions([ConflictResolutionRequest(**item) for item in raw_resolutions]) _enforce_connector_policy(source_provenance_json, connector_policy_json, operation="import") metadata = _source_metadata_from_form(source_provenance_json, source_revision) for upload in files: filename = upload.filename or "file" content_type = upload.content_type or None upload_limit = settings.file_upload_zip_max_bytes if unpack_zip and filename.lower().endswith(".zip") else settings.file_upload_max_bytes if unpack_zip and filename.lower().endswith(".zip"): zip_path = await _spool_limited_upload_to_temp(upload, max_bytes=upload_limit, suffix=".zip") try: extracted = extract_zip_upload( session, tenant_id=principal.tenant_id, owner_type=owner_type, owner_id=target_owner, user_id=principal.user.id, zip_data=zip_path, folder=path, campaign_id=campaign_id, conflict_strategy=conflict_strategy, conflict_resolutions=upload_resolutions, metadata=metadata, is_admin=_is_admin(principal), max_file_bytes=settings.file_upload_max_bytes, max_total_bytes=settings.file_upload_zip_max_bytes, ) finally: _cleanup_temp_file(zip_path) uploaded_assets.extend(item.asset for item in extracted) continue data = await _read_limited_upload(upload, max_bytes=upload_limit) stored = create_file_asset( session, tenant_id=principal.tenant_id, owner_type=owner_type, owner_id=target_owner, user_id=principal.user.id, filename=filename, data=data, folder=path, content_type=content_type, campaign_id=campaign_id, conflict_strategy=conflict_strategy, conflict_resolutions=upload_resolutions, metadata=metadata, is_admin=_is_admin(principal), ) uploaded_assets.append(stored.asset) _audit_connector_imports(session, principal, uploaded_assets) session.commit() except ConnectorPolicyDenied as exc: session.rollback() raise _connector_policy_error(exc) from exc except (FileStorageError, UnsafeFilePathError, ValueError) as exc: session.rollback() raise _http_error(exc) from exc return FileUploadResponse(files=[_asset_response(session, asset, include_shares=True) for asset in uploaded_assets]) @router.post("/upload-zip", response_model=FileUploadResponse) async def upload_zip( file: UploadFile = FastAPIFile(...), owner_type: Literal["user", "group"] = Form(default="user"), owner_id: str | None = Form(default=None), path: str = Form(default=""), campaign_id: str | None = Form(default=None), conflict_strategy: Literal["reject", "overwrite", "rename"] = Form(default="reject"), conflict_resolutions_json: str | None = Form(default=None), source_provenance_json: str | None = Form(default=None), source_revision: str | None = Form(default=None), connector_policy_json: str | None = Form(default=None), session: Session = Depends(get_session), principal: ApiPrincipal = Depends(require_scope("files:file:upload")), ): target_owner = owner_id or principal.user.id zip_path: str | None = None try: raw_resolutions = json.loads(conflict_resolutions_json) if conflict_resolutions_json else [] upload_resolutions = _conflict_resolutions([ConflictResolutionRequest(**item) for item in raw_resolutions]) _enforce_connector_policy(source_provenance_json, connector_policy_json, operation="import") metadata = _source_metadata_from_form(source_provenance_json, source_revision) zip_path = await _spool_limited_upload_to_temp(file, max_bytes=settings.file_upload_zip_max_bytes, suffix=".zip") extracted = extract_zip_upload( session, tenant_id=principal.tenant_id, owner_type=owner_type, owner_id=target_owner, user_id=principal.user.id, zip_data=zip_path, folder=path, campaign_id=campaign_id, conflict_strategy=conflict_strategy, conflict_resolutions=upload_resolutions, metadata=metadata, is_admin=_is_admin(principal), max_file_bytes=settings.file_upload_max_bytes, max_total_bytes=settings.file_upload_zip_max_bytes, ) _audit_connector_imports(session, principal, [item.asset for item in extracted]) session.commit() except ConnectorPolicyDenied as exc: session.rollback() raise _connector_policy_error(exc) from exc except (FileStorageError, UnsafeFilePathError, ValueError) as exc: session.rollback() raise _http_error(exc) from exc finally: if zip_path: _cleanup_temp_file(zip_path) return FileUploadResponse(files=[_asset_response(session, item.asset, include_shares=True) for item in extracted]) @router.post("/connector-policy/evaluate", response_model=FileConnectorPolicyEvaluateResponse) def evaluate_connector_policy( payload: FileConnectorPolicyEvaluateRequest, principal: ApiPrincipal = Depends(require_any_scope("files:file:read", "files:file:upload", "files:file:download")), ): del principal sources = connector_policy_sources_from_payload([item.model_dump(mode="json") for item in payload.policy_sources]) request = ConnectorAccessRequest.from_provenance(payload.source_provenance.model_dump(mode="json", exclude_none=True), operation=payload.operation) return FileConnectorPolicyEvaluateResponse(decision=connector_policy_decision(request, sources).to_dict()) def _full_file_connector_settings_delta_response( session: Session, principal: ApiPrincipal, *, scope_type: str, scope_id: str | None, provider: str | None, campaign_id: str | None, include_disabled: bool, include_inactive: bool, owner_type: Literal["user", "group"] | None, owner_id: str | None, ) -> FileConnectorSettingsDeltaResponse: if campaign_id: _ensure_campaign_file_access(session, principal, campaign_id) profiles = _visible_connector_profiles( session, principal, provider=provider, campaign_id=campaign_id, include_disabled=include_disabled and _can_read_disabled_connector_profiles(principal), include_admin_scopes=_can_read_disabled_connector_profiles(principal), include_effective_policy=False, ) credentials = _visible_connector_credentials( session, principal, provider=provider, include_disabled=include_disabled, ) spaces = _visible_connector_spaces( session, principal, owner_type=owner_type, owner_id=owner_id, include_inactive=include_inactive, ) return FileConnectorSettingsDeltaResponse( profiles=[FileConnectorProfileResponse(**profile.to_response()) for profile in profiles], credentials=[FileConnectorCredentialResponse(**credential.to_response()) for credential in credentials], spaces=[_connector_space_response(space) for space in spaces], policy=FileConnectorPolicyResponse(**connector_policy_response(session, tenant_id=principal.tenant_id, scope_type=scope_type, scope_id=scope_id)), changed_sections=["profiles", "credentials", "spaces", "policy"], deleted=[], watermark=_file_connector_settings_watermark(session, tenant_id=principal.tenant_id), has_more=False, full=True, ) @router.get("/connectors/settings/delta", response_model=FileConnectorSettingsDeltaResponse) def connector_settings_delta( scope_type: str = Query(default="tenant"), scope_id: str | None = Query(default=None), provider: str | None = None, campaign_id: str | None = None, include_disabled: bool = False, include_inactive: bool = False, owner_type: Literal["user", "group"] | None = None, owner_id: str | None = None, since: str | None = None, limit: int = Query(default=100, ge=1, le=500), session: Session = Depends(get_session), principal: ApiPrincipal = Depends(require_any_scope("files:file:admin", "system:settings:read", "admin:settings:read")), ): scope_type = scope_type.strip().casefold() _require_connector_policy_read(principal, scope_type) try: if since is None: return _full_file_connector_settings_delta_response( session, principal, scope_type=scope_type, scope_id=scope_id, provider=provider, campaign_id=campaign_id, include_disabled=include_disabled, include_inactive=include_inactive, owner_type=owner_type, owner_id=owner_id, ) entries, has_more = _file_connector_settings_entries(session, tenant_id=principal.tenant_id, since=since, limit=limit) if entries is None: return _full_file_connector_settings_delta_response( session, principal, scope_type=scope_type, scope_id=scope_id, provider=provider, campaign_id=campaign_id, include_disabled=include_disabled, include_inactive=include_inactive, owner_type=owner_type, owner_id=owner_id, ) changed_profile_ids = { entry.resource_id for entry in entries if entry.collection == FILES_CONNECTOR_PROFILES_COLLECTION and entry.resource_type == FILES_CONNECTOR_PROFILE_RESOURCE } changed_credential_ids = { entry.resource_id for entry in entries if entry.collection == FILES_CONNECTOR_CREDENTIALS_COLLECTION and entry.resource_type == FILES_CONNECTOR_CREDENTIAL_RESOURCE } changed_space_ids = { entry.resource_id for entry in entries if entry.collection == FILES_CONNECTOR_SPACES_COLLECTION and entry.resource_type == FILES_CONNECTOR_SPACE_RESOURCE } policy_changed = any(entry.collection == FILES_CONNECTOR_POLICIES_COLLECTION for entry in entries) profiles = _visible_connector_profiles( session, principal, provider=provider, campaign_id=campaign_id, include_disabled=include_disabled and _can_read_disabled_connector_profiles(principal), include_admin_scopes=_can_read_disabled_connector_profiles(principal), include_effective_policy=False, ) visible_profiles = { profile.id: profile for profile in profiles if profile.id in changed_profile_ids or (profile.credential_profile_id and profile.credential_profile_id in changed_credential_ids) } credentials = _visible_connector_credentials( session, principal, provider=provider, include_disabled=include_disabled, ) visible_credentials = { credential.id: credential for credential in credentials if credential.id in changed_credential_ids } spaces = _visible_connector_spaces( session, principal, owner_type=owner_type, owner_id=owner_id, include_inactive=include_inactive, ) visible_spaces = { space.id: space for space in spaces if space.id in changed_space_ids } changed_sections = [] if changed_profile_ids or any(profile.credential_profile_id and profile.credential_profile_id in changed_credential_ids for profile in profiles): changed_sections.append("profiles") if changed_credential_ids: changed_sections.append("credentials") if changed_space_ids: changed_sections.append("spaces") if policy_changed: changed_sections.append("policy") deleted = [ _connector_deleted_item(entry) for entry in entries if ( entry.resource_type == FILES_CONNECTOR_PROFILE_RESOURCE and entry.resource_id not in visible_profiles ) or ( entry.resource_type == FILES_CONNECTOR_CREDENTIAL_RESOURCE and entry.resource_id not in visible_credentials ) or ( entry.resource_type == FILES_CONNECTOR_SPACE_RESOURCE and entry.resource_id not in visible_spaces ) ] return FileConnectorSettingsDeltaResponse( profiles=[FileConnectorProfileResponse(**profile.to_response()) for profile in visible_profiles.values()], credentials=[FileConnectorCredentialResponse(**credential.to_response()) for credential in visible_credentials.values()], spaces=[_connector_space_response(space) for space in visible_spaces.values()], policy=FileConnectorPolicyResponse(**connector_policy_response(session, tenant_id=principal.tenant_id, scope_type=scope_type, scope_id=scope_id)) if policy_changed else None, changed_sections=changed_sections, deleted=deleted, watermark=_file_connector_settings_response_watermark(session, tenant_id=principal.tenant_id, entries=entries, has_more=has_more), has_more=has_more, full=False, ) except (FileStorageError, ValueError, json.JSONDecodeError) as exc: raise _http_error(exc) from exc @router.get("/connectors/providers", response_model=FileConnectorProvidersResponse) def list_connector_providers( principal: ApiPrincipal = Depends(require_any_scope("files:file:read", "files:file:upload", "files:file:admin")), ): del principal return FileConnectorProvidersResponse(providers=[FileConnectorProviderResponse(**item.to_response()) for item in connector_provider_descriptors()]) @router.post("/connectors/discover", response_model=FileConnectorDiscoveryResponse) def discover_connector_endpoint( payload: FileConnectorDiscoveryRequest, principal: ApiPrincipal = Depends(require_any_scope("files:file:admin", "system:settings:write", "admin:settings:write")), ): if payload.provider not in {"webdav", "nextcloud"}: return FileConnectorDiscoveryResponse( provider=payload.provider, endpoint_url=None, base_path=payload.base_path, status="unsupported", message=f"Discovery is not implemented for {payload.provider} connectors yet", ) candidates: list[dict[str, str]] = [] for endpoint_url in _webdav_discovery_candidates(payload): profile = _discovery_profile_from_payload(payload, endpoint_url, principal=principal) try: browse_connector_profile(profile, path=payload.base_path or "") except (ConnectorBrowseError, ConnectorBrowseUnsupported, OSError, ValueError, json.JSONDecodeError) as exc: message = str(exc) if "credentials were rejected" in message.casefold(): if payload.require_valid_credentials: candidates.append({"endpoint_url": endpoint_url, "status": "credentials_rejected", "message": "The endpoint exists, but the credentials were rejected."}) return FileConnectorDiscoveryResponse( provider=payload.provider, endpoint_url=endpoint_url, base_path=payload.base_path, status="credentials_rejected", message="The endpoint was found, but login failed with these credentials.", candidates=candidates, metadata={**payload.metadata, "discovered_by": "webdav-auth-challenge"}, ) candidates.append({"endpoint_url": endpoint_url, "status": "found", "message": "The endpoint exists, but credentials are required or were rejected."}) return FileConnectorDiscoveryResponse( provider=payload.provider, endpoint_url=endpoint_url, base_path=payload.base_path, status="found", message="The endpoint was found. Add working credentials before saving or testing the connection.", candidates=candidates, metadata={**payload.metadata, "discovered_by": "webdav-auth-challenge"}, ) candidates.append({"endpoint_url": endpoint_url, "status": "failed", "message": message}) continue status_value = "usable" if _same_endpoint(endpoint_url, payload.endpoint_url) else "found" message = "The supplied URL is directly usable." if status_value == "usable" else "A usable connector endpoint was discovered." candidates.append({"endpoint_url": endpoint_url, "status": status_value, "message": message}) return FileConnectorDiscoveryResponse( provider=payload.provider, endpoint_url=endpoint_url, base_path=payload.base_path, status=status_value, message=message, candidates=candidates, metadata={**payload.metadata, "discovered_by": "webdav-propfind"}, ) return FileConnectorDiscoveryResponse( provider=payload.provider, endpoint_url=None, base_path=payload.base_path, status="not_found", message="No usable WebDAV endpoint was found for this server URL.", candidates=candidates, ) @router.get("/connectors/credentials", response_model=FileConnectorCredentialsResponse) def list_connector_credentials( provider: str | None = None, include_disabled: bool = False, session: Session = Depends(get_session), principal: ApiPrincipal = Depends(require_any_scope("files:file:admin", "system:settings:read", "admin:settings:read")), ): provider_norm = provider.strip().casefold() if provider else None credentials = list_database_connector_credentials( session, tenant_id=principal.tenant_id, include_disabled=include_disabled and _can_read_disabled_connector_profiles(principal), ) return FileConnectorCredentialsResponse( credentials=[ FileConnectorCredentialResponse(**credential.to_response()) for credential in credentials if provider_norm is None or credential.provider in {None, provider_norm} ] ) @router.get("/connectors/policies/{scope_type}", response_model=FileConnectorPolicyResponse) def read_connector_policy( scope_type: str, scope_id: str | None = Query(default=None), session: Session = Depends(get_session), principal: ApiPrincipal = Depends(require_any_scope("files:file:admin", "system:settings:read", "admin:settings:read")), ): _require_connector_policy_read(principal, scope_type) try: return FileConnectorPolicyResponse(**connector_policy_response(session, tenant_id=principal.tenant_id, scope_type=scope_type, scope_id=scope_id)) except (FileStorageError, ValueError, json.JSONDecodeError) as exc: raise _http_error(exc, not_found=True) from exc @router.put("/connectors/policies/{scope_type}", response_model=FileConnectorPolicyResponse) def write_connector_policy( scope_type: str, payload: FileConnectorPolicyUpdateRequest, scope_id: str | None = Query(default=None), session: Session = Depends(get_session), principal: ApiPrincipal = Depends(require_any_scope("files:file:admin", "system:settings:write", "admin:settings:write")), ): _require_connector_profile_write(principal, scope_type) try: set_connector_policy( session, tenant_id=principal.tenant_id, user_id=principal.user.id, scope_type=scope_type, scope_id=scope_id, policy=payload.policy, ) _record_connector_settings_change( session, collection=FILES_CONNECTOR_POLICIES_COLLECTION, resource_type=FILES_CONNECTOR_POLICY_RESOURCE, resource_id=_file_connector_policy_resource_id(scope_type, scope_id), operation="updated", principal=principal, tenant_id=None if scope_type.strip().casefold() == "system" else principal.tenant_id, payload={"scope_type": scope_type.strip().casefold(), "scope_id": scope_id}, ) session.commit() return FileConnectorPolicyResponse(**connector_policy_response(session, tenant_id=principal.tenant_id, scope_type=scope_type, scope_id=scope_id)) except (FileStorageError, ValueError, json.JSONDecodeError) as exc: session.rollback() raise _http_error(exc) from exc @router.post("/connectors/credentials", response_model=FileConnectorCredentialResponse, status_code=status.HTTP_201_CREATED) def create_connector_credential( payload: FileConnectorCredentialCreateRequest, session: Session = Depends(get_session), principal: ApiPrincipal = Depends(require_any_scope("files:file:admin", "system:settings:write", "admin:settings:write")), ): _require_connector_credential_write(principal, payload.scope_type) credentials = payload.credentials try: _ensure_connector_credential_configuration_allowed( session, principal, credential_id=payload.id, provider=payload.provider, scope_type=payload.scope_type, scope_id=payload.scope_id, operation="configure_credentials", ) _ensure_connector_local_policy_allowed( session, principal, scope_type=payload.scope_type, scope_id=payload.scope_id, policy=payload.policy, ) row = create_connector_credential_row( session, tenant_id=principal.tenant_id, user_id=principal.user.id, credential_id=payload.id, label=payload.label, provider=payload.provider, scope_type=payload.scope_type, scope_id=payload.scope_id, enabled=payload.enabled, credential_mode=payload.credential_mode, username=credentials.username, password=credentials.password, token=credentials.token, password_env=credentials.password_env, token_env=credentials.token_env, secret_ref=credentials.secret_ref, policy=payload.policy, metadata=payload.metadata, ) _record_connector_settings_change( session, collection=FILES_CONNECTOR_CREDENTIALS_COLLECTION, resource_type=FILES_CONNECTOR_CREDENTIAL_RESOURCE, resource_id=row.id, operation="created", principal=principal, tenant_id=row.tenant_id, payload={"scope_type": row.scope_type, "scope_id": row.scope_id, "provider": row.provider}, ) session.commit() session.refresh(row) return _connector_credential_response(row) except ConnectorPolicyDenied as exc: session.rollback() raise _connector_policy_error(exc) from exc except (FileStorageError, ValueError, json.JSONDecodeError) as exc: session.rollback() raise _http_error(exc) from exc @router.get("/connectors/credentials/{credential_id}", response_model=FileConnectorCredentialResponse) def get_connector_credential( credential_id: str, include_disabled: bool = False, session: Session = Depends(get_session), principal: ApiPrincipal = Depends(require_any_scope("files:file:admin", "system:settings:read", "admin:settings:read")), ): try: row = get_connector_credential_row( session, tenant_id=principal.tenant_id, credential_id=credential_id, include_disabled=include_disabled and _can_read_disabled_connector_profiles(principal), ) return _connector_credential_response(row) except FileStorageError as exc: raise _http_error(exc, not_found=True) from exc @router.patch("/connectors/credentials/{credential_id}", response_model=FileConnectorCredentialResponse) def update_connector_credential( credential_id: str, payload: FileConnectorCredentialUpdateRequest, session: Session = Depends(get_session), principal: ApiPrincipal = Depends(require_any_scope("files:file:admin", "system:settings:write", "admin:settings:write")), ): try: row = get_connector_credential_row(session, tenant_id=principal.tenant_id, credential_id=credential_id, include_disabled=True) except FileStorageError as exc: raise _http_error(exc, not_found=True) from exc _require_connector_credential_write(principal, row.scope_type) credentials = payload.credentials try: provider = payload.provider if payload.provider is not None else row.provider _ensure_connector_credential_configuration_allowed( session, principal, credential_id=row.id, provider=provider, scope_type=row.scope_type, scope_id=row.scope_id, operation="configure_credentials", ) if payload.policy is not None: _ensure_connector_local_policy_allowed( session, principal, scope_type=row.scope_type, scope_id=row.scope_id, policy=payload.policy, ) update_connector_credential_row( session, row, user_id=principal.user.id, label=payload.label, provider=payload.provider, enabled=payload.enabled, credential_mode=payload.credential_mode, username=credentials.username if credentials else None, password=credentials.password if credentials else None, token=credentials.token if credentials else None, password_env=credentials.password_env if credentials else None, token_env=credentials.token_env if credentials else None, secret_ref=credentials.secret_ref if credentials else None, clear_password=payload.clear_password, clear_token=payload.clear_token, policy=payload.policy, metadata=payload.metadata, ) _record_connector_settings_change( session, collection=FILES_CONNECTOR_CREDENTIALS_COLLECTION, resource_type=FILES_CONNECTOR_CREDENTIAL_RESOURCE, resource_id=row.id, operation="updated", principal=principal, tenant_id=row.tenant_id, payload={"scope_type": row.scope_type, "scope_id": row.scope_id, "provider": row.provider}, ) session.commit() session.refresh(row) return _connector_credential_response(row) except ConnectorPolicyDenied as exc: session.rollback() raise _connector_policy_error(exc) from exc except (FileStorageError, ValueError, json.JSONDecodeError) as exc: session.rollback() raise _http_error(exc) from exc @router.delete("/connectors/credentials/{credential_id}", response_model=FileConnectorCredentialResponse) def deactivate_connector_credential( credential_id: str, session: Session = Depends(get_session), principal: ApiPrincipal = Depends(require_any_scope("files:file:admin", "system:settings:write", "admin:settings:write")), ): try: row = get_connector_credential_row(session, tenant_id=principal.tenant_id, credential_id=credential_id, include_disabled=True) except FileStorageError as exc: raise _http_error(exc, not_found=True) from exc _require_connector_credential_write(principal, row.scope_type) try: deactivate_connector_credential_row(session, row, user_id=principal.user.id) _record_connector_settings_change( session, collection=FILES_CONNECTOR_CREDENTIALS_COLLECTION, resource_type=FILES_CONNECTOR_CREDENTIAL_RESOURCE, resource_id=row.id, operation="deleted", principal=principal, tenant_id=row.tenant_id, payload={"scope_type": row.scope_type, "scope_id": row.scope_id, "provider": row.provider}, ) session.commit() session.refresh(row) return _connector_credential_response(row) except FileStorageError as exc: session.rollback() raise _http_error(exc) from exc @router.get("/connectors/profiles", response_model=FileConnectorProfilesResponse) def list_connector_profiles( provider: str | None = None, campaign_id: str | None = None, include_disabled: bool = False, session: Session = Depends(get_session), principal: ApiPrincipal = Depends(require_any_scope("files:file:read", "files:file:upload", "files:file:download", "files:file:admin", "system:settings:read", "admin:settings:read")), ): try: if campaign_id: _ensure_campaign_file_access(session, principal, campaign_id) profiles = _visible_connector_profiles( session, principal, provider=provider, campaign_id=campaign_id, include_disabled=include_disabled and _can_read_disabled_connector_profiles(principal), include_admin_scopes=_can_read_disabled_connector_profiles(principal), include_effective_policy=False, ) except (OSError, ValueError, json.JSONDecodeError) as exc: raise _http_error(exc) from exc return FileConnectorProfilesResponse(profiles=[FileConnectorProfileResponse(**profile.to_response()) for profile in profiles]) @router.post("/connectors/profiles", response_model=FileConnectorProfileResponse, status_code=status.HTTP_201_CREATED) def create_connector_profile( payload: FileConnectorProfileCreateRequest, session: Session = Depends(get_session), principal: ApiPrincipal = Depends(require_any_scope("files:file:admin", "system:settings:write", "admin:settings:write")), ): _require_connector_profile_write(principal, payload.scope_type) credentials = payload.credentials try: credential_row = _credential_row_for_profile( session, principal, credential_profile_id=payload.credential_profile_id, provider=payload.provider, ) _ensure_connector_configuration_allowed( session, principal, connector_id=payload.id, credential_id=payload.credential_profile_id, provider=payload.provider, endpoint_url=payload.endpoint_url, base_path=payload.base_path, scope_type=payload.scope_type, scope_id=payload.scope_id, operation="configure", ) _ensure_connector_local_policy_allowed( session, principal, scope_type=payload.scope_type, scope_id=payload.scope_id, policy=payload.policy, ) row = create_connector_profile_row( session, tenant_id=principal.tenant_id, user_id=principal.user.id, profile_id=payload.id, label=payload.label, provider=payload.provider, scope_type=payload.scope_type, scope_id=payload.scope_id, endpoint_url=payload.endpoint_url, base_path=payload.base_path, enabled=payload.enabled, credential_profile_id=payload.credential_profile_id, credential_mode=payload.credential_mode, username=credentials.username, password=credentials.password, token=credentials.token, password_env=credentials.password_env, token_env=credentials.token_env, secret_ref=credentials.secret_ref, capabilities=payload.capabilities, policy=payload.policy, metadata=payload.metadata, ) _record_connector_settings_change( session, collection=FILES_CONNECTOR_PROFILES_COLLECTION, resource_type=FILES_CONNECTOR_PROFILE_RESOURCE, resource_id=row.id, operation="created", principal=principal, tenant_id=row.tenant_id, payload={"scope_type": row.scope_type, "scope_id": row.scope_id, "provider": row.provider}, ) session.commit() session.refresh(row) return FileConnectorProfileResponse(**connector_profile_from_row(row, credential_row=credential_row).to_response()) except ConnectorPolicyDenied as exc: session.rollback() raise _connector_policy_error(exc) from exc except (FileStorageError, ValueError, json.JSONDecodeError) as exc: session.rollback() raise _http_error(exc) from exc @router.post("/connectors/profiles/{profile_id}/import", response_model=FileUploadResponse) def import_connector_file( profile_id: str, payload: FileConnectorImportRequest, session: Session = Depends(get_session), principal: ApiPrincipal = Depends(require_scope("files:file:upload")), ): try: if payload.campaign_id: _ensure_campaign_file_access(session, principal, payload.campaign_id) profile = _visible_connector_profile(session, principal, profile_id, campaign_id=payload.campaign_id) _source_path, downloaded, metadata = _download_connector_payload(profile, payload, operation="import") target_owner = payload.owner_id or principal.user.id stored = create_file_asset( session, tenant_id=principal.tenant_id, owner_type=payload.owner_type, owner_id=target_owner, user_id=principal.user.id, filename=downloaded.filename, data=downloaded.data, folder=payload.target_folder, display_path=payload.target_path, content_type=downloaded.content_type, metadata=metadata, campaign_id=payload.campaign_id, conflict_strategy=payload.conflict_strategy, is_admin=_is_admin(principal), ) _audit_connector_imports(session, principal, [stored.asset]) session.commit() except ConnectorPolicyDenied as exc: session.rollback() raise _connector_policy_error(exc) from exc except ConnectorImportUnsupported as exc: session.rollback() raise HTTPException(status_code=status.HTTP_501_NOT_IMPLEMENTED, detail=str(exc)) from exc except (ConnectorImportError, FileStorageError, UnsafeFilePathError, ValueError, json.JSONDecodeError) as exc: session.rollback() raise _http_error(exc) from exc return FileUploadResponse(files=[_asset_response(session, stored.asset, include_shares=True)]) @router.post("/connectors/profiles/{profile_id}/sync", response_model=FileConnectorSyncResponse) def sync_connector_file( profile_id: str, payload: FileConnectorImportRequest, session: Session = Depends(get_session), principal: ApiPrincipal = Depends(require_scope("files:file:upload")), ): try: if payload.campaign_id: _ensure_campaign_file_access(session, principal, payload.campaign_id) profile = _visible_connector_profile(session, principal, profile_id, campaign_id=payload.campaign_id) _source_path, downloaded, metadata = _download_connector_payload(profile, payload, operation="sync") target_owner = payload.owner_id or principal.user.id stored, sync_action, previous_version_id = sync_file_asset_from_source( session, tenant_id=principal.tenant_id, owner_type=payload.owner_type, owner_id=target_owner, user_id=principal.user.id, filename=downloaded.filename, data=downloaded.data, folder=payload.target_folder, display_path=payload.target_path, content_type=downloaded.content_type, metadata=metadata, campaign_id=payload.campaign_id, conflict_strategy=payload.conflict_strategy, is_admin=_is_admin(principal), ) _audit_connector_sync(session, principal, stored.asset, sync_action=sync_action, previous_version_id=previous_version_id) session.commit() except ConnectorPolicyDenied as exc: session.rollback() raise _connector_policy_error(exc) from exc except ConnectorImportUnsupported as exc: session.rollback() raise HTTPException(status_code=status.HTTP_501_NOT_IMPLEMENTED, detail=str(exc)) from exc except (ConnectorImportError, FileStorageError, UnsafeFilePathError, ValueError, json.JSONDecodeError) as exc: session.rollback() raise _http_error(exc) from exc return FileConnectorSyncResponse( file=_asset_response(session, stored.asset, include_shares=True), action=sync_action, previous_version_id=previous_version_id, current_version_id=stored.version.id, ) @router.get("/connectors/profiles/{profile_id}/browse", response_model=FileConnectorBrowseResponse) def browse_connector_profile_items( profile_id: str, path: str | None = None, library_id: str | None = None, campaign_id: str | None = None, session: Session = Depends(get_session), principal: ApiPrincipal = Depends(require_any_scope("files:file:read", "files:file:upload", "files:file:download", "files:file:admin", "system:settings:read", "admin:settings:read")), ): try: profile = _visible_connector_profile(session, principal, profile_id, campaign_id=campaign_id) browse_path = normalize_connector_browse_path(path) decision = connector_policy_decision( ConnectorAccessRequest( connector_id=profile.id, credential_id=profile.credential_profile_id, provider=profile.provider, external_path=browse_path, external_url=profile.endpoint_url, operation="browse", ), profile.policy_sources, ) if not decision.allowed: raise ConnectorPolicyDenied(decision) items = browse_connector_profile(profile, path=browse_path, library_id=library_id) except ConnectorPolicyDenied as exc: raise _connector_policy_error(exc) from exc except ConnectorBrowseUnsupported as exc: raise HTTPException(status_code=status.HTTP_501_NOT_IMPLEMENTED, detail=str(exc)) from exc except (ConnectorBrowseError, OSError, ValueError, json.JSONDecodeError) as exc: raise _http_error(exc) from exc return FileConnectorBrowseResponse( profile_id=profile.id, provider=profile.provider, path=browse_path, library_id=library_id, decision=decision.to_dict(), items=[FileConnectorBrowseItem(**item.to_response()) for item in items], ) @router.get("/connectors/profiles/{profile_id}", response_model=FileConnectorProfileResponse) def get_connector_profile( profile_id: str, campaign_id: str | None = None, include_disabled: bool = False, session: Session = Depends(get_session), principal: ApiPrincipal = Depends(require_any_scope("files:file:read", "files:file:upload", "files:file:download", "files:file:admin", "system:settings:read", "admin:settings:read")), ): try: profile = _visible_connector_profile( session, principal, profile_id, campaign_id=campaign_id, include_disabled=include_disabled and _can_read_disabled_connector_profiles(principal), include_effective_policy=False, ) except (OSError, ValueError, json.JSONDecodeError) as exc: raise _http_error(exc) from exc return FileConnectorProfileResponse(**profile.to_response()) @router.patch("/connectors/profiles/{profile_id}", response_model=FileConnectorProfileResponse) def update_connector_profile( profile_id: str, payload: FileConnectorProfileUpdateRequest, session: Session = Depends(get_session), principal: ApiPrincipal = Depends(require_any_scope("files:file:admin", "system:settings:write", "admin:settings:write")), ): try: row = get_connector_profile_row(session, tenant_id=principal.tenant_id, profile_id=profile_id, include_disabled=True) except FileStorageError as exc: raise _http_error(exc, not_found=True) from exc _require_connector_profile_write(principal, row.scope_type) credentials = payload.credentials try: credential_profile_id = payload.credential_profile_id if payload.credential_profile_id is not None else row.credential_profile_id provider = payload.provider if payload.provider is not None else row.provider credential_row = _credential_row_for_profile( session, principal, credential_profile_id=credential_profile_id, provider=provider, include_disabled=True, ) _ensure_connector_configuration_allowed( session, principal, connector_id=row.id, credential_id=credential_profile_id, provider=provider, endpoint_url=payload.endpoint_url if payload.endpoint_url is not None else row.endpoint_url, base_path=payload.base_path if payload.base_path is not None else row.base_path, scope_type=row.scope_type, scope_id=row.scope_id, operation="configure", ) if payload.policy is not None: _ensure_connector_local_policy_allowed( session, principal, scope_type=row.scope_type, scope_id=row.scope_id, policy=payload.policy, ) update_connector_profile_row( session, row, user_id=principal.user.id, label=payload.label, provider=payload.provider, endpoint_url=payload.endpoint_url, base_path=payload.base_path, enabled=payload.enabled, credential_profile_id=payload.credential_profile_id, credential_mode=payload.credential_mode, username=credentials.username if credentials else None, password=credentials.password if credentials else None, token=credentials.token if credentials else None, password_env=credentials.password_env if credentials else None, token_env=credentials.token_env if credentials else None, secret_ref=credentials.secret_ref if credentials else None, clear_password=payload.clear_password, clear_token=payload.clear_token, capabilities=payload.capabilities, policy=payload.policy, metadata=payload.metadata, ) _record_connector_settings_change( session, collection=FILES_CONNECTOR_PROFILES_COLLECTION, resource_type=FILES_CONNECTOR_PROFILE_RESOURCE, resource_id=row.id, operation="updated", principal=principal, tenant_id=row.tenant_id, payload={"scope_type": row.scope_type, "scope_id": row.scope_id, "provider": row.provider}, ) session.commit() session.refresh(row) return FileConnectorProfileResponse(**connector_profile_from_row(row, credential_row=credential_row).to_response()) except ConnectorPolicyDenied as exc: session.rollback() raise _connector_policy_error(exc) from exc except (FileStorageError, ValueError, json.JSONDecodeError) as exc: session.rollback() raise _http_error(exc) from exc @router.delete("/connectors/profiles/{profile_id}", response_model=FileConnectorProfileResponse) def deactivate_connector_profile( profile_id: str, session: Session = Depends(get_session), principal: ApiPrincipal = Depends(require_any_scope("files:file:admin", "system:settings:write", "admin:settings:write")), ): try: row = get_connector_profile_row(session, tenant_id=principal.tenant_id, profile_id=profile_id, include_disabled=True) except FileStorageError as exc: raise _http_error(exc, not_found=True) from exc _require_connector_profile_write(principal, row.scope_type) try: deactivate_connector_profile_row(session, row, user_id=principal.user.id) _record_connector_settings_change( session, collection=FILES_CONNECTOR_PROFILES_COLLECTION, resource_type=FILES_CONNECTOR_PROFILE_RESOURCE, resource_id=row.id, operation="deleted", principal=principal, tenant_id=row.tenant_id, payload={"scope_type": row.scope_type, "scope_id": row.scope_id, "provider": row.provider}, ) session.commit() session.refresh(row) return FileConnectorProfileResponse(**connector_profile_from_row(row).to_response()) except FileStorageError as exc: session.rollback() raise _http_error(exc) from exc @router.get("/{file_id}", response_model=FileAssetResponse) def get_file( file_id: str, session: Session = Depends(get_session), principal: ApiPrincipal = Depends(require_scope("files:file:read")), ): try: asset = get_asset_for_user(session, tenant_id=principal.tenant_id, user_id=principal.user.id, asset_id=file_id, is_admin=_is_admin(principal)) return _asset_response(session, asset, include_shares=True) except FileStorageError as exc: raise _http_error(exc, not_found=True) from exc @router.get("/{file_id}/download") def download_file( file_id: str, session: Session = Depends(get_session), principal: ApiPrincipal = Depends(require_scope("files:file:download")), ): try: asset = get_asset_for_user(session, tenant_id=principal.tenant_id, user_id=principal.user.id, asset_id=file_id, is_admin=_is_admin(principal)) data, version, blob = read_asset_bytes(session, asset) _audit_connector_event( session, principal, action="files.connector.accessed", asset=asset, version=version, blob=blob, operation="download", commit=True, ) except FileStorageError as exc: raise _http_error(exc, not_found=True) from exc headers = {"Content-Disposition": _attachment_disposition(asset.filename)} return StreamingResponse(BytesIO(data), media_type=blob.content_type or "application/octet-stream", headers=headers) @router.delete("/{file_id}", response_model=BulkDeleteResponse) def delete_file( file_id: str, session: Session = Depends(get_session), principal: ApiPrincipal = Depends(require_scope("files:file:delete")), ): try: asset = get_asset_for_user(session, tenant_id=principal.tenant_id, user_id=principal.user.id, asset_id=file_id, require_write=True, is_admin=_is_admin(principal)) count = soft_delete_assets(session, [asset]) session.commit() return BulkDeleteResponse(deleted_count=count) except FileStorageError as exc: session.rollback() raise _http_error(exc, not_found=True) from exc @router.post("/bulk-delete", response_model=BulkDeleteResponse) def bulk_delete_files( payload: BulkDeleteRequest, session: Session = Depends(get_session), principal: ApiPrincipal = Depends(require_scope("files:file:delete")), ): try: assets = [ get_asset_for_user(session, tenant_id=principal.tenant_id, user_id=principal.user.id, asset_id=file_id, require_write=True, is_admin=_is_admin(principal)) for file_id in payload.file_ids ] count = soft_delete_assets(session, assets) session.commit() return BulkDeleteResponse(deleted_count=count) except FileStorageError as exc: session.rollback() raise _http_error(exc) from exc @router.post("/{file_id}/shares", response_model=FileShareResponse) def create_share( file_id: str, payload: FileShareRequest, session: Session = Depends(get_session), principal: ApiPrincipal = Depends(require_scope("files:file:share")), ): try: asset = get_asset_for_user(session, tenant_id=principal.tenant_id, user_id=principal.user.id, asset_id=file_id, require_write=True, is_admin=_is_admin(principal)) share = share_file( session, tenant_id=principal.tenant_id, asset=asset, target_type=payload.target_type, target_id=payload.target_id, permission=payload.permission, user_id=principal.user.id, ) session.commit() return FileShareResponse( id=share.id, target_type=share.target_type, target_id=share.target_id, permission=share.permission, created_at=share.created_at.isoformat(), revoked_at=share.revoked_at.isoformat() if share.revoked_at else None, ) except FileStorageError as exc: session.rollback() raise _http_error(exc) from exc @router.post("/bulk-shares", response_model=BulkFileShareResponse) def create_bulk_shares( payload: BulkFileShareRequest, session: Session = Depends(get_session), principal: ApiPrincipal = Depends(require_scope("files:file:share")), ): try: file_ids = list(dict.fromkeys(payload.file_ids)) assets = [ get_asset_for_user(session, tenant_id=principal.tenant_id, user_id=principal.user.id, asset_id=file_id, require_write=True, is_admin=_is_admin(principal)) for file_id in file_ids ] shares = share_files( session, tenant_id=principal.tenant_id, assets=assets, target_type=payload.target_type, target_id=payload.target_id, permission=payload.permission, user_id=principal.user.id, ) session.commit() return BulkFileShareResponse( shared_count=len(shares), shares=[ FileShareResponse( id=share.id, target_type=share.target_type, target_id=share.target_id, permission=share.permission, created_at=share.created_at.isoformat(), revoked_at=share.revoked_at.isoformat() if share.revoked_at else None, ) for share in shares ], ) except FileStorageError as exc: session.rollback() raise _http_error(exc) from exc @router.post("/bulk-rename", response_model=RenameResponse) def bulk_rename( payload: RenameRequest, session: Session = Depends(get_session), principal: ApiPrincipal = Depends(require_scope("files:file:organize")), ): try: plan = rename_selection( session, tenant_id=principal.tenant_id, user_id=principal.user.id, file_ids=payload.file_ids, folder_paths=payload.folder_paths, owner_type=payload.owner_type, owner_id=payload.owner_id, mode=payload.mode, new_name=payload.new_name, find=payload.find, replacement=payload.replacement, prefix=payload.prefix, suffix=payload.suffix, recursive=payload.recursive, dry_run=payload.dry_run, is_admin=_is_admin(principal), ) if not payload.dry_run: session.commit() return RenameResponse( dry_run=payload.dry_run, items=[ RenamePreviewItem( kind=item.kind, id=item.id, file_id=item.id if item.kind == "file" else None, folder_path=item.old_path if item.kind == "folder" else None, old_path=item.old_path, new_path=item.new_path, ) for item in plan ], ) except (FileStorageError, UnsafeFilePathError, ValueError) as exc: session.rollback() raise _http_error(exc) from exc @router.post("/transfer", response_model=TransferResponse) def transfer_files( payload: TransferRequest, session: Session = Depends(get_session), principal: ApiPrincipal = Depends(require_scope("files:file:organize")), ): try: files, folders = transfer_selection( session, tenant_id=principal.tenant_id, user_id=principal.user.id, operation=payload.operation, file_ids=payload.file_ids, folder_paths=payload.folder_paths, source_owner_type=payload.source_owner_type, source_owner_id=payload.source_owner_id, target_owner_type=payload.target_owner_type, target_owner_id=payload.target_owner_id, target_folder=payload.target_folder, conflict_strategy=payload.conflict_strategy, conflict_resolutions=_conflict_resolutions(payload.conflict_resolutions), is_admin=_is_admin(principal), ) session.commit() return TransferResponse(operation=payload.operation, files=files, folders=folders) except (FileStorageError, UnsafeFilePathError, ValueError) as exc: session.rollback() raise _http_error(exc) from exc @router.post("/archive.zip") def download_archive( payload: ArchiveRequest, session: Session = Depends(get_session), principal: ApiPrincipal = Depends(require_scope("files:file:download")), ): try: assets = [ get_asset_for_user(session, tenant_id=principal.tenant_id, user_id=principal.user.id, asset_id=file_id, is_admin=_is_admin(principal)) for file_id in payload.file_ids ] tmp = tempfile.NamedTemporaryFile(prefix="multimailer-files-", suffix=".zip", delete=False) tmp_path = tmp.name tmp.close() try: create_zip_file(session, assets, tmp_path) except Exception: _cleanup_temp_file(tmp_path) raise _audit_connector_access(session, principal, assets, operation="archive") except FileStorageError as exc: raise _http_error(exc) from exc filename = filename_from_path(normalize_logical_path(payload.filename, fallback_filename="files.zip")) headers = {"Content-Disposition": _attachment_disposition(filename)} return FileResponse(tmp_path, media_type="application/zip", headers=headers, background=BackgroundTask(_cleanup_temp_file, tmp_path)) @router.post("/resolve-patterns", response_model=PatternResolveResponse) def resolve_file_patterns( payload: PatternResolveRequest, session: Session = Depends(get_session), principal: ApiPrincipal = Depends(require_scope("files:file:read")), ): _ensure_list_owner_access(session, principal, payload.owner_type, payload.owner_id) _ensure_campaign_file_access(session, principal, payload.campaign_id) try: assets = list_assets_for_user( session, tenant_id=principal.tenant_id, user_id=principal.user.id, owner_type=payload.owner_type, owner_id=payload.owner_id, campaign_id=payload.campaign_id, path_prefix=payload.path_prefix, is_admin=_is_admin(principal), ) resolved, unmatched = resolve_patterns(assets, payload.patterns, base_path=payload.path_prefix, case_sensitive=payload.case_sensitive) return PatternResolveResponse( patterns=[PatternMatchResponse(pattern=item.pattern, matches=[_asset_response(session, asset) for asset in item.matches]) for item in resolved], unmatched=[_asset_response(session, asset) for asset in unmatched] if payload.include_unmatched else [], ) except (FileStorageError, UnsafeFilePathError, ValueError) as exc: raise _http_error(exc) from exc