deerflow-code/offline-backend-20260512/backend/packages/harness/deerflow/persistence/roundtable_jobs/model.py
2026-09-07 18:24:55 +08:00

113 lines
7.8 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

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