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

98 lines
5.7 KiB
Python

"""ORM models for user scheduled tasks."""
from __future__ import annotations
from datetime import UTC, datetime
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, PortableJSON
class SchedulerThreadRow(Base):
__tablename__ = "scheduler_threads"
user_id: Mapped[str] = mapped_column(String(64), primary_key=True)
thread_id: Mapped[str] = mapped_column(String(64), unique=True, index=True)
created_at: Mapped[datetime] = mapped_column(BeijingDateTime(), default=lambda: datetime.now(UTC))
updated_at: Mapped[datetime] = mapped_column(BeijingDateTime(), default=lambda: datetime.now(UTC), onupdate=lambda: datetime.now(UTC))
class ScheduledTaskRow(Base):
__tablename__ = "scheduled_tasks"
__table_args__ = (UniqueConstraint("user_id", "name", name="uq_scheduled_tasks_user_name"),)
task_id: Mapped[str] = mapped_column(String(64), primary_key=True)
user_id: Mapped[str] = mapped_column(String(64), index=True)
scheduler_thread_id: Mapped[str] = mapped_column(String(64), index=True)
source_thread_id: Mapped[str | None] = mapped_column(String(64), index=True)
source_agent_name: Mapped[str | None] = mapped_column(String(128), index=True)
execution_agent_name: Mapped[str | None] = mapped_column(String(128), index=True)
name: Mapped[str] = mapped_column(String(127))
description: Mapped[str | None] = mapped_column(Text)
prompt: Mapped[str] = mapped_column(Text)
schedule_text: Mapped[str] = mapped_column(String(256))
cron_expr: Mapped[str] = mapped_column(String(128))
timezone: Mapped[str] = mapped_column(String(64), default="UTC")
enabled: Mapped[bool] = mapped_column(Boolean, default=True)
next_run_at: Mapped[datetime | None] = mapped_column(BeijingDateTime(), index=True)
last_run_at: Mapped[datetime | None] = mapped_column(BeijingDateTime())
last_status: Mapped[str | None] = mapped_column(String(32))
execution_context_json: Mapped[dict] = mapped_column(PortableJSON(), default=dict)
published: Mapped[bool] = mapped_column(Boolean, default=False, index=True)
published_at: Mapped[datetime | None] = mapped_column(BeijingDateTime())
created_at: Mapped[datetime] = mapped_column(BeijingDateTime(), default=lambda: datetime.now(UTC))
updated_at: Mapped[datetime] = mapped_column(BeijingDateTime(), default=lambda: datetime.now(UTC), onupdate=lambda: datetime.now(UTC))
class ScheduledTaskSubscriptionRow(Base):
__tablename__ = "scheduled_task_subscriptions"
__table_args__ = (UniqueConstraint("task_id", "user_id", name="uq_scheduled_task_subscriptions_task_user"),)
id: Mapped[str] = mapped_column(String(64), primary_key=True)
task_id: Mapped[str] = mapped_column(String(64), index=True)
user_id: Mapped[str] = mapped_column(String(64), index=True)
notify_element: Mapped[bool] = mapped_column(Boolean, default=False)
# Source of this subscription: ``"user"`` (subscribed manually from the
# 广场) or ``"position"`` (auto-granted by the user's 岗位 via the sync
# engine). Sync only ever touches ``"position"`` rows.
origin: Mapped[str] = mapped_column(String(16), nullable=False, default="user", server_default="user")
position_id: Mapped[str | None] = mapped_column(String(64), nullable=True, index=True)
created_at: Mapped[datetime] = mapped_column(BeijingDateTime(), default=lambda: datetime.now(UTC))
updated_at: Mapped[datetime] = mapped_column(BeijingDateTime(), default=lambda: datetime.now(UTC), onupdate=lambda: datetime.now(UTC))
class ScheduledTaskRunRow(Base):
__tablename__ = "scheduled_task_runs"
__table_args__ = (UniqueConstraint("task_id", "scheduled_for", name="uq_scheduled_task_runs_task_scheduled_for"),)
id: Mapped[str] = mapped_column(String(64), primary_key=True)
task_id: Mapped[str] = mapped_column(String(64), index=True)
user_id: Mapped[str] = mapped_column(String(64), index=True)
scheduler_thread_id: Mapped[str] = mapped_column(String(64), index=True)
agent_run_id: Mapped[str | None] = mapped_column(String(64), index=True)
scheduled_for: Mapped[datetime] = mapped_column(BeijingDateTime(), index=True)
started_at: Mapped[datetime | None] = mapped_column(BeijingDateTime())
finished_at: Mapped[datetime | None] = mapped_column(BeijingDateTime())
status: Mapped[str] = mapped_column(String(32), default="pending")
error: Mapped[str | None] = mapped_column(Text)
result_json: Mapped[dict | None] = mapped_column(PortableJSON(), default=dict)
# Number of recovery attempts made after the first try. A run is retried
# when it overruns the per-attempt timeout or is found orphaned in the
# "running" state (e.g. after a service restart). See ScheduledTaskService.
retry_count: Mapped[int] = mapped_column(Integer, default=0, server_default="0")
created_at: Mapped[datetime] = mapped_column(BeijingDateTime(), default=lambda: datetime.now(UTC))
updated_at: Mapped[datetime] = mapped_column(BeijingDateTime(), default=lambda: datetime.now(UTC), onupdate=lambda: datetime.now(UTC))
class ScheduledTaskDeliveryProfileRow(Base):
__tablename__ = "scheduled_task_delivery_profiles"
user_id: Mapped[str] = mapped_column(String(64), primary_key=True)
element_user_id: Mapped[str | None] = mapped_column(String(191), index=True)
element_room_id: Mapped[str | None] = mapped_column(String(256), nullable=True)
created_at: Mapped[datetime] = mapped_column(BeijingDateTime(), default=lambda: datetime.now(UTC))
updated_at: Mapped[datetime] = mapped_column(BeijingDateTime(), default=lambda: datetime.now(UTC), onupdate=lambda: datetime.now(UTC))