from __future__ import annotations import json from typing import Any from supply_infra.db.models.demand_feedback import DemandFeedback from supply_infra.db.session import get_session from supply_infra.pipeline.dates import china_now _PROCESSING_STATUSES = {"accepted", "rejected", "needs_evidence"} def process_feedback( feedback_id: int, *, status: str, processor: str, reason: str, impact: dict[str, Any] | None = None, ) -> dict[str, Any]: if status not in _PROCESSING_STATUSES: raise ValueError(f"unsupported feedback status: {status}") if not processor.strip() or not reason.strip(): raise ValueError("processor and reason are required") with get_session() as session: row = session.get(DemandFeedback, int(feedback_id), with_for_update=True) if row is None: raise LookupError(f"feedback not found: {feedback_id}") if row.processing_status in {"accepted", "rejected"}: if row.processing_status != status: raise RuntimeError( f"feedback already finalized as {row.processing_status}" ) return { "feedback_id": int(row.id), "processing_status": row.processing_status, "idempotent_replay": True, } row.processing_status = status row.processed_by = processor[:128] row.processed_at = china_now() row.resolution_reason = reason row.impact_json = ( json.dumps(impact, ensure_ascii=False, sort_keys=True) if impact is not None else None ) session.flush() return { "feedback_id": int(row.id), "processing_status": row.processing_status, "processed_by": row.processed_by, "processed_at": row.processed_at.isoformat(), "resolution_reason": row.resolution_reason, "impact": impact, "idempotent_replay": False, }