Files
leocrm/app/services/attachment_service.py
Agent Zero abbe7a18fc fix(audit): P0-P3 audit fixes — 838 ruff errors → 0, 30 F821 bugs fixed, 118 files changed
- P0: hooks.py 3-tuple fix, trigger_dispatcher Contract, contacts/plugin unregister_actions_by_owner
- P0: 5 test files — check_permission mocks removed, hardcoded DB credential → env var
- P1: attachment_service DmsFile via Contract helper, restore_registry/history_hooks dedup
- P1: mail/plugin restore unregister, mcp_client datetime.now(UTC), saved_views/filters patterns
- P1: ProtectedRoute fail-closed, 13 test assertion fixes (bcrypt, DB-URLs, SECRET_KEYs)
- P2: deprecated notifications → post_system_message (3 files), forgejo Base, report_generator lazy import
- P2: webhooks permissions, deps.py/roles.py plugin perms removed, import_export default
- P2: address/tags/entity_links patterns removed, worker.py Contract-Umgehungen fixed
- P2: 28 frontend TODOs (hardcoded constants, deprecated notification API)
- P3: dead code, duplicates, deprecated imports, private attr, __import__ inline
- P3: 8 frontend TODOs (LucideIcons, inline styles, XSS, i18n)
- ruff: 838 → 0 (612 auto-fix + 246 manual + 27 F821 regression fix)
- F821: 30 → 0 (AutomationDefinition, DmsFile, user_id, Path, Any, String)
- Contract-Umgehungen: 2 neue gefunden (worker.py:169, worker.py:280) und gefixt
2026-08-16 01:17:18 +02:00

338 lines
11 KiB
Python

"""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