feedback.py 2.0 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758
  1. from __future__ import annotations
  2. import json
  3. from typing import Any
  4. from supply_infra.db.models.demand_feedback import DemandFeedback
  5. from supply_infra.db.session import get_session
  6. from supply_infra.pipeline.dates import china_now
  7. _PROCESSING_STATUSES = {"accepted", "rejected", "needs_evidence"}
  8. def process_feedback(
  9. feedback_id: int,
  10. *,
  11. status: str,
  12. processor: str,
  13. reason: str,
  14. impact: dict[str, Any] | None = None,
  15. ) -> dict[str, Any]:
  16. if status not in _PROCESSING_STATUSES:
  17. raise ValueError(f"unsupported feedback status: {status}")
  18. if not processor.strip() or not reason.strip():
  19. raise ValueError("processor and reason are required")
  20. with get_session() as session:
  21. row = session.get(DemandFeedback, int(feedback_id), with_for_update=True)
  22. if row is None:
  23. raise LookupError(f"feedback not found: {feedback_id}")
  24. if row.processing_status in {"accepted", "rejected"}:
  25. if row.processing_status != status:
  26. raise RuntimeError(
  27. f"feedback already finalized as {row.processing_status}"
  28. )
  29. return {
  30. "feedback_id": int(row.id),
  31. "processing_status": row.processing_status,
  32. "idempotent_replay": True,
  33. }
  34. row.processing_status = status
  35. row.processed_by = processor[:128]
  36. row.processed_at = china_now()
  37. row.resolution_reason = reason
  38. row.impact_json = (
  39. json.dumps(impact, ensure_ascii=False, sort_keys=True)
  40. if impact is not None
  41. else None
  42. )
  43. session.flush()
  44. return {
  45. "feedback_id": int(row.id),
  46. "processing_status": row.processing_status,
  47. "processed_by": row.processed_by,
  48. "processed_at": row.processed_at.isoformat(),
  49. "resolution_reason": row.resolution_reason,
  50. "impact": impact,
  51. "idempotent_replay": False,
  52. }