"""Generating print files in the background, one order item at a time. A paid order queues one job per item. The worker claims a job, commits the claim, and renders outside any transaction, so a long render never holds a row lock or a database connection. A claim older than CLAIM_TIMEOUT is assumed to belong to a worker that died and is taken again. The result is an ordinary upload row owned by an identity derived from the order, already marked clean: its only inputs are artwork that passed the malware scan, and the bytes are written here. The operator still decides whether it becomes the final file; generation never approves anything. """ import logging import os import tempfile import time from datetime import timedelta from uuid import NAMESPACE_URL, UUID, uuid4, uuid5 from psycopg.types.json import Jsonb from .core.auth import audit from .core.db import connect from .printfile import Unsupported, render CLAIM_TIMEOUT = timedelta(minutes=15) MAX_ATTEMPTS = 3 def generated_identity(order_id): return uuid5(NAMESPACE_URL, f'dtf-print-file:{order_id}') def queue(c, order_id, items, only_missing=False): """Queue (or re-queue) generation for every item of an order.""" for index in range(items): if only_missing: c.execute('''INSERT INTO dtf_local.print_files(id,order_id,item_index) VALUES(%s,%s,%s) ON CONFLICT(order_id,item_index) DO UPDATE SET status='pending', attempts=0, claimed_at=NULL, finished_at=NULL, detail='{}' WHERE dtf_local.print_files.status IN ('failed','manual')''', (uuid4(), order_id, index)) else: c.execute('''INSERT INTO dtf_local.print_files(id,order_id,item_index) VALUES(%s,%s,%s) ON CONFLICT(order_id,item_index) DO NOTHING''', (uuid4(), order_id, index)) def claim(c): return c.execute('''SELECT p.*, o.snapshot, o.number FROM dtf_local.print_files p JOIN dtf_local.orders o ON o.id=p.order_id WHERE p.status='pending' OR (p.status='rendering' AND p.claimed_at < now()-%s) ORDER BY p.created_at FOR UPDATE OF p SKIP LOCKED LIMIT 1''', (CLAIM_TIMEOUT,)).fetchone() def render_one(storage): with connect() as c: job = claim(c) if not job: return False c.execute('''UPDATE dtf_local.print_files SET status='rendering', claimed_at=now(), attempts=attempts+1 WHERE id=%s''', (job['id'],)) item = job['snapshot']['items'][job['item_index']] uploads = c.execute('''SELECT id,name,object_key,scan_state,purged_at,expires_at, (expires_at<=now()) AS expired FROM dtf_local.uploads WHERE id=ANY(%s)''', ([UUID(u) for u in item['uploads']],)).fetchall() # Every generated file shares its order's artwork retention deadline. expiry = c.execute('SELECT min(created_at)+interval \'30 days\' AS e FROM dtf_local.uploads WHERE id=ANY(%s)', ([UUID(u) for u in item['uploads']],)).fetchone()['e'] by_id = {str(row['id']): row for row in uploads} try: result = produce(storage, job, item, by_id) except Unsupported as reason: finish(job, 'manual', {'reason': str(reason)}) audit('print_file_manual', order=str(job['order_id']), item=job['item_index']) return True except Exception as exc: logging.exception('Print file generation failed') final = job['attempts'] + 1 >= MAX_ATTEMPTS finish(job, 'failed' if final else 'pending', {'error': type(exc).__name__}) return True path, name, size, detail = result try: uid = uuid4() key = f'originals/{uid}' storage.store(key, path, 'application/pdf') finally: os.unlink(path) with connect() as c: c.execute('''INSERT INTO dtf_local.uploads(id,owner,name,size,object_key,multipart_id,complete, expires_at,scan_state,scan_reason,scanned_at) VALUES(%s,%s,%s,%s,%s,'',true,%s,'clean','generated from scanned artwork',now())''', (uid, generated_identity(job['order_id']), name, size, key, expiry)) c.execute('''UPDATE dtf_local.print_files SET status='ready', upload_id=%s, detail=%s, finished_at=now() WHERE id=%s''', (uid, Jsonb(detail), job['id'])) audit('print_file_ready', order=str(job['order_id']), item=job['item_index']) return True def produce(storage, job, item, uploads): """Fetch the item's artwork and render it. Returns (pdf path, name, size, evidence).""" for upload_id in item['uploads']: row = uploads.get(upload_id) if not row or row['scan_state'] != 'clean': raise Unsupported('an original file is missing or not cleared by the malware scan') if row['purged_at'] or row['expired']: raise Unsupported('an original file has passed its retention period') number = job['number'] name = f'pedido-{number}-item-{job["item_index"] + 1}.pdf' with tempfile.TemporaryDirectory(prefix='print-') as scratch: files = {} for index, upload_id in enumerate(item['uploads']): local = os.path.join(scratch, f'source-{index}') storage.fetch(uploads[upload_id]['object_key'], local) files[index] = (local, uploads[upload_id]['name']) handle, output = tempfile.mkstemp(prefix='print-', suffix='.pdf') try: with os.fdopen(handle, 'wb') as out: detail = render(item, files, out, f'Pedido {number} - item {job["item_index"] + 1}') except BaseException: os.unlink(output) raise return output, name, os.path.getsize(output), detail def finish(job, status, detail): with connect() as c: c.execute('''UPDATE dtf_local.print_files SET status=%s, detail=%s, finished_at=CASE WHEN %s='pending' THEN NULL ELSE now() END, claimed_at=NULL WHERE id=%s''', (status, Jsonb(detail), status, job['id'])) def render_loop(storage): while True: try: if render_one(storage): continue except Exception: logging.exception('Print render worker tick failed') time.sleep(2)