Files
leocrm/app/services/attachment_service.py
T
Agent Zero 29d55cb187
Check Cross-Plugin Imports / check (push) Has been cancelled
Phase 6: DMS & Attachments — Streaming, Deduplikation, API-Bereinigung
6.4 Upload streamen:
- attachment_service.save_attachment: Streamt in 1MB Chunks statt await file.read()
- routes/attachments.py: Uebergibt UploadFile direkt statt bytes

6.5 Download streamen:
- DMS preview_file: FileResponse fuer LocalStorage (automatisches Streaming)
- Kein storage.read() mehr fuer LocalStorage

6.6 Tenantlokale Deduplikation:
- DMS Upload: Prueft content_hash vor Erstellung, wiederverwendet existierendes File
- attachment_service: Dedup bereits vorhanden, jetzt mit Streaming kompatibel
- Migration 0098: Partial Unique Index (tenant_id, content_hash) WHERE content_hash IS NOT NULL AND deleted_at IS NULL

6.7 API-Ausgabe bereinigt:
- attachment_service: storage_path und content_hash aus API-Ausgaben entfernt
- DMS routes: content_hash aus 4 API-Endpunkten entfernt

Tests: 54/54 bestanden (17 Workspace + 13 API Token + 24 Command)
2026-08-03 14:21:43 +02:00

316 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, func
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
from app.plugins.builtins.dms.models import File as DmsFile
# 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: DmsFile | None = 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.
"""
import hashlib
from app.core.storage import get_storage_backend
# 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()
# Check for existing DMS file with same hash in same tenant (deduplication)
existing_file = await db.execute(
select(DmsFile).where(
DmsFile.tenant_id == tenant_id,
DmsFile.content_hash == content_hash,
DmsFile.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 = DmsFile(
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)."""
q = (
select(EntityAttachment, DmsFile)
.join(DmsFile, EntityAttachment.dms_file_id == DmsFile.id)
.where(
EntityAttachment.tenant_id == tenant_id,
EntityAttachment.entity_type == entity_type,
EntityAttachment.entity_id == entity_id,
EntityAttachment.deleted_at.is_(None),
DmsFile.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)."""
q = (
select(EntityAttachment, DmsFile)
.join(DmsFile, EntityAttachment.dms_file_id == DmsFile.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."""
q = (
select(EntityAttachment, DmsFile)
.join(DmsFile, EntityAttachment.dms_file_id == DmsFile.id)
.where(
EntityAttachment.id == attachment_id,
EntityAttachment.tenant_id == tenant_id,
EntityAttachment.deleted_at.is_(None),
DmsFile.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