refactor: 完成会话持久化与事务机制重构,清理冗余代码

本次提交是一次大型架构重构,核心变更包括:
1. 调整持久化适配器为无状态实现,通过session_factory按需获取会话
2. 重构事务共享与透传机制,统一使用_session_scope管理会话生命周期
3. 移除OutboxEntry聚合根内的版本自增逻辑,由持久化层统一管理
4. 优化微信插件会话类型枚举与配置对齐
5. 简化出站管道阶段依赖注入与流程逻辑
6. 删除健康检查自动释放会话的冗余代码
7. 重构多个定时任务处理器,移除显式会话工厂创建逻辑
8. 修复会话延迟加载异常问题,新增时区转换工具函数
This commit is contained in:
Kris 2026-07-10 04:10:33 +08:00
parent 41e4385708
commit bb1934023e
36 changed files with 4041 additions and 3858 deletions

View File

@ -6,11 +6,15 @@
装配策略分两类§4.12
- **请求级 / 无状态适配器**14 ``create_driven_adapters(db,
outbox_config, redis_client, arq_pool, execution_port)`` factory 构造
需要共享 ``AsyncSession`` 的适配器``ChannelPersistenceAdapter`` /
``ConversationAdapter`` / ``SqlAlchemyTransactionAdapter``共享同一 ``db``
无状态适配器Logger / Tracer使用模块级全局实例``identity_resolver``
默认为 ``None``由插件注册时注入
session_factory, outbox_config, redis_client, arq_pool, execution_port)``
factory 构造
``ChannelPersistenceAdapter`` 为无状态适配器注入 ``session_factory``
``Callable[[], AsyncSession]``通过 ``_session_scope(tx)`` 按需获取
session``tx`` 非空时复用应用层主事务 session``tx`` ``None`` 时自主
创建并提交``ConversationAdapter`` / ``SqlAlchemyTransactionAdapter``
为请求级事务边界适配器共享同一 ``db`` 会话以保证事务一致性无状态适配器
Logger / Tracer使用模块级全局实例``identity_resolver`` 默认为
``None``由插件注册时注入
- **独立装配路径适配器**4 不经过工厂不进入 ``DrivenAdapters`` 聚合
``ContentReviewRepositoryAdapter`` / ``DefaultContentModerationAdapter``
应用级单例 ``factory.create_host_bootstrap`` 构造并注册 DI
@ -26,7 +30,7 @@
from __future__ import annotations
from collections.abc import Mapping
from collections.abc import Callable, Mapping
from arq import ArqRedis
from redis.asyncio import Redis
@ -72,6 +76,7 @@ _tracer_adapter = InMemoryTracerAdapter(_logger_adapter)
def create_driven_adapters(
db: AsyncSession,
session_factory: Callable[[], AsyncSession],
outbox_config: OutboxConfig,
redis_client: Redis,
arq_pool: ArqRedis,
@ -81,20 +86,25 @@ def create_driven_adapters(
) -> DrivenAdapters:
"""创建被驱动适配器聚合实例。
请求级适配器``ChannelPersistenceAdapter`` / ``ConversationAdapter`` /
``SqlAlchemyTransactionAdapter````
共享同一 ``db`` 会话以保证事务一致性无状态适配器Logger / Tracer
使用模块级全局实例其余适配器按需构造``ConversationAdapter`` 仅依赖
共享 ``db`` 会话合并策略开关决策已上移至 usecase ``ConfigPort`` /
``PersistencePort`` 死依赖已移除``RedisConfigAdapter`` /
``RedisCacheAdapter`` 接受 ``redis_client`` 执行配置读写与缓存操作
``ARQQueueAdapter`` 接受 ``arq_pool`` 入队``AgentRunAdapter`` 接受
``execution_port`` 委托 Agent 运行技术操作构造注入INV-5 / ADP-001 /
ADP-005 / ADP-022``identity_resolver`` 默认为 ``None``由插件注册时注入
``ChannelPersistenceAdapter`` 为无状态适配器注入 ``session_factory``
``Callable[[], AsyncSession]``通过 ``_session_scope(tx)`` 按需获取
session``tx`` 非空时复用应用层主事务 sessionC-I1 透传``tx``
``None`` 时通过 ``session_factory()`` 创建独立 session 并自主提交
``ConversationAdapter`` / ``SqlAlchemyTransactionAdapter`` 为请求级事务
边界适配器共享同一 ``db`` 会话以保证事务一致性无状态适配器
Logger / Tracer使用模块级全局实例其余适配器按需构造
``RedisConfigAdapter`` / ``RedisCacheAdapter`` 接受 ``redis_client`` 执行
配置读写与缓存操作``ARQQueueAdapter`` 接受 ``arq_pool`` 入队
``AgentRunAdapter`` 接受 ``execution_port`` 委托 Agent 运行技术操作
构造注入INV-5 / ADP-001 / ADP-005 / ADP-022``identity_resolver``
默认为 ``None``由插件注册时注入
Args:
db: SQLAlchemy 异步会话由框架层注入persistence / conversation /
transaction 适配器共享此会话
db: SQLAlchemy 异步会话请求级由框架层注入conversation /
transaction 适配器共享此会话以保证事务一致性
session_factory: ``Callable[[], AsyncSession]``注入到
``ChannelPersistenceAdapter`` ``_session_scope(None)`` 创建
独立 session无状态适配器不持有请求级 ``db``
outbox_config: 发件箱配置注入到 ``ChannelPersistenceAdapter`` 替代
直接读取全局 ``app_config``§6.1 应用服务层禁止依赖具体技术适配器
redis_client: ``redis.asyncio.Redis`` 客户端实例注入到
@ -120,7 +130,7 @@ def create_driven_adapters(
key_to_scope_map=key_to_scope_map,
declared_keys=declared_keys,
)
persistence = ChannelPersistenceAdapter(db, outbox_config=outbox_config, logger=_logger_adapter)
persistence = ChannelPersistenceAdapter(session_factory, outbox_config=outbox_config, logger=_logger_adapter)
# 服务账号适配器单例:仅依赖 logger共享注入到 AgentRunAdapter / 管道阶段 / AccountHandler
service_account = ServiceAccountAdapter(logger=_logger_adapter)
# Agent 渠道访问配置适配器单例:仅依赖 logger注入到 agent-run-enqueue 阶段

View File

@ -435,7 +435,7 @@ class AgentRunAdapter(AgentRunPort):
委托 ``execution_port.streamAgentRunEvents`` 输出流式事件并转换为
``StreamEvent`` DTO``current_uid`` 通过
``execution_port.getRunUid`` run 记录查询获取创建 run 时由
``channel_sender_id`` 写入确保 ``streamAgentRunEvents`` 内部的
``service_account.uid`` 写入确保 ``streamAgentRunEvents`` 内部的
权限校验通过
Args:
@ -451,10 +451,13 @@ class AgentRunAdapter(AgentRunPort):
"""
try:
# 通过 execution_port 查询 run 记录的 uid创建时由
# channel_sender_id 写入),确保 streamAgentRunEvents 内部
# service_account.uid 写入),确保 streamAgentRunEvents 内部
# get_run_for_user 权限校验通过。current_uid="" 会导致查不到
# 记录而报错。
current_uid = await self._execution_port.getRunUid(run_id.value, None)
# 记录而报错。需通过 _session_scope 获取独立会话,不能传 None
# 否则 AgentRunRepository(None).get_run() 会在 None.execute()
# 处抛 AttributeError。
async with self._session_scope(None) as session:
current_uid = await self._execution_port.getRunUid(run_id.value, session)
async for event in self._execution_port.streamAgentRunEvents(
run_id=run_id.value,

View File

@ -28,6 +28,8 @@ yuxi.repositories.channels仓储层、sqlalchemy、typing。
from __future__ import annotations
from collections.abc import AsyncIterator, Callable
from contextlib import asynccontextmanager
from datetime import datetime
from typing import TYPE_CHECKING, Any
@ -232,27 +234,46 @@ class ContentReviewRepositoryAdapter(ContentReviewRepositoryPort):
def __init__(
self,
db: AsyncSession,
session_factory: Callable[[], AsyncSession],
logger: LoggerPort,
) -> None:
"""初始化仓储适配器。
"""初始化仓储适配器,注入 session 工厂(不持有 session 实例)。
适配器为无状态协议转换器可安全注册为应用级单例每次方法调用
通过 ``_session_scope(tx)`` 按需获取 session
Args:
db: SQLAlchemy 异步会话所有读写操作复用该会话以保证事务一致性
session_factory: 可调用对象调用返回 ``AsyncSession`` 实例
logger: 日志被驱动端口用于记录故障时的异常信息
"""
self._db = db
self._session_factory = session_factory
self._logger = logger
self._repo = ChannelContentReviewRecordRepository(db)
def _should_commit(self, tx: TransactionContext | None) -> bool:
"""判断是否由适配器自主提交事务。
@asynccontextmanager
async def _session_scope(
self, tx: TransactionContext | None
) -> AsyncIterator[tuple[AsyncSession, ChannelContentReviewRecordRepository, bool]]:
"""统一 session 解析与生命周期管理。
事务边界由应用层控制``tx`` 非空时加入应用层事务
适配器 **不得** 自主提交返回 ``False````tx`` ``None``
按单方法提交返回 ``True``向后兼容
- ``tx`` 非空且 ``get_session()`` 返回 session复用应用层主事务
sessionC-I1 透传``commit=False`` close
- ``tx`` ``None`` 或无 session创建独立 session``commit=True``
异常 rollbackfinally close
"""
return tx is None
if tx is not None:
session = tx.get_session()
if session is not None:
yield session, ChannelContentReviewRecordRepository(session), False
return
session = self._session_factory()
try:
try:
yield session, ChannelContentReviewRecordRepository(session), True
except Exception:
await session.rollback()
raise
finally:
await session.close()
def _translate_db_error(self, exc: Exception, resource: str) -> Error:
"""将数据库异常翻译为契约层错误。
@ -277,43 +298,35 @@ class ContentReviewRepositoryAdapter(ContentReviewRepositoryPort):
``categories`` / ``detail`` tuple 转为 JSON 兼容的 list
事务策略
- ``tx`` 非空时加入应用层事务COOPERATIVE使用共享的
``AsyncSession``**不得** 自主提交
- ``tx`` ``None`` 使用共享会话按单方法提交INDEPENDENT
本端口默认 Strong异常时回滚
- ``tx`` 非空时加入应用层事务COOPERATIVE复用主事务
session**不得** 自主提交
- ``tx`` ``None`` 创建独立 session 按单方法提交
INDEPENDENT本端口默认 Strong异常时回滚
Raises:
ConflictError: ``review_id`` 已存在唯一约束冲突
DependencyError: 数据库故障
"""
commit = self._should_commit(tx)
try:
data = _record_to_orm_data(record)
await self._repo.create(data, commit=commit)
except IntegrityError as exc:
if commit:
await self._db.rollback()
raise self._translate_db_error(exc, "content_review_record") from exc
except SQLAlchemyError as exc:
if commit:
await self._db.rollback()
raise self._translate_db_error(exc, "content_review_record") from exc
except (ConflictError, DependencyError):
if commit:
await self._db.rollback()
raise
except Exception as exc:
if commit:
await self._db.rollback()
await self._logger.error(
"content_review_repository_failed",
resource="content_review_repository",
error=str(exc),
)
raise DependencyError(
"content_review_repository",
Error(str(exc)),
) from exc
data = _record_to_orm_data(record)
async with self._session_scope(tx) as (_, repo, commit):
try:
await repo.create(data, commit=commit)
except IntegrityError as exc:
raise self._translate_db_error(exc, "content_review_record") from exc
except SQLAlchemyError as exc:
raise self._translate_db_error(exc, "content_review_record") from exc
except (ConflictError, DependencyError):
raise
except Exception as exc:
await self._logger.error(
"content_review_repository_failed",
resource="content_review_repository",
error=str(exc),
)
raise DependencyError(
"content_review_repository",
Error(str(exc)),
) from exc
async def queryReviewHistory(
self,
@ -332,35 +345,36 @@ class ContentReviewRepositoryAdapter(ContentReviewRepositoryPort):
Raises:
DependencyError: 数据库故障
"""
try:
orms = await self._repo.list(
channel_type=filter.channel_type if filter.channel_type else None,
account_id=filter.account_id,
verdict=filter.verdict.value if filter.verdict else None,
resource_type=filter.resource_type.value if filter.resource_type else None,
trace_id=filter.trace_id,
start_time=_to_naive_utc(filter.start_time),
end_time=_to_naive_utc(filter.end_time),
limit=limit,
offset=offset,
)
return tuple(_orm_to_history_item(orm) for orm in orms)
except IntegrityError as exc:
raise self._translate_db_error(exc, "content_review_record") from exc
except SQLAlchemyError as exc:
raise self._translate_db_error(exc, "content_review_record") from exc
except (ConflictError, DependencyError):
raise
except Exception as exc:
await self._logger.error(
"content_review_repository_failed",
resource="content_review_repository",
error=str(exc),
)
raise DependencyError(
"content_review_repository",
Error(str(exc)),
) from exc
async with self._session_scope(None) as (_, repo, _):
try:
orms = await repo.list(
channel_type=filter.channel_type if filter.channel_type else None,
account_id=filter.account_id,
verdict=filter.verdict.value if filter.verdict else None,
resource_type=filter.resource_type.value if filter.resource_type else None,
trace_id=filter.trace_id,
start_time=_to_naive_utc(filter.start_time),
end_time=_to_naive_utc(filter.end_time),
limit=limit,
offset=offset,
)
return tuple(_orm_to_history_item(orm) for orm in orms)
except IntegrityError as exc:
raise self._translate_db_error(exc, "content_review_record") from exc
except SQLAlchemyError as exc:
raise self._translate_db_error(exc, "content_review_record") from exc
except (ConflictError, DependencyError):
raise
except Exception as exc:
await self._logger.error(
"content_review_repository_failed",
resource="content_review_repository",
error=str(exc),
)
raise DependencyError(
"content_review_repository",
Error(str(exc)),
) from exc
async def countReviewHistory(
self,
@ -374,32 +388,33 @@ class ContentReviewRepositoryAdapter(ContentReviewRepositoryPort):
Raises:
DependencyError: 数据库故障
"""
try:
return await self._repo.count(
channel_type=filter.channel_type if filter.channel_type else None,
account_id=filter.account_id,
verdict=filter.verdict.value if filter.verdict else None,
resource_type=filter.resource_type.value if filter.resource_type else None,
trace_id=filter.trace_id,
start_time=_to_naive_utc(filter.start_time),
end_time=_to_naive_utc(filter.end_time),
)
except IntegrityError as exc:
raise self._translate_db_error(exc, "content_review_record") from exc
except SQLAlchemyError as exc:
raise self._translate_db_error(exc, "content_review_record") from exc
except (ConflictError, DependencyError):
raise
except Exception as exc:
await self._logger.error(
"content_review_repository_failed",
resource="content_review_repository",
error=str(exc),
)
raise DependencyError(
"content_review_repository",
Error(str(exc)),
) from exc
async with self._session_scope(None) as (_, repo, _):
try:
return await repo.count(
channel_type=filter.channel_type if filter.channel_type else None,
account_id=filter.account_id,
verdict=filter.verdict.value if filter.verdict else None,
resource_type=filter.resource_type.value if filter.resource_type else None,
trace_id=filter.trace_id,
start_time=_to_naive_utc(filter.start_time),
end_time=_to_naive_utc(filter.end_time),
)
except IntegrityError as exc:
raise self._translate_db_error(exc, "content_review_record") from exc
except SQLAlchemyError as exc:
raise self._translate_db_error(exc, "content_review_record") from exc
except (ConflictError, DependencyError):
raise
except Exception as exc:
await self._logger.error(
"content_review_repository_failed",
resource="content_review_repository",
error=str(exc),
)
raise DependencyError(
"content_review_repository",
Error(str(exc)),
) from exc
async def getReviewDetail(
self,
@ -415,27 +430,28 @@ class ContentReviewRepositoryAdapter(ContentReviewRepositoryPort):
Raises:
DependencyError: 数据库故障
"""
try:
orm = await self._repo.get_by_review_id(review_id)
if orm is None:
return None
return _orm_to_detail(orm)
except IntegrityError as exc:
raise self._translate_db_error(exc, "content_review_record") from exc
except SQLAlchemyError as exc:
raise self._translate_db_error(exc, "content_review_record") from exc
except (ConflictError, DependencyError):
raise
except Exception as exc:
await self._logger.error(
"content_review_repository_failed",
resource="content_review_repository",
error=str(exc),
)
raise DependencyError(
"content_review_repository",
Error(str(exc)),
) from exc
async with self._session_scope(None) as (_, repo, _):
try:
orm = await repo.get_by_review_id(review_id)
if orm is None:
return None
return _orm_to_detail(orm)
except IntegrityError as exc:
raise self._translate_db_error(exc, "content_review_record") from exc
except SQLAlchemyError as exc:
raise self._translate_db_error(exc, "content_review_record") from exc
except (ConflictError, DependencyError):
raise
except Exception as exc:
await self._logger.error(
"content_review_repository_failed",
resource="content_review_repository",
error=str(exc),
)
raise DependencyError(
"content_review_repository",
Error(str(exc)),
) from exc
async def getReviewAnalytics(
self,
@ -456,38 +472,39 @@ class ContentReviewRepositoryAdapter(ContentReviewRepositoryPort):
Raises:
DependencyError: 数据库故障
"""
try:
data = await self._repo.query_analytics(
start_time=_to_naive_utc(query.start_time),
end_time=_to_naive_utc(query.end_time),
channel_type=query.channel_type if query.channel_type else None,
granularity=query.granularity,
)
return ContentReviewAnalyticsResult(
total_reviews=data["total_reviews"],
block_count=data["block_count"],
block_rate=data["block_rate"],
by_category=tuple(
CategoryStat(category=item["category"], count=item["count"]) for item in data["by_category"]
),
trend=tuple(data["trend"]),
)
except IntegrityError as exc:
raise self._translate_db_error(exc, "content_review_record") from exc
except SQLAlchemyError as exc:
raise self._translate_db_error(exc, "content_review_record") from exc
except (ConflictError, DependencyError):
raise
except Exception as exc:
await self._logger.error(
"content_review_repository_failed",
resource="content_review_repository",
error=str(exc),
)
raise DependencyError(
"content_review_repository",
Error(str(exc)),
) from exc
async with self._session_scope(None) as (_, repo, _):
try:
data = await repo.query_analytics(
start_time=_to_naive_utc(query.start_time),
end_time=_to_naive_utc(query.end_time),
channel_type=query.channel_type if query.channel_type else None,
granularity=query.granularity,
)
return ContentReviewAnalyticsResult(
total_reviews=data["total_reviews"],
block_count=data["block_count"],
block_rate=data["block_rate"],
by_category=tuple(
CategoryStat(category=item["category"], count=item["count"]) for item in data["by_category"]
),
trend=tuple(data["trend"]),
)
except IntegrityError as exc:
raise self._translate_db_error(exc, "content_review_record") from exc
except SQLAlchemyError as exc:
raise self._translate_db_error(exc, "content_review_record") from exc
except (ConflictError, DependencyError):
raise
except Exception as exc:
await self._logger.error(
"content_review_repository_failed",
resource="content_review_repository",
error=str(exc),
)
raise DependencyError(
"content_review_repository",
Error(str(exc)),
) from exc
async def getReviewStats(
self,
@ -510,52 +527,53 @@ class ContentReviewRepositoryAdapter(ContentReviewRepositoryPort):
Raises:
DependencyError: 数据库故障
"""
try:
data = await self._repo.query_stats(
channel_type=query.channel_type if query.channel_type else None,
account_id=query.account_id,
start_time=_to_naive_utc(query.start_time),
end_time=_to_naive_utc(query.end_time),
granularity=query.granularity,
)
trend = tuple(
ReviewTrendPoint(
timestamp=item["timestamp"],
pass_count=item["pass_count"],
block_count=item["block_count"],
async with self._session_scope(None) as (_, repo, _):
try:
data = await repo.query_stats(
channel_type=query.channel_type if query.channel_type else None,
account_id=query.account_id,
start_time=_to_naive_utc(query.start_time),
end_time=_to_naive_utc(query.end_time),
granularity=query.granularity,
)
for item in data["trend"]
)
return ContentReviewStatsResult(
total_reviews=data["total_reviews"],
pass_count=data["pass_count"],
review_count=data["review_count"],
block_count=data["block_count"],
pass_rate=data["pass_rate"],
block_rate=data["block_rate"],
manual_intervention_rate=data["manual_intervention_rate"],
avg_decision_seconds=data["avg_decision_seconds"],
by_category=tuple(
CategoryStat(category=item["category"], count=item["count"]) for item in data["by_category"]
),
trend=trend,
)
except IntegrityError as exc:
raise self._translate_db_error(exc, "content_review_record") from exc
except SQLAlchemyError as exc:
raise self._translate_db_error(exc, "content_review_record") from exc
except (ConflictError, DependencyError):
raise
except Exception as exc:
await self._logger.error(
"content_review_repository_failed",
resource="content_review_repository",
error=str(exc),
)
raise DependencyError(
"content_review_repository",
Error(str(exc)),
) from exc
trend = tuple(
ReviewTrendPoint(
timestamp=item["timestamp"],
pass_count=item["pass_count"],
block_count=item["block_count"],
)
for item in data["trend"]
)
return ContentReviewStatsResult(
total_reviews=data["total_reviews"],
pass_count=data["pass_count"],
review_count=data["review_count"],
block_count=data["block_count"],
pass_rate=data["pass_rate"],
block_rate=data["block_rate"],
manual_intervention_rate=data["manual_intervention_rate"],
avg_decision_seconds=data["avg_decision_seconds"],
by_category=tuple(
CategoryStat(category=item["category"], count=item["count"]) for item in data["by_category"]
),
trend=trend,
)
except IntegrityError as exc:
raise self._translate_db_error(exc, "content_review_record") from exc
except SQLAlchemyError as exc:
raise self._translate_db_error(exc, "content_review_record") from exc
except (ConflictError, DependencyError):
raise
except Exception as exc:
await self._logger.error(
"content_review_repository_failed",
resource="content_review_repository",
error=str(exc),
)
raise DependencyError(
"content_review_repository",
Error(str(exc)),
) from exc
async def updateReviewVerdict(
self,
@ -577,7 +595,7 @@ class ContentReviewRepositoryAdapter(ContentReviewRepositoryPort):
事务策略
- ``tx`` 非空时加入应用层事务COOPERATIVE**不得** 自主提交
- ``tx`` ``None`` 使用共享会话按单方法提交INDEPENDENT
- ``tx`` ``None`` 创建独立 session 按单方法提交INDEPENDENT
Returns:
更新后的 ``ContentReviewDetail``记录不存在时返回 ``None``
@ -585,38 +603,28 @@ class ContentReviewRepositoryAdapter(ContentReviewRepositoryPort):
Raises:
DependencyError: 数据库故障
"""
commit = self._should_commit(tx)
try:
orm = await self._repo.update_verdict(review_id, verdict, reviewer, commit=commit)
if orm is None:
if commit:
await self._db.rollback()
return None
return _orm_to_detail(orm)
except IntegrityError as exc:
if commit:
await self._db.rollback()
raise self._translate_db_error(exc, "content_review_record") from exc
except SQLAlchemyError as exc:
if commit:
await self._db.rollback()
raise self._translate_db_error(exc, "content_review_record") from exc
except (ConflictError, DependencyError):
if commit:
await self._db.rollback()
raise
except Exception as exc:
if commit:
await self._db.rollback()
await self._logger.error(
"content_review_repository_failed",
resource="content_review_repository",
error=str(exc),
)
raise DependencyError(
"content_review_repository",
Error(str(exc)),
) from exc
async with self._session_scope(tx) as (_, repo, commit):
try:
orm = await repo.update_verdict(review_id, verdict, reviewer, commit=commit)
if orm is None:
return None
return _orm_to_detail(orm)
except IntegrityError as exc:
raise self._translate_db_error(exc, "content_review_record") from exc
except SQLAlchemyError as exc:
raise self._translate_db_error(exc, "content_review_record") from exc
except (ConflictError, DependencyError):
raise
except Exception as exc:
await self._logger.error(
"content_review_repository_failed",
resource="content_review_repository",
error=str(exc),
)
raise DependencyError(
"content_review_repository",
Error(str(exc)),
) from exc
async def deleteOldReviewRecords(
self,
@ -634,29 +642,26 @@ class ContentReviewRepositoryAdapter(ContentReviewRepositoryPort):
Raises:
DependencyError: 数据库故障
"""
try:
return await self._repo.delete_old_records(
_to_naive_utc(before) or before,
limit=limit,
commit=True,
)
except IntegrityError as exc:
await self._db.rollback()
raise self._translate_db_error(exc, "content_review_record") from exc
except SQLAlchemyError as exc:
await self._db.rollback()
raise self._translate_db_error(exc, "content_review_record") from exc
except (ConflictError, DependencyError):
await self._db.rollback()
raise
except Exception as exc:
await self._db.rollback()
await self._logger.error(
"content_review_repository_failed",
resource="content_review_repository",
error=str(exc),
)
raise DependencyError(
"content_review_repository",
Error(str(exc)),
) from exc
async with self._session_scope(None) as (_, repo, commit):
try:
return await repo.delete_old_records(
_to_naive_utc(before) or before,
limit=limit,
commit=commit,
)
except IntegrityError as exc:
raise self._translate_db_error(exc, "content_review_record") from exc
except SQLAlchemyError as exc:
raise self._translate_db_error(exc, "content_review_record") from exc
except (ConflictError, DependencyError):
raise
except Exception as exc:
await self._logger.error(
"content_review_repository_failed",
resource="content_review_repository",
error=str(exc),
)
raise DependencyError(
"content_review_repository",
Error(str(exc)),
) from exc

View File

@ -511,6 +511,12 @@ class ConversationAdapter(ConversationPort):
)
if updated is None:
raise NotFoundError("message", cmd.message_id)
# update_channel_status 用 db.get 加载 Message不预加载
# conversation relationshipcommit 后expire_on_commit=True
# 关系过期_orm_to_message 访问 orm.conversation.channel_type 会
# 触发 lazy-load 抛 MissingGreenletSQLAlchemyError 子类)。需显式
# refresh 加载 relationship与 saveMessage / _createMessage 一致。
await self._db.refresh(updated, ["conversation"])
return self._orm_to_message(updated)
except IntegrityError as exc:
raise self._translate_db_error(exc, "message") from exc

View File

@ -51,6 +51,22 @@ from yuxi.storage.postgres.models_channels import (
from yuxi.utils.crypto import decrypt_sensitive_fields, encrypt_sensitive_fields
def _to_aware_utc(value: dt.datetime | None) -> dt.datetime | None:
"""将 ORM 读出的 naive UTC datetime 转换为 aware UTC datetime。
Postgres ``TIMESTAMP WITHOUT TIME ZONE`` 列由 SQLAlchemy 映射为
offset-naive datetime领域层统一使用 ``datetime.now(UTC)``
offset-aware datetime 运算 helper adaptercore 边界处转换
避免 ``can't subtract offset-naive and offset-aware datetimes`` 等
时区混用错误已是 aware 的值原样返回
"""
if value is None:
return None
if value.tzinfo is None:
return value.replace(tzinfo=dt.UTC)
return value
def _enum[T, V](value: V, cls: type[T], resource: str) -> T:
"""安全构造枚举,将非法值翻译为 DependencyError。
@ -319,19 +335,19 @@ def orm_to_outbox_entry(
durability_policy=_enum(orm.durability_policy, MessageDurabilityPolicy, "outbox_entry"),
retry_count=orm.retry_count,
max_retry=orm.max_retry,
next_retry_at=orm.next_retry_at,
next_retry_at=_to_aware_utc(orm.next_retry_at),
last_error=orm.last_error,
created_at=orm.created_at,
updated_at=orm.updated_at,
expires_at=orm.expires_at,
created_at=_to_aware_utc(orm.created_at),
updated_at=_to_aware_utc(orm.updated_at),
expires_at=_to_aware_utc(orm.expires_at),
channel_msg_id=orm.channel_msg_id,
version=orm.version,
channel_session_id=str(orm.channel_session_id) if orm.channel_session_id is not None else None,
# 聚合视图域 analytics 子域字段(延迟分布 / 漏斗聚合查询)
latency_ms=orm.latency_ms,
funnel_node=orm.funnel_node,
sent_at=orm.sent_at,
last_retry_at=orm.last_retry_at,
sent_at=_to_aware_utc(orm.sent_at),
last_retry_at=_to_aware_utc(orm.last_retry_at),
channel_type=ChannelType(account_orm.channel_type),
# 投递原子语义与降级追踪字段O-02/O-09/O-10/O-11/O-16
idempotency_key=orm.idempotency_key,

View File

@ -2,12 +2,16 @@
实现 ``TransactionPort`` 契约基于 SQLAlchemy ``AsyncSession`` 提供事务
边界控制能力事务边界由应用层管道或用例编排器显式调用被驱动适配
器通过构造时共享的 ``AsyncSession`` 加入同一事务**不得** 自主提交
器通过 ``TransactionContext`` 透传C-I1加入同一事务**不得** 自主提交
事务共享机制``create_driven_adapters(db)`` 将同一 ``AsyncSession``
注入到 ``ChannelPersistenceAdapter`` / ``ConversationAdapter`` /
``SqlAlchemyTransactionAdapter````begin()`` 在共享 session 上开启
事务所有适配器的写操作自动加入
事务共享机制``create_driven_adapters(db, session_factory, ...)`` 将同一
请求级 ``AsyncSession`` 注入到 ``ConversationAdapter`` /
``SqlAlchemyTransactionAdapter``请求级事务边界适配器共享 ``db`` 保证
事务一致性``ChannelPersistenceAdapter`` 为无状态适配器注入
``session_factory``通过 ``_session_scope(tx)`` 按需获取 session
``tx`` 非空时复用 ``tx.get_session()`` 返回的请求级 ``db``加入应用层
事务``tx`` ``None`` 时通过 ``session_factory()`` 创建独立 session
并自主提交
"""
from __future__ import annotations

View File

@ -5,8 +5,9 @@
设计要点
- **会话隔离**通过构造函数注入 ``session_factory``每次执行创建独立 db
会话与适配器避免 worker 进程长生命周期会话问题
- **无状态仓储**通过构造函数注入无状态仓储端口实例适配器内部通过
``_session_scope(tx)`` 按需获取 session``tx`` None 时自主创建并提交
避免长生命周期会话问题
- **批量处理**单次执行最多删除 ``batch_size`` 避免长事务锁竞争
- **分布式锁**``execute`` 通过 ``acquireAdvisoryLock`` 串行化多 worker
并发执行lock_key=``scheduler:{class_name}``, TTL=300s锁获取失败时
@ -19,12 +20,8 @@
from __future__ import annotations
from collections.abc import Callable
from contextlib import AbstractAsyncContextManager
from datetime import timedelta
from sqlalchemy.ext.asyncio import AsyncSession
from yuxi.channels.contract.errors import DependencyError
from yuxi.channels.contract.ports.driven.audit_log_repository_port import (
AuditLogRepositoryPort,
@ -47,11 +44,8 @@ class ChannelAuditLogRetentionHandler:
batch_size: 单次执行最多删除的记录数默认 1000
依赖注入
session_factory: 调用返回 ``AsyncSession`` 的异步上下文管理器工厂
通常为 ``pg_manager.get_async_session_context``每次执行创建独立
会话避免长生命周期会话问题
audit_log_repo_factory: 审计日志仓储端口工厂接收 ``AsyncSession``
返回 ``AuditLogRepositoryPort`` 实例
audit_log_repo: 审计日志仓储端口实例无状态内部通过
``_session_scope(tx)`` 按需获取 session
cache_port: 缓存端口用于获取分布式锁串行化多 worker 并发执行
logger: 日志端口实例
"""
@ -66,13 +60,11 @@ class ChannelAuditLogRetentionHandler:
def __init__(
self,
session_factory: Callable[[], AbstractAsyncContextManager[AsyncSession]],
audit_log_repo_factory: Callable[[AsyncSession], AuditLogRepositoryPort],
audit_log_repo: AuditLogRepositoryPort,
cache_port: CachePort,
logger: LoggerPort,
) -> None:
self._session_factory = session_factory
self._audit_log_repo_factory = audit_log_repo_factory
self._audit_log_repo = audit_log_repo
self._cache_port = cache_port
self._logger = logger
@ -109,9 +101,8 @@ class ChannelAuditLogRetentionHandler:
流程
1. ctx.payload 读取 retention_days batch_size
2. 计算保留截止时间当前时间 - retention_days
3. 通过 session_factory 创建独立 session + 适配器
4. 调用 ``deleteOldAuditLogs`` 物理删除过期日志最多 batch_size
5. 记录结构化日志返回 TaskResult(success=True, output={...})
3. 调用 ``deleteOldAuditLogs`` 物理删除过期日志最多 batch_size
4. 记录结构化日志返回 TaskResult(success=True, output={...})
"""
retention_days: int = ctx.payload.get(
"retention_days",
@ -124,12 +115,10 @@ class ChannelAuditLogRetentionHandler:
cutoff = utc_now_naive() - timedelta(days=retention_days)
try:
async with self._session_factory() as db:
audit_log_repo = self._audit_log_repo_factory(db)
deleted_count = await audit_log_repo.deleteOldAuditLogs(
cutoff,
limit=batch_size,
)
deleted_count = await self._audit_log_repo.deleteOldAuditLogs(
cutoff,
limit=batch_size,
)
if deleted_count > 0:
await self._logger.info(

View File

@ -6,8 +6,9 @@
设计要点
- **会话隔离**通过构造函数注入 ``session_factory``每次执行创建独立 db
会话与适配器避免 worker 进程长生命周期会话问题
- **无状态仓储**通过构造函数注入无状态仓储端口实例适配器内部通过
``_session_scope(tx)`` 按需获取 session``tx`` None 时自主创建并提交
避免长生命周期会话问题
- **批量处理**单次执行最多删除 ``batch_size`` 避免长事务锁竞争
- **分布式锁**``execute`` 通过 ``acquireAdvisoryLock`` 串行化多 worker
并发执行lock_key=``scheduler:{class_name}``, TTL=300s锁获取失败时
@ -20,12 +21,8 @@
from __future__ import annotations
from collections.abc import Callable
from contextlib import AbstractAsyncContextManager
from datetime import timedelta
from sqlalchemy.ext.asyncio import AsyncSession
from yuxi.channels.contract.errors import DependencyError
from yuxi.channels.contract.ports.driven.cache_port import CachePort
from yuxi.channels.contract.ports.driven.content_review_repository_port import (
@ -49,11 +46,8 @@ class ChannelContentReviewRetentionHandler:
batch_size: 单次执行最多删除的记录数默认 1000
依赖注入
session_factory: 调用返回 ``AsyncSession`` 的异步上下文管理器工厂
通常为 ``pg_manager.get_async_session_context``每次执行创建独立
会话避免长生命周期会话问题
content_review_repo_factory: 内容审核仓储端口工厂接收 ``AsyncSession``
返回 ``ContentReviewRepositoryPort`` 实例
content_review_repo: 内容审核仓储端口实例无状态内部通过
``_session_scope(tx)`` 按需获取 session
cache_port: 缓存端口用于获取分布式锁串行化多 worker 并发执行
logger: 日志端口实例
"""
@ -68,13 +62,11 @@ class ChannelContentReviewRetentionHandler:
def __init__(
self,
session_factory: Callable[[], AbstractAsyncContextManager[AsyncSession]],
content_review_repo_factory: Callable[[AsyncSession], ContentReviewRepositoryPort],
content_review_repo: ContentReviewRepositoryPort,
cache_port: CachePort,
logger: LoggerPort,
) -> None:
self._session_factory = session_factory
self._content_review_repo_factory = content_review_repo_factory
self._content_review_repo = content_review_repo
self._cache_port = cache_port
self._logger = logger
@ -111,9 +103,8 @@ class ChannelContentReviewRetentionHandler:
流程
1. ctx.payload 读取 retention_days batch_size
2. 计算保留截止时间当前时间 - retention_days
3. 通过 session_factory 创建独立 session + 适配器
4. 调用 ``deleteOldReviewRecords`` 物理删除过期记录最多 batch_size
5. 记录结构化日志返回 TaskResult(success=True, output={...})
3. 调用 ``deleteOldReviewRecords`` 物理删除过期记录最多 batch_size
4. 记录结构化日志返回 TaskResult(success=True, output={...})
"""
retention_days: int = ctx.payload.get(
"retention_days",
@ -126,12 +117,10 @@ class ChannelContentReviewRetentionHandler:
cutoff = utc_now_naive() - timedelta(days=retention_days)
try:
async with self._session_factory() as db:
review_repo = self._content_review_repo_factory(db)
deleted_count = await review_repo.deleteOldReviewRecords(
cutoff,
limit=batch_size,
)
deleted_count = await self._content_review_repo.deleteOldReviewRecords(
cutoff,
limit=batch_size,
)
if deleted_count > 0:
await self._logger.info(

View File

@ -6,8 +6,9 @@
设计要点
- **会话隔离**通过构造函数注入 ``session_factory``每次执行创建独立 db
会话与适配器避免 worker 进程长生命周期会话问题
- **无状态仓储**通过构造函数注入无状态仓储端口实例适配器内部通过
``_session_scope(tx)`` 按需获取 session``tx`` None 时自主创建并提交
避免长生命周期会话问题
- **分布式锁**``execute`` 通过 ``acquireAdvisoryLock`` 串行化多 worker
并发执行lock_key=``scheduler:{class_name}``, TTL=300s锁获取失败时
skip返回 ``TaskResult(skipped=True)`` fail-closed下一周期重试
@ -19,11 +20,6 @@
from __future__ import annotations
from collections.abc import Callable
from contextlib import AbstractAsyncContextManager
from sqlalchemy.ext.asyncio import AsyncSession
from yuxi.channels.contract.errors import DependencyError
from yuxi.channels.contract.ports.driven.cache_port import CachePort
from yuxi.channels.contract.ports.driven.idempotency_repository_port import (
@ -42,11 +38,8 @@ class ChannelIdempotencyCleanupHandler:
``system:channel_idempotency_cleanup`` 任务配置为每小时执行一次
依赖注入
session_factory: 调用返回 ``AsyncSession`` 的异步上下文管理器工厂
通常为 ``pg_manager.get_async_session_context``每次执行创建独立
会话避免长生命周期会话问题
idempotency_repo_factory: 幂等记录仓储端口工厂接收 ``AsyncSession``
返回 ``IdempotencyRepositoryPort`` 实例
idempotency_repo: 幂等记录仓储端口实例无状态内部通过
``_session_scope(tx)`` 按需获取 session
cache_port: 缓存端口用于获取分布式锁串行化多 worker 并发执行
logger: 日志端口实例
"""
@ -58,13 +51,11 @@ class ChannelIdempotencyCleanupHandler:
def __init__(
self,
session_factory: Callable[[], AbstractAsyncContextManager[AsyncSession]],
idempotency_repo_factory: Callable[[AsyncSession], IdempotencyRepositoryPort],
idempotency_repo: IdempotencyRepositoryPort,
cache_port: CachePort,
logger: LoggerPort,
) -> None:
self._session_factory = session_factory
self._idempotency_repo_factory = idempotency_repo_factory
self._idempotency_repo = idempotency_repo
self._cache_port = cache_port
self._logger = logger
@ -100,15 +91,12 @@ class ChannelIdempotencyCleanupHandler:
流程
1. 以当前时间作为截止时间
2. 通过 session_factory 创建独立 session + 适配器
3. 调用 ``deleteExpiredRecords`` 物理删除 ``expires_at`` 已过期记录
4. 记录结构化日志返回 TaskResult(success=True, output={...})
2. 调用 ``deleteExpiredRecords`` 物理删除 ``expires_at`` 已过期记录
3. 记录结构化日志返回 TaskResult(success=True, output={...})
"""
now = utc_now_naive()
try:
async with self._session_factory() as db:
idempotency_repo = self._idempotency_repo_factory(db)
deleted_count = await idempotency_repo.deleteExpiredRecords(now)
deleted_count = await self._idempotency_repo.deleteExpiredRecords(now)
if deleted_count > 0:
await self._logger.info(

View File

@ -10,8 +10,9 @@
设计要点
- **会话隔离**通过构造函数注入 ``session_factory``每次执行创建独立 db
会话与适配器避免 worker 进程长生命周期会话问题
- **无状态仓储**通过构造函数注入无状态仓储端口实例适配器内部通过
``_session_scope(tx)`` 按需获取 session``tx`` None 时自主创建并提交
避免长生命周期会话问题
- **批量处理**单次执行最多处理 ``batch_size`` 避免长事务锁竞争
- **分布式锁**``execute`` 通过 ``acquireAdvisoryLock`` 串行化多 worker
并发执行lock_key=``scheduler:{class_name}``, TTL=300s锁获取失败时
@ -26,12 +27,8 @@
from __future__ import annotations
from collections.abc import Callable
from contextlib import AbstractAsyncContextManager
from datetime import datetime, timedelta
from sqlalchemy.ext.asyncio import AsyncSession
from yuxi.channels.contract.dtos.audit import AuditOperationType
from yuxi.channels.contract.dtos.outbox import OutboxConfig, OutboxEntry, OutboxStatus
from yuxi.channels.contract.dtos.persistence import SaveAuditLogCmd
@ -65,13 +62,9 @@ class ChannelOutboxRecoveryHandler:
默认 300
依赖注入
session_factory: 调用返回 ``AsyncSession`` 的异步上下文管理器工厂
通常为 ``pg_manager.get_async_session_context``每次执行创建独立
会话避免长生命周期会话问题
outbox_repo_factory: outbox 仓储端口工厂接收 ``AsyncSession`` 返回
``OutboxRepositoryPort`` 实例
audit_log_repo_factory: 审计日志仓储端口工厂接收 ``AsyncSession``
返回 ``AuditLogRepositoryPort`` 实例
outbox_repo: outbox 仓储端口实例无状态内部通过
``_session_scope(tx)`` 按需获取 session
audit_log_repo: 审计日志仓储端口实例无状态同上
queue_port: 队列被驱动端口用于入队 ARQ 重试任务
outbox_config: 发件箱配置提供 ``ttl_seconds`` / ``max_retry`` /
``retry_backoff_schedule`` 等参数
@ -93,17 +86,15 @@ class ChannelOutboxRecoveryHandler:
def __init__(
self,
session_factory: Callable[[], AbstractAsyncContextManager[AsyncSession]],
outbox_repo_factory: Callable[[AsyncSession], OutboxRepositoryPort],
audit_log_repo_factory: Callable[[AsyncSession], AuditLogRepositoryPort],
outbox_repo: OutboxRepositoryPort,
audit_log_repo: AuditLogRepositoryPort,
queue_port: QueuePort,
outbox_config: OutboxConfig,
cache_port: CachePort,
logger: LoggerPort,
) -> None:
self._session_factory = session_factory
self._outbox_repo_factory = outbox_repo_factory
self._audit_log_repo_factory = audit_log_repo_factory
self._outbox_repo = outbox_repo
self._audit_log_repo = audit_log_repo
self._queue = queue_port
self._outbox_config = outbox_config
self._cache_port = cache_port
@ -143,12 +134,11 @@ class ChannelOutboxRecoveryHandler:
流程
1. ctx.payload 读取 batch_size / scan_ttl_seconds /
sent_unconfirmed_timeout_seconds
2. 通过 session_factory 创建独立 session + 适配器
3. 扫描 PENDING 条目死信标记 DEAD其余入队重试
4. 扫描 SENT_UNCONFIRMED 条目死信标记 DEAD无回执超时入队重试
5. 扫描 FAILED 条目O-04未达最大重试且已到重试时间则重置为
2. 扫描 PENDING 条目死信标记 DEAD其余入队重试
3. 扫描 SENT_UNCONFIRMED 条目死信标记 DEAD无回执超时入队重试
4. 扫描 FAILED 条目O-04未达最大重试且已到重试时间则重置为
PENDING 并入队满足死信条件标记 DEAD
6. 记录结构化日志返回 TaskResult(success=True, output={...})
5. 记录结构化日志返回 TaskResult(success=True, output={...})
"""
batch_size: int = ctx.payload.get(
"batch_size",
@ -168,51 +158,47 @@ class ChannelOutboxRecoveryHandler:
processed_ids: set[str] = set()
try:
async with self._session_factory() as db:
outbox_repo = self._outbox_repo_factory(db)
audit_log_repo = self._audit_log_repo_factory(db)
await self._scan_pending(
outbox_repo,
audit_log_repo,
scan_start,
processed_ids,
batch_size,
scan_ttl_seconds,
await self._scan_pending(
self._outbox_repo,
self._audit_log_repo,
scan_start,
processed_ids,
batch_size,
scan_ttl_seconds,
)
if self._is_scan_expired(scan_start, scan_ttl_seconds):
return TaskResult(
success=True,
output={
"processed_count": len(processed_ids),
"reason": "scan_ttl_reached_after_pending",
},
)
if self._is_scan_expired(scan_start, scan_ttl_seconds):
return TaskResult(
success=True,
output={
"processed_count": len(processed_ids),
"reason": "scan_ttl_reached_after_pending",
},
)
await self._scan_sent_unconfirmed(
outbox_repo,
audit_log_repo,
scan_start,
processed_ids,
batch_size,
scan_ttl_seconds,
sent_unconfirmed_timeout,
)
if self._is_scan_expired(scan_start, scan_ttl_seconds):
return TaskResult(
success=True,
output={
"processed_count": len(processed_ids),
"reason": "scan_ttl_reached_after_sent_unconfirmed",
},
)
await self._scan_failed(
outbox_repo,
audit_log_repo,
scan_start,
processed_ids,
batch_size,
scan_ttl_seconds,
await self._scan_sent_unconfirmed(
self._outbox_repo,
self._audit_log_repo,
scan_start,
processed_ids,
batch_size,
scan_ttl_seconds,
sent_unconfirmed_timeout,
)
if self._is_scan_expired(scan_start, scan_ttl_seconds):
return TaskResult(
success=True,
output={
"processed_count": len(processed_ids),
"reason": "scan_ttl_reached_after_sent_unconfirmed",
},
)
await self._scan_failed(
self._outbox_repo,
self._audit_log_repo,
scan_start,
processed_ids,
batch_size,
scan_ttl_seconds,
)
await self._logger.info(
"outbox 恢复扫描完成",

View File

@ -6,8 +6,9 @@ SENT 条目保留 30 天(默认),保留期内可供审计与排查。
设计要点
- **会话隔离**通过构造函数注入 ``session_factory``每次执行创建独立 db
会话与适配器避免 worker 进程长生命周期会话问题
- **无状态仓储**通过构造函数注入无状态仓储端口实例适配器内部通过
``_session_scope(tx)`` 按需获取 session``tx`` None 时自主创建并提交
避免长生命周期会话问题
- **批量处理**单次执行最多清理 ``batch_size`` 避免长事务锁竞争
- **分布式锁**``execute`` 通过 ``acquireAdvisoryLock`` 串行化多 worker
并发执行lock_key=``scheduler:{class_name}``, TTL=300s锁获取失败时
@ -23,12 +24,8 @@ SENT 条目保留 30 天(默认),保留期内可供审计与排查。
from __future__ import annotations
from collections.abc import Callable
from contextlib import AbstractAsyncContextManager
from datetime import UTC, datetime, timedelta
from sqlalchemy.ext.asyncio import AsyncSession
from yuxi.channels.contract.dtos.event import OutboxEntryPurgedEvent
from yuxi.channels.contract.errors import DependencyError
from yuxi.channels.contract.ports.driven.cache_port import CachePort
@ -55,12 +52,9 @@ class ChannelOutboxTerminalCleanupHandler:
batch_size: 单次执行每种状态最多清理的条目数默认 1000
依赖注入
session_factory: 调用返回 ``AsyncSession`` 的异步上下文管理器工厂
通常为 ``pg_manager.get_async_session_context``每次执行创建独立
会话避免长生命周期会话问题
outbox_repo_factory: outbox 仓储端口工厂接收 ``AsyncSession`` 返回
``OutboxRepositoryPort`` 实例 infrastructure 层装配具体适配器
``SqlAlchemyOutboxRepositoryAdapter``应用层不直接依赖
outbox_repo: outbox 仓储端口实例 infrastructure 层装配具体适配器
``ChannelPersistenceAdapter``无状态内部通过
``_session_scope(tx)`` 按需获取 session应用层不直接依赖
适配器实现INV-1 依赖方向合规
event_publisher: 事件发布端口实例或 Noneworker 进程未注入时为 None
cache_port: 缓存端口用于获取分布式锁串行化多 worker 并发执行
@ -74,14 +68,12 @@ class ChannelOutboxTerminalCleanupHandler:
def __init__(
self,
session_factory: Callable[[], AbstractAsyncContextManager[AsyncSession]],
outbox_repo_factory: Callable[[AsyncSession], OutboxRepositoryPort],
outbox_repo: OutboxRepositoryPort,
event_publisher: EventPublisherPort | None,
cache_port: CachePort,
logger: LoggerPort,
) -> None:
self._session_factory = session_factory
self._outbox_repo_factory = outbox_repo_factory
self._outbox_repo = outbox_repo
self._event_publisher = event_publisher
self._cache_port = cache_port
self._logger = logger
@ -121,11 +113,10 @@ class ChannelOutboxTerminalCleanupHandler:
1. ctx.payload 读取 dead_retention_hours / sent_retention_hours /
batch_size
2. 计算 dead_before sent_before 截止时间
3. 通过 session_factory 创建独立 session + 适配器
4. 调用 cleanupOldDeadEntries 清理 DEAD 条目
5. 调用 cleanupOldSentEntries 清理 SENT 条目
6. best-effort 发布 OutboxEntryPurgedEventevent_publisher 非空时
7. 记录结构化日志返回 TaskResult(success=True, output={...})
3. 调用 cleanupOldDeadEntries 清理 DEAD 条目
4. 调用 cleanupOldSentEntries 清理 SENT 条目
5. best-effort 发布 OutboxEntryPurgedEventevent_publisher 非空时
6. 记录结构化日志返回 TaskResult(success=True, output={...})
"""
dead_retention_hours: int = ctx.payload.get("dead_retention_hours", 168)
sent_retention_hours: int = ctx.payload.get("sent_retention_hours", 720)
@ -135,17 +126,14 @@ class ChannelOutboxTerminalCleanupHandler:
sent_before = datetime.now(UTC) - timedelta(hours=sent_retention_hours)
try:
async with self._session_factory() as db:
outbox_repo = self._outbox_repo_factory(db)
dead_ids = await outbox_repo.cleanupOldDeadEntries(
before=dead_before,
limit=batch_size,
)
sent_ids = await outbox_repo.cleanupOldSentEntries(
before=sent_before,
limit=batch_size,
)
dead_ids = await self._outbox_repo.cleanupOldDeadEntries(
before=dead_before,
limit=batch_size,
)
sent_ids = await self._outbox_repo.cleanupOldSentEntries(
before=sent_before,
limit=batch_size,
)
# best-effort 发布事件,携带被清理条目的完整 outbox_id 列表
if self._event_publisher is not None:

View File

@ -6,8 +6,9 @@ PENDING 状态的配对记录,通过聚合根方法 ``PairingApproval.expireIf
设计要点
- **会话隔离**通过构造函数注入 ``session_factory``每次执行创建独立 db
会话与适配器避免 worker 进程长生命周期会话问题
- **无状态仓储**通过构造函数注入无状态仓储端口实例适配器内部通过
``_session_scope(tx)`` 按需获取 session``tx`` None 时自主创建并提交
避免长生命周期会话问题
- **批量处理**单次执行最多处理 ``batch_size`` 避免长事务锁竞争
- **分布式锁**``execute`` 通过 ``acquireAdvisoryLock`` 串行化多 worker
并发执行lock_key=``scheduler:{class_name}``, TTL=300s锁获取失败时
@ -20,11 +21,6 @@ PENDING 状态的配对记录,通过聚合根方法 ``PairingApproval.expireIf
from __future__ import annotations
from collections.abc import Callable
from contextlib import AbstractAsyncContextManager
from sqlalchemy.ext.asyncio import AsyncSession
from yuxi.channels.contract.dtos.pairing import PairingStatus
from yuxi.channels.contract.errors import DependencyError
from yuxi.channels.contract.ports.driven.cache_port import CachePort
@ -48,11 +44,8 @@ class ChannelPairingExpirationHandler:
batch_size: 单次执行最多扫描的记录数默认 100
依赖注入
session_factory: 调用返回 ``AsyncSession`` 的异步上下文管理器工厂
通常为 ``pg_manager.get_async_session_context``每次执行创建独立
会话避免长生命周期会话问题
pairing_repo_factory: 配对仓储端口工厂接收 ``AsyncSession`` 返回
``PairingRepositoryPort`` 实例
pairing_repo: 配对仓储端口实例无状态内部通过
``_session_scope(tx)`` 按需获取 session
cache_port: 缓存端口用于获取分布式锁串行化多 worker 并发执行
logger: 日志端口实例
"""
@ -66,13 +59,11 @@ class ChannelPairingExpirationHandler:
def __init__(
self,
session_factory: Callable[[], AbstractAsyncContextManager[AsyncSession]],
pairing_repo_factory: Callable[[AsyncSession], PairingRepositoryPort],
pairing_repo: PairingRepositoryPort,
cache_port: CachePort,
logger: LoggerPort,
) -> None:
self._session_factory = session_factory
self._pairing_repo_factory = pairing_repo_factory
self._pairing_repo = pairing_repo
self._cache_port = cache_port
self._logger = logger
@ -109,11 +100,10 @@ class ChannelPairingExpirationHandler:
流程
1. ctx.payload 读取 batch_size
2. 通过 session_factory 创建独立 session + 适配器
3. 拉取 ``expires_at`` 早于当前时间的 PENDING 配对
4. 重建 ``PairingApproval`` 聚合根并调用 ``expireIfOverdue``
5. 返回 ``True`` 时持久化为 EXPIRED
6. 记录结构化日志返回 TaskResult(success=True, output={...})
2. 拉取 ``expires_at`` 早于当前时间的 PENDING 配对
3. 重建 ``PairingApproval`` 聚合根并调用 ``expireIfOverdue``
4. 返回 ``True`` 时持久化为 EXPIRED
5. 记录结构化日志返回 TaskResult(success=True, output={...})
"""
batch_size: int = ctx.payload.get(
"batch_size",
@ -122,37 +112,35 @@ class ChannelPairingExpirationHandler:
now = utc_now_naive()
try:
async with self._session_factory() as db:
pairing_repo = self._pairing_repo_factory(db)
records = await pairing_repo.listExpiredPendingPairings(
before=now,
limit=batch_size,
records = await self._pairing_repo.listExpiredPendingPairings(
before=now,
limit=batch_size,
)
if not records:
return TaskResult(
success=True,
output={"expired_count": 0, "scanned_count": 0},
)
if not records:
return TaskResult(
success=True,
output={"expired_count": 0, "scanned_count": 0},
)
expired_count = 0
for record in records:
try:
approval = PairingApproval.fromRecord(record)
if approval.expireIfOverdue():
await pairing_repo.updatePairingStatus(
record.pairing_id,
PairingStatus.EXPIRED,
expected_version=record.version,
updated_by="system",
)
expired_count += 1
except Exception as exc:
await self._logger.exception(
"配对过期标记失败,跳过",
exc_info=exc,
pairing_id=record.pairing_id,
expired_count = 0
for record in records:
try:
approval = PairingApproval.fromRecord(record)
if approval.expireIfOverdue():
await self._pairing_repo.updatePairingStatus(
record.pairing_id,
PairingStatus.EXPIRED,
expected_version=record.version,
updated_by="system",
)
expired_count += 1
except Exception as exc:
await self._logger.exception(
"配对过期标记失败,跳过",
exc_info=exc,
pairing_id=record.pairing_id,
)
if expired_count > 0:
await self._logger.info(

View File

@ -6,8 +6,9 @@
设计要点
- **会话隔离**通过构造函数注入 ``session_factory``每次执行创建独立 db
会话与适配器避免 worker 进程长生命周期会话问题
- **无状态仓储**通过构造函数注入无状态仓储端口实例适配器内部通过
``_session_scope(tx)`` 按需获取 session``tx`` None 时自主创建并提交
避免长生命周期会话问题
- **批量处理**单次执行最多清理 ``batch_size`` 避免长事务锁竞争
- **分布式锁**``execute`` 通过 ``acquireAdvisoryLock`` 串行化多 worker
并发执行lock_key=``scheduler:{class_name}``, TTL=300s锁获取失败时
@ -20,12 +21,8 @@
from __future__ import annotations
from collections.abc import Callable
from contextlib import AbstractAsyncContextManager
from datetime import UTC, datetime, timedelta
from sqlalchemy.ext.asyncio import AsyncSession
from yuxi.channels.contract.errors import DependencyError
from yuxi.channels.contract.ports.driven.cache_port import CachePort
from yuxi.channels.contract.ports.driven.logger_port import LoggerPort
@ -47,12 +44,9 @@ class ChannelPairingTerminalCleanupHandler:
batch_size: 单次执行最多清理的记录数默认 500
依赖注入
session_factory: 调用返回 ``AsyncSession`` 的异步上下文管理器工厂
通常为 ``pg_manager.get_async_session_context``每次执行创建独立
会话避免长生命周期会话问题
pairing_repo_factory: 配对仓储端口工厂接收 ``AsyncSession`` 返回
``PairingRepositoryPort`` 实例 infrastructure 层装配具体适配器
``SqlAlchemyPairingRepositoryAdapter``应用层不直接依赖
pairing_repo: 配对仓储端口实例 infrastructure 层装配具体适配器
``ChannelPersistenceAdapter``无状态内部通过
``_session_scope(tx)`` 按需获取 session应用层不直接依赖
适配器实现INV-1 依赖方向合规
cache_port: 缓存端口用于获取分布式锁串行化多 worker 并发执行
logger: 日志端口实例
@ -65,13 +59,11 @@ class ChannelPairingTerminalCleanupHandler:
def __init__(
self,
session_factory: Callable[[], AbstractAsyncContextManager[AsyncSession]],
pairing_repo_factory: Callable[[AsyncSession], PairingRepositoryPort],
pairing_repo: PairingRepositoryPort,
cache_port: CachePort,
logger: LoggerPort,
) -> None:
self._session_factory = session_factory
self._pairing_repo_factory = pairing_repo_factory
self._pairing_repo = pairing_repo
self._cache_port = cache_port
self._logger = logger
@ -108,9 +100,8 @@ class ChannelPairingTerminalCleanupHandler:
流程
1. ctx.payload 读取 retention_days batch_size
2. 计算 before 截止时间
3. 通过 session_factory 创建独立 session + 适配器
4. 调用 cleanupOldTerminalPairings 清理终态配对记录
5. 记录结构化日志返回 TaskResult(success=True, output={...})
3. 调用 cleanupOldTerminalPairings 清理终态配对记录
4. 记录结构化日志返回 TaskResult(success=True, output={...})
"""
retention_days: int = ctx.payload.get("retention_days", 30)
batch_size: int = ctx.payload.get("batch_size", 500)
@ -118,13 +109,10 @@ class ChannelPairingTerminalCleanupHandler:
before = datetime.now(UTC) - timedelta(days=retention_days)
try:
async with self._session_factory() as db:
pairing_repo = self._pairing_repo_factory(db)
cleaned_count = await pairing_repo.cleanupOldTerminalPairings(
before=before,
limit=batch_size,
)
cleaned_count = await self._pairing_repo.cleanupOldTerminalPairings(
before=before,
limit=batch_size,
)
await self._logger.info(
"终态配对记录物理清理完成",

View File

@ -6,8 +6,9 @@
设计要点
- **会话隔离**通过构造函数注入 ``session_factory``每次执行创建独立 db
会话与适配器避免 worker 进程长生命周期会话问题
- **无状态仓储**通过构造函数注入无状态仓储端口实例适配器内部通过
``_session_scope(tx)`` 按需获取 session``tx`` None 时自主创建并提交
避免长生命周期会话问题
- **批量处理**单次执行最多清理 ``batch_size`` 避免长事务锁竞争
- **分布式锁**``execute`` 通过 ``acquireAdvisoryLock`` 串行化多 worker
并发执行lock_key=``scheduler:{class_name}``, TTL=300s锁获取失败时
@ -20,12 +21,8 @@
from __future__ import annotations
from collections.abc import Callable
from contextlib import AbstractAsyncContextManager
from datetime import UTC, datetime, timedelta
from sqlalchemy.ext.asyncio import AsyncSession
from yuxi.channels.contract.errors import DependencyError
from yuxi.channels.contract.ports.driven.cache_port import CachePort
from yuxi.channels.contract.ports.driven.channel_session_repository_port import (
@ -50,13 +47,10 @@ class ChannelSessionInactiveCleanupHandler:
batch_size: 单次执行最多清理的会话数默认 100
依赖注入
session_factory: 调用返回 ``AsyncSession`` 的异步上下文管理器工厂
通常为 ``pg_manager.get_async_session_context``每次执行创建独立
会话避免长生命周期会话问题
session_repo_factory: 渠道会话仓储端口工厂接收 ``AsyncSession`` 返回
``ChannelSessionRepositoryPort`` 实例 infrastructure 层装配具体
适配器 ``SqlAlchemyChannelSessionRepositoryAdapter``应用层
不直接依赖适配器实现INV-1 依赖方向合规
session_repo: 渠道会话仓储端口实例 infrastructure 层装配具体
适配器 ``ChannelPersistenceAdapter``无状态内部通过
``_session_scope(tx)`` 按需获取 session应用层不直接依赖
适配器实现INV-1 依赖方向合规
cache_port: 缓存端口用于获取分布式锁串行化多 worker 并发执行
logger: 日志端口实例
"""
@ -68,14 +62,12 @@ class ChannelSessionInactiveCleanupHandler:
def __init__(
self,
session_factory: Callable[[], AbstractAsyncContextManager[AsyncSession]],
session_repo_factory: Callable[[AsyncSession], ChannelSessionRepositoryPort],
session_repo: ChannelSessionRepositoryPort,
cache_port: CachePort,
logger: LoggerPort,
realtime_metrics: RealtimeMetricsPort | None = None,
) -> None:
self._session_factory = session_factory
self._session_repo_factory = session_repo_factory
self._session_repo = session_repo
self._cache_port = cache_port
self._logger = logger
self._realtime_metrics = realtime_metrics
@ -113,10 +105,10 @@ class ChannelSessionInactiveCleanupHandler:
流程
1. ctx.payload 读取 inactive_threshold_minutes batch_size
2. 计算 inactive_before 截止时间
3. 通过 session_factory 创建独立 session + 适配器
4. 调用 listInactiveTemporarySessions 查询待清理会话
5. 无会话时返回 successcount=0
6. 提取 session_id 列表调用 cleanupInactiveSessions 批量关闭
3. 调用 listInactiveTemporarySessions 查询待清理会话
4. 无会话时返回 successcount=0
5. 提取 session_id 列表调用 cleanupInactiveSessions 批量关闭
6. 释放实时指标中的活跃会话DSB-REALTIME
7. 记录结构化日志返回 TaskResult(success=True, output={...})
"""
inactive_threshold_minutes: int = ctx.payload.get("inactive_threshold_minutes", 60)
@ -125,24 +117,21 @@ class ChannelSessionInactiveCleanupHandler:
inactive_before = datetime.now(UTC) - timedelta(minutes=inactive_threshold_minutes)
try:
async with self._session_factory() as db:
session_repo = self._session_repo_factory(db)
sessions = await session_repo.listInactiveTemporarySessions(
inactive_before=inactive_before,
limit=batch_size,
sessions = await self._session_repo.listInactiveTemporarySessions(
inactive_before=inactive_before,
limit=batch_size,
)
if not sessions:
return TaskResult(
success=True,
output={
"cleaned_count": 0,
"inactive_threshold_minutes": inactive_threshold_minutes,
},
)
if not sessions:
return TaskResult(
success=True,
output={
"cleaned_count": 0,
"inactive_threshold_minutes": inactive_threshold_minutes,
},
)
session_ids = [s.session_id for s in sessions]
cleaned_count = await session_repo.cleanupInactiveSessions(session_ids)
session_ids = [s.session_id for s in sessions]
cleaned_count = await self._session_repo.cleanupInactiveSessions(session_ids)
# 释放实时指标中的活跃会话DSB-REALTIME避免清理后的会话
# 仍计入 active_sessions / active_conversations。

View File

@ -292,12 +292,6 @@ class HealthAggregator:
degraded_component_details=self._buildDegradedComponentDetails([self.HEALTH_CHECK_TIMEOUT_COMPONENT]),
status_message="系统部分降级:健康检查超时,聚合结果可能不完整,建议稍后重试或导出诊断包排查。",
)
finally:
# 健康检查在应用级共享 session 上执行 ping / getConnectionPoolStatus /
# listChannelAccountsSQLAlchemy 2.0 autobegin 会开启隐式事务并
# 持有连接直至显式提交。每轮检查结束(含超时/异常路径)调用
# releaseSession 归还连接,避免健康检查反复触发导致连接池耗尽。
await self._releasePersistenceSession()
if not query.with_probe:
# 缓存写入前序列化为 JSON-safe dict避免 dataclass 被
@ -311,30 +305,6 @@ class HealthAggregator:
return snapshot
async def _releasePersistenceSession(self) -> None:
"""释放持久化 session 占用的连接。
``persistence_port`` 实际注入的是 ``PersistencePort``同时满足
``PersistenceHealthPort`` ``ChannelAccountRepositoryPort``
调用 ``releaseSession`` 归还应用级共享 ``AsyncSession`` 占用的连接
SQLAlchemy 2.0 autobegin 隐式事务在只读查询后持有连接
释放失败仅记 WARN 日志不向上抛出健康检查不应因连接归还失败
而中断返回快照``persistence_port`` 未实现 ``releaseSession``
no-op测试场景注入的 mock 可能不实现该方法
"""
releaser = getattr(self._persistence_port, "releaseSession", None)
if releaser is None:
return
try:
await releaser()
except Exception as exc:
if self._logger is not None:
await self._logger.warn(
"健康检查释放会话连接失败",
error=str(exc),
)
async def _aggregate(self, query: HealthQuery) -> HealthSnapshot:
"""执行健康聚合逻辑(不含缓存读写,由 ``checkHealth`` 包裹超时)。

View File

@ -494,7 +494,6 @@ class DeliverStage:
latest_aggregate.partial_failure = True
latest_aggregate.last_error = f"partial failure on parts: {failed_parts}"
latest_aggregate.updated_at = datetime.now(UTC)
latest_aggregate.version += 1
await self.persistence_port.updateOutboxEntry(
latest_aggregate,
expected_status=OutboxStatus.SENT_UNCONFIRMED,

View File

@ -1,18 +1,23 @@
"""出站管道 load-build 阶段。
从上下文构建 OutboundPayload关联 Agent 运行 ID 并携带流式分块供下游
阶段格式化与投递本阶段无领域服务调用仅做数据装配
阶段格式化与投递持久化模式下AgentRun 响应内容由本阶段通过阻塞消费
streamAgentRun 事件流加载stream_chunk_stage 在持久化模式下被跳过
"""
from __future__ import annotations
from collections.abc import Callable
from contextlib import aclosing
from typing import Any
from yuxi.channels.application.context.outbound_context import OutboundContext
from yuxi.channels.contract.dtos.agent_run import AgentRunId
from yuxi.channels.contract.dtos.common import Attachment
from yuxi.channels.contract.dtos.option import Some
from yuxi.channels.contract.dtos.outbound import OutboundPayload, RichMessageFields
from yuxi.channels.contract.plugin.extension_point import FailureStrategy
from yuxi.channels.contract.ports.driven.agent_run_port import AgentRunPort
__all__ = ["LoadBuildStage"]
@ -21,15 +26,18 @@ class LoadBuildStage:
"""出站管道 load-build 阶段。
职责构建出站负载 OutboundPayload关联 Agent 运行 ID 与流式分块
携带富消息字段FR-08与附件列表FR-48
携带富消息字段FR-08与附件列表FR-48持久化模式下负责加载
AgentRun 完整响应内容stream_chunk_stage 被跳过内容加载职责由
本阶段承接
关联 FRFR-13流式分块装配FR-08富消息字段填充
FR-48附件装配
StageContract
- id="load-build"
- reads=("agent_run_id", "stream_chunks", "rich_message", "attachments")
- writes=("outbound_payload",)
- reads=("agent_run_id", "stream_chunks", "rich_message", "attachments",
"delivery_mode")
- writes=("outbound_payload", "stream_chunks")
- idempotent=True
- thread_safe=True
- failure=TERMINATE
@ -38,30 +46,45 @@ class LoadBuildStage:
"""
id: str = "load-build"
reads: tuple[str, ...] = ("agent_run_id", "stream_chunks", "rich_message", "attachments")
writes: tuple[str, ...] = ("outbound_payload",)
reads: tuple[str, ...] = (
"agent_run_id",
"stream_chunks",
"rich_message",
"attachments",
"delivery_mode",
)
writes: tuple[str, ...] = ("outbound_payload", "stream_chunks")
idempotent: bool = True
thread_safe: bool = True
failure: FailureStrategy = FailureStrategy.TERMINATE
compensate: str | None = None
condition: Callable[[Any], bool] | None = None
def __init__(self) -> None:
"""初始化 load-build 阶段。"""
pass
def __init__(self, agent_run_port: AgentRunPort) -> None:
"""初始化 load-build 阶段。
参数
agent_run_port: Agent 运行被驱动端口持久化模式下通过
``streamAgentRun`` 阻塞消费事件流等待 AgentRun 完成
再通过 ``getAgentRunFinalOutput`` 获取完整输出文本
"""
self._agent_run_port = agent_run_port
async def process(self, context: OutboundContext) -> bool:
"""构建出站负载。
处理步骤
1. ctx.stream_chunks 构造分块元组
2. ctx.rich_message 构造富消息字段FR-08
3. 装配附件列表FR-48
1. 持久化模式下 agent_run_id 非空且 stream_chunks 为空
通过 streamAgentRun 阻塞消费至 AgentRun 完成再调用
getAgentRunFinalOutput 获取完整文本填充 stream_chunks
2. ctx.stream_chunks 构造分块元组
3. ctx.rich_message 构造富消息字段FR-08
4. 装配附件列表FR-48
- ctx.attachments 直接并入
- rich_message.image_url 转换为 Attachment(type="image") 并入
- rich_message.video_url 转换为 Attachment(type="video") 并入
4. 构造 OutboundPayload关联 agent_run_id
5. 设置 ctx.outbound_payload
5. 构造 OutboundPayload关联 agent_run_id
6. 设置 ctx.outbound_payload
参数
context: 出站管道上下文
@ -71,6 +94,17 @@ class LoadBuildStage:
抛出
"""
# 持久化模式下AgentRun 响应内容尚未加载stream_chunk_stage 被
# condition 跳过)。阻塞消费事件流等待 AgentRun 完成后加载完整输出,
# 由下游 outbox-persist/deliver 投递。流式模式下 stream_chunk_stage
# 负责加载+投递,此处不介入。
if (
context.delivery_mode == "persistent"
and context.agent_run_id
and not context.stream_chunks
):
await self._loadAgentRunOutput(context)
rich_message_fields = None
attachments: list[Attachment] = list(context.attachments)
@ -90,3 +124,27 @@ class LoadBuildStage:
attachments=tuple(attachments),
)
return True
async def _loadAgentRunOutput(self, context: OutboundContext) -> None:
"""持久化模式下加载 AgentRun 完整输出。
两步
1. 消费 ``streamAgentRun`` 事件流阻塞至 AgentRun 完成``end``
事件或 DB 终态兜底触发迭代结束
2. 调用 ``getAgentRunFinalOutput`` 一次性获取完整文本复用
已有的 payload 提取逻辑避免重复解析 StreamEvent payload
嵌套结构
参数
context: 出站管道上下文读取 ``agent_run_id``写入
``stream_chunks``
"""
run_id = AgentRunId(value=context.agent_run_id)
# 步骤 1阻塞消费至完成不提取内容仅等待 end 信号)
async with aclosing(self._agent_run_port.streamAgentRun(run_id)) as aiterator:
async for _ in aiterator:
pass
# 步骤 2一次性获取完整文本
result = await self._agent_run_port.getAgentRunFinalOutput(run_id)
if isinstance(result, Some):
context.stream_chunks.append(result.value)

View File

@ -220,7 +220,7 @@ class OutboundPipeline(Pipeline):
conversation_port=conversation_port,
logger=logger,
)
load_build_stage = LoadBuildStage()
load_build_stage = LoadBuildStage(agent_run_port=agent_run_port)
capability_verify_stage = CapabilityVerifyStage(
capability_verifier=capability_verifier,
plugin_registry=plugin_registry,
@ -244,6 +244,7 @@ class OutboundPipeline(Pipeline):
streaming_adapter_registry=streaming_adapter_registry,
config_port=config_port,
agent_run_port=agent_run_port,
outbound_adapter_registry=outbound_adapter_registry,
logger=logger,
)
truncation_check_stage = TruncationCheckStage(

View File

@ -84,7 +84,7 @@ class OutboxMarkFailed:
处理步骤
1. ctx.outbox_entry 获取 outbox_id
2. **重新查询最新 outbox 条目**补偿幂等性保证deliver 阶段
可能已通过 markSentUnconfirmed 递增 version 并持久化
可能已通过 updateOutboxEntry 递增 DB version 并持久化
``context.outbox_entry`` version 可能已过期直接用会
导致乐观锁 ConflictError重新查询保证读到最新 version
3. 终态检查SENT / FAILED / DEAD / SUPPRESSED 等终态条目
@ -111,8 +111,8 @@ class OutboxMarkFailed:
return True
# 补偿幂等性保证重新查询最新快照。deliver 阶段在发送前会调用
# markSentUnconfirmed + updateOutboxEntry 递增 version若发送失败
# 触发补偿时 context.outbox_entry.version 仍是旧值deliver 失败路径
# updateOutboxEntry 递增 DB version若发送失败触发补偿时
# context.outbox_entry.version 仍是旧值deliver 失败路径
# 不保证同步最新 DTO直接用旧 version 更新会触发乐观锁
# ConflictError。重新查询确保读到 DB 当前 version。
latest = await self.persistence_port.getOutboxEntry(dto.outbox_id)

View File

@ -40,7 +40,7 @@ class PrefixStage:
writes: tuple[str, ...] = ("final_message",)
idempotent: bool = True
thread_safe: bool = True
failure: FailureStrategy = FailureStrategy.SKIP
failure: FailureStrategy = FailureStrategy.TERMINATE
compensate: str | None = None
# 流式模式下 content 由 stream-chunk 阶段独立投递,不构造 FinalMessage
# 跳过 prefix 避免对流式尚未到达的 content 做非空校验

View File

@ -19,26 +19,36 @@ TTL 安全阀FR-13 / AC-62累计耗时超过 ``streaming_ttl_ms`` 时
发送分块标记 ``stream_aborted_at_chunk`` 并降级为持久化模式 typing-stop
阶段强制结束流式会话``endStreaming``deliver 阶段检测中断标记后调用
``sendMessageContinuation`` 续发剩余内容H-2
降级时重新格式化TTL 超时或异常降级为持久化模式后``format_stage``
流式模式下设置的空 content ``formatted_message`` 已过期本阶段在降级时
用已消费的 ``stream_chunks`` 重新调用 ``formatOutbound`` 更新
``formatted_message``确保后续 ``outbox_persist``/``deliver`` 阶段读到
非空 content
"""
from __future__ import annotations
import asyncio
import time
from collections.abc import AsyncIterator, Callable
from collections.abc import AsyncIterator, Callable, Iterator
from contextlib import aclosing
from typing import Any
from yuxi.channels.application.context.outbound_context import OutboundContext
from yuxi.channels.contract.dtos.agent_run import AgentRunId
from yuxi.channels.contract.dtos.channel import ChannelType
from yuxi.channels.contract.dtos.common import MessageContent, MessageFormat
from yuxi.channels.contract.dtos.config import ConfigScope
from yuxi.channels.contract.dtos.outbound import FormattedMessage
from yuxi.channels.contract.dtos.streaming import (
ChunkResult,
StreamChunk,
StreamingCompleted,
)
from yuxi.channels.contract.dtos.stream_event import StreamEvent
from yuxi.channels.contract.errors.domain import ChannelDegradedError
from yuxi.channels.contract.plugin.adapters.outbound_adapter import OutboundAdapter
from yuxi.channels.contract.plugin.adapters.streaming_adapter import StreamingAdapter
from yuxi.channels.contract.plugin.extension_point import FailureStrategy
from yuxi.channels.contract.ports.driven.agent_run_port import AgentRunPort
@ -48,19 +58,58 @@ from yuxi.channels.contract.ports.driven.logger_port import LoggerPort
__all__ = ["StreamChunkStage"]
def _extractTextFromStreamEvent(event: StreamEvent) -> Iterator[str]:
"""从 StreamEvent 提取文本内容。
``AgentRunAdapter.getAgentRunFinalOutput`` 的提取逻辑一致
1. 仅处理 ``event_type == "messages"`` 的事件跳过 metadata/custom/
error/end 等非内容事件
2. 穿透两层 payload``event.payload``envelope
``envelope["payload"]``inner ``inner["items"]``
3. 对每个 item 优先取 ``stream_event.content``语义增量
回退 ``response``直接文本避免重复拼接
参数
event: StreamEvent DTO
生成
文本内容字符串
"""
if event.event_type != "messages":
return
envelope = event.payload or {}
inner_payload = envelope.get("payload") or {}
items = inner_payload.get("items") or []
for item in items:
if not isinstance(item, dict):
continue
stream_event = item.get("stream_event")
if isinstance(stream_event, dict):
content = stream_event.get("content")
if isinstance(content, str) and content:
yield content
continue
response = item.get("response")
if isinstance(response, str) and response:
yield response
class StreamChunkStage:
"""出站管道 stream-chunk 阶段。
职责流式分块投递FR-13消费 Agent Run 事件流并按序将分块
发送至渠道侧实现实时流式输出AC-15
发送至渠道侧实现实时流式输出AC-15降级时重新格式化
``formatted_message``确保持久化路径读到非空 content
关联 FRFR-13流式输出
StageContract
- id="stream-chunk"
- reads=("channel_type", "account_id", "peer_id", "agent_run_id", "stream_chunks", "delivery_mode", "trace_id")
- reads=("channel_type", "account_id", "peer_id", "agent_run_id",
"stream_chunks", "delivery_mode", "trace_id", "formatted_message")
- writes=("chunk_results", "stream_chunks", "streaming_completed",
"delivery_mode", "degraded", "degraded_reason")
"delivery_mode", "degraded", "degraded_reason",
"stream_aborted_at_chunk", "formatted_message")
- idempotent=False
- thread_safe=False
- failure=DEGRADE
@ -77,6 +126,7 @@ class StreamChunkStage:
"stream_chunks",
"delivery_mode",
"trace_id",
"formatted_message",
)
writes: tuple[str, ...] = (
"chunk_results",
@ -86,6 +136,7 @@ class StreamChunkStage:
"degraded",
"degraded_reason",
"stream_aborted_at_chunk",
"formatted_message",
)
idempotent: bool = False
thread_safe: bool = False
@ -98,6 +149,7 @@ class StreamChunkStage:
streaming_adapter_registry: dict[ChannelType, StreamingAdapter],
config_port: ConfigPort,
agent_run_port: AgentRunPort,
outbound_adapter_registry: dict[ChannelType, OutboundAdapter],
logger: LoggerPort | None = None,
) -> None:
"""初始化 stream-chunk 阶段。
@ -106,12 +158,15 @@ class StreamChunkStage:
streaming_adapter_registry: 流式适配器注册表按渠道类型索引
config_port: 配置被驱动端口用于读取最小分块间隔与 TTL
agent_run_port: Agent 运行被驱动端口用于流式消费 Agent 运行事件
outbound_adapter_registry: 出站适配器注册表降级时调用
``formatOutbound`` 重新格式化 ``formatted_message``
logger: 可选日志端口用于 TTL 安全阀告警 ``None`` 时跳过
告警日志输出
"""
self.streaming_adapter_registry = streaming_adapter_registry
self.config_port = config_port
self.agent_run_port = agent_run_port
self.outbound_adapter_registry = outbound_adapter_registry
self.logger = logger
async def process(self, context: OutboundContext) -> bool:
@ -123,12 +178,13 @@ class StreamChunkStage:
3. 消费 Agent Run 事件流或预填充分块逐块投递至渠道侧
AC-15用户看到逐步更新
4. 累计耗时超过 TTL 时停止发送标记 stream_aborted_at_chunk
降级为持久化模式记录告警日志AC-62H-2
降级为持久化模式重新格式化 formatted_message记录告警日志
AC-62H-2
5. 构造 StreamingCompleted设置 ctx.streaming_completed
适配器不存在或投递异常时 delivery_mode 降级为持久化模式并显式
标记 context.degraded由后续持久化投递阶段outbox-persist /
deliver继续执行
标记 context.degraded重新格式化 formatted_message由后续持久化
投递阶段outbox-persist / deliver继续执行
参数
context: 出站管道上下文
@ -147,6 +203,7 @@ class StreamChunkStage:
context.channel_type,
trace_id=context.trace_id,
)
await self._reformatAfterDegrade(context)
return True
streaming_config = await self.config_port.getStreamingConfig(ConfigScope.ACCOUNT, context.account_id)
@ -177,6 +234,9 @@ class StreamChunkStage:
# sendMessageContinuation 续发剩余内容H-2
context.stream_aborted_at_chunk = len(results)
context.delivery_mode = "persistent"
# 降级后重新格式化 formatted_message确保后续
# outbox-persist/deliver 阶段读到非空 content
await self._reformatAfterDegrade(context)
break
if results and min_interval_s > 0:
@ -198,6 +258,9 @@ class StreamChunkStage:
# sendMessageContinuation 续发剩余内容H-2
context.stream_aborted_at_chunk = len(results)
context.delivery_mode = "persistent"
# 降级后重新格式化 formatted_message在抛出前执行因 DEGRADE
# 策略会 continue 后续阶段,需确保 formatted_message 已更新)
await self._reformatAfterDegrade(context)
raise ChannelDegradedError(
context.channel_type,
trace_id=context.trace_id,
@ -222,6 +285,10 @@ class StreamChunkStage:
- ``agent_run_id`` 为空时如管理员消息遍历上游预填充的
``stream_chunks`` 列表
``StreamEvent`` 提取文本内容时 ``event_type == "messages"``
过滤并穿透两层 payloadenvelope inner items
``getAgentRunFinalOutput`` 的提取逻辑一致
参数
context: 出站管道上下文
@ -232,12 +299,50 @@ class StreamChunkStage:
run_id = AgentRunId(value=context.agent_run_id)
aiterator = self.agent_run_port.streamAgentRun(run_id)
try:
async for chunk in aiterator:
content = getattr(chunk, "content", str(chunk))
context.stream_chunks.append(content)
yield content
async for event in aiterator:
for content in _extractTextFromStreamEvent(event):
context.stream_chunks.append(content)
yield content
finally:
await aiterator.aclose()
else:
for content in context.stream_chunks:
yield content
async def _reformatAfterDegrade(self, context: OutboundContext) -> None:
"""降级后重新格式化 ``formatted_message``。
流式模式下 ``format_stage`` 先执行时 ``stream_chunks`` 为空
``formatted_message.content`` 为空字符串降级为持久化模式后
用已消费的 ``stream_chunks`` 重新调用 ``formatOutbound`` 更新
``formatted_message``确保后续 ``outbox_persist``/``deliver``
阶段读到非空 content 与正确的 ``rich_message``
参数
context: 出站管道上下文
"""
content = "".join(context.stream_chunks)
if not content:
return
adapter = self.outbound_adapter_registry.get(context.channel_type)
if adapter is None:
# 理论上不会发生:所有渠道均注册出站适配器。直接构造 TEXT
# 格式的 formatted_message 保证管道不中断。
context.formatted_message = FormattedMessage(
content=content,
format=MessageFormat.TEXT,
attachments=context.formatted_message.attachments
if context.formatted_message is not None
else (),
)
return
message_content = MessageContent(text=content, format=MessageFormat.TEXT)
formatted = await adapter.formatOutbound(message_content)
context.formatted_message = FormattedMessage(
content=formatted.content,
format=formatted.format,
rich_message=formatted.rich_message,
attachments=context.formatted_message.attachments
if context.formatted_message is not None
else (),
)

View File

@ -438,30 +438,31 @@ class TransportManager:
source=source,
)
started = False
# 先持久化插件运行态为 runningFR-32 诊断字段best-effort 写入),
# 必须在 start_account 之前调用start_account 内部通过 asyncio.create_task
# 调度的 _runAccountLoop 后台任务会使用同一共享 AsyncSession 调用
# getChannelAccount若 updatePluginStatus 与 task 并发执行会触发
# SQLAlchemy AsyncSession 并发访问异常(该异常在 channel_persistence_adapter
# 的 except SQLAlchemyError 分支被静默翻译为 DependencyError无原始异常日志
# 提前完成 updatePluginStatus 的 commit 可确保 task 启动时 session 处于干净状态。
# plugin_status 为诊断字段,即使后续 start_account 因 circuit breaker open
# 等原因未真正启动 task状态轻微不一致可接受下次状态变化时纠正
await self._touchPluginStatus(channel_type, account_id, "running", trace_id)
if transport_mode == "pull":
if puller_adapter is not None and self._puller_worker is not None:
await self._puller_worker.start_account(channel_type, account_id, puller_adapter)
started = True
elif transport_mode == "stream":
if stream_adapter is not None and self._stream_worker is not None:
await self._stream_worker.start_account(channel_type, account_id, stream_adapter)
started = True
else:
# Task 11.2: both 模式优先 StreamPuller 作为降级。
# Stream 适配器可用时仅启动 StreamStream 健康时不 poll
# Stream 适配器不可用时降级启动 Puller。
if stream_adapter is not None and self._stream_worker is not None:
await self._stream_worker.start_account(channel_type, account_id, stream_adapter)
started = True
elif puller_adapter is not None and self._puller_worker is not None:
await self._puller_worker.start_account(channel_type, account_id, puller_adapter)
started = True
# 传输任务成功启动后持久化插件运行态FR-32。诊断字段写入失败
# 仅告警不中止(与 last_health_check_at 一致的 best-effort 语义)。
if started:
await self._touchPluginStatus(channel_type, account_id, "running", trace_id)
async def _restoreOnlineAccounts(self, trace_id: str) -> None:
"""重启恢复:扫描 DB 中 ACTIVE 状态账号,启动传输任务。

View File

@ -250,15 +250,11 @@ class IdentityMergeService:
f"identity {cmd.merged_identity_id} is not pending review",
)
# 3. 记录 DB 版本merge / clearPendingReview / unbindUser 仅递增内存 version
# 3. 记录 DB 版本merge 仅递增内存 version,持久化时用 DB 版本校验
canonical_db_version = canonical.version
# 4. 执行合并 + 清除 pending_review + 清除 sibling 的 user_idM-7
# 清除 user_id 与 _clearSiblingAndMigrateSessions 对齐,避免 sibling
# 仍被 listUserIdentitiesByUserId 返回导致后续 bindUserToSession 重复发现
# 4. 执行合并canonical
canonical.merge(sibling, trigger="approveMerge", operator_id=operator.user_id)
sibling.clearPendingReview()
sibling.unbindUser()
# 5. 持久化 canonical 的 bindings/merged_from乐观锁expected_version 用 DB 版本)
await self._user_identity_repo.updateUserIdentityBindings(
@ -270,10 +266,22 @@ class IdentityMergeService:
tx=tx,
)
# 6. 持久化 sibling 的 pending_review 清除与 user_id 清除(乐观锁)
sibling_updated_dto = _aggregate_to_dto(sibling, sibling_dto)
# 6. 持久化 sibling 的状态变更clearPendingReview 与 unbindUser 拆分为两次持久化。
# updateUserIdentity 适配器以 ``identity.version - 1`` 作为期望旧版本号
# (假设单次状态转换),多次状态转换需拆分持久化并在中间重建聚合根
# 同步 DB version否则乐观锁校验失败M+2-1=M+1 ≠ DB M
sibling.clearPendingReview()
sibling_dto_after = _aggregate_to_dto(sibling, sibling_dto)
updated_sibling_dto = await self._user_identity_repo.updateUserIdentity(
sibling_dto_after,
operator_id=operator.user_id,
tx=tx,
)
sibling = UserIdentityAggregate.from_dto(updated_sibling_dto)
sibling.unbindUser()
sibling_dto_after = _aggregate_to_dto(sibling, updated_sibling_dto)
await self._user_identity_repo.updateUserIdentity(
sibling_updated_dto,
sibling_dto_after,
operator_id=operator.user_id,
tx=tx,
)

View File

@ -62,15 +62,18 @@ class DrivenAdapters:
聚合 11 个被驱动适配器实例 1 个单实例 ``identity_resolver``
22 个渠道插件适配器列表 ``channel_context_providers`` 端口列表
供框架层与应用服务层统一注入需要共享 ``AsyncSession`` 的适配器
``persistence`` / ``conversation`` / ``transaction`` factory 中共享同一 ``db``无状态适配器
供框架层与应用服务层统一注入``persistence`` 为无状态适配器注入
``session_factory``通过 ``_session_scope(tx)`` 按需获取 session不持有
请求级 ``db````conversation`` / ``transaction`` 为请求级事务边界适配器
factory 中共享同一 ``db`` 会话以保证事务一致性无状态适配器
使用全局实例``identity_resolver`` 默认为 ``None``由插件注册时注入
其余 22 类渠道适配器inbound/outbound/status 默认为空列表由插件
通过 ``PluginHost.registerAdapter`` 注入``create_channel_use_cases``
``PluginRegistry`` 读取后填充管道注册表
字段
persistence: 持久化适配器PersistencePort
persistence: 持久化适配器PersistencePort无状态注入
``session_factory``通过 ``_session_scope(tx)`` 按需获取 session
conversation: 会话适配器ConversationPort
agent_run: Agent 运行适配器AgentRunPort
queue: 任务队列适配器QueuePort
@ -162,14 +165,16 @@ class DrivenAdapters:
CloseablePort)`` 显式判断是否实现关闭契约INV-3仅对实现方
调用 ``aclose`` 释放资源
``ChannelPersistenceAdapter`` 持有共享 ``AsyncSession`` 并实现
``CloseablePort``关闭它即释放 persistence / conversation / transaction
三者共享的会话``ConversationAdapter`` ``SqlAlchemyTransactionAdapter``
不实现 ``CloseablePort``共享 session persistence 统一关闭避免
重复关闭``AgentRunAdapter`` ``ARQQueueAdapter`` 不持有 ``db``
每次调用内部获取临时会话同样不实现 ``CloseablePort``
``ChannelPersistenceAdapter`` 为无状态适配器注入 ``session_factory``
不持有请求级 ``db``不再实现 ``CloseablePort``关闭时跳过
``ConversationAdapter`` ``SqlAlchemyTransactionAdapter`` 为请求级
事务边界适配器共享同一 ``db`` 会话但其生命周期由请求级框架管理
请求结束时由调用方 ``db.close()``不实现 ``CloseablePort``
``AgentRunAdapter`` ``ARQQueueAdapter`` 不持有 ``db``每次调用
内部获取临时会话同样不实现 ``CloseablePort``
未遍历的字段说明
- ``persistence``无状态适配器不持有需释放资源
- ``config`` / ``cache`` / ``logger`` / ``tracer`` / ``masking``
无状态或由框架管理生命周期无需关闭
- ``identity_resolver`` / ``service_account`` / ``agent_access_config``
@ -185,7 +190,6 @@ class DrivenAdapters:
"""
closed: set[int] = set()
for adapter in (
self.persistence,
self.conversation,
self.agent_run,
self.queue,

View File

@ -1,10 +1,12 @@
"""持久化健康检查被驱动端口。
定义核心层对数据库可用性连接池诊断与会话生命周期管理的依赖契约
application 层调用framework 层实现 ``HealthAggregator``
``DiagnosticsExporter`` 聚合健康状态与填充诊断包并供后台扫描器在每轮
扫描结束时释放隐式事务占用的连接避免应用级共享 ``AsyncSession`` 长期
持有连接导致连接池耗尽
定义核心层对数据库可用性与连接池诊断的依赖契约 application 层调用
framework 层实现 ``HealthAggregator`` ``DiagnosticsExporter`` 聚合
健康状态与填充诊断包
适配器为无状态协议转换器不持有 session每次方法调用通过
``_session_scope(tx)`` 按需获取 session 并在调用结束后归还连接无需
额外的会话生命周期管理方法
"""
from __future__ import annotations
@ -19,12 +21,13 @@ if TYPE_CHECKING:
class PersistenceHealthPort(Protocol):
"""持久化健康检查被驱动端口。
覆盖数据库可用性探测连接池状态查询与会话生命周期管理用例方法遵循
领域语义命名入参/出参均为契约层不可变值对象
覆盖数据库可用性探测与连接池状态查询方法遵循领域语义命名入参/出参
均为契约层不可变值对象
事务边界§10.1
- 健康检查方法无 ``tx`` 参数按只读单次查询执行不参与应用层
事务
事务适配器通过 ``_session_scope(None)`` 创建独立 session
查询结束后自动 close 归还连接无需显式释放
约束
- 复用约束INV-I3必须复用现有 pg_manager不得引入新连接池
@ -62,21 +65,3 @@ class PersistenceHealthPort(Protocol):
- DependencyError数据库故障
"""
...
async def releaseSession(self) -> None:
"""释放当前会话占用的连接FR-35 会话生命周期管理)。
对应用级共享 ``AsyncSession`` 执行 ``commit`` 以结束隐式事务并归还
连接至池SQLAlchemy 2.0 autobegin 只读查询会开启隐式事务并
持有连接直至显式提交/回滚后台扫描器在每轮扫描结束时调用本方法
避免应用级共享会话长期占用连接导致连接池耗尽
语义约束
- 无活动事务时为 no-op``commit`` 对干净 session 安全
- 不关闭会话本身``aclose`` 负责最终关闭仅归还连接
- 异常翻译为 ``DependencyError``禁止原生异常穿透至核心层
@failure
- DependencyError提交事务时数据库故障
"""
...

View File

@ -4,18 +4,13 @@
实现事务边界 **必须** 由应用层管道或用例编排器显式开启与提交
被驱动适配器 **不得** 自主开启跨调用的事务§10.1
事务共享机制被驱动适配器 ``ChannelPersistenceAdapter`` /
``ConversationAdapter``在构造时通过 ``create_driven_adapters(db)``
``SqlAlchemyTransactionAdapter`` 共享同一 ``AsyncSession``应用层通过
``TransactionPort.begin()`` 在共享 session 上开启事务所有适配器的写
操作自动加入同一事务``tx`` 参数仅作为"是否自主提交"的标志位``tx``
非空时 ``commit=False``由应用层统一提交
事务透传机制C-I1被驱动适配器 ``AgentRunAdapter``在构造时
**** ``SqlAlchemyTransactionAdapter`` 共享 session而是通过
``pg_manager`` 自主获取为让其加入应用层主事务``TransactionContext``
暴露 ``get_session()`` 返回底层共享会话适配器通过 ``tx.get_session()``
复用主事务的 session避免独立提交产生孤儿记录C-I1
事务透传机制C-I1被驱动适配器 ``ChannelPersistenceAdapter`` /
``ContentReviewRepositoryAdapter`` / ``AgentRunAdapter``为无状态协议
转换器构造时仅注入 ``session_factory````Callable[[], AsyncSession]``
不持有 session 实例适配器通过 ``_session_scope(tx)`` 统一管理 session
``tx`` 非空时通过 ``tx.get_session()`` 复用应用层主事务 sessionC-I1
透传``commit=False``由应用层统一提交``tx`` ``None`` 时通过
``session_factory`` 创建独立 session 并自主提交``commit=True``
"""
from __future__ import annotations
@ -73,13 +68,11 @@ class TransactionContext(Protocol):
``TransactionPort.begin()`` 创建上下文管理器退出时自动提交
无异常或回滚有异常
事务共享机制被驱动适配器在构造时与事务适配器共享同一底层会话
SQLAlchemy ``AsyncSession````begin()`` 在共享会话上开启事务
所有适配器的写操作自动加入
事务透传机制C-I1未在构造时共享 session 的适配器
``AgentRunAdapter``通过 ``get_session()`` 获取底层共享会话复用
主事务的 session避免独立提交产生孤儿记录
事务透传机制C-I1无状态适配器通过 ``get_session()`` 获取底层
共享会话复用主事务的 session避免独立提交产生孤儿记录适配器
``_session_scope(tx)`` 中判断 ``tx`` 非空且 ``get_session()``
返回 session 时使用该 session``commit=False``否则创建独立
session``commit=True``
被驱动适配器 **不得** 自主调用 ``commit`` / ``rollback``
@ -89,10 +82,10 @@ class TransactionContext(Protocol):
"""
def get_session(self) -> AsyncSession | None:
"""返回底层共享会话,供未在构造时共享 session 的适配器复用主事务
"""返回底层共享会话,供无状态适配器复用主事务 session
SQL 实现返回真实 ``AsyncSession`` SQL 实现返回 ``None``
调用方回退到自主获取 session 的路径
适配器回退到 ``session_factory`` 创建独立 session 的路径
@pre
- 事务已开启``__aenter__`` 已调用

View File

@ -1,7 +1,9 @@
"""发件箱条目聚合根。
实现持久化投递记录的状态机管理封装投递成功抑制失败未确认与死信
规则聚合根为可变 Python 业务方法直接修改内部状态并递增乐观锁版本
规则聚合根为可变 Python 业务方法直接修改内部状态乐观锁版本号
``version`` 由持久化层``updateOutboxEntry``统一管理聚合根不递增
以支持单次持久化前的多次内存状态转换 ``markFailed`` + ``markDead``
"""
from __future__ import annotations
@ -161,7 +163,6 @@ class OutboxEntry:
if channel_msg_id is not None:
self.channel_msg_id = channel_msg_id
self.updated_at = datetime.now(UTC)
self.version += 1
# 首次调用写入 latency_ms 与 sent_at重试成功不覆盖
if self.latency_ms is None:
self.latency_ms = int((self.updated_at - self.created_at).total_seconds() * 1000)
@ -182,7 +183,6 @@ class OutboxEntry:
self.status = OutboxStatus.SUPPRESSED
self.last_error = reason
self.updated_at = datetime.now(UTC)
self.version += 1
self.funnel_node = "suppressed"
def markFailed(self, error: str) -> None:
@ -208,7 +208,6 @@ class OutboxEntry:
self.last_error = error
self.updated_at = datetime.now(UTC)
self.last_retry_at = self.updated_at
self.version += 1
self.funnel_node = "failed"
def markDead(self, reason: str) -> None:
@ -225,7 +224,6 @@ class OutboxEntry:
self.next_retry_at = None
self.last_error = reason
self.updated_at = datetime.now(UTC)
self.version += 1
self.funnel_node = "dead"
def markDeliveryUnconfirmedFailed(self, reason: str) -> None:
@ -259,7 +257,6 @@ class OutboxEntry:
self.degraded_reason = reason
self.updated_at = datetime.now(UTC)
self.last_retry_at = self.updated_at
self.version += 1
self.funnel_node = "failed"
def markPartialFailure(
@ -306,11 +303,9 @@ class OutboxEntry:
self.degraded_reason = self.last_error
self.updated_at = datetime.now(UTC)
self.last_retry_at = self.updated_at
self.version += 1
self.funnel_node = "failed"
return
self.updated_at = datetime.now(UTC)
self.version += 1
def requeueForRetry(self) -> None:
"""将 FAILED 条目重新入队重试FAILED → PENDING
@ -327,7 +322,6 @@ class OutboxEntry:
self.status = OutboxStatus.PENDING
self.next_retry_at = None
self.updated_at = datetime.now(UTC)
self.version += 1
self.funnel_node = "enter"
def markSentUnconfirmed(self, channel_msg_id: str | None) -> None:
@ -344,7 +338,6 @@ class OutboxEntry:
if channel_msg_id is not None:
self.channel_msg_id = channel_msg_id
self.updated_at = datetime.now(UTC)
self.version += 1
def canRetry(self) -> bool:
"""判断发件箱条目是否可重试。
@ -375,7 +368,7 @@ class OutboxEntry:
调用
状态转换DEAD FAILEDretry_count 重置为 0last_error 清空
next_retry_at 设为当前时间立即可重投version 递增
next_retry_at 设为当前时间立即可重投
抛出
RuleViolationError: 当前状态非 DEAD仅死信可复活
@ -387,7 +380,6 @@ class OutboxEntry:
self.last_error = None
self.next_retry_at = datetime.now(UTC)
self.updated_at = datetime.now(UTC)
self.version += 1
self.funnel_node = "failed"
def _applyStatusTransition(

View File

@ -6,7 +6,7 @@
external_systems ``UseCases`` + ``create_use_cases_from_db(db)`` 模式
装配顺序
1. ``create_driven_adapters(db)`` 创建被驱动适配器聚合
1. ``create_driven_adapters(db, session_factory, ...)`` 创建被驱动适配器聚合
2. DI 容器获取应用级单例注册中心等
3. 按依赖图拓扑顺序构造领域服务
4. 通过 ``Pipeline.create()`` 创建管道
@ -25,12 +25,9 @@ from typing import Any
from arq import ArqRedis
from redis.asyncio import Redis
from sqlalchemy.ext.asyncio import AsyncSession
from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker
from yuxi.channels.adapters import DrivenAdapters, create_driven_adapters
from yuxi.channels.adapters.content_review_repository_adapter import (
ContentReviewRepositoryAdapter,
)
from yuxi.channels.adapters.default_content_moderation_adapter import (
DefaultContentModerationAdapter,
)
@ -276,7 +273,7 @@ __all__ = [
class ChannelUseCases:
"""渠道用例服务聚合。
串联 ``create_driven_adapters(db)`` 领域服务 管道 用例服务
串联 ``create_driven_adapters(db, session_factory, ...)`` 领域服务 管道 用例服务
请求级创建共享同一 ``db`` session
对齐 external_systems ``UseCases`` dataclass 模式字段类型使用驱动
@ -418,7 +415,7 @@ def create_channel_use_cases(
"""请求级 factory创建用例服务聚合。
装配顺序对齐 external_systems ``create_use_cases_from_db`` 模式
1. ``create_driven_adapters(db)`` 创建被驱动适配器聚合
1. ``create_driven_adapters(db, session_factory, ...)`` 创建被驱动适配器聚合
2. DI 容器获取应用级单例注册中心等
3. 创建领域服务按依赖图拓扑顺序
4. 通过 ``Pipeline.create()`` 创建管道
@ -453,14 +450,19 @@ def create_channel_use_cases(
# 配置变更后发布领域事件CON-030 / ADP-029事件由应用层发布
# F-02ConfigScopeRegistry 从 DI 容器解析,提取 key_to_scope_map
# 注入到请求级 RedisConfigAdapter 供 _build_key 校验作用域一致性。
# session_factoryasync_sessionmaker从 DI 容器解析,注入到
# ChannelPersistenceAdapter 供 _session_scope(tx) 按需创建独立 session
# (无状态适配器不持有请求级 db
redis_client: Redis = di_container.resolve(Redis)
arq_pool: ArqRedis = di_container.resolve(ArqRedis)
execution_port: AgentRunExecutionPort = di_container.resolve(AgentRunExecutionPort)
event_publisher: EventPublisherPort = di_container.resolve(EventPublisherPort)
queue_port: QueuePort = di_container.resolve(QueuePort)
config_scope_registry: ConfigScopeRegistry = di_container.resolve(ConfigScopeRegistry)
session_factory: async_sessionmaker = di_container.resolve(async_sessionmaker)
adapters: DrivenAdapters = create_driven_adapters(
db,
session_factory=session_factory,
outbox_config=outbox_config,
redis_client=redis_client,
arq_pool=arq_pool,
@ -484,10 +486,10 @@ def create_channel_use_cases(
# - ContentReviewRepositoryPort审核历史仓储端口绑定注入 DispatchStage
# (写路径)与 ContentReviewQueryService读路径
default_moderation_adapter: DefaultContentModerationAdapter = di_container.resolve(DefaultContentModerationAdapter)
# 请求级构造 ContentReviewRepositoryAdapterDI 单例持有共享 AsyncSession
# 预览 INSERT 失败后事务进入 aborted 状态且无法回滚,污染后续所有请求。
# 改为复用请求级 ``db`` 会话,失败时由应用层事务统一回滚
review_repository: ContentReviewRepositoryPort = ContentReviewRepositoryAdapter(db, logger=adapters.logger)
# ContentReviewRepositoryPort 从 DI 容器解析(无状态单例,由
# create_host_bootstrap 注册)。适配器通过 _session_scope(tx) 按需获取
# sessiontx 非空时复用应用层主事务 sessiontx 为 None 时自主创建并提交
review_repository: ContentReviewRepositoryPort = di_container.resolve(ContentReviewRepositoryPort)
# PluginLoader 应用级单例P1 PLG-INSTALL/UNINSTALL 文件操作),由
# create_host_bootstrap 注册到 DI 容器,注入 DispatchStage。
plugin_loader: PluginLoader = di_container.resolve(PluginLoader)
@ -1079,21 +1081,23 @@ def create_channel_session_inactive_cleanup_handler_dependencies(
调用传入 worker 进程级 ``session_factory`` ``cache_port``返回
``ChannelSessionInactiveCleanupHandler.__init__`` 所需的依赖字典
依赖方向合规INV-1``session_repo_factory`` infrastructure
装配具体适配器 ``ChannelPersistenceAdapter``fat adapter实现
依赖方向合规INV-1``session_repo`` infrastructure 装配具体
适配器 ``ChannelPersistenceAdapter``fat adapter实现
``ChannelSessionRepositoryPort``应用层 handler 仅依赖端口抽象
不直接 import 适配器实现
适配器无状态注入 ``session_factory``不持有 session 实例通过
``_session_scope(tx)`` 按需获取 session可安全在 handler 长生命周期内复用
Args:
session_factory: 调用返回 ``AsyncSession`` 的异步上下文管理器工厂
通常为 ``pg_manager.get_async_session_context``
session_factory: 可调用对象调用返回 ``AsyncSession`` 实例
通常为 ``pg_manager.AsyncSession``
cache_port: worker 进程级 CachePort 实例RedisCacheAdapter
handler 分布式锁使用
Returns:
依赖字典key ``ChannelSessionInactiveCleanupHandler.__init__``
参数名一致``{"session_factory": ..., "session_repo_factory": ...,
"cache_port": ..., "logger": ...}``
参数名一致``{"session_repo": ..., "cache_port": ..., "logger": ...}``
"""
from yuxi.channels.adapters.channel_persistence_adapter import (
ChannelPersistenceAdapter,
@ -1102,12 +1106,12 @@ def create_channel_session_inactive_cleanup_handler_dependencies(
logger_port = _create_worker_logger_port()
outbox_config = _create_outbox_config()
def session_repo_factory(db: AsyncSession) -> ChannelSessionRepositoryPort:
return ChannelPersistenceAdapter(db, outbox_config=outbox_config, logger=logger_port)
session_repo: ChannelSessionRepositoryPort = ChannelPersistenceAdapter(
session_factory, outbox_config=outbox_config, logger=logger_port
)
return {
"session_factory": session_factory,
"session_repo_factory": session_repo_factory,
"session_repo": session_repo,
"cache_port": cache_port,
"logger": logger_port,
}
@ -1124,7 +1128,8 @@ def create_channel_outbox_terminal_cleanup_handler_dependencies(
WARN 不影响清理主流程
Args:
session_factory: 调用返回 ``AsyncSession`` 的异步上下文管理器工厂
session_factory: 可调用对象调用返回 ``AsyncSession`` 实例
通常为 ``pg_manager.AsyncSession``
cache_port: worker 进程级 CachePort 实例RedisCacheAdapter
handler 分布式锁使用
@ -1139,12 +1144,12 @@ def create_channel_outbox_terminal_cleanup_handler_dependencies(
logger_port = _create_worker_logger_port()
outbox_config = _create_outbox_config()
def outbox_repo_factory(db: AsyncSession) -> OutboxRepositoryPort:
return ChannelPersistenceAdapter(db, outbox_config=outbox_config, logger=logger_port)
outbox_repo: OutboxRepositoryPort = ChannelPersistenceAdapter(
session_factory, outbox_config=outbox_config, logger=logger_port
)
return {
"session_factory": session_factory,
"outbox_repo_factory": outbox_repo_factory,
"outbox_repo": outbox_repo,
"event_publisher": None, # worker 进程未注入事件发布端口best-effort
"cache_port": cache_port,
"logger": logger_port,
@ -1158,7 +1163,8 @@ def create_channel_pairing_terminal_cleanup_handler_dependencies(
"""构造 ChannelPairingTerminalCleanupHandler 依赖字典worker 进程装配用)。
Args:
session_factory: 调用返回 ``AsyncSession`` 的异步上下文管理器工厂
session_factory: 可调用对象调用返回 ``AsyncSession`` 实例
通常为 ``pg_manager.AsyncSession``
cache_port: worker 进程级 CachePort 实例RedisCacheAdapter
handler 分布式锁使用
@ -1173,12 +1179,12 @@ def create_channel_pairing_terminal_cleanup_handler_dependencies(
logger_port = _create_worker_logger_port()
outbox_config = _create_outbox_config()
def pairing_repo_factory(db: AsyncSession) -> PairingRepositoryPort:
return ChannelPersistenceAdapter(db, outbox_config=outbox_config, logger=logger_port)
pairing_repo: PairingRepositoryPort = ChannelPersistenceAdapter(
session_factory, outbox_config=outbox_config, logger=logger_port
)
return {
"session_factory": session_factory,
"pairing_repo_factory": pairing_repo_factory,
"pairing_repo": pairing_repo,
"cache_port": cache_port,
"logger": logger_port,
}
@ -1191,13 +1197,19 @@ def create_channel_outbox_recovery_handler_dependencies(
) -> dict[str, Any]:
"""构造 ChannelOutboxRecoveryHandler 依赖字典worker 进程装配用)。
依赖方向合规INV-1``outbox_repo_factory`` / ``audit_log_repo_factory``
infrastructure 层装配具体适配器 ``ChannelPersistenceAdapter``fat adapter
分别实现 ``OutboxRepositoryPort`` / ``AuditLogRepositoryPort``应用层 handler
依赖方向合规INV-1``outbox_repo`` / ``audit_log_repo`` infrastructure
层装配具体适配器 ``ChannelPersistenceAdapter``fat adapter分别实现
``OutboxRepositoryPort`` / ``AuditLogRepositoryPort``应用层 handler
仅依赖端口抽象不直接 import 适配器实现
适配器无状态注入 ``session_factory``不持有 session 实例通过
``_session_scope(tx)`` 按需获取 session恢复扫描中每条 entry
``updateOutboxEntry`` / ``saveAuditLog`` 各自独立提交单条失败不影响其他
条目与原共享 session 模式相比避免了 session 异常状态污染后续操作
Args:
session_factory: 调用返回 ``AsyncSession`` 的异步上下文管理器工厂
session_factory: 可调用对象调用返回 ``AsyncSession`` 实例
通常为 ``pg_manager.AsyncSession``
arq_pool: worker 进程级 ARQ 连接池 ``run_worker.py``
``_worker_startup`` 提前初始化后传入用于构造 ``ARQQueueAdapter``
cache_port: worker 进程级 CachePort 实例RedisCacheAdapter
@ -1214,18 +1226,18 @@ def create_channel_outbox_recovery_handler_dependencies(
logger_port = _create_worker_logger_port()
outbox_config = _create_outbox_config()
def outbox_repo_factory(db: AsyncSession) -> OutboxRepositoryPort:
return ChannelPersistenceAdapter(db, outbox_config=outbox_config, logger=logger_port)
def audit_log_repo_factory(db: AsyncSession) -> AuditLogRepositoryPort:
return ChannelPersistenceAdapter(db, outbox_config=outbox_config, logger=logger_port)
outbox_repo: OutboxRepositoryPort = ChannelPersistenceAdapter(
session_factory, outbox_config=outbox_config, logger=logger_port
)
audit_log_repo: AuditLogRepositoryPort = ChannelPersistenceAdapter(
session_factory, outbox_config=outbox_config, logger=logger_port
)
queue_port = ARQQueueAdapter(arq_pool, logger_port)
return {
"session_factory": session_factory,
"outbox_repo_factory": outbox_repo_factory,
"audit_log_repo_factory": audit_log_repo_factory,
"outbox_repo": outbox_repo,
"audit_log_repo": audit_log_repo,
"queue_port": queue_port,
"outbox_config": outbox_config,
"cache_port": cache_port,
@ -1239,13 +1251,17 @@ def create_channel_pairing_expiration_handler_dependencies(
) -> dict[str, Any]:
"""构造 ChannelPairingExpirationHandler 依赖字典worker 进程装配用)。
依赖方向合规INV-1``pairing_repo_factory`` infrastructure 层装配具体
依赖方向合规INV-1``pairing_repo`` infrastructure 层装配具体
适配器 ``ChannelPersistenceAdapter``fat adapter实现
``PairingRepositoryPort``应用层 handler 仅依赖端口抽象不直接 import
适配器实现
适配器无状态注入 ``session_factory``不持有 session 实例通过
``_session_scope(tx)`` 按需获取 session
Args:
session_factory: 调用返回 ``AsyncSession`` 的异步上下文管理器工厂
session_factory: 可调用对象调用返回 ``AsyncSession`` 实例
通常为 ``pg_manager.AsyncSession``
cache_port: worker 进程级 CachePort 实例RedisCacheAdapter
handler 分布式锁使用
@ -1259,12 +1275,12 @@ def create_channel_pairing_expiration_handler_dependencies(
logger_port = _create_worker_logger_port()
outbox_config = _create_outbox_config()
def pairing_repo_factory(db: AsyncSession) -> PairingRepositoryPort:
return ChannelPersistenceAdapter(db, outbox_config=outbox_config, logger=logger_port)
pairing_repo: PairingRepositoryPort = ChannelPersistenceAdapter(
session_factory, outbox_config=outbox_config, logger=logger_port
)
return {
"session_factory": session_factory,
"pairing_repo_factory": pairing_repo_factory,
"pairing_repo": pairing_repo,
"cache_port": cache_port,
"logger": logger_port,
}
@ -1276,13 +1292,17 @@ def create_channel_audit_log_retention_handler_dependencies(
) -> dict[str, Any]:
"""构造 ChannelAuditLogRetentionHandler 依赖字典worker 进程装配用)。
依赖方向合规INV-1``audit_log_repo_factory`` infrastructure 层装配具体
依赖方向合规INV-1``audit_log_repo`` infrastructure 层装配具体
适配器 ``ChannelPersistenceAdapter``fat adapter实现
``AuditLogRepositoryPort``应用层 handler 仅依赖端口抽象不直接 import
适配器实现
适配器无状态注入 ``session_factory``不持有 session 实例通过
``_session_scope(tx)`` 按需获取 session
Args:
session_factory: 调用返回 ``AsyncSession`` 的异步上下文管理器工厂
session_factory: 可调用对象调用返回 ``AsyncSession`` 实例
通常为 ``pg_manager.AsyncSession``
cache_port: worker 进程级 CachePort 实例RedisCacheAdapter
handler 分布式锁使用
@ -1296,12 +1316,12 @@ def create_channel_audit_log_retention_handler_dependencies(
logger_port = _create_worker_logger_port()
outbox_config = _create_outbox_config()
def audit_log_repo_factory(db: AsyncSession) -> AuditLogRepositoryPort:
return ChannelPersistenceAdapter(db, outbox_config=outbox_config, logger=logger_port)
audit_log_repo: AuditLogRepositoryPort = ChannelPersistenceAdapter(
session_factory, outbox_config=outbox_config, logger=logger_port
)
return {
"session_factory": session_factory,
"audit_log_repo_factory": audit_log_repo_factory,
"audit_log_repo": audit_log_repo,
"cache_port": cache_port,
"logger": logger_port,
}
@ -1313,12 +1333,16 @@ def create_channel_content_review_retention_handler_dependencies(
) -> dict[str, Any]:
"""构造 ChannelContentReviewRetentionHandler 依赖字典worker 进程装配用)。
依赖方向合规INV-1``content_review_repo_factory`` infrastructure
依赖方向合规INV-1``content_review_repo`` infrastructure
层装配具体适配器 ``ContentReviewRepositoryAdapter``应用层 handler 仅依赖
``ContentReviewRepositoryPort`` 端口抽象不直接 import 适配器实现
适配器无状态注入 ``session_factory``不持有 session 实例通过
``_session_scope(tx)`` 按需获取 session
Args:
session_factory: 调用返回 ``AsyncSession`` 的异步上下文管理器工厂
session_factory: 可调用对象调用返回 ``AsyncSession`` 实例
通常为 ``pg_manager.AsyncSession``
cache_port: worker 进程级 CachePort 实例RedisCacheAdapter
handler 分布式锁使用
@ -1332,12 +1356,12 @@ def create_channel_content_review_retention_handler_dependencies(
logger_port = _create_worker_logger_port()
def content_review_repo_factory(db: AsyncSession) -> ContentReviewRepositoryPort:
return ContentReviewRepositoryAdapter(db, logger=logger_port)
content_review_repo: ContentReviewRepositoryPort = ContentReviewRepositoryAdapter(
session_factory, logger=logger_port
)
return {
"session_factory": session_factory,
"content_review_repo_factory": content_review_repo_factory,
"content_review_repo": content_review_repo,
"cache_port": cache_port,
"logger": logger_port,
}
@ -1349,13 +1373,17 @@ def create_channel_idempotency_cleanup_handler_dependencies(
) -> dict[str, Any]:
"""构造 ChannelIdempotencyCleanupHandler 依赖字典worker 进程装配用)。
依赖方向合规INV-1``idempotency_repo_factory`` infrastructure
依赖方向合规INV-1``idempotency_repo`` infrastructure
装配具体适配器 ``ChannelPersistenceAdapter``fat adapter实现
``IdempotencyRepositoryPort``应用层 handler 仅依赖端口抽象不直接
import 适配器实现
适配器无状态注入 ``session_factory``不持有 session 实例通过
``_session_scope(tx)`` 按需获取 session
Args:
session_factory: 调用返回 ``AsyncSession`` 的异步上下文管理器工厂
session_factory: 可调用对象调用返回 ``AsyncSession`` 实例
通常为 ``pg_manager.AsyncSession``
cache_port: worker 进程级 CachePort 实例RedisCacheAdapter
handler 分布式锁使用
@ -1369,12 +1397,12 @@ def create_channel_idempotency_cleanup_handler_dependencies(
logger_port = _create_worker_logger_port()
outbox_config = _create_outbox_config()
def idempotency_repo_factory(db: AsyncSession) -> IdempotencyRepositoryPort:
return ChannelPersistenceAdapter(db, outbox_config=outbox_config, logger=logger_port)
idempotency_repo: IdempotencyRepositoryPort = ChannelPersistenceAdapter(
session_factory, outbox_config=outbox_config, logger=logger_port
)
return {
"session_factory": session_factory,
"idempotency_repo_factory": idempotency_repo_factory,
"idempotency_repo": idempotency_repo,
"cache_port": cache_port,
"logger": logger_port,
}

View File

@ -7,7 +7,7 @@
设计原则
- framework 层组件注册中心EventBus熔断器等为应用级单例跨请求复用
- 请求级组件被驱动适配器管道用例服务由各自 factory 按请求创建不进入容器
``create_driven_adapters(db, outbox_config, redis_client, arq_pool, execution_port)``
``create_driven_adapters(db, session_factory, outbox_config, redis_client, arq_pool, execution_port)``
``InboundPipeline.create()`` / ``ControlPlanePipeline.create()`` /
``OutboundPipeline.create()``
- 容器仅负责应用级单例的注册与解析保持职责单一

View File

@ -14,7 +14,7 @@ from typing import Any, Protocol
from arq import ArqRedis
from redis.asyncio import Redis
from sqlalchemy.ext.asyncio import AsyncSession
from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker
from yuxi.channels.adapters import DrivenAdapters, create_driven_adapters
from yuxi.channels.adapters.agent_run_execution_adapter import (
@ -319,6 +319,9 @@ async def _assemble_plugin_loading_core(
# driven_adapters_factory 为请求级工厂,每次调用创建独立 AsyncSession
# 与 DrivenAdapters 实例,供 PluginHostImpl 获取端口。
# ChannelPersistenceAdapter 为无状态适配器,注入 ensure_schema.AsyncSession
# session_factory而非请求级 db 实例,通过 _session_scope(tx) 按需获取
# session。
# F-02传入 config_scope_registry.key_to_scope_map 供 RedisConfigAdapter
# 校验作用域一致性。§13.1:传入 config_scope_registry._declared_keys 供
# _assertDeclared 校验键合法性。registry 通过 register() 原地更新内部
@ -327,6 +330,7 @@ async def _assemble_plugin_loading_core(
db = ensure_schema.AsyncSession()
return create_driven_adapters(
db,
session_factory=ensure_schema.AsyncSession,
outbox_config=outbox_config,
redis_client=redis_client,
arq_pool=arq_pool,
@ -573,7 +577,7 @@ def _register_di_singletons(
core: _PluginLoadingCore,
*,
persistence_port: PersistencePort,
persistence_db: AsyncSession,
session_factory: Callable[[], AsyncSession],
queue_port: QueuePort,
config_port: ConfigPort,
stage_slot_injector: StageSlotInjector,
@ -600,16 +604,12 @@ def _register_di_singletons(
- ``DefaultContentModerationAdapter``未注册插件渠道时的兜底
审核适配器 ``create_channel_use_cases`` resolve 并注入
DispatchStage
- ``ContentReviewRepositoryAdapter``审核历史仓储适配器复用
``persistence_db`` 共享 AsyncSession ChannelPersistenceAdapter
同一会话§10.1 共享会话约定``tx`` 参数仅作为"是否自主提交"
标志位同时注册为 ``ContentReviewRepositoryPort`` 端口绑定
INV-3 契约显式化
``persistence_db`` 的生命周期由 ``HostShutdown`` 通过
``persistence_port.aclose()`` ``self._db.close()`` 统一管理
CloseablePort 契约不注册为 ``AsyncSession`` 单例请求级资源
注册为单例存在并发污染隐患
- ``ContentReviewRepositoryAdapter``审核历史仓储适配器无状态
注入 ``session_factory``不持有 session 实例可安全注册为
应用级单例``tx`` 非空时通过 ``tx.get_session()`` 复用应用层主
事务 sessionC-I1 透传``tx`` ``None`` 时通过
``session_factory`` 创建独立 session 并自主提交同时注册为
``ContentReviewRepositoryPort`` 端口绑定INV-3 契约显式化
"""
# 1. 注册中心与共享依赖
di_container.registerSingleton(PluginRegistry, core.plugin_registry)
@ -685,6 +685,10 @@ def _register_di_singletons(
di_container.registerSingleton(IdentityResolverRegistry, identity_resolver_registry)
# 3. 应用级端口与编排组件
# session_factoryasync_sessionmaker 实例)注册为 DI 单例,供
# create_channel_use_cases 解析并注入到请求级 create_driven_adapters
# 使 ChannelPersistenceAdapter 通过 _session_scope(tx) 按需创建独立 session。
di_container.registerSingleton(async_sessionmaker, session_factory)
di_container.registerSingleton(StageSlotInjector, stage_slot_injector)
# HealthAggregator 预注册HostBootstrap._registerCoreSingletons 也会注册,
# 此处预注册保证 bootstrap 失败后 DI 容器状态完整(与 EventBus 等同理)。
@ -706,7 +710,7 @@ def _register_di_singletons(
logger=logger,
)
review_repository = ContentReviewRepositoryAdapter(
db=persistence_db,
session_factory=session_factory,
logger=logger,
)
di_container.registerSingleton(DefaultContentModerationAdapter, default_moderation_adapter)
@ -737,8 +741,9 @@ async def create_host_bootstrap(
Args:
ensure_schema: 渠道 schema 初始化端口``PostgresManager`` 满足此
Protocol提供 ``ensure_channel_schema`` 方法同时作为数据库
会话工厂来源通过 ``ensure_schema.AsyncSession()`` 创建应用级
与请求级 ``AsyncSession``
会话工厂来源``ensure_schema.AsyncSession``可调用对象注入到
``ChannelPersistenceAdapter`` / ``ContentReviewRepositoryAdapter``
适配器按需调用创建独立 ``AsyncSession``不在应用级持有共享会话
logger: 日志被驱动端口供所有 framework 层组件注入
plugin_dir: 插件目录路径 ``PluginLifecycleManager.discover`` 扫描
di_container: 依赖注入容器由调用方创建并传入HostBootstrap
@ -763,11 +768,11 @@ async def create_host_bootstrap(
core = await _assemble_plugin_loading_core(logger, ensure_schema, plugin_dir=plugin_dir)
# 2. 构造应用级 persistence_port 与 queue_port
# persistence_port 需要 AsyncSession创建应用级会话供 TransportManager /
# HealthAggregator / DiagnosticsExporter 共用queue_port 无需会话。
persistence_db = ensure_schema.AsyncSession()
# persistence_port 注入 session_factoryCallable[[], AsyncSession]
# 适配器无状态,每次方法调用通过 _session_scope(tx) 按需获取 session
# queue_port 无需会话。
persistence_port: PersistencePort = ChannelPersistenceAdapter(
persistence_db, outbox_config=core.outbox_config, logger=logger
ensure_schema.AsyncSession, outbox_config=core.outbox_config, logger=logger
)
queue_port: QueuePort = ARQQueueAdapter(arq_pool=core.arq_pool, logger=logger)
@ -867,7 +872,7 @@ async def create_host_bootstrap(
di_container=di_container,
core=core,
persistence_port=persistence_port,
persistence_db=persistence_db,
session_factory=ensure_schema.AsyncSession,
queue_port=queue_port,
config_port=config_port,
stage_slot_injector=stage_slot_injector,
@ -881,6 +886,8 @@ async def create_host_bootstrap(
# 供 ``_loadConfig`` 加载渠道配置、``_initDrivenAdapters`` 对应用级
# 被驱动适配器执行连通性检查driving_adapter_unregistrar 用于启动
# 失败回滚时注销驱动适配器)
# persistence_port 保留在元组中供启动期 ping 连通性检查PingablePort
# 但不再实现 CloseablePort无状态适配器HostShutdown 关停时跳过 close。
app_driven_adapters = (
core.cache_port,
persistence_port,
@ -923,8 +930,10 @@ def create_host_shutdown(
INF-017同时从 DI 容器解析 4 个应用级被驱动适配器
``cache_port`` / ``persistence_port`` / ``queue_port`` / ``config_port``
传入 ``app_driven_adapters`` 参数 ``HostShutdown`` 在关停步骤 6
显式调用 ``close()`` 释放资源
传入 ``app_driven_adapters`` 参数 ``HostBootstrap`` 启动期执行
ping 连通性检查``HostShutdown`` 关停期对实现 ``CloseablePort``
适配器调用 ``close()`` 释放资源``persistence_port`` 为无状态适配器
不持有 session关停时通过 ``isinstance(CloseablePort)`` 跳过
INF-006依赖来源由模块级全局状态改为 DI 容器
消除模块级可变状态INV-5
@ -947,6 +956,8 @@ def create_host_shutdown(
执行关停编排
"""
# 解析应用级被驱动适配器HostBootstrap 构造时注册到 DI 容器)
# persistence_port 保留在元组中(启动期 ping 检查需要),关停时由
# HostShutdown 通过 isinstance(CloseablePort) 跳过(无状态适配器无需 close
app_driven_adapters = (
di_container.resolve(CachePort),
di_container.resolve(PersistencePort),

View File

@ -688,10 +688,12 @@ class HostBootstrap:
async def _initDrivenAdapters(self) -> None:
"""加载并初始化被驱动适配器INF-010应用级适配器连通性检查
被驱动适配器为请求级组件 ``create_driven_adapters(db)`` 按请求创建
共享同一 ``db`` session本步骤对 ``HostBootstrap`` 构造时传入的应用级
被驱动适配器``cache_port`` / ``persistence_port`` / ``queue_port`` /
``config_port`` 等跨请求共享实例执行连通性检查确保下游依赖可达
被驱动适配器为请求级组件 ``create_driven_adapters(db, session_factory, ...)``
按请求创建其中 ``conversation`` / ``transaction`` 共享同一 ``db`` session
``persistence`` 为无状态适配器注入 ``session_factory``本步骤对
``HostBootstrap`` 构造时传入的应用级被驱动适配器``cache_port`` /
``persistence_port`` / ``queue_port`` / ``config_port`` 等跨请求共享实例
执行连通性检查确保下游依赖可达
任一失败必须抛 ``DependencyError`` 终止启动INV-I5避免启动后请求期
才发现依赖故障

View File

@ -442,18 +442,21 @@ class HostShutdown:
应用级被驱动适配器``cache_port`` / ``persistence_port`` /
``queue_port`` / ``config_port`` ``HostBootstrap`` 构造时创建
跨请求共享需在关停时显式关闭以释放 Redis 连接DB 会话等资源
插件被驱动适配器由 ``PluginRegistry`` 管理通过
``listPluginAdapters`` 遍历关闭
跨请求共享需在关停时显式关闭以释放 Redis 连接等资源
``persistence_port`` 为无状态适配器不持有 session不实现
``CloseablePort``遍历时跳过插件被驱动适配器由 ``PluginRegistry``
管理通过 ``listPluginAdapters`` 遍历关闭
单个适配器关闭异常记录告警并继续确保后续适配器仍能关闭避免
资源泄漏INV-9 回退路径完整性
"""
# 1. 关闭应用级被驱动适配器cache / persistence / queue / config
# persistence_port 为无状态适配器(不持有 session不实现
# CloseablePortisinstance 检查为 False 时跳过DEBUG 级别记录)。
for adapter in self._app_driven_adapters:
try:
if not isinstance(adapter, CloseablePort):
await self._logger.info(
await self._logger.debug(
"应用级被驱动适配器未实现 CloseablePort跳过",
adapter_type=type(adapter).__name__,
)

View File

@ -92,7 +92,11 @@ async def register_scheduler_handlers(registry: HandlerRegistry) -> None:
create_worker_cache_port,
)
session_factory = pg_manager.get_async_session_context
# ``pg_manager.AsyncSession`` 为 ``async_sessionmaker`` 实例callable
# 调用返回 ``AsyncSession``),符合 ``Callable[[], AsyncSession]`` 协议。
# 由 ``run_worker._worker_startup`` 在调用本函数前通过 ``pg_manager.initialize()``
# 初始化完成。
session_factory = pg_manager.AsyncSession
arq_pool = await get_arq_pool()
cache_port = await create_worker_cache_port()

View File

@ -76,8 +76,8 @@ _GROUP_TALKER_SUFFIX = "@chatroom"
# 拒绝含空格、斜杠、中文等异常字符的 talker/sender防止 bridge 数据污染下游。
_WXID_RE = re.compile(r"^[A-Za-z0-9_@-]+$")
# P2-12: session_type 合法枚举值bridge 返回的会话类型
_VALID_SESSION_TYPES = frozenset({"single", "group"})
# P2-12: session_type 合法枚举值bridge 返回的会话类型,与 session_adapter 的 chat_type 保持一致
_VALID_SESSION_TYPES = frozenset({"p2p", "group"})
class WeChatWocInboundAdapter:

View File

@ -643,9 +643,12 @@ class WocBridgeClient:
if last_exc is not None:
raise last_exc
except httpx.HTTPError as exc:
await self._logger.exception(
# 仅记录关键诊断字段,不打印原始 httpx 堆栈(易误判为异常穿透)。
# 翻译后的契约层异常OperationTimeoutError/DependencyError 等)由
# translate_call_error 抛出,调用方 catch 时自行记录完整堆栈。
await self._logger.error(
"woc bridge HTTP request failed",
exc_info=exc,
error_type=type(exc).__name__,
method=method,
url=_redact_url(url),
trace_id=trace_id,