Files
group_fqcd_jr/app/service/profile_governance_service.py
T
张胜宇 e239eb778b docs: 品牌全量口径统一为「南方基金」+ 作废文档清理
1) 客服 Agent 四份交付文档 + 构建脚手架:品牌由包装占位 XX科技 / 旧名 南方财富
   统一为南方基金(热线 400-889-8899 / 官网 nffund.com),系统名改为「智能服务系统」;
   同步追加 §0.4 修订记录行,工程记录行保留原占位字面以支撑硬编码扫描验收。
2) 开发文档:清理 28 份已作废/残留文档(14 份移出归档 + 14 份仓库副本),
   新增《文档规整方案与开发前待决事项-2026-09-17》。
3) 客服agent 四份交付文档首次纳入本分支。
2026-09-17 15:15:22 +08:00

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