"""Profile-tag drift review and advisory-operation gating.""" from datetime import UTC, datetime from typing import Any, Protocol from uuid import uuid4 from sqlalchemy import update from sqlalchemy.ext.asyncio import AsyncSession from app.core.contracts import RequestContext from app.core.errors import GenericResourceNotFoundError, InvalidStateError from app.core.profile_governance_contracts import ProfileDriftReviewRequest from app.infrastructure.db import SessionFactory from app.model.audit import InteractionAudit from app.model.memory import MemorySyncOutbox from app.model.profile_tag import AdvisorProfileDriftReview, AdvisorProfileTag from app.model.risk_questionnaire import ProfileSnapshot from app.repository.risk_questionnaire_repository import RiskQuestionnaireRepository from app.service.api_transaction_service import ApiTransactionService, digest from app.service.authorization_service import AuthorizationService class SessionFactoryLike(Protocol): def __call__(self) -> AsyncSession: ... class ProfileGovernanceService: """Internal-only profile governance; no customer-facing projection is returned.""" def __init__(self, session_factory: SessionFactoryLike = SessionFactory) -> None: self.session_factory = session_factory async def require_operable(self, customer_id: int) -> None: """Block advisory decisions while a changed profile is awaiting review.""" async with self.session_factory() as session: pending = await RiskQuestionnaireRepository(session).pending_drift_review(customer_id) if pending is not None: raise InvalidStateError("画像标签漂移正在复核,暂不能执行投顾分析") async def tags(self, customer_id: int, context: RequestContext) -> dict[str, object]: await AuthorizationService.require(context, "profile-governance:read", admin=True) async with self.session_factory() as session: rows = await RiskQuestionnaireRepository(session).tags(customer_id) return { "data": [self._tag_view(row) for row in rows], "meta": {"trace_id": context.trace_id}, } async def pending_reviews(self, context: RequestContext) -> dict[str, object]: await AuthorizationService.require(context, "profile-governance:read", admin=True) async with self.session_factory() as session: rows = await RiskQuestionnaireRepository(session).pending_reviews(limit=100) return { "data": [self._review_view(row) for row in rows], "meta": {"trace_id": context.trace_id}, } async def review( self, review_id: int, payload: ProfileDriftReviewRequest, context: RequestContext, key: str | None, ) -> dict[str, object]: await AuthorizationService.require(context, "profile-governance:review", admin=True) async def operation(session: AsyncSession) -> dict[str, Any]: repository = RiskQuestionnaireRepository(session) review = await repository.drift_review(review_id, lock=True) if review is None: raise GenericResourceNotFoundError("画像漂移复核记录不存在") if review.status != "pending_review": raise InvalidStateError("画像漂移复核记录当前不能审核") now = datetime.now(UTC).replace(tzinfo=None) tags = await repository.tags_for_review(review_id, lock=True) if payload.decision == "approved": await self._approve(repository, review, tags, now) else: await session.execute( update(AdvisorProfileTag) .where(AdvisorProfileTag.drift_review_id == review_id) .values(status="rejected", active_customer_tag=None, updated_at=now) ) review.status = payload.decision review.reviewer_user_id = int(context.user_id) review.reviewed_at = now review.review_comment = payload.comment review.updated_at = now session.add(InteractionAudit( actor_type="user", actor_id=int(context.user_id), target_customer_id=review.customer_id, portal="admin", action_type="advisor.profile_drift_reviewed", detail={"review_id": review_id, "decision": payload.decision, "trace_id": context.trace_id}, created_at=now, )) await session.flush() return { "data": {"review_id": str(review.id), "status": review.status}, "meta": {"trace_id": context.trace_id}, } return await ApiTransactionService().execute( context, f"advisor:profile-drift-review:{review_id}", key, payload.model_dump(mode="json"), operation, ) async def _approve( self, repository: RiskQuestionnaireRepository, review: AdvisorProfileDriftReview, tags: list[AdvisorProfileTag], now: datetime, ) -> None: await repository.deactivate_current_profile(review.customer_id, now) active = await repository.active_tags(review.customer_id, lock=True) await repository.supersede_active_tags( review.customer_id, tuple(tag.tag_key for tag in active), now ) repository.add_profile(ProfileSnapshot( profile_uuid=review.candidate_profile_uuid, customer_id=review.customer_id, version=review.candidate_profile_version, snapshot=review.candidate_snapshot, generation_basis=review.candidate_generation_basis, snapshot_hash=digest(review.candidate_snapshot), is_current=True, generated_at=now, created_at=now, updated_at=now, )) tag_keys = {tag.tag_key for tag in tags} if tag_keys: await repository.session.execute( update(AdvisorProfileTag) .where( AdvisorProfileTag.drift_review_id == review.id, AdvisorProfileTag.status == "pending_review", ) .values( status="active", active_customer_tag=( # A single SQL value cannot vary by tag; update individually below. None ), updated_at=now, ) ) for tag in tags: tag.status = "active" tag.active_customer_tag = f"{review.customer_id}:{tag.tag_key}" tag.updated_at = now self._add_sync_events(repository, review, now) @staticmethod def _add_sync_events( repository: RiskQuestionnaireRepository, review: AdvisorProfileDriftReview, now: datetime, ) -> None: # The same event UUID across targets is intentional; the baseline unique key is # (event_uuid, target_store), so both projections share one logical change. event_uuid = str(uuid4()) payload = { "customer_id": str(review.customer_id), "profile_uuid": review.candidate_profile_uuid, "version": review.candidate_profile_version, "profile": review.candidate_snapshot, } for target_store in ("milvus", "neo4j"): repository.add_sync_event(MemorySyncOutbox( event_uuid=event_uuid, aggregate_type="profile", aggregate_uuid=review.candidate_profile_uuid, aggregate_version=review.candidate_profile_version, target_store=target_store, operation="upsert", payload=payload, status="pending", retry_count=0, next_retry_at=None, last_error=None, created_at=now, processed_at=None, )) @staticmethod def _tag_view(row: AdvisorProfileTag) -> dict[str, object]: return { "tag_id": str(row.id), "customer_id": str(row.customer_id), "tag_key": row.tag_key, "tag_value": row.tag_value, "confidence": str(row.confidence), "source_type": row.source_type, "source_reference": row.source_reference, "source_confidence": str(row.source_confidence), "profile_version": row.profile_version, "status": row.status, "drift_review_id": str(row.drift_review_id) if row.drift_review_id else None, } @staticmethod def _review_view(row: AdvisorProfileDriftReview) -> dict[str, object]: return { "review_id": str(row.id), "drift_no": row.drift_no, "customer_id": str(row.customer_id), "candidate_profile_version": row.candidate_profile_version, "changed_tags": row.changed_tags, "status": row.status, "created_at": row.created_at.isoformat() + "Z", }