"""Attachment service — unified through DMS. All files are stored in the DMS (files table). entity_attachments just references the DMS file with entity_type/entity_id. This unifies: one upload path, one download path, one permission model, one deduplication (content_hash), one storage backend. """ from __future__ import annotations import hashlib import os import uuid from datetime import UTC, datetime from typing import Any from sqlalchemy import select from sqlalchemy.ext.asyncio import AsyncSession from app.core.audit import log_audit from app.core.storage import get_storage_backend from app.core.visibility import apply_visibility_filter, check_single_entity_access from app.models.entity_attachment import EntityAttachment # DMS File model accessed via contract to avoid Core→Plugin dependency (P1-9 fix) def _get_dms_file_model(): """Get DMS File model via contract registry, or None if DMS plugin inactive.""" from app.plugins.builtins.contracts import get_contract dms_contract = get_contract("dms") if dms_contract is not None: return dms_contract.dms_file return None # File size limit: 50MB MAX_FILE_SIZE = 50 * 1024 * 1024 def _generate_unique_filename(original_filename: str) -> str: """Generate a unique filename using UUID + original extension.""" ext = os.path.splitext(original_filename)[1] return f"{uuid.uuid4().hex}{ext}" def _entity_attachment_to_dict(ea: EntityAttachment, dms_file: Any = None) -> dict[str, Any]: """Serialize an EntityAttachment + DMS File to dict.""" return { "id": str(ea.id), "entity_type": ea.entity_type, "entity_id": str(ea.entity_id), "dms_file_id": str(ea.dms_file_id), "category": ea.category, "display_name": ea.display_name, "filename": dms_file.name if dms_file else (ea.display_name or "unknown"), "mime_type": dms_file.mime_type if dms_file else "application/octet-stream", "file_size": dms_file.size_bytes if dms_file else 0, "uploaded_by": str(ea.created_by) if ea.created_by else None, "owner_id": str(ea.owner_id) if ea.owner_id else None, "created_at": ea.created_at.isoformat() if ea.created_at else None, "updated_at": ea.updated_at.isoformat() if ea.updated_at else None, } async def save_attachment( db: AsyncSession, tenant_id: uuid.UUID, user_id: uuid.UUID, entity_type: str, entity_id: uuid.UUID, filename: str, file: Any, # UploadFile or async iterator of chunks mime_type: str, is_system_admin: bool = False, ) -> dict[str, Any]: """Save a file to DMS and create an entity_attachments reference. Streams the file in chunks to avoid loading entire file into RAM. """ # Generate unique filename and storage path unique_filename = _generate_unique_filename(filename) storage_path = f"attachments/{entity_type}/{entity_id}/{unique_filename}" # Stream file to storage — compute hash and size during streaming sha256 = hashlib.sha256() file_size = 0 chunk_size = 1024 * 1024 # 1MB chunks async def chunk_stream(): nonlocal file_size if hasattr(file, 'read'): # UploadFile object while True: chunk = await file.read(chunk_size) if not chunk: break file_size += len(chunk) sha256.update(chunk) yield chunk else: # Already bytes (backward compat) nonlocal_bytes = file if isinstance(file, bytes) else b''.join([c async for c in file]) file_size = len(nonlocal_bytes) sha256.update(nonlocal_bytes) yield nonlocal_bytes # Check file size limit during streaming # (we check after streaming — for true streaming we'd need a wrapper) # For now, stream and check size after storage = get_storage_backend() await storage.save_stream(storage_path, chunk_stream()) if file_size > MAX_FILE_SIZE: await storage.delete(storage_path) raise ValueError(f"File too large: {file_size} bytes (max {MAX_FILE_SIZE})") # Check for blocked file types import os as _os _blocked = {".exe", ".bat", ".cmd", ".sh", ".jar", ".com", ".scr", ".msi", ".dll", ".vbs", ".ps1", ".app", ".bin", ".reg", ".inf"} _ext = _os.path.splitext(filename)[1].lower() if _ext in _blocked: await storage.delete(storage_path) raise ValueError(f"File type not allowed: {_ext}") # Check access on parent entity if not is_system_admin: has_access = await check_single_entity_access( db, entity_type, entity_id, user_id, tenant_id, "write", is_system_admin ) if not has_access: await storage.delete(storage_path) raise PermissionError(f"No write access to {entity_type} {entity_id}") content_hash = sha256.hexdigest() # Get DMS File model via contract (avoid Core→Plugin dependency) dms_file = _get_dms_file_model() if dms_file is None: await storage.delete(storage_path) raise RuntimeError("DMS plugin not available — cannot save attachment") # Check for existing DMS file with same hash in same tenant (deduplication) existing_file = await db.execute( select(dms_file).where( dms_file.tenant_id == tenant_id, dms_file.content_hash == content_hash, dms_file.deleted_at.is_(None), ).limit(1) ) existing_dms_file = existing_file.scalar_one_or_none() if existing_dms_file: # Deduplicate: reuse existing DMS file, just create new reference dms_file = existing_dms_file await storage.delete(storage_path) # Remove the duplicate we just saved else: # File already streamed to storage — create DMS File record dms_file = dms_file( tenant_id=tenant_id, name=filename, folder_id=None, # Attachments don't go in DMS folders uploaded_by=user_id, mime_type=mime_type, size_bytes=file_size, storage_path=storage_path, content_hash=content_hash, owner_id=user_id, ) db.add(dms_file) await db.flush() await db.refresh(dms_file) # Create entity_attachments reference entity_attachment = EntityAttachment( tenant_id=tenant_id, entity_type=entity_type, entity_id=entity_id, dms_file_id=dms_file.id, category=None, display_name=filename, owner_id=user_id, created_by=user_id, ) db.add(entity_attachment) await db.flush() await db.refresh(entity_attachment) await log_audit( db, tenant_id, user_id, "upload", "attachment", entity_attachment.id, changes={"filename": filename, "entity_type": entity_type, "entity_id": str(entity_id), "dms_file_id": str(dms_file.id)}, ) return _entity_attachment_to_dict(entity_attachment, dms_file) async def list_attachments( db: AsyncSession, tenant_id: uuid.UUID, entity_type: str, entity_id: uuid.UUID, user_id: uuid.UUID | None = None, is_system_admin: bool = False, ) -> dict[str, Any]: """List attachments for a specific entity (via DMS files).""" dms_file = _get_dms_file_model() if dms_file is None: return {"items": [], "total": 0} q = ( select(EntityAttachment, dms_file) .join(dms_file, EntityAttachment.dms_file_id == dms_file.id) .where( EntityAttachment.tenant_id == tenant_id, EntityAttachment.entity_type == entity_type, EntityAttachment.entity_id == entity_id, EntityAttachment.deleted_at.is_(None), dms_file.deleted_at.is_(None), ) .order_by(EntityAttachment.created_at.desc()) ) if user_id and not is_system_admin: q = await apply_visibility_filter( db, q, "entity_attachment", EntityAttachment, user_id, tenant_id, is_system_admin ) result = await db.execute(q) rows = result.all() return { "items": [_entity_attachment_to_dict(ea, df) for ea, df in rows], "total": len(rows), } async def get_attachment( db: AsyncSession, tenant_id: uuid.UUID, attachment_id: uuid.UUID, user_id: uuid.UUID | None = None, is_system_admin: bool = False, ) -> dict[str, Any] | None: """Get a single attachment by ID (with DMS file info).""" dms_file = _get_dms_file_model() if dms_file is None: return None q = ( select(EntityAttachment, dms_file) .join(dms_file, EntityAttachment.dms_file_id == dms_file.id) .where( EntityAttachment.id == attachment_id, EntityAttachment.tenant_id == tenant_id, EntityAttachment.deleted_at.is_(None), ) ) result = await db.execute(q) row = result.first() if row is None: return None ea, dms_file = row if user_id and not is_system_admin: has_access = await check_single_entity_access( db, "entity_attachment", ea.id, user_id, tenant_id, "read", is_system_admin ) if not has_access: raise PermissionError("No access") return _entity_attachment_to_dict(ea, dms_file) async def get_attachment_download_path( db: AsyncSession, tenant_id: uuid.UUID, attachment_id: uuid.UUID, user_id: uuid.UUID | None = None, is_system_admin: bool = False, ) -> str | None: """Get the storage path for downloading an attachment's DMS file.""" dms_file = _get_dms_file_model() if dms_file is None: return None q = ( select(EntityAttachment, dms_file) .join(dms_file, EntityAttachment.dms_file_id == dms_file.id) .where( EntityAttachment.id == attachment_id, EntityAttachment.tenant_id == tenant_id, EntityAttachment.deleted_at.is_(None), dms_file.deleted_at.is_(None), ) ) result = await db.execute(q) row = result.first() if row is None: return None ea, dms_file = row if user_id and not is_system_admin: has_access = await check_single_entity_access( db, "entity_attachment", ea.id, user_id, tenant_id, "read", is_system_admin ) if not has_access: raise PermissionError("No access") return dms_file.storage_path async def delete_attachment( db: AsyncSession, tenant_id: uuid.UUID, user_id: uuid.UUID, attachment_id: uuid.UUID, is_system_admin: bool = False, ) -> bool: """Soft-delete an entity_attachments reference. The DMS file is NOT deleted because other entities may reference it. DMS file cleanup happens via DMS's own deletion workflow. """ q = select(EntityAttachment).where( EntityAttachment.id == attachment_id, EntityAttachment.tenant_id == tenant_id, EntityAttachment.deleted_at.is_(None), ) result = await db.execute(q) ea = result.scalar_one_or_none() if ea is None: return False if not is_system_admin: has_access = await check_single_entity_access( db, "entity_attachment", ea.id, user_id, tenant_id, "admin", is_system_admin ) if not has_access: raise PermissionError("No access") ea.deleted_at = datetime.now(UTC) await db.flush() await log_audit( db, tenant_id, user_id, "delete", "attachment", attachment_id, changes={"display_name": ea.display_name, "dms_file_id": str(ea.dms_file_id)}, ) return True