diff --git a/backend/app/task/crud/crud_scheduler.py b/backend/app/task/crud/crud_scheduler.py index 55e4acc5..2cb46b4b 100644 --- a/backend/app/task/crud/crud_scheduler.py +++ b/backend/app/task/crud/crud_scheduler.py @@ -85,7 +85,7 @@ class CRUDTaskScheduler(CRUDPlus[TaskScheduler]): TaskScheduler.no_changes = False return 1 - async def set_status(self, db: AsyncSession, pk: int, *, status: bool) -> int: + async def set_status(self, db: AsyncSession, pk: int, *, status: int) -> int: """ 设置任务调度状态 @@ -95,7 +95,7 @@ class CRUDTaskScheduler(CRUDPlus[TaskScheduler]): :return: """ task_scheduler = await self.get(db, pk) - task_scheduler.enabled = status + task_scheduler.status = status TaskScheduler.no_changes = False return 1 diff --git a/backend/app/task/model/scheduler.py b/backend/app/task/model/scheduler.py index d8e36a9c..7b1fbb47 100644 --- a/backend/app/task/model/scheduler.py +++ b/backend/app/task/model/scheduler.py @@ -8,6 +8,7 @@ import sqlalchemy as sa from sqlalchemy import event from sqlalchemy.orm import Mapped, mapped_column +from backend.common.enums import StatusType from backend.common.exception import errors from backend.common.model import Base, TimeZone, UniversalText, id_key from backend.core.conf import settings @@ -40,7 +41,7 @@ class TaskScheduler(Base): interval_period: Mapped[str | None] = mapped_column(sa.String(256), comment='任务运行之间的周期类型') crontab: Mapped[str | None] = mapped_column(sa.String(64), default='* * * * *', comment='Crontab 表达式') one_off: Mapped[bool] = mapped_column(default=False, comment='是否仅运行一次') - enabled: Mapped[bool] = mapped_column(default=True, comment='是否启用任务') + status: Mapped[int] = mapped_column(default=StatusType.enable.value, comment='状态(0停用 1正常)') total_run_count: Mapped[int] = mapped_column(default=0, comment='任务触发的总次数') last_run_time: Mapped[datetime | None] = mapped_column(TimeZone, default=None, comment='任务最后触发的时间') remark: Mapped[str | None] = mapped_column(UniversalText, default=None, comment='备注') diff --git a/backend/app/task/schema/scheduler.py b/backend/app/task/schema/scheduler.py index 3efc7a62..b47616a1 100644 --- a/backend/app/task/schema/scheduler.py +++ b/backend/app/task/schema/scheduler.py @@ -4,6 +4,7 @@ from pydantic import ConfigDict, Field from pydantic.types import JsonValue from backend.app.task.enums import PeriodType, TaskSchedulerType +from backend.common.enums import StatusType from backend.common.schema import SchemaBase @@ -42,7 +43,7 @@ class GetTaskSchedulerDetail(TaskSchedulerSchemaBase): model_config = ConfigDict(from_attributes=True) id: int = Field(description='任务调度 ID') - enabled: bool = Field(description='是否启用任务') + status: StatusType = Field(description='状态') total_run_count: int = Field(description='已运行总次数') last_run_time: datetime | None = Field(None, description='最后运行时间') created_time: datetime = Field(description='创建时间') diff --git a/backend/app/task/service/scheduler_service.py b/backend/app/task/service/scheduler_service.py index 68b5fb32..b0d31333 100644 --- a/backend/app/task/service/scheduler_service.py +++ b/backend/app/task/service/scheduler_service.py @@ -12,6 +12,7 @@ from backend.app.task.enums import TaskSchedulerType from backend.app.task.model import TaskScheduler from backend.app.task.schema.scheduler import CreateTaskSchedulerParam, UpdateTaskSchedulerParam from backend.app.task.utils.tzcrontab import crontab_verify +from backend.common.enums import StatusType from backend.common.exception import errors from backend.common.pagination import paging_data @@ -110,7 +111,8 @@ class TaskSchedulerService: task_scheduler = await task_scheduler_dao.get(db, pk) if not task_scheduler: raise errors.NotFoundError(msg='任务调度不存在') - count = await task_scheduler_dao.set_status(db, pk, status=not task_scheduler.enabled) + next_status = StatusType.disable if task_scheduler.status == StatusType.enable else StatusType.enable + count = await task_scheduler_dao.set_status(db, pk, status=next_status) return count @staticmethod diff --git a/backend/app/task/utils/schedulers.py b/backend/app/task/utils/schedulers.py index 46616369..340f3506 100644 --- a/backend/app/task/utils/schedulers.py +++ b/backend/app/task/utils/schedulers.py @@ -19,6 +19,7 @@ from backend.app.task.enums import PeriodType, TaskSchedulerType from backend.app.task.model.scheduler import TaskScheduler from backend.app.task.schema.scheduler import CreateTaskSchedulerParam from backend.app.task.utils.tzcrontab import TzAwareCrontab, crontab_verify +from backend.common.enums import StatusType from backend.common.exception import errors from backend.core.conf import settings from backend.database.db import async_db_session @@ -91,22 +92,24 @@ class ModelEntry(ScheduleEntry): self.last_run_at = timezone.from_datetime(model.last_run_time) self.options['periodic_task_name'] = model.name self.model = model + self.enabled = model.status == StatusType.enable async def _disable(self, model: TaskScheduler) -> None: """禁用任务""" model.no_changes = True - self.model.enabled = self.enabled = model.enabled = False + self.model.status = model.status = StatusType.disable + self.enabled = False async with async_db_session.begin() as db: stmt = select(TaskScheduler).where(TaskScheduler.id == model.id, TaskScheduler.deleted == 0) query = await db.execute(stmt) task = query.scalars().first() if task: task.no_changes = True - task.enabled = False + task.status = StatusType.disable def is_due(self) -> tuple[bool, int | float | datetime]: """任务到期状态""" - if not self.model.enabled: + if self.model.status != StatusType.enable: # 重新启用时延迟 5 秒 return schedules.schedstate(is_due=False, next=5) @@ -119,11 +122,11 @@ class ModelEntry(ScheduleEntry): return schedules.schedstate(is_due=False, next=delay) # 一次性任务 - if self.model.one_off and self.model.enabled and self.model.total_run_count > 0: - self.model.enabled = False + if self.model.one_off and self.model.status == StatusType.enable and self.model.total_run_count > 0: + self.model.status = StatusType.disable self.model.total_run_count = 0 self.model.no_changes = False - save_fields = ('enabled',) + save_fields = ('status',) run_await(self.save)(save_fields) return schedules.schedstate(is_due=False, next=1000000000) # 高延迟,避免重新检查 @@ -237,6 +240,9 @@ class ModelEntry(ScheduleEntry): **cls._unpack_options(**options or {}), **entry, ) + if 'enabled' in model_dict: + enabled = model_dict.pop('enabled') + model_dict['status'] = StatusType.enable if enabled else StatusType.disable return model_dict @classmethod @@ -359,7 +365,7 @@ class DatabaseScheduler(Scheduler): try: for name, entry_fields in beat_dict.items(): entry = run_await(self.Entry.from_entry)(name, app=self.app, **entry_fields) - if entry.model.enabled: + if entry.model.status == StatusType.enable: s[name] = entry except Exception: logger.error(f'添加任务 {name} 到数据库失败') @@ -389,7 +395,7 @@ class DatabaseScheduler(Scheduler): async with async_db_session() as db: logger.debug('DatabaseScheduler: Fetching database schedule') stmt = select(TaskScheduler).where( - TaskScheduler.enabled.is_(True), + TaskScheduler.status == StatusType.enable, TaskScheduler.deleted == 0, ) query = await db.execute(stmt)