"""Human feedback for demand-summary records.""" from __future__ import annotations import json from typing import Any from pydantic import ValidationError from sqlalchemy.exc import IntegrityError from api.schemas.demand_feedback import CreateDemandFeedbackBody from api.services.video_discovery import _parse_video_ids from supply_infra.db.models.demand_feedback import DemandFeedback from supply_infra.db.repositories.demand_feedback_repo import DemandFeedbackRepository from supply_infra.db.repositories.demand_grade_repo import DemandGradeRepository from supply_infra.db.repositories.demand_video_expansion_repo import ( DemandVideoExpansionRepository, ) from supply_infra.db.repositories.multi_demand_video_detail_repo import ( MultiDemandVideoDetailRepository, ) from supply_infra.db.session import get_session _SOURCE_VIDEO_GRADES = frozenset({"B", "C", "D"}) class FeedbackTargetNotFoundError(Exception): """Raised when the requested demand, video or hit content does not exist.""" class FeedbackTargetConflictError(Exception): """Raised when target identifiers do not belong to the same record.""" class FeedbackRequestConflictError(Exception): """Raised when an idempotency key is reused for a different request.""" def _feedback_summary(count: int) -> dict[str, int]: return {"count": int(count)} def _serialize_feedback(row: DemandFeedback) -> dict[str, Any]: try: target_snapshot = json.loads(row.target_snapshot_json) except (TypeError, ValueError): target_snapshot = {} return { "id": int(row.id), "target_type": row.target_type, "biz_dt": row.biz_dt, "demand_grade_id": int(row.demand_grade_id), "video_id": row.video_id, "demand_video_expansion_id": ( int(row.demand_video_expansion_id) if row.demand_video_expansion_id is not None else None ), "feedback_action": row.feedback_action, "reason_code": row.reason_code, "content": row.content, "target_snapshot": target_snapshot, "feedback_user": { "id": int(row.feedback_user_id), "username": row.feedback_username_snapshot, "display_name": row.feedback_user_name_snapshot, }, "created_at": row.created_at.isoformat() if row.created_at else None, } def _same_request( row: DemandFeedback, body: CreateDemandFeedbackBody, feedback_user_id: int, ) -> bool: return ( row.feedback_user_id == feedback_user_id and row.target_type == body.target_type and row.demand_grade_id == body.demand_grade_id and row.video_id == body.video_id and row.demand_video_expansion_id == body.demand_video_expansion_id and row.feedback_action == body.feedback_action and row.reason_code == body.reason_code and row.content == body.content ) def _resolve_target( session: Any, body: CreateDemandFeedbackBody, ) -> tuple[Any, dict[str, Any]]: grade = DemandGradeRepository(session).get_by_id(body.demand_grade_id) if grade is None: raise FeedbackTargetNotFoundError("需求不存在") snapshot: dict[str, Any] = { "demand_name": grade.demand_name, "grade": grade.grade, } if body.target_type == "demand": snapshot["reason"] = grade.reason return grade, snapshot video_id = body.video_id or "" grade_code = str(grade.grade or "").upper() expansion_repo = DemandVideoExpansionRepository(session) if grade_code in _SOURCE_VIDEO_GRADES: video_exists = video_id in _parse_video_ids(grade.video_list) video_source = "pool" else: video_exists = expansion_repo.has_video( biz_dt=str(grade.biz_dt), source_demand_grade_id=int(grade.id), video_id=video_id, ) video_source = "expansion" if not video_exists: raise FeedbackTargetConflictError("视频不属于当前需求") detail = MultiDemandVideoDetailRepository(session).list_by_vids([video_id]).get( video_id ) snapshot.update( { "video_id": video_id, "video_title": detail.title if detail else None, "video_source": video_source, } ) if body.target_type == "video": return grade, snapshot expansion = expansion_repo.get_active_by_id(body.demand_video_expansion_id or 0) if expansion is None: raise FeedbackTargetNotFoundError("命中内容不存在") if ( int(expansion.source_demand_grade_id) != int(grade.id) or str(expansion.biz_dt) != str(grade.biz_dt) or str(expansion.video_id) != video_id ): raise FeedbackTargetConflictError("命中内容不属于当前需求和视频") snapshot.update( { "point_type": expansion.point_type, "expanded_text": expansion.expanded_text, "point_desc": expansion.point_desc, "reason": expansion.reason, } ) return grade, snapshot def create_demand_feedback( body: CreateDemandFeedbackBody, current_user: dict[str, Any], ) -> dict[str, Any]: feedback_user_id = int(current_user["id"]) feedback_username = str(current_user["username"]) feedback_user_name = str( current_user.get("display_name") or current_user.get("username") or feedback_user_id ) with get_session() as session: repo = DemandFeedbackRepository(session) existing = repo.get_by_client_request_id(body.client_request_id) if existing is not None: if not _same_request(existing, body, feedback_user_id): raise FeedbackRequestConflictError("请求标识已用于其他反馈") _, video_counts, expansion_counts = repo.count_for_demand( body.demand_grade_id ) count = ( repo.count_demands([body.demand_grade_id]).get(body.demand_grade_id, 0) if body.target_type == "demand" else video_counts.get(body.video_id or "", 0) if body.target_type == "video" else expansion_counts.get(body.demand_video_expansion_id or 0, 0) ) return { "item": _serialize_feedback(existing), "feedback_summary": _feedback_summary(count), } grade, target_snapshot = _resolve_target(session, body) feedback = DemandFeedback( client_request_id=body.client_request_id, target_type=body.target_type, biz_dt=str(grade.biz_dt), demand_grade_id=int(grade.id), video_id=body.video_id, demand_video_expansion_id=body.demand_video_expansion_id, feedback_action=body.feedback_action, reason_code=body.reason_code, content=body.content, target_snapshot_json=json.dumps( target_snapshot, ensure_ascii=False, separators=(",", ":"), ), feedback_user_id=feedback_user_id, feedback_user_name_snapshot=feedback_user_name[:128], feedback_username_snapshot=feedback_username[:64], ) try: repo.add(feedback) session.refresh(feedback) except IntegrityError: session.rollback() existing = repo.get_by_client_request_id(body.client_request_id) if existing is None or not _same_request(existing, body, feedback_user_id): raise FeedbackRequestConflictError("请求标识已用于其他反馈") from None feedback = existing rows, total = repo.list_for_target( target_type=body.target_type, demand_grade_id=body.demand_grade_id, video_id=body.video_id, demand_video_expansion_id=body.demand_video_expansion_id, limit=1, offset=0, ) del rows return { "item": _serialize_feedback(feedback), "feedback_summary": _feedback_summary(total), } def list_demand_feedback( *, target_type: str, demand_grade_id: int, video_id: str | None, demand_video_expansion_id: int | None, limit: int, offset: int, ) -> dict[str, Any]: try: body = CreateDemandFeedbackBody( client_request_id="history-query", target_type=target_type, demand_grade_id=demand_grade_id, video_id=video_id, demand_video_expansion_id=demand_video_expansion_id, feedback_action="support", ) except ValidationError as exc: raise FeedbackTargetConflictError("反馈目标参数不完整") from exc with get_session() as session: repo = DemandFeedbackRepository(session) rows, total = repo.list_for_target( target_type=target_type, demand_grade_id=demand_grade_id, video_id=video_id, demand_video_expansion_id=demand_video_expansion_id, limit=limit, offset=offset, ) if total == 0: _resolve_target(session, body) return { "items": [_serialize_feedback(row) for row in rows], "total": total, "limit": limit, "offset": offset, }