infrasynth-backend-kit/infrasynth/files/processing.py
jcv-dev 551b42eab5 feat: production-hardening pass across the kit
Close the gaps between the documented contract (API-STANDARD, TENANCY,
ENTITLEMENTS) and the implementation, and remove committed build artifacts.

Security:
- verify + process inbound webhooks (HMAC/handler verify, size limit,
  timestamp tolerance, idempotency via InboundEvent.external_id)
- real 2FA login flow (pre-auth challenge; tokens only after verify/recovery)
- wire HybridPermission into security/audit views; add API-key rotate and
  users/<id>/permissions|roles endpoints
- tenant-scoped throttling on by default; webhook replay protection
- verify MercadoPago webhook signatures
- login brute-force guard, configurable password policy, real ALTCHA PoW

Correctness:
- apply verified billing webhooks idempotently (subscription/entitlement/
  invoice/PaymentTransaction); scheduled payment lifecycle jobs
- capture audit update diffs automatically; add audit retention purge
- working notification retries, per-channel rate limits, log retention
- pluggable virus scanner, upload-size limit, pipeline toggle
- feature rollout %/environment targeting; settings-driven registrations
- workflow guards (instance cap, route depth, self-assignment, clone on re-entry)
- wire every previously-dead INFRASYNTH_* setting; drop truly dead ones

Delivery:
- README + CHANGELOG; CI format check + coverage gate
- keep test media out of the tree; untrack .coverage, __pycache__,
  egg-info, docs/ and invoice artifacts
2026-09-24 10:41:21 -05:00

156 lines
6 KiB
Python

import io
import logging
from celery import shared_task
from django.utils import timezone
logger = logging.getLogger(__name__)
class PipelineExecutor:
"""Executes a processing pipeline over a stored file, step by step."""
def execute(self, execution):
from .models import PipelineExecution
execution.status = PipelineExecution.Status.RUNNING
execution.started_at = timezone.now()
execution.save(update_fields=["status", "started_at"])
try:
data, mime_type = self._read_source(execution.file)
pipeline = execution.pipeline
steps = pipeline.steps if pipeline else []
for step in steps:
step_type = step.get("type")
params = step.get("params", {})
handler = getattr(self, f"_step_{step_type}", None)
if handler is None:
raise ValueError(f"Unknown pipeline step type: '{step_type}'")
data, mime_type = handler(data, mime_type, params)
output = self._store_output(execution, data, mime_type)
execution.status = PipelineExecution.Status.COMPLETED
execution.output_file = output
execution.completed_at = timezone.now()
execution.error = ""
execution.save(update_fields=["status", "output_file", "completed_at", "error"])
except Exception as exc: # noqa: BLE001
logger.exception("Pipeline execution %s failed", execution.pk)
execution.status = PipelineExecution.Status.FAILED
execution.completed_at = timezone.now()
execution.error = str(exc)
execution.save(update_fields=["status", "completed_at", "error"])
from .signals import file_processed
file_processed.send(
sender=PipelineExecution,
file_id=execution.file_id,
pipeline_name=execution.pipeline.slug if execution.pipeline else "",
output_file_id=execution.output_file_id,
status=execution.status,
)
return execution
def _read_source(self, stored_file):
from .storage import get_storage_backend
backend = get_storage_backend(stored_file.storage_backend)
fh = backend.open(stored_file.storage_key, "rb")
return fh.read(), stored_file.mime_type
def _step_resize(self, data: bytes, mime_type: str, params: dict):
from PIL import Image
img = Image.open(io.BytesIO(data))
width = int(params.get("width", 800))
height = params.get("height")
if height:
img.thumbnail((width, int(height)))
else:
img.thumbnail((width, width))
output = io.BytesIO()
img.save(output, format=img.format or "PNG")
return output.getvalue(), mime_type
def _step_optimize(self, data: bytes, mime_type: str, params: dict):
from PIL import Image
quality = int(params.get("quality", 80))
img = Image.open(io.BytesIO(data))
fmt = img.format or "PNG"
if fmt.upper() == "PNG":
img = img.convert("P", palette=Image.Palette.ADAPTIVE, colors=256) # type: ignore[assignment]
output = io.BytesIO()
img.save(output, format="PNG", optimize=True)
else:
output = io.BytesIO()
img.save(output, format=fmt, quality=quality, optimize=True)
return output.getvalue(), mime_type
def _step_watermark(self, data: bytes, mime_type: str, params: dict):
from PIL import Image, ImageDraw, ImageFont
text = params.get("text", "Confidential")
img = Image.open(io.BytesIO(data)).convert("RGBA")
layer = Image.new("RGBA", img.size, (0, 0, 0, 0))
draw = ImageDraw.Draw(layer)
try:
font = ImageFont.load_default(size=48)
except TypeError:
font = ImageFont.load_default()
width, height = img.size
draw.text((width // 4, height // 2), text, fill=(255, 255, 255, 120), font=font)
out = Image.alpha_composite(img, layer)
output = io.BytesIO()
out.save(output, format="PNG")
return output.getvalue(), "image/png"
def _step_scan(self, data: bytes, mime_type: str, params: dict):
from .scanner import get_scanner
result = get_scanner().scan(data)
if not result.clean:
raise ValueError(f"Virus scan rejected the file: {result.threat}")
logger.info("Virus scan passed for %d bytes (scanner=%s)", len(data), result.scanner)
return data, mime_type
def _store_output(self, execution, data: bytes, mime_type: str):
from django.core.files.base import ContentFile
from .services import FileService
original = execution.file
name = f"processed_{execution.pipeline.slug}_{original.original_filename}"
content_file = ContentFile(data, name=name)
content_file.content_type = mime_type # type: ignore[attr-defined]
service = FileService()
output = service.upload(
content_file,
filename=name,
user=original.uploaded_by,
metadata={"source_file_id": original.id, "pipeline": execution.pipeline.slug},
)
return output
@shared_task(name="infrasynth.files.run_pipeline_execution", bind=True, max_retries=3)
def run_pipeline_execution(self, execution_id, tenant_id=None):
"""Celery task wrapper around the pipeline executor."""
from infrasynth.tenancy.context import tenant_context
from infrasynth.tenancy.models import Tenant
from .models import PipelineExecution
try:
execution = PipelineExecution.all_objects.select_related("file", "pipeline").get(pk=execution_id)
except PipelineExecution.DoesNotExist:
logger.warning("Pipeline execution %s not found", execution_id)
return None
tenant = Tenant.objects.filter(pk=tenant_id).first() if tenant_id else execution.tenant
with tenant_context(tenant):
return PipelineExecutor().execute(execution)