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

163 lines
9.3 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.

"""Persistence models for the standalone position-collaboration workspace.
This deliberately does **not** reuse ``roundtable_drafts`` or
``roundtable_jobs``. A session captures the selected business chain and the
confirmed task intent; every chain seat then has one durable node record.
LangGraph threads remain the execution source, while these tables also retain
database-owned recovery copies of conversations and text deliverables so the
workspace survives replacement of container-local checkpoint/output volumes.
"""
from __future__ import annotations
from datetime import UTC, datetime
from sqlalchemy import Boolean, ForeignKey, Index, Integer, String, UniqueConstraint
from sqlalchemy.orm import Mapped, mapped_column
from deerflow.persistence.base import Base
from deerflow.persistence.types import BeijingDateTime, PortableLongText
class PositionRoundtableSessionRow(Base):
__tablename__ = "position_roundtable_sessions"
__table_args__ = (
# History is read as "one user's conversations for one task, newest
# first". `user_id` partitions visibility; a shared task workspace no
# longer exposes other users' collaboration records.
Index(
"ix_position_roundtable_sessions_task_updated",
"external_task_id",
"updated_at",
),
)
id: Mapped[str] = mapped_column(String(64), primary_key=True)
# Owner field: every read/write path filters on it. `task_scoped` still
# marks rows created from an explicit route taskId (locking the snapshot
# task id and keeping them out of the unscoped personal list), but it no
# longer grants cross-user access.
user_id: Mapped[str] = mapped_column(String(64), nullable=False, index=True)
external_task_id: Mapped[str | None] = mapped_column(String(128), nullable=True, index=True)
task_scoped: Mapped[bool] = mapped_column(Boolean, nullable=False, default=False, server_default="0")
task_snapshot: Mapped[str | None] = mapped_column(PortableLongText, nullable=True)
chain_id: Mapped[str | None] = mapped_column(String(64), nullable=True, index=True)
chain_snapshot: Mapped[str | None] = mapped_column(PortableLongText, nullable=True)
intent_thread_id: Mapped[str | None] = mapped_column(String(128), nullable=True)
intent_snapshot: Mapped[str | None] = mapped_column(PortableLongText, nullable=True)
# Database-owned recovery copy for the complete Step-1 conversation. The
# normal LangGraph checkpoint remains the execution source, while this
# snapshot lets the position workspace recover after a checkpoint volume
# or container-local store has been replaced.
intent_conversation: Mapped[str | None] = mapped_column(PortableLongText, nullable=True)
# The two built-in closing agents own normal LangGraph threads too. Keep
# their ids on the position session so a refresh can continue the same
# report/planning dialogue instead of silently starting a fresh context.
summary_thread_id: Mapped[str | None] = mapped_column(String(128), nullable=True)
summary_snapshot: Mapped[str | None] = mapped_column(PortableLongText, nullable=True)
summary_conversation: Mapped[str | None] = mapped_column(PortableLongText, nullable=True)
action_plan_thread_id: Mapped[str | None] = mapped_column(String(128), nullable=True)
action_plan_snapshot: Mapped[str | None] = mapped_column(PortableLongText, nullable=True)
action_plan_conversation: Mapped[str | None] = mapped_column(PortableLongText, nullable=True)
# intent_pending -> active -> completed / archived. The router owns legal
# transitions; the DB keeps it as a concise, portable string.
status: Mapped[str] = mapped_column(String(32), nullable=False, default="intent_pending", index=True)
# 乐观并发控制版本号:保护「确定意图 / 激活链 / 归档」等 session 级变更。
# 每次 session 写操作自增;传入 expected_version 时做 ``WHERE version=?`` CAS,
# 命中 0 行即并发冲突(消除两个用户同时激活/归档同一 session 的竞态)。
version: Mapped[int] = mapped_column(Integer, nullable=False, default=0, server_default="0")
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),
)
class PositionRoundtableNodeRow(Base):
__tablename__ = "position_roundtable_nodes"
__table_args__ = (
UniqueConstraint("session_id", "node_key", name="uq_position_roundtable_nodes_session_key"),
)
id: Mapped[str] = mapped_column(String(64), primary_key=True)
session_id: Mapped[str] = mapped_column(
String(64),
ForeignKey("position_roundtable_sessions.id", ondelete="CASCADE"),
nullable=False,
index=True,
)
# Stable inside one frozen chain snapshot; e.g. ``stage-1-seat-2-agent``.
node_key: Mapped[str] = mapped_column(String(256), nullable=False)
stage_index: Mapped[int] = mapped_column(Integer, nullable=False)
seat_index: Mapped[int] = mapped_column(Integer, nullable=False)
agent_id: Mapped[str] = mapped_column(String(128), nullable=False, index=True)
position_id: Mapped[str | None] = mapped_column(String(64), nullable=True, index=True)
# The LangGraph/checkpoint store remains the execution source for messages.
thread_id: Mapped[str | None] = mapped_column(String(128), nullable=True, unique=True)
# locked / ready / running / done / stale / rejected / error
status: Mapped[str] = mapped_column(String(32), nullable=False, default="locked", index=True)
latest_answer: Mapped[str | None] = mapped_column(PortableLongText, nullable=True)
artifact_manifest: Mapped[str | None] = mapped_column(PortableLongText, nullable=True)
# A complete, JSON-serialisable copy of the node chat plus its composer
# draft. It is intentionally colocated with the node so deleting a
# position Session removes every recovery snapshot through the FK cascade.
conversation_snapshot: Mapped[str | None] = mapped_column(PortableLongText, nullable=True)
revision: Mapped[int] = mapped_column(Integer, nullable=False, default=0)
# 乐观并发控制版本号:每次 update_node 自增;传入 expected_version 时做
# ``WHERE version=?`` 校验,命中 0 行即并发冲突。与 revision(业务交付计数)
# 刻意解耦——revision 仅在真正交付新 Markdown 报告时 +1,version 在任意字段
# 变更时 +1,专用于消除「complete_node_turn 多步状态机」并发盲写丢失更新。
version: Mapped[int] = mapped_column(Integer, nullable=False, default=0, server_default="0")
upstream_revision: Mapped[str | None] = mapped_column(PortableLongText, nullable=True)
rejection_count: Mapped[int] = mapped_column(Integer, nullable=False, default=0, server_default="0")
last_rejection: Mapped[str | None] = mapped_column(PortableLongText, nullable=True)
invalidated_by: Mapped[str | None] = mapped_column(String(256), nullable=True)
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),
)
class PositionRoundtableCommandRow(Base):
"""命令幂等记录:同一次客户端操作(由 command_id 标识)只执行一次。
网络重试 / 双击 / beforeunload beacon 可能把同一命令送达多次。靠
``UNIQUE(session_id, node_key, command_id)`` 让数据库天然去重:第二次到达时
命中已有记录,直接返回 ``result_json``(上次结果),绝不重复推进状态机。
``payload_hash`` 防止客户端误复用 command_id 发送不同参数——同 command_id 但
payload 不同 → 409。session 级命令(activate/archive)node_key 存空串。
"""
__tablename__ = "position_roundtable_commands"
__table_args__ = (
UniqueConstraint(
"session_id", "node_key", "command_id",
name="uq_position_roundtable_commands_session_node_cmd",
),
)
id: Mapped[str] = mapped_column(String(64), primary_key=True)
session_id: Mapped[str] = mapped_column(String(64), nullable=False, index=True)
# 命令作用的节点;session 级命令(activate/archive/delete)存空串 ""。
node_key: Mapped[str] = mapped_column(String(256), nullable=False, default="")
# 客户端生成的幂等键(一次点击/重试共用)。
command_id: Mapped[str] = mapped_column(String(64), nullable=False)
# complete / bind / reject / conversation / activate / archive / delete
command_type: Mapped[str] = mapped_column(String(32), nullable=False)
# 规范化 payload 的 sha256;同 command_id 但 hash 不同 → 409(防误复用)。
payload_hash: Mapped[str] = mapped_column(String(64), nullable=False)
# 上次执行结果的 JSON 序列化,重放时原样返回。
result_json: Mapped[str | None] = mapped_column(PortableLongText, nullable=True)
created_at: Mapped[datetime] = mapped_column(
BeijingDateTime(), nullable=False, default=lambda: datetime.now(UTC)
)