refactor(storage/postgres): 清理废弃字段与冗余模型

1. 移除ChannelAccount表的capabilities、config_version、last_message_at字段
2. 扩大service_user_uid字段长度至225
3. 移除AuditExportTask模型并清理相关代码
4. 更新ChannelOutboxEntry、ChannelIdempotency、ChannelContentReviewRecord的注释与索引
5. 为ChannelPairing添加SQLAlchemy乐观锁支持
6. 新增ChannelContentReviewRecord的决策时长统计字段
This commit is contained in:
Kris 2026-07-08 03:58:34 +08:00
parent ff07cef881
commit 16f7e2d6ab

View File

@ -71,9 +71,8 @@ class ChannelAccount(Base):
account_id = Column(String(128), nullable=False, comment="渠道账户 ID全局唯一业务标识")
display_name = Column(String(128), nullable=False, comment="展示名称")
# 配置与能力
# 配置
config = Column(JSON_VALUE, nullable=False, default=dict, comment="渠道配置(敏感字段加密存储)")
capabilities = Column(JSON_VALUE, nullable=False, default=dict, comment="能力声明ChannelCapabilities 序列化)")
# 状态
enabled = Column(
@ -86,7 +85,6 @@ class ChannelAccount(Base):
default="stopped",
comment="插件生命周期状态FR-32discovered/parsed/loaded/initialized/starting/running/paused/resuming/stopping/stopped/unloaded/failed",
)
config_version = Column(Integer, nullable=False, default=1, comment="配置版本号FR-37 热更新冲突检测)")
# 接入态状态机(与运行态 status 解耦pending/configured/verified/online/offline/failed
onboarding_status = Column(
String(32),
@ -99,13 +97,12 @@ class ChannelAccount(Base):
credential_ref = Column(String(128), nullable=True, comment="凭证引用,指向凭证存储中的版本化凭证记录")
credential_version = Column(Integer, nullable=False, default=0, server_default="0", comment="凭证版本号")
# 关联服务账号user_type='service'
service_user_uid = Column(String(64), nullable=True, comment="关联服务账号 User.uid")
service_user_uid = Column(String(225), nullable=True, comment="关联服务账号 User.uid")
last_error = Column(Text, nullable=True, comment="最近一次错误信息FR-35 健康检查详情)")
last_error_at = Column(DateTime, nullable=True, comment="最近一次错误时间")
last_health_check_at = Column(
DateTime, nullable=True, comment="最近一次健康检查时间FR-35无论成功失败均更新用于诊断检查任务是否存活"
)
last_message_at = Column(DateTime, nullable=True, comment="最近一次消息时间(冗余字段)")
transport_cursor = Column(
String(256), nullable=False, default="", server_default="", comment="传输游标Puller类型使用"
)
@ -134,11 +131,9 @@ class ChannelAccount(Base):
"account_id": self.account_id,
"display_name": self.display_name,
"config": self.config or {},
"capabilities": self.capabilities or {},
"enabled": bool(self.enabled),
"status": self.status,
"plugin_status": self.plugin_status,
"config_version": self.config_version,
"onboarding_status": self.onboarding_status,
"credential_ref": self.credential_ref,
"credential_version": self.credential_version,
@ -146,7 +141,6 @@ class ChannelAccount(Base):
"last_error": self.last_error,
"last_error_at": format_utc_datetime(self.last_error_at),
"last_health_check_at": format_utc_datetime(self.last_health_check_at),
"last_message_at": format_utc_datetime(self.last_message_at),
"transport_cursor": self.transport_cursor or "",
"last_rotated_at": format_utc_datetime(self.last_rotated_at),
"version": self.version,
@ -403,6 +397,11 @@ class ChannelPairing(Base):
),
)
# 启用 SQLAlchemy 乐观锁ORM 实例 flush 时自动追加
# ``WHERE id = ? AND version = ?`` 并自增 version并发状态转换冲突时
# 抛 ``StaleDataError``,由适配器层翻译为 ``ConflictError``。
__mapper_args__ = {"version_id_col": version}
def to_dict(self) -> dict[str, Any]:
"""序列化为字典。"""
return {
@ -506,69 +505,6 @@ class ChannelAuditLog(Base):
}
class AuditExportTask(Base):
"""审计日志导出任务。
对应聚合根 ``AuditExportTask````yuxi.channels.core.model.audit_export_task``
持久化异步导出任务的状态与结果配合 ARQ worker 实现审计日志的异步导出
任务状态机pending running completed | failed终态不可变更
"""
__tablename__ = "audit_export_tasks"
# 主键
id = Column(BigInteger, primary_key=True, autoincrement=True)
# 业务标识
task_id = Column(String(64), nullable=False, unique=True, index=True, comment="导出任务 IDUUID业务标识")
status = Column(
String(32), nullable=False, index=True, comment="任务状态pending / running / completed / failed / cancelled"
)
# 导出参数与结果
query_snapshot = Column(JSON_VALUE, nullable=False, comment="查询条件快照AuditQuery 序列化)")
format = Column(String(16), nullable=False, default="json", comment="导出格式(当前仅支持 json")
total_records = Column(Integer, nullable=True, comment="导出记录总数completed 状态时填充)")
file_path = Column(String(512), nullable=True, comment="导出报告文件路径completed 状态时填充)")
partial_file_path = Column(String(512), nullable=True, comment="部分导出文件路径cancelled 状态时填充)")
error = Column(Text, nullable=True, comment="失败原因failed 状态时填充)")
# 时间与操作人
created_by = Column(String(64), nullable=False, index=True, comment="创建者用户 ID")
created_at = Column(DateTime(timezone=True), nullable=False, comment="任务创建时间UTC")
completed_at = Column(
DateTime(timezone=True), nullable=True, comment="任务完成时间completed / failed 状态时填充)"
)
# 审计字段
is_deleted = Column(Integer, nullable=False, default=0)
deleted_at = Column(DateTime(timezone=True), nullable=True)
__table_args__ = (
Index("ix_audit_export_tasks_status_time", "status", "created_at"),
Index("ix_audit_export_tasks_created_by_time", "created_by", "created_at"),
)
def to_dict(self) -> dict[str, Any]:
"""序列化为字典。"""
return {
"id": self.id,
"task_id": self.task_id,
"status": self.status,
"query_snapshot": self.query_snapshot or {},
"format": self.format,
"total_records": self.total_records,
"file_path": self.file_path,
"partial_file_path": self.partial_file_path,
"error": self.error,
"created_by": self.created_by,
"created_at": format_utc_datetime(self.created_at),
"completed_at": format_utc_datetime(self.completed_at),
"is_deleted": self.is_deleted,
"deleted_at": format_utc_datetime(self.deleted_at),
}
class ChannelOutboxEntry(Base):
"""持久化投递队列,存储出站消息的持久化投递状态。
@ -576,10 +512,10 @@ class ChannelOutboxEntry(Base):
仅存消息引用不存完整消息内容FR-22conversation_id / channel_session_id
persistence 适配器从管道上下文注入聚合根无这两个字段
并发投递锁 worker 拉取 pending 消息时通过 ``locked_by`` / ``locked_at``
实现显式租约worker 取走消息后写入 ``locked_by`` ``locked_at``处理完成
后清空并更新 ``status`` worker 崩溃``locked_at`` 超时后由回收任务重置
锁并重新投递保证 at-least-once 语义
并发投递锁已废弃``locked_by`` / ``locked_at`` 列保留兼容历史数据
但并发控制已由 Redis advisory lock``OutboxRetryWorker.acquireAdvisoryLock``
实现DB 级租约机制不再使用``delivery_mode`` / ``ack_policy`` 列保留默认值
投递模式为运行时管道决策不持久化``ack_policy`` 为入站概念不在出站流程读取
"""
__tablename__ = "channel_outbox_entries"
@ -723,6 +659,13 @@ class ChannelOutboxEntry(Base):
"last_error": self.last_error,
"channel_msg_id": self.channel_msg_id,
"sent_at": format_utc_datetime(self.sent_at),
"latency_ms": self.latency_ms,
"funnel_node": self.funnel_node,
"idempotency_key": self.idempotency_key,
"channel_request_id": self.channel_request_id,
"partial_failure": self.partial_failure,
"stream_aborted_at_chunk": self.stream_aborted_at_chunk,
"degraded_reason": self.degraded_reason,
"expires_at": format_utc_datetime(self.expires_at),
"trace_id": self.trace_id,
"version": self.version,
@ -786,6 +729,9 @@ class ChannelRouteBinding(Base):
"match_value",
unique=True,
postgresql_where=text("is_deleted = 0"),
# PG15+NULL 视为相等,使 account/default 等 match_value 为 NULL 的
# tier 也能在 (channel_type, account_id, match_source) 维度上唯一。
postgresql_nulls_not_distinct=True,
),
Index(
"ix_channel_route_bindings_list",
@ -800,7 +746,8 @@ class ChannelIdempotency(Base):
"""写操作幂等记录,支持重复请求重放首次响应。
不复用现有 ``scheduled_task_idempotency`` 限界上下文隔离
is_deleted 恒为 0不设 version 字段过期后由定时任务硬删除
is_deleted 恒为 0不设 version 字段过期后由定时任务
``ChannelIdempotencyCleanupHandler`` 硬删除
幂等处理流程
1. 请求到达 ``idempotency_key`` 查询记录
@ -808,11 +755,12 @@ class ChannelIdempotency(Base):
持有处理权处理完成后更新为 ``status='completed'`` 并写入 ``response_body``
处理失败则更新为 ``status='failed'``可由后续请求覆盖重试
3. 存在且 ``status='in_progress'``重复请求应拒绝409 Conflict或阻塞等待
避免并发执行同一操作
4. 存在且 ``status='completed'``校验 ``request_hash``匹配则重放
``response_body``不匹配说明同 key 不同参数拒绝409 Conflict以防
掩盖客户端 bug
5. 已过期``expires_at < now``硬删除后按新请求处理
避免并发执行同一操作超时的 in_progress超过 ``in_progress_started_at``
阈值视为崩溃遗留删除重建
4. 存在且 ``status='completed'``重放 ``response_body``管理员消息路径
或拒绝重复创建 AgentRun入站消息路径
5. 过期记录由定时任务 ``ChannelIdempotencyCleanupHandler`` ``expires_at``
硬删除不阻塞请求处理
"""
__tablename__ = "channel_idempotency"
@ -825,18 +773,7 @@ class ChannelIdempotency(Base):
operation = Column(
String(32),
nullable=False,
index=True,
comment="操作类型admin_message_send/pairing_approve/pairing_reject/config_update/whitelist_modify/session_merge",
)
account_id = Column(
Integer,
ForeignKey("channel_accounts.id"),
nullable=True,
index=True,
comment="关联渠道账户(可选,部分操作如全局配置变更无账户)",
)
request_hash = Column(
String(64), nullable=True, comment="请求参数 SHA256用于检测同 key 不同参数的冲突;无参数操作可留空)"
comment="操作类型inbound_message入站消息/ message/send管理员消息发送",
)
status = Column(
String(16), nullable=False, default="in_progress", index=True, comment="处理状态in_progress/completed/failed"
@ -862,11 +799,7 @@ class ChannelIdempotency(Base):
is_deleted = Column(Integer, nullable=False, default=0, index=True)
deleted_at = Column(DateTime, nullable=True)
__table_args__ = (
Index("ix_channel_idempotency_created_at", "created_at"),
Index("ix_channel_idempotency_expires_at", "expires_at"),
Index("ix_channel_idempotency_status_started", "status", "in_progress_started_at"),
)
__table_args__ = (Index("ix_channel_idempotency_expires_at", "expires_at"),)
def to_dict(self) -> dict[str, Any]:
"""序列化为字典。"""
@ -874,8 +807,6 @@ class ChannelIdempotency(Base):
"id": self.id,
"idempotency_key": self.idempotency_key,
"operation": self.operation,
"account_id": self.account_id,
"request_hash": self.request_hash,
"status": self.status,
"in_progress_started_at": format_utc_datetime(self.in_progress_started_at),
"response_body": self.response_body,
@ -930,11 +861,17 @@ class ChannelContentReviewRecord(Base):
source = Column(
String(32),
nullable=False,
index=True,
comment="审核来源manual_preview / inbound_pipeline / outbound_pipeline",
)
trace_id = Column(String(64), nullable=True, comment="链路追踪 ID")
# 决策时长统计CR-STATS-01 avg_decision_seconds 真实化)
# decision_started_at审核请求开始时间首次写入时与 reviewed_at 同值)
# decision_completed_at审核结论最终确定时间人工 update_verdict 时刷新)
# 仅对有 completed_at 的记录参与 avg 计算(首次审核即终态的不计入)
decision_started_at = Column(DateTime, nullable=True, comment="审核决策开始时间UTC")
decision_completed_at = Column(DateTime, nullable=True, comment="审核决策完成时间UTC人工覆盖时刷新")
# 审计字段(与现有渠道模型统一)
created_by = Column(String(64), nullable=True)
updated_by = Column(String(64), nullable=True)
@ -944,6 +881,8 @@ class ChannelContentReviewRecord(Base):
deleted_at = Column(DateTime, nullable=True)
# 复合索引(按渠道+账户查询、按时间范围查询)
# 注source 不作为查询过滤条件(仅聚合 case when 使用),不建索引;
# reviewed_at 是高频时间范围扫描字段,单列索引直接支持 stats/analytics 趋势分桶。
__table_args__ = (
Index(
"ix_channel_content_review_channel_account",
@ -951,8 +890,7 @@ class ChannelContentReviewRecord(Base):
"account_id",
),
Index(
"ix_channel_content_review_source_reviewed",
"source",
"ix_channel_content_review_reviewed_at",
"reviewed_at",
),
)
@ -974,6 +912,8 @@ class ChannelContentReviewRecord(Base):
"reviewer": self.reviewer,
"source": self.source,
"trace_id": self.trace_id,
"decision_started_at": format_utc_datetime(self.decision_started_at),
"decision_completed_at": format_utc_datetime(self.decision_completed_at),
"created_by": self.created_by,
"updated_by": self.updated_by,
"created_at": format_utc_datetime(self.created_at),
@ -981,71 +921,3 @@ class ChannelContentReviewRecord(Base):
"is_deleted": self.is_deleted,
"deleted_at": format_utc_datetime(self.deleted_at),
}
class ChannelReport(Base):
"""渠道报告 ORM 模型。
存储一次性报告的元数据与内容报告内容由 scheduler handler
异步生成后写入 ``content`` 字段JSONB
设计依据对齐19-聚合视图域-reports-router-设计方案.md§8.4
"""
__tablename__ = "channel_reports"
# 主键
id = Column(Integer, primary_key=True, autoincrement=True)
# 业务标识
report_id = Column(String(64), nullable=False, unique=True, index=True, comment="报告 IDUUID业务标识")
task_id = Column(String(64), nullable=False, unique=True, index=True, comment="关联调度任务 ID业务标识")
report_type = Column(String(32), nullable=False, comment="报告类型")
status = Column(String(16), nullable=False, default="pending", comment="状态pending/generating/ready/failed")
# 内容与参数
params = Column(JSON_VALUE, nullable=False, default=dict, comment="报告生成参数(序列化)")
content = Column(JSON_VALUE, nullable=True, comment="报告内容ready 状态时填充JSONB")
error_message = Column(String(512), nullable=True, comment="失败原因failed 状态时填充)")
# 时间与来源
created_at = Column(DateTime, nullable=False, default=utc_now_naive, comment="创建时间")
ready_at = Column(DateTime, nullable=True, comment="报告就绪时间ready 状态时填充)")
created_by = Column(String(64), nullable=False, default="system", comment="创建者(用户 ID 或 system")
# 重试关联RPT-ONEOFF-RETRY
retried_from = Column(String(64), nullable=True, comment="重试来源报告的 task_id仅重试生成的新报告非空")
retried_at = Column(DateTime, nullable=True, comment="原报告被重试的时间(原报告标记字段)")
# 审计字段(与 channels 其他表一致)
updated_by = Column(String(64), nullable=True)
updated_at = Column(DateTime, nullable=False, default=utc_now_naive, onupdate=utc_now_naive)
is_deleted = Column(Integer, nullable=False, default=0, index=True)
deleted_at = Column(DateTime, nullable=True)
__table_args__ = (
Index("ix_channel_reports_status_created", "status", "created_at"),
Index("ix_channel_reports_type_created", "report_type", "created_at"),
)
def to_dict(self) -> dict[str, Any]:
"""序列化为字典。"""
return {
"id": self.id,
"report_id": self.report_id,
"task_id": self.task_id,
"report_type": self.report_type,
"status": self.status,
"params": self.params or {},
"content": self.content,
"error_message": self.error_message,
"created_at": format_utc_datetime(self.created_at),
"ready_at": format_utc_datetime(self.ready_at),
"created_by": self.created_by,
"retried_from": self.retried_from,
"retried_at": format_utc_datetime(self.retried_at),
"updated_by": self.updated_by,
"updated_at": format_utc_datetime(self.updated_at),
"is_deleted": self.is_deleted,
"deleted_at": format_utc_datetime(self.deleted_at),
}