| 12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758 |
- 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,
- }
|