98 lines
5.7 KiB
Python
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))
|