113 lines
7.8 KiB
Python
113 lines
7.8 KiB
Python
"""ORM model for roundtable background-orchestration jobs.
|
||
|
||
A "job" = one **backend** run of the roundtable orchestration pipeline kicked off
|
||
from Step 2「后台挂起」: it drives the leader→dispatch→sub-agents loop to consensus
|
||
and then Step 3 report drawing, entirely server-side, surviving page refresh/close.
|
||
|
||
Bound 1:1 to a ``roundtable_drafts`` row (``draft_id``) so the history dropdown can
|
||
show "this task is running" and the dispatch-chain modal can render live progress.
|
||
|
||
Big JSON columns (``dispatch_chain`` / ``thread_ids`` / ``chain``) are stored as
|
||
serialized strings in ``PortableLongText`` (MySQL ``LONGTEXT``) — same approach as
|
||
``ai_writing_sessions`` / ``roundtable_drafts`` — to avoid MySQL's 64 KB ``TEXT`` cap.
|
||
|
||
See ``frontend-web/docs/multi-agent-background-run-dev.md`` §2.5 for the field
|
||
contract and §Phase 1 for the build steps.
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
from datetime import UTC, datetime
|
||
|
||
import sqlalchemy as sa
|
||
from sqlalchemy import Boolean, Integer, String, Text, UniqueConstraint
|
||
from sqlalchemy.orm import Mapped, mapped_column
|
||
|
||
from deerflow.persistence.base import Base
|
||
from deerflow.persistence.types import BeijingDateTime, PortableLongText
|
||
|
||
|
||
class RoundtableJobRow(Base):
|
||
__tablename__ = "roundtable_jobs"
|
||
__table_args__ = (
|
||
# 「同一 scope 至多一个活跃作业」由该唯一键强制(Phase 3):活跃时存
|
||
# sha256(draft_id + task_id|user_id),进入终态时清 NULL;MySQL/SQLite 唯一索引
|
||
# 都允许多个 NULL,故终态历史作业不占键。INSERT 冲突即幂等命中 —— 一致性来自
|
||
# DB 唯一键,不再依赖 GET_LOCK 命名锁。
|
||
UniqueConstraint("active_dedupe_key", name="uq_roundtable_jobs_active_dedupe_key"),
|
||
)
|
||
|
||
id: Mapped[str] = mapped_column(String(64), primary_key=True)
|
||
# The draft/session this job runs for. Indexed for "active job of draft" lookups.
|
||
draft_id: Mapped[str | None] = mapped_column(String(64), nullable=True, index=True)
|
||
user_id: Mapped[str | None] = mapped_column(String(64), nullable=True, index=True)
|
||
# 外部任务 id(taskId 深链会商)。非空 = 这是一条 task 作业:读取按 task_id 共享、
|
||
# **不按 user 分权**(任何人打开同一 task 都能看在跑的研讨);其 draft_id 指向独立的
|
||
# ``roundtable_task_drafts`` 表。为空 = 普通个人作业(按 user 分权)。
|
||
task_id: Mapped[str | None] = mapped_column(String(128), nullable=True, index=True)
|
||
# ── Phase 3:唯一活跃键 + 租约 + dispatcher ──────────────────────────────
|
||
# 前端启动请求的幂等键(同键重试返回同一作业)与入参指纹(审计/排查)。
|
||
request_id: Mapped[str | None] = mapped_column(String(64), nullable=True)
|
||
request_hash: Mapped[str | None] = mapped_column(String(64), nullable=True)
|
||
# 活跃去重键:活跃(非终态)时 = sha256(draft_id + task_id|user_id),终态清 NULL。
|
||
# 唯一索引(__table_args__)保证「同一 scope 至多一个活跃作业」——INSERT 冲突即幂等。
|
||
active_dedupe_key: Mapped[str | None] = mapped_column(String(64), nullable=True)
|
||
# 启动入参完整 JSON 快照:dispatcher 在任意 worker 领取 queued/租约过期作业时
|
||
# 据此重建 StartParams(跨进程可接管的前提;resume 续跑也重写本列)。
|
||
input_snapshot: Mapped[str | None] = mapped_column(PortableLongText, nullable=True)
|
||
# 租约:条件 UPDATE 抢占(queued 或租约过期)+ 心跳续租;worker 崩溃后租约自然过期,
|
||
# 其它 worker 接管。恢复能力来自 DB 而非进程内存。
|
||
lease_owner: Mapped[str | None] = mapped_column(String(64), nullable=True)
|
||
lease_until: Mapped[datetime | None] = mapped_column(BeijingDateTime(), nullable=True)
|
||
# 被 dispatcher 领取执行的次数(接管计数)。
|
||
attempt: Mapped[int] = mapped_column(Integer, nullable=False, default=0, server_default=sa.text("0"))
|
||
# 行级 CAS 版本:条件 UPDATE 乐观锁,防止并发写互相覆盖。
|
||
version: Mapped[int] = mapped_column(Integer, nullable=False, default=0, server_default=sa.text("0"))
|
||
# Job status machine (§2.1): queued / running / awaiting_input / done / error / cancelled.
|
||
status: Mapped[str] = mapped_column(String(32), nullable=False, default="queued", index=True)
|
||
# Fine-grained phase (§2.2): initializing / leader_thinking / dispatching /
|
||
# consensus / report_drawing / report_done.
|
||
phase: Mapped[str] = mapped_column(String(32), nullable=False, default="initializing")
|
||
# Current orchestration cycle (0..MAX_CYCLES).
|
||
cycle: Mapped[int] = mapped_column(Integer, nullable=False, default=0)
|
||
# Dispatch chain nodes (§2.3): JSON array of
|
||
# { agentId, name, avatarType, order, state, lastSpokeCycle? }.
|
||
dispatch_chain: Mapped[str | None] = mapped_column(PortableLongText, nullable=True)
|
||
# Currently-speaking seat agent_id (null = coordinator / none).
|
||
active_agent_id: Mapped[str | None] = mapped_column(String(128), nullable=True)
|
||
consensus_percentage: Mapped[int] = mapped_column(Integer, nullable=False, default=0)
|
||
# When status=awaiting_input: the coordinator's clarification question text.
|
||
pending_clarification: Mapped[str | None] = mapped_column(Text, nullable=True)
|
||
# ThreadIdMap: JSON object { "<agent_id>": "<thread_id>", ... }.
|
||
thread_ids: Mapped[str | None] = mapped_column(PortableLongText, nullable=True)
|
||
coordinator_name: Mapped[str | None] = mapped_column(String(128), nullable=True)
|
||
# "recommend"(总控自由派活)| "chain"(按业务链条顺序派活)| "dag"(分层 DAG:stage 间串行、stage 内并行).
|
||
orchestration_mode: Mapped[str] = mapped_column(String(16), nullable=False, default="recommend")
|
||
# chain 模式来源链条快照:JSON object { id, title }(recommend 模式为 null).
|
||
chain: Mapped[str | None] = mapped_column(PortableLongText, nullable=True)
|
||
# dag 模式分层编排计划:JSON object { mode, stages:[{id,agentIds,goal?}], finalSynthesis,
|
||
# coordinatorPrompt? }(对齐前端 OrchestrationPlan)。仅 dag 模式有值;驱动后端 _run_dag
|
||
# (stage 串行 / stage 内并行)+ 前端大弹窗按 stage 分层画流程图 + resume 时重建计划。
|
||
orchestration_plan: Mapped[str | None] = mapped_column(PortableLongText, nullable=True)
|
||
# 研讨对话记录:JSON array of { id, role('leader'|'seat'), agentId, name, content, cycle }。
|
||
# 让用户跑完 / 进行中都能看到真实研讨内容(leader 说明 + 各席位交付),而非只有流程图。
|
||
dialogues: Mapped[str | None] = mapped_column(PortableLongText, nullable=True)
|
||
# Step3 报告快照:JSON object { html, generatedAt, model, summary }。done 时写入,
|
||
# 使前端从单个 job 即可拿到完整结果(对话 + 报告)。
|
||
step3: Mapped[str | None] = mapped_column(PortableLongText, nullable=True)
|
||
# Terminal error detail when status=error.
|
||
error: Mapped[str | None] = mapped_column(Text, nullable=True)
|
||
# 草稿回写 pending 标记:``_write_run_to_draft`` 乐观锁重试耗尽后置 1(绝不退化
|
||
# last-write-wins 盲写);由 outbox 补偿器投影结果到草稿后清 0。一致性来自 DB 标记。
|
||
draft_persist_pending: Mapped[bool] = mapped_column(
|
||
Boolean, nullable=False, default=False, server_default=sa.false()
|
||
)
|
||
# 时间列必须用 BeijingDateTime(项目规范,与 ai_writing_sessions / roundtable_drafts 对齐)。
|
||
created_at: Mapped[datetime] = mapped_column(BeijingDateTime(), nullable=False, default=lambda: datetime.now(UTC))
|
||
updated_at: Mapped[datetime] = mapped_column(
|
||
BeijingDateTime(),
|
||
nullable=False,
|
||
default=lambda: datetime.now(UTC),
|
||
onupdate=lambda: datetime.now(UTC),
|
||
)
|