200 lines
8.8 KiB
Python
200 lines
8.8 KiB
Python
"""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",
|
|
}
|