merge: integrate ZSY customer service and profile capabilities

This commit is contained in:
张胜宇
2026-09-11 22:31:51 +08:00
94 changed files with 7933 additions and 77 deletions
+56 -1
View File
@@ -1,11 +1,17 @@
import argparse
import asyncio
import logging
from typing import Any
from app.core.config import get_settings
from app.core.errors import RecoverableAgentError
from app.infrastructure.db import engine
from app.service.agent.bootstrap import get_relationship_service
from app.infrastructure.milvus_profile_projection import MilvusProfileProjection
from app.infrastructure.neo4j_profile_projection import Neo4jProfileProjection
from app.service.agent.bootstrap import get_memory_embedding_service, get_relationship_service
from app.service.model_gateway import DatabaseModelEndpointResolver
from app.service.projection_cleanup_service import ProjectionCleanupService
from app.worker.memory_sync_outbox_worker import MemorySyncOutboxWorker
from app.worker.offsite_mail_worker import OffsiteMailWorker
from app.worker.runtime import WorkerRuntime
@@ -28,6 +34,46 @@ async def serve(*, once: bool = False) -> None:
relationships=relationships
),
)
# 画像投影是 MySQL 审核结果的异步派生写入;任一外部存储未配置时保持事件 pending。
neo4j_driver: Any | None = None
milvus_client: Any | None = None
memory_sync_handlers: dict[str, Any] = {}
if settings.neo4j_password:
from neo4j import AsyncGraphDatabase
neo4j_driver = AsyncGraphDatabase.driver(
settings.neo4j_uri,
auth=(settings.neo4j_username, settings.neo4j_password),
)
memory_sync_handlers["neo4j"] = Neo4jProfileProjection(neo4j_driver).upsert
else:
logger.warning("Neo4j password not configured; profile projection remains pending")
if settings.resolved_milvus_uri and settings.knowledge_embedding_endpoint_code:
from pymilvus import AsyncMilvusClient # type: ignore[import-untyped]
milvus_client = AsyncMilvusClient(
uri=settings.resolved_milvus_uri,
token=settings.milvus_token or None,
)
async def embed_profile(text: str) -> list[float]:
endpoints = await DatabaseModelEndpointResolver().resolve(
agent_type="memory_projection", task_type="embedding"
)
if not endpoints:
raise RecoverableAgentError("没有可用的 embedding 端点")
return (await get_memory_embedding_service().embed(endpoints, text)).vector
memory_sync_handlers["milvus"] = MilvusProfileProjection(
milvus_client, embed_profile
).upsert
else:
logger.warning(
"Milvus profile projection not configured; profile projection remains pending"
)
memory_sync_worker = (
MemorySyncOutboxWorker(memory_sync_handlers) if memory_sync_handlers else None
)
# 场外收件 Worker 必须与底座 Worker 同进程同入口:2026-09-11 01:45 的一次批量
# 文件覆盖把这处接线删掉了,导致邮件 Worker 完全不再运行、邮箱无人收取。
offsite_worker = OffsiteMailWorker(settings)
@@ -35,6 +81,11 @@ async def serve(*, once: bool = False) -> None:
while True:
try:
worked = await runtime.run_once()
if memory_sync_worker is not None:
for target_store in ("neo4j", "milvus"):
worked = await memory_sync_worker.run_once(
target_store=target_store
) or worked
worked = await offsite_worker.run_once() or worked
except Exception:
# 常驻 Worker 不能因为"某一轮"的异常就整体退出:数据库抖动、
@@ -52,6 +103,10 @@ async def serve(*, once: bool = False) -> None:
await asyncio.sleep(settings.worker_poll_seconds)
finally:
await offsite_worker.close()
if neo4j_driver is not None:
await neo4j_driver.close()
if milvus_client is not None:
await milvus_client.close()
await engine.dispose()
@@ -0,0 +1,50 @@
"""客服对话画像候选消费者。
该消费者只接收已登录用户的候选事件,并将脱敏后的模型抽取结果写为
``memory_unit.status='candidate'``。候选不会进入客服召回,也不会修改正式画像。
"""
from datetime import datetime
from typing import Any
from sqlalchemy.ext.asyncio import AsyncSession
from app.service.memory_extraction_service import (
MemoryExtractionService,
get_memory_extraction_service,
)
from app.service.memory_service import CacheDeleteAdapter
from app.worker.memory_extraction_worker import MemoryExtractionWorker
class CustomerProfileCandidateWorker:
"""消费已登录客服会话的画像候选事件;访客事件失败关闭。"""
def __init__(
self,
session: AsyncSession,
*,
extractor: MemoryExtractionService | None = None,
cache: CacheDeleteAdapter | None = None,
) -> None:
self.worker = MemoryExtractionWorker(
session,
extractor=extractor or get_memory_extraction_service(),
cache=cache,
memory_status="candidate",
source_type="AI对话提取",
event_type="customer_profile.candidate_requested",
sanitize_source=True,
)
async def handle(
self,
payload: dict[str, Any],
*,
event_id: str | None = None,
occurred_at: datetime | None = None,
) -> bool:
"""只允许受理服务标记的已登录客户事件,防止访客写入画像候选。"""
if payload.get("actor_type") != "authenticated_customer":
return False
return await self.worker.handle(payload, event_id=event_id, occurred_at=occurred_at)
+15 -3
View File
@@ -6,6 +6,7 @@ from uuid import uuid4
from sqlalchemy import select
from sqlalchemy.ext.asyncio import AsyncSession
from app.core.conversation_privacy import sanitize_customer_service_message
from app.model.conversation import ConversationMessage
from app.model.memory import MemoryEvidence
from app.model.platform import AgentRun, DomainEventOutbox
@@ -42,6 +43,10 @@ class MemoryExtractionWorker:
*,
extractor: MemoryExtractionService | None = None,
cache: CacheDeleteAdapter | None = None,
memory_status: str = "active",
source_type: str = SOURCE_TYPE,
event_type: str = "memory.extraction_requested",
sanitize_source: bool = False,
) -> None:
self.session = session
# 默认走生产装配(与业务 Agent 同一个 ModelGenerationService);验收探针
@@ -49,6 +54,10 @@ class MemoryExtractionWorker:
self.extractor = extractor if extractor is not None else get_memory_extraction_service()
# 召回热缓存适配器:写入生效后必须失效,否则新记忆在 TTL 内召回不到。
self.cache = cache
self.memory_status = memory_status
self.source_type = source_type
self.event_type = event_type
self.sanitize_source = sanitize_source
async def handle(
self, payload: dict[str, Any], *, event_id: str | None = None,
@@ -74,7 +83,7 @@ class MemoryExtractionWorker:
# 没有可追溯的事件 id 就不能建立幂等边界,重复消费将无法去重。
logger.warning("memory extraction skipped: event id not found run_id=%s", run_id)
return False
idempotency_key = f"memory.extraction_requested:{event_id}"
idempotency_key = f"{self.event_type}:{event_id}"
seen = await self.session.scalar(
select(MemoryEvidence.id).where(MemoryEvidence.idempotency_key == idempotency_key))
if seen is not None:
@@ -92,7 +101,8 @@ class MemoryExtractionWorker:
extracted.value,
memory_type=extracted.memory_type,
confidence=extracted.confidence,
source_type=SOURCE_TYPE,
source_type=self.source_type,
status=self.memory_status,
structured_value={
"memory_key": extracted.memory_key,
"value": extracted.value,
@@ -146,7 +156,7 @@ class MemoryExtractionWorker:
async def _event_id(self, run_id: str, result_message_id: int) -> str | None:
"""按 payload 回查本事件的事件 id,作为重复消费的幂等边界。"""
conditions = [DomainEventOutbox.event_type == "memory.extraction_requested"]
conditions = [DomainEventOutbox.event_type == self.event_type]
if run_id:
conditions.append(DomainEventOutbox.aggregate_id == run_id)
else:
@@ -179,6 +189,8 @@ class MemoryExtractionWorker:
content = (message.content or "").strip()
if not content:
return None, message.id
if self.sanitize_source:
content = sanitize_customer_service_message(content)
return content, message.id
async def _request_message(
+92
View File
@@ -0,0 +1,92 @@
"""画像投影 Outbox 消费器。
外部存储通过 handler 注入;本模块只负责领取、重试、死信和 MySQL 状态更新。
"""
from collections.abc import Awaitable, Callable
from datetime import UTC, datetime, timedelta
from typing import Any
from sqlalchemy import select
from sqlalchemy.ext.asyncio import AsyncSession
from app.infrastructure.db import SessionFactory
from app.model.memory import MemorySyncOutbox
ProjectionHandler = Callable[[dict[str, Any]], Awaitable[Any]]
MAX_RETRY_COUNT = 5
class MemorySyncOutboxWorker:
"""按目标存储独立消费画像投影事件。"""
def __init__(
self,
handlers: dict[str, ProjectionHandler],
*,
session_factory: Callable[[], AsyncSession] = SessionFactory,
) -> None:
self.handlers = handlers
self.session_factory = session_factory
async def run_once(
self, *, target_store: str | None = None, event_uuid: str | None = None
) -> bool:
"""领取并处理一条到期事件;没有可处理事件时返回 False。"""
if not self.handlers:
return False
async with self.session_factory() as session:
now = datetime.now(UTC).replace(tzinfo=None)
conditions: list[Any] = [
MemorySyncOutbox.status.in_({"pending", "failed"}),
MemorySyncOutbox.next_retry_at.is_(None)
| (MemorySyncOutbox.next_retry_at <= now),
]
if target_store is not None:
conditions.append(MemorySyncOutbox.target_store == target_store)
if event_uuid is not None:
conditions.append(MemorySyncOutbox.event_uuid == event_uuid)
event = await session.scalar(
select(MemorySyncOutbox)
.where(*conditions)
.order_by(MemorySyncOutbox.id)
.limit(1)
.with_for_update(skip_locked=True)
)
if event is None:
await session.rollback()
return False
handler = self.handlers.get(event.target_store)
if handler is None:
self._fail(event, "target_handler_not_configured", now, dead=True)
await session.commit()
return True
try:
await handler(event.payload)
except Exception as exc:
self._fail(event, type(exc).__name__, now)
else:
event.status = "processed"
event.processed_at = datetime.now(UTC).replace(tzinfo=None)
event.last_error = None
event.next_retry_at = None
await session.commit()
return True
@staticmethod
def _fail(
event: MemorySyncOutbox,
reason: str,
now: datetime,
*,
dead: bool = False,
) -> None:
"""写入可重试失败或死信状态,不吞掉失败事实。"""
event.retry_count = int(event.retry_count) + 1
event.last_error = reason[:500]
event.status = "dead" if dead or event.retry_count >= MAX_RETRY_COUNT else "failed"
event.next_retry_at = (
None
if event.status == "dead"
else now + timedelta(seconds=min(300, 2 ** int(event.retry_count)))
)
+135 -16
View File
@@ -10,12 +10,17 @@ from uuid import uuid4
from sqlalchemy import select, update
from app.core.config import Settings, get_settings
from app.core.contracts import AgentRequest, AgentResult, RequestContext
from app.core.contracts import AgentRequest, AgentRequestMetadata, AgentResult, RequestContext
from app.core.errors import AgentError, RecoverableAgentError, RunLeaseLostError
from app.infrastructure.db import SessionFactory
from app.model.audit import InteractionAudit
from app.model.conversation import ConversationMessage
from app.model.platform import AgentRun, DomainEventOutbox, RequestIdempotency
from app.model.platform import (
AgentRun,
DomainEventOutbox,
HandoverTicket,
RequestIdempotency,
)
from app.repository.agent_run_repository import AgentRunRepository
from app.service.agent.bootstrap import (
get_agent_factory,
@@ -27,6 +32,11 @@ from app.service.agent.bootstrap import (
from app.service.agent.executor import AgentExecutor
from app.service.agent.factory import AgentFactory
from app.service.agent_persistence_service import AgentPersistenceService
from app.service.customer_service_session_memory_service import (
CustomerServiceSessionMemory,
CustomerServiceSessionTurn,
build_customer_service_session_memory,
)
from app.service.identity_service import IdentityService
from app.service.memory_extraction_service import (
ExtractionEndpointResolver,
@@ -37,6 +47,7 @@ from app.service.memory_recall_service import MemoryRecallService
from app.service.memory_service import CacheDeleteAdapter, MemoryService
from app.service.memory_taxonomy import BUSINESS_EVENT_TYPES
from app.service.model_gateway import DatabaseModelEndpointResolver, ModelGenerationService
from app.worker.customer_profile_candidate_worker import CustomerProfileCandidateWorker
from app.worker.episode_worker import (
EpisodeConsumptionResult,
EpisodeExtractionConsumer,
@@ -91,6 +102,7 @@ class WorkerRuntime:
knowledge_writer: Any = _UNSET,
knowledge_embedder: Any = _UNSET,
knowledge_endpoint_resolver: Any = _UNSET,
session_memory: CustomerServiceSessionMemory | None = None,
) -> None:
self.factory = factory if factory is not None else get_agent_factory()
self.settings = settings or get_settings()
@@ -112,6 +124,11 @@ class WorkerRuntime:
# 这里不兜底。组件内部给默认实现会把"尚未装配"这一事实悄悄盖住——而"未注入即显式
# 降级并留痕"是本模块刻意保留的语义(有单测守着),因此默认值保持 None。
self.projection_cleaner = projection_cleaner
# 客服短期会话 Redis 仅用于当前会话上下文,不参与长期画像召回。
self.session_memory = (
session_memory if session_memory is not None
else build_customer_service_session_memory()
)
# 图关系服务:画像投影用它写入节点与关系(投顾的多跳推荐、风控的关系网络都读它)。
# 默认取生产装配;图库不可用时该值为 None,投影如实降级而不是失败。
if relationships is not None:
@@ -141,6 +158,48 @@ class WorkerRuntime:
# episode 聚合是低频批处理,按轮次节流而不是每轮都查。
self._episode_rounds = 0
async def restore_context(
self, *, actor_type: str, actor_id: str, trace_id: str
) -> RequestContext:
"""按受理事件中的可信身份恢复最小执行权限。"""
identity = RequestContext(user_id=actor_id, trace_id=trace_id)
if actor_type == "visitor":
return identity.model_copy(update={
"roles": ("visitor",),
"permissions": ("agent:run", "knowledge:query"),
"data_scope": "public",
})
return await self.resolve_identity(identity)
@staticmethod
def should_request_memory_extraction(
*, agent_type: str, context: RequestContext, message: str,
result: AgentResult, business_events: tuple[str, ...] | list[str],
) -> bool:
"""长期记忆抽取只接收非客服、非访客的明确业务事实。"""
if agent_type == "customer_service" or "visitor" in context.roles:
return False
return MemoryService.should_extract_memory(
conversation_content=message,
role="user",
tool_result=any(call.status == "succeeded" for call in result.result.tool_calls),
event_type=business_events[0] if business_events else None,
signals=MemoryService.detect_memory_signals(message),
)
@staticmethod
def should_request_profile_candidate(
*, agent_type: str, context: RequestContext, message: str,
) -> bool:
"""客服仅为已登录且 self 范围内的用户生成待确认画像候选。"""
if agent_type != "customer_service" or "visitor" in context.roles:
return False
if not {"customer", "authenticated_user"}.intersection(context.roles):
return False
if context.data_scope != "self":
return False
return bool(MemoryService.detect_memory_signals(message))
async def dispatch_one(self, *, run_id: str | None = None) -> bool:
# Outbox acknowledges a durable SQL queue entry, not an in-memory task.
async with SessionFactory() as session:
@@ -158,6 +217,13 @@ class WorkerRuntime:
session, extractor=self.memory_extraction, cache=self.memory_cache
).handle(payload)
async def dispatch_profile_candidate(payload: dict[str, Any]) -> None:
if "message_id" not in payload or "customer_id" not in payload:
raise OutboxHandlerError("profile candidate payload is incomplete")
await CustomerProfileCandidateWorker(
session, extractor=self.memory_extraction, cache=self.memory_cache
).handle(payload)
async def dispatch_run_completed(payload: dict[str, Any]) -> None:
# 结果消息与审计已由 complete_run 同事务落库,此事件只承担
# "运行已完成"的对外通知职责。当前没有独立外部消费者,
@@ -196,6 +262,33 @@ class WorkerRuntime:
logger.warning("graph projection degraded customer_id=%s reason=%s",
customer_id, result.reason)
async def dispatch_handover_queue_ready(payload: dict[str, Any]) -> None:
"""记录转人工队列已就绪;不向客户承诺已接单或处理时限。"""
ticket_no = str(payload.get("ticket_no", "")).strip()
if not ticket_no:
raise OutboxHandlerError(
"conversation.transfer_requested payload is incomplete"
)
ticket = await session.scalar(
select(HandoverTicket).where(HandoverTicket.ticket_no == ticket_no)
)
if ticket is None:
raise OutboxHandlerError("handover ticket not found")
session.add(InteractionAudit(
actor_type="system", actor_id=None,
target_customer_id=ticket.customer_id,
session_id=ticket.session_id, portal="worker",
action_type="handover.queue_ready",
detail={
"ticket_no": ticket.ticket_no,
"source_agent": ticket.source_agent,
"reason_code": ticket.reason_code,
"ticket_status": ticket.status,
},
created_at=datetime.now(UTC).replace(tzinfo=None),
))
await session.flush()
async def dispatch_projection_cleanup(payload: dict[str, Any]) -> None:
# memory.invalidated / memory.deleted 由 MemoryLifecycleService 按
# memory_uuid 写入,这里做幂等的投影清理(Milvus 向量、Neo4j 关系)。
@@ -220,6 +313,7 @@ class WorkerRuntime:
handlers: dict[str, Callable[[dict[str, Any]], Awaitable[None]]] = {
"agent.run_requested": dispatch,
"memory.extraction_requested": dispatch_memory_extraction,
"customer_profile.candidate_requested": dispatch_profile_candidate,
"agent.run_completed": dispatch_run_completed,
"config.cache_invalidate_requested": dispatch_cache_invalidate,
"memory.deletion_requested": dispatch_memory_deletion,
@@ -228,6 +322,7 @@ class WorkerRuntime:
"memory.deleted": dispatch_projection_cleanup,
# 画像重建:记忆写入后自动触发,使「记忆 → 画像 → 图」全链路无需手工介入
"profile.rebuild_requested": dispatch_profile_rebuild,
"conversation.transfer_requested": dispatch_handover_queue_ready,
}
# 知识向量同步/删除:Task 5 交付了 handler 与写适配器,但先前没有任何生产装配
# 调用它们 —— 事件类型不在上面的白名单里,`OutboxWorker.publish_one` 的
@@ -558,12 +653,19 @@ class WorkerRuntime:
request = AgentRequest(
agent_type=run.agent_type, message=message.content, session_id=run.session_id,
idempotency_key=idem.idempotency_key,
metadata=event.payload.get("metadata", {}) if event else {},
metadata=AgentRequestMetadata.model_validate(
event.payload.get("metadata", {}) if event else {}
),
history=history,
)
identity = RequestContext(user_id=str(run.user_id), trace_id=run.trace_id)
actor_type = (
str(event.payload.get("actor_type", "authenticated"))
if event else "authenticated"
)
# Re-check account and permissions at execution time, including delayed jobs.
context = await self.resolve_identity(identity)
context = await self.restore_context(
actor_type=actor_type, actor_id=str(run.user_id), trace_id=run.trace_id
)
result: AgentResult | None = None
async for event_data in AgentExecutor(self.factory).execute(
request.agent_type, request, context, run_id
@@ -585,19 +687,36 @@ class WorkerRuntime:
async with SessionFactory() as session:
await AgentPersistenceService(session).complete_run(
run_id, result, worker_id=worker_id,
memory_extraction_requested=MemoryService.should_extract_memory(
conversation_content=request.message,
role="user",
# 工具产出的权威事实同样构成持久记忆(工具调用记录来自终态结果)。
tool_result=any(
call.status == "succeeded" for call in result.result.tool_calls
),
# 本 run 落库的业务事件(风险评估完成、交易完成等)。
event_type=business_events[0] if business_events else None,
# 用户明确陈述的偏好/约束/身份/目标,命中才触发抽取。
signals=MemoryService.detect_memory_signals(request.message),
memory_extraction_requested=self.should_request_memory_extraction(
agent_type=run.agent_type, context=context, message=request.message,
result=result, business_events=business_events,
),
profile_candidate_requested=self.should_request_profile_candidate(
agent_type=run.agent_type, context=context, message=request.message,
),
)
await self._append_customer_service_session_memory(
agent_type=run.agent_type, actor_id=str(run.user_id), session_id=run.session_id,
request_message=request.message, response_message=result.result.text,
)
async def _append_customer_service_session_memory(
self, *, agent_type: str, actor_id: str, session_id: str,
request_message: str, response_message: str,
) -> None:
"""成功落库后追加短期会话;Redis 故障不影响主事务。"""
if agent_type != "customer_service":
return
try:
await self.session_memory.append(
actor_id=actor_id, session_id=session_id,
turns=(
CustomerServiceSessionTurn(role="user", content=request_message),
CustomerServiceSessionTurn(role="assistant", content=response_message),
),
)
except Exception:
logger.warning("客服短期会话写入降级,不影响已完成的客服运行", exc_info=True)
async def _failure(
self, run_id: str, worker_id: str, error_code: str, *, retryable: bool