refactor: extract turntable branch and exhausted-retry handler from render_order_line_task
_render_turntable: Option B — resolved objects as params (render_invocation, step_path, output_path, template, emit, pl). Session and PipelineLogger stay in the caller so no second DB connection is opened and log steps roll up to the main task. _handle_render_task_exhausted: extracted 68-line mark-as-failed block; retry/raise logic with Celery self.retry stays in render_order_line_task since it needs the bound-task context. Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
This commit is contained in:
@@ -533,3 +533,6 @@ Der Admin-Settings-Endpunkt (`GET /api/admin/settings`) erfordert `global_admin`
|
|||||||
|
|
||||||
### 2026-07-22 | Workflow-Editor | _legacy_dispatch umging cancelled/rejected Pre-Check
|
### 2026-07-22 | Workflow-Editor | _legacy_dispatch umging cancelled/rejected Pre-Check
|
||||||
`_legacy_dispatch` in `dispatch_service.py` rief `render_order_line_task.delay()` direkt auf und umging damit den Pre-Check in `dispatch_order_line_render` (der cancelled/rejected Lines überspringt). Alle Legacy-Dispatch-Pfade (auch im Graph-Fallback) liefen so durch, auch für bereits gecancelte Jobs. **Lösung:** `_legacy_dispatch` ruft jetzt `dispatch_order_line_render.delay()` auf statt `render_order_line_task.delay()` direkt.
|
`_legacy_dispatch` in `dispatch_service.py` rief `render_order_line_task.delay()` direkt auf und umging damit den Pre-Check in `dispatch_order_line_render` (der cancelled/rejected Lines überspringt). Alle Legacy-Dispatch-Pfade (auch im Graph-Fallback) liefen so durch, auch für bereits gecancelte Jobs. **Lösung:** `_legacy_dispatch` ruft jetzt `dispatch_order_line_render.delay()` auf statt `render_order_line_task.delay()` direkt.
|
||||||
|
|
||||||
|
### 2026-07-22 | Refactoring | Turntable-Branch aus render_order_line_task extrahiert
|
||||||
|
`render_order_line_task` war ein 427-Zeilen-Monolith. Der Turntable-Branch (~55 Zeilen) wurde in `_render_turntable(*, render_invocation, step_path, output_path, template, order_line_id, emit, pl)` ausgelagert — resolved objects als Parameter (Option B), damit Session und PipelineLogger im Main Task verbleiben und kein doppeltes DB-Lookup entsteht. Der 68-Zeilen Exception-Handler (Mark-as-failed bei max_retries) wurde in `_handle_render_task_exhausted(order_line_id, exc, tenant_id)` extrahiert; die Retry-Logik (`self.retry`) bleibt im Main Task, da sie den Celery `self`-Context benötigt. Beide Helper stehen in `render_order_line.py` vor den Task-Definitionen.
|
||||||
|
|||||||
@@ -53,6 +53,154 @@ def dispatch_order_line_render(order_line_id: str):
|
|||||||
render_order_line_task.apply_async(args=[order_line_id], queue=target_queue)
|
render_order_line_task.apply_async(args=[order_line_id], queue=target_queue)
|
||||||
|
|
||||||
|
|
||||||
|
def _render_turntable(
|
||||||
|
*,
|
||||||
|
render_invocation,
|
||||||
|
step_path,
|
||||||
|
output_path: str,
|
||||||
|
template,
|
||||||
|
order_line_id: str,
|
||||||
|
emit,
|
||||||
|
pl: PipelineLogger,
|
||||||
|
) -> tuple[bool, dict]:
|
||||||
|
"""Execute the Blender turntable render for one order line.
|
||||||
|
|
||||||
|
Session and PipelineLogger are owned by the caller (render_order_line_task).
|
||||||
|
"""
|
||||||
|
from pathlib import Path as _Path
|
||||||
|
from app.services.render_blender import is_blender_available, render_turntable_to_file
|
||||||
|
from app.services.step_processor import _get_all_settings
|
||||||
|
|
||||||
|
render_width = render_invocation.width
|
||||||
|
render_height = render_invocation.height
|
||||||
|
render_engine = render_invocation.engine
|
||||||
|
render_samples = render_invocation.samples
|
||||||
|
cycles_device_val = render_invocation.cycles_device
|
||||||
|
frame_count = render_invocation.frame_count
|
||||||
|
fps = render_invocation.fps
|
||||||
|
tmpl_info = f" template={template.name}" if template else ""
|
||||||
|
|
||||||
|
emit(order_line_id, f"Starting turntable render: {frame_count} frames @ {fps}fps, {render_width or 1920}x{render_height or 1920}{tmpl_info}")
|
||||||
|
pl.step_start("blender_turntable", {"frame_count": frame_count, "fps": fps})
|
||||||
|
|
||||||
|
if not is_blender_available():
|
||||||
|
raise RuntimeError("Blender not available on this worker")
|
||||||
|
|
||||||
|
_sys = _get_all_settings()
|
||||||
|
try:
|
||||||
|
turntable_kwargs = render_invocation.as_turntable_renderer_kwargs(
|
||||||
|
step_path=step_path,
|
||||||
|
output_path=_Path(output_path),
|
||||||
|
default_width=1920,
|
||||||
|
default_height=1920,
|
||||||
|
default_engine=_sys.get("blender_engine", "cycles"),
|
||||||
|
default_samples=int(
|
||||||
|
_sys.get(
|
||||||
|
f"blender_{render_engine or _sys.get('blender_engine', 'cycles')}_samples",
|
||||||
|
128,
|
||||||
|
)
|
||||||
|
),
|
||||||
|
smooth_angle=int(_sys.get("blender_smooth_angle", 30)),
|
||||||
|
)
|
||||||
|
service_data = render_turntable_to_file(**turntable_kwargs)
|
||||||
|
render_log = {
|
||||||
|
"renderer": "blender",
|
||||||
|
"type": "turntable",
|
||||||
|
"format": "mp4",
|
||||||
|
"engine": render_engine or _sys.get("blender_engine", "cycles"),
|
||||||
|
"engine_used": service_data.get("engine_used", "cycles"),
|
||||||
|
"samples": render_samples,
|
||||||
|
"cycles_device": cycles_device_val,
|
||||||
|
"width": render_width or 1920,
|
||||||
|
"height": render_height or 1920,
|
||||||
|
"frame_count": service_data.get("frame_count", frame_count),
|
||||||
|
"fps": fps,
|
||||||
|
"total_duration_s": service_data.get("total_duration_s"),
|
||||||
|
"stl_duration_s": service_data.get("stl_duration_s"),
|
||||||
|
"render_duration_s": service_data.get("render_duration_s"),
|
||||||
|
"ffmpeg_duration_s": service_data.get("ffmpeg_duration_s"),
|
||||||
|
"stl_size_bytes": service_data.get("stl_size_bytes"),
|
||||||
|
"output_size_bytes": service_data.get("output_size_bytes"),
|
||||||
|
"log_lines": service_data.get("log_lines", []),
|
||||||
|
}
|
||||||
|
if template:
|
||||||
|
render_log["template"] = template.blend_file_path
|
||||||
|
pl.step_done("blender_turntable")
|
||||||
|
return True, render_log
|
||||||
|
except Exception as exc:
|
||||||
|
render_log = {"renderer": "blender", "type": "turntable", "error": str(exc)[:500]}
|
||||||
|
pl.step_error("blender_turntable", str(exc), exc)
|
||||||
|
logger.error("Turntable render failed for %s: %s", order_line_id, exc)
|
||||||
|
return False, render_log
|
||||||
|
|
||||||
|
|
||||||
|
def _handle_render_task_exhausted(
|
||||||
|
order_line_id: str,
|
||||||
|
exc: Exception,
|
||||||
|
tenant_id: str | None,
|
||||||
|
) -> None:
|
||||||
|
"""Mark order line as failed and emit notifications after all retries are exhausted."""
|
||||||
|
from sqlalchemy import create_engine, update as sql_update2, select as sel
|
||||||
|
from sqlalchemy.orm import Session as SyncSession
|
||||||
|
from app.config import settings as app_settings
|
||||||
|
from app.models.order_line import OrderLine as OL2
|
||||||
|
from app.core.tenant_context import set_tenant_context_sync
|
||||||
|
from datetime import datetime as dt2
|
||||||
|
|
||||||
|
sync_url = app_settings.database_url.replace("+asyncpg", "")
|
||||||
|
|
||||||
|
eng2 = create_engine(sync_url)
|
||||||
|
with SyncSession(eng2) as s2:
|
||||||
|
set_tenant_context_sync(s2, tenant_id)
|
||||||
|
s2.execute(
|
||||||
|
sql_update2(OL2).where(OL2.id == order_line_id)
|
||||||
|
.values(
|
||||||
|
render_status="failed",
|
||||||
|
render_completed_at=dt2.utcnow(),
|
||||||
|
render_log={"error": str(exc)[:500]},
|
||||||
|
)
|
||||||
|
)
|
||||||
|
s2.commit()
|
||||||
|
eng2.dispose()
|
||||||
|
|
||||||
|
from app.services.order_status_service import check_order_completion
|
||||||
|
eng3 = create_engine(sync_url)
|
||||||
|
with SyncSession(eng3) as s3:
|
||||||
|
set_tenant_context_sync(s3, tenant_id)
|
||||||
|
row = s3.execute(sel(OL2.order_id).where(OL2.id == order_line_id)).scalar_one_or_none()
|
||||||
|
if row:
|
||||||
|
check_order_completion(str(row))
|
||||||
|
eng3.dispose()
|
||||||
|
|
||||||
|
try:
|
||||||
|
from sqlalchemy import select as sel2
|
||||||
|
from app.models.order import Order as OrderModel2
|
||||||
|
from app.domains.rendering.workflow_runtime_services import emit_order_line_render_notifications
|
||||||
|
eng4 = create_engine(sync_url)
|
||||||
|
with SyncSession(eng4) as s4:
|
||||||
|
set_tenant_context_sync(s4, tenant_id)
|
||||||
|
order_row2 = s4.execute(
|
||||||
|
sel2(OrderModel2.created_by, OrderModel2.order_number)
|
||||||
|
.join(OL2, OL2.order_id == OrderModel2.id)
|
||||||
|
.where(OL2.id == order_line_id)
|
||||||
|
).one_or_none()
|
||||||
|
eng4.dispose()
|
||||||
|
if order_row2:
|
||||||
|
emit_order_line_render_notifications(
|
||||||
|
success=False,
|
||||||
|
order_line_id=order_line_id,
|
||||||
|
order_number=order_row2[1],
|
||||||
|
order_creator_id=str(order_row2[0]),
|
||||||
|
product_name="unknown",
|
||||||
|
output_type_name="unknown",
|
||||||
|
render_log={"error": str(exc)},
|
||||||
|
emit_websocket=False,
|
||||||
|
activity_entity_id=None,
|
||||||
|
)
|
||||||
|
except Exception:
|
||||||
|
logger.exception("Failed to emit render failure activity event")
|
||||||
|
|
||||||
|
|
||||||
@celery_app.task(bind=True, name="app.tasks.step_tasks.render_order_line_task", queue="asset_pipeline", max_retries=3)
|
@celery_app.task(bind=True, name="app.tasks.step_tasks.render_order_line_task", queue="asset_pipeline", max_retries=3)
|
||||||
def render_order_line_task(self, order_line_id: str):
|
def render_order_line_task(self, order_line_id: str):
|
||||||
"""Render a specific output type for an order line.
|
"""Render a specific output type for an order line.
|
||||||
@@ -240,59 +388,15 @@ def render_order_line_task(self, order_line_id: str):
|
|||||||
logger.error("Cinematic render failed for %s: %s", order_line_id, exc)
|
logger.error("Cinematic render failed for %s: %s", order_line_id, exc)
|
||||||
elif is_animation:
|
elif is_animation:
|
||||||
# ── Turntable animation path ────────────────────────────────
|
# ── Turntable animation path ────────────────────────────────
|
||||||
emit(order_line_id, f"Starting turntable render: {frame_count} frames @ {fps}fps, {render_width or 1920}x{render_height or 1920}{tmpl_info}")
|
success, render_log = _render_turntable(
|
||||||
pl.step_start("blender_turntable", {"frame_count": frame_count, "fps": fps})
|
render_invocation=render_invocation,
|
||||||
from app.services.render_blender import is_blender_available, render_turntable_to_file
|
step_path=step_path,
|
||||||
if not is_blender_available():
|
output_path=output_path,
|
||||||
raise RuntimeError("Blender not available on this worker")
|
template=template,
|
||||||
|
order_line_id=order_line_id,
|
||||||
from app.services.step_processor import _get_all_settings
|
emit=emit,
|
||||||
_sys = _get_all_settings()
|
pl=pl,
|
||||||
try:
|
)
|
||||||
turntable_kwargs = render_invocation.as_turntable_renderer_kwargs(
|
|
||||||
step_path=step_path,
|
|
||||||
output_path=_Path(output_path),
|
|
||||||
default_width=1920,
|
|
||||||
default_height=1920,
|
|
||||||
default_engine=_sys.get("blender_engine", "cycles"),
|
|
||||||
default_samples=int(
|
|
||||||
_sys.get(
|
|
||||||
f"blender_{render_engine or _sys.get('blender_engine', 'cycles')}_samples",
|
|
||||||
128,
|
|
||||||
)
|
|
||||||
),
|
|
||||||
smooth_angle=int(_sys.get("blender_smooth_angle", 30)),
|
|
||||||
)
|
|
||||||
service_data = render_turntable_to_file(**turntable_kwargs)
|
|
||||||
success = True
|
|
||||||
render_log = {
|
|
||||||
"renderer": "blender",
|
|
||||||
"type": "turntable",
|
|
||||||
"format": "mp4",
|
|
||||||
"engine": render_engine or _sys.get("blender_engine", "cycles"),
|
|
||||||
"engine_used": service_data.get("engine_used", "cycles"),
|
|
||||||
"samples": render_samples,
|
|
||||||
"cycles_device": cycles_device_val,
|
|
||||||
"width": render_width or 1920,
|
|
||||||
"height": render_height or 1920,
|
|
||||||
"frame_count": service_data.get("frame_count", frame_count),
|
|
||||||
"fps": fps,
|
|
||||||
"total_duration_s": service_data.get("total_duration_s"),
|
|
||||||
"stl_duration_s": service_data.get("stl_duration_s"),
|
|
||||||
"render_duration_s": service_data.get("render_duration_s"),
|
|
||||||
"ffmpeg_duration_s": service_data.get("ffmpeg_duration_s"),
|
|
||||||
"stl_size_bytes": service_data.get("stl_size_bytes"),
|
|
||||||
"output_size_bytes": service_data.get("output_size_bytes"),
|
|
||||||
"log_lines": service_data.get("log_lines", []),
|
|
||||||
}
|
|
||||||
if template:
|
|
||||||
render_log["template"] = template.blend_file_path
|
|
||||||
pl.step_done("blender_turntable")
|
|
||||||
except Exception as exc:
|
|
||||||
success = False
|
|
||||||
render_log = {"renderer": "blender", "type": "turntable", "error": str(exc)[:500]}
|
|
||||||
pl.step_error("blender_turntable", str(exc), exc)
|
|
||||||
logger.error("Turntable render failed for %s: %s", order_line_id, exc)
|
|
||||||
else:
|
else:
|
||||||
# ── Still image path ────────────────────────────────────────
|
# ── Still image path ────────────────────────────────────────
|
||||||
_render_path_label = "USD → Blender" if usd_path else "STEP → GLB → Blender"
|
_render_path_label = "USD → Blender" if usd_path else "STEP → GLB → Blender"
|
||||||
@@ -358,69 +462,10 @@ def render_order_line_task(self, order_line_id: str):
|
|||||||
|
|
||||||
except Exception as exc:
|
except Exception as exc:
|
||||||
logger.error(f"render_order_line_task failed for {order_line_id}: {exc}")
|
logger.error(f"render_order_line_task failed for {order_line_id}: {exc}")
|
||||||
# If retries exhausted, mark as failed so the line doesn't stay stuck
|
|
||||||
if self.request.retries >= self.max_retries:
|
if self.request.retries >= self.max_retries:
|
||||||
logger.error(f"Max retries reached for {order_line_id}, marking as failed")
|
logger.error(f"Max retries reached for {order_line_id}, marking as failed")
|
||||||
try:
|
try:
|
||||||
from sqlalchemy import create_engine, update as sql_update2
|
_handle_render_task_exhausted(order_line_id, exc, _tenant_id)
|
||||||
from sqlalchemy.orm import Session as SyncSession
|
|
||||||
from app.config import settings as app_settings
|
|
||||||
from app.models.order_line import OrderLine as OL2
|
|
||||||
sync_url2 = app_settings.database_url.replace("+asyncpg", "")
|
|
||||||
eng2 = create_engine(sync_url2)
|
|
||||||
with SyncSession(eng2) as s2:
|
|
||||||
set_tenant_context_sync(s2, _tenant_id)
|
|
||||||
from datetime import datetime as dt2
|
|
||||||
s2.execute(
|
|
||||||
sql_update2(OL2).where(OL2.id == order_line_id)
|
|
||||||
.values(
|
|
||||||
render_status="failed",
|
|
||||||
render_completed_at=dt2.utcnow(),
|
|
||||||
render_log={"error": str(exc)[:500]},
|
|
||||||
)
|
|
||||||
)
|
|
||||||
s2.commit()
|
|
||||||
eng2.dispose()
|
|
||||||
from app.services.order_status_service import check_order_completion
|
|
||||||
# Try to get order_id from DB
|
|
||||||
eng3 = create_engine(sync_url2)
|
|
||||||
with SyncSession(eng3) as s3:
|
|
||||||
set_tenant_context_sync(s3, _tenant_id)
|
|
||||||
from sqlalchemy import select as sel
|
|
||||||
row = s3.execute(sel(OL2.order_id).where(OL2.id == order_line_id)).scalar_one_or_none()
|
|
||||||
if row:
|
|
||||||
check_order_completion(str(row))
|
|
||||||
eng3.dispose()
|
|
||||||
# Notify the order creator about the failure
|
|
||||||
try:
|
|
||||||
from sqlalchemy import select as sel2
|
|
||||||
from app.models.order import Order as OrderModel2
|
|
||||||
from app.domains.rendering.workflow_runtime_services import (
|
|
||||||
emit_order_line_render_notifications,
|
|
||||||
)
|
|
||||||
eng4 = create_engine(sync_url2)
|
|
||||||
with SyncSession(eng4) as s4:
|
|
||||||
set_tenant_context_sync(s4, _tenant_id)
|
|
||||||
order_row2 = s4.execute(
|
|
||||||
sel2(OrderModel2.created_by, OrderModel2.order_number)
|
|
||||||
.join(OL2, OL2.order_id == OrderModel2.id)
|
|
||||||
.where(OL2.id == order_line_id)
|
|
||||||
).one_or_none()
|
|
||||||
eng4.dispose()
|
|
||||||
if order_row2:
|
|
||||||
emit_order_line_render_notifications(
|
|
||||||
success=False,
|
|
||||||
order_line_id=order_line_id,
|
|
||||||
order_number=order_row2[1],
|
|
||||||
order_creator_id=str(order_row2[0]),
|
|
||||||
product_name="unknown",
|
|
||||||
output_type_name="unknown",
|
|
||||||
render_log={"error": str(exc)},
|
|
||||||
emit_websocket=False,
|
|
||||||
activity_entity_id=None,
|
|
||||||
)
|
|
||||||
except Exception:
|
|
||||||
logger.exception("Failed to emit render failure activity event")
|
|
||||||
except Exception:
|
except Exception:
|
||||||
logger.exception(f"Failed to mark {order_line_id} as failed in DB")
|
logger.exception(f"Failed to mark {order_line_id} as failed in DB")
|
||||||
raise
|
raise
|
||||||
|
|||||||
Reference in New Issue
Block a user