feat(schedule): 任务支持「每天定时 / 每周定时」,用选择器而不是手写 cron
- schedule.py: 新增调度计算模块(纯函数,便于单测)。三种模式:interval(每 N 分钟)/ daily(选钟点)/ weekly(选星期 + 钟点) - 钟点模式是「固定时刻」而非「固定延迟」——从日历重算,所以某轮跑晚了不会把之后每一轮都拖晚。interval 保持原语义:从上一轮开始计时 - 抖动只加给 interval。给「每天 9:00」也加抖动就成了 9:00–9:01 随机触发,操作者选的时间被悄悄改掉,只会像 bug - 模型 / schema / service: 新增 schedule_mode / schedule_hours / schedule_days / schedule_minute。时钟字段存逗号分隔文本——几个小整数、永远整体读写,单开一张表只会换来 join。interval_minutes 保留且仍是默认值,已有任务不受影响 - service: 改动任何调度字段都按合并后的状态重算 next_run_at。重新启用也算改动,否则停用一个月再打开会带着一个月前的 next_run_at,一保存就立即触发 - 前端: 运行方式三选一 + 小时/星期胶囊多选 + 分钟下拉,并实时预览结果句子。任务卡片改显示后端拼好的 schedule_label,避免列表和编辑器对同一计划给出两种说法 - 校验: 钟点模式至少选一个时间,按周至少选一个星期 tests/test_schedule.py 新增 24 个用例,含「恰好等于当前时刻的档位归属下一天」这个会让调度器自循环的边界。
This commit is contained in:
+93
-5
@@ -28,7 +28,7 @@ from sqlalchemy.ext.asyncio import AsyncSession
|
||||
|
||||
from tools.time_util import get_current_timestamp
|
||||
|
||||
from . import app_settings, platforms
|
||||
from . import app_settings, platforms, schedule
|
||||
from .db import get_session
|
||||
from .platforms import PLATFORM_XHS
|
||||
from .models import (
|
||||
@@ -53,6 +53,28 @@ _background_runs: set[asyncio.Task] = set()
|
||||
MIN_INTERVAL_MINUTES = 30
|
||||
MAX_INTERVAL_MINUTES = 7 * 24 * 60
|
||||
|
||||
# Changing any of these invalidates the pending run slot.
|
||||
SCHEDULE_FIELDS = {
|
||||
"interval_minutes",
|
||||
"schedule_mode",
|
||||
"schedule_hours",
|
||||
"schedule_days",
|
||||
"schedule_minute",
|
||||
}
|
||||
|
||||
|
||||
def _require_clock_fields(mode: str, hours: list, days: list) -> None:
|
||||
"""A clock schedule with no clock time can never fire.
|
||||
|
||||
The create schema already rejects that shape, but the update path merges
|
||||
partial fields and therefore has no schema-level view of the result -- so the
|
||||
check lives here, where both paths meet.
|
||||
"""
|
||||
if mode in schedule.CLOCK_MODES and not hours:
|
||||
raise ValueError("按钟点调度至少要选一个时间")
|
||||
if mode == schedule.MODE_WEEKLY and not days:
|
||||
raise ValueError("按周调度至少要选一个星期")
|
||||
|
||||
_CREATOR_URL_RE = re.compile(r"xiaohongshu\.com/user/profile/([A-Za-z0-9_-]+)")
|
||||
_NOTE_URL_RE = re.compile(r"xiaohongshu\.com/(?:explore|discovery/item)/([A-Za-z0-9_-]+)")
|
||||
# XHS user ids and note ids are 24-char hex; allow a slightly wider range so a
|
||||
@@ -139,7 +161,12 @@ async def create_task(session: AsyncSession, payload: Dict[str, Any]) -> Monitor
|
||||
# the Settings page actually governs new tasks.
|
||||
defaults = await app_settings.defaults(session, platform)
|
||||
interval_minutes = payload.get("interval_minutes") or defaults["interval_minutes"]
|
||||
interval_ms = int(interval_minutes) * 60_000
|
||||
|
||||
schedule_mode = payload.get("schedule_mode") or schedule.MODE_INTERVAL
|
||||
schedule_hours = list(payload.get("schedule_hours") or [])
|
||||
schedule_days = list(payload.get("schedule_days") or [])
|
||||
schedule_minute = int(payload.get("schedule_minute") or 0)
|
||||
_require_clock_fields(schedule_mode, schedule_hours, schedule_days)
|
||||
|
||||
task = MonitorTask(
|
||||
name=payload["name"],
|
||||
@@ -147,12 +174,23 @@ async def create_task(session: AsyncSession, payload: Dict[str, Any]) -> Monitor
|
||||
mode=mode,
|
||||
enabled=payload.get("enabled", True),
|
||||
interval_minutes=interval_minutes,
|
||||
schedule_mode=schedule_mode,
|
||||
schedule_hours=schedule.format_hours(schedule_hours),
|
||||
schedule_days=schedule.format_days(schedule_days),
|
||||
schedule_minute=schedule_minute,
|
||||
max_notes_count=payload.get("max_notes_count") or defaults["max_notes_count"],
|
||||
enable_comments=payload.get("enable_comments", True),
|
||||
max_comments_count=payload.get("max_comments_count") or defaults["max_comments_count"],
|
||||
run_timeout_seconds=payload.get("run_timeout_seconds", 3600),
|
||||
notify_enabled=payload.get("notify_enabled", False),
|
||||
next_run_at=now + interval_ms,
|
||||
next_run_at=schedule.next_occurrence(
|
||||
mode=schedule_mode,
|
||||
interval_minutes=interval_minutes,
|
||||
hours=schedule_hours,
|
||||
days=schedule_days,
|
||||
minute=schedule_minute,
|
||||
after_ms=now,
|
||||
),
|
||||
last_status="idle",
|
||||
created_at=now,
|
||||
updated_at=now,
|
||||
@@ -189,10 +227,14 @@ async def update_task(session: AsyncSession, task_id: int, payload: Dict[str, An
|
||||
if task is None:
|
||||
raise ValueError(f"Task {task_id} not found")
|
||||
|
||||
was_enabled = task.enabled
|
||||
|
||||
for field in (
|
||||
"name",
|
||||
"enabled",
|
||||
"interval_minutes",
|
||||
"schedule_mode",
|
||||
"schedule_minute",
|
||||
"max_notes_count",
|
||||
"enable_comments",
|
||||
"max_comments_count",
|
||||
@@ -202,6 +244,13 @@ async def update_task(session: AsyncSession, task_id: int, payload: Dict[str, An
|
||||
if field in payload and payload[field] is not None:
|
||||
setattr(task, field, payload[field])
|
||||
|
||||
# The clock lists are stored as comma-separated text, so they cannot go
|
||||
# through the generic loop above.
|
||||
if payload.get("schedule_hours") is not None:
|
||||
task.schedule_hours = schedule.format_hours(payload["schedule_hours"])
|
||||
if payload.get("schedule_days") is not None:
|
||||
task.schedule_days = schedule.format_days(payload["schedule_days"])
|
||||
|
||||
# Replacing targets resets the baseline implicitly: a note set that now
|
||||
# includes new ids will simply report them as new on the next run.
|
||||
if payload.get("targets") is not None:
|
||||
@@ -227,8 +276,23 @@ async def update_task(session: AsyncSession, task_id: int, payload: Dict[str, An
|
||||
)
|
||||
)
|
||||
|
||||
if "interval_minutes" in payload and payload["interval_minutes"]:
|
||||
task.next_run_at = get_current_timestamp() + payload["interval_minutes"] * 60_000
|
||||
# Any change to when the task runs invalidates the pending slot, so recompute
|
||||
# it from the merged state rather than working out which field moved.
|
||||
# Re-enabling counts as a change too: otherwise a task switched off for a
|
||||
# month comes back holding a next_run_at a month in the past and fires the
|
||||
# instant it is saved.
|
||||
if (SCHEDULE_FIELDS & set(payload)) or (task.enabled and not was_enabled):
|
||||
hours = schedule.parse_hours(task.schedule_hours)
|
||||
days = schedule.parse_days(task.schedule_days)
|
||||
_require_clock_fields(task.schedule_mode, hours, days)
|
||||
task.next_run_at = schedule.next_occurrence(
|
||||
mode=task.schedule_mode,
|
||||
interval_minutes=task.interval_minutes,
|
||||
hours=hours,
|
||||
days=days,
|
||||
minute=task.schedule_minute,
|
||||
after_ms=get_current_timestamp(),
|
||||
)
|
||||
|
||||
task.updated_at = get_current_timestamp()
|
||||
await session.flush()
|
||||
@@ -623,6 +687,7 @@ async def list_tasks(
|
||||
"mode": task.mode,
|
||||
"enabled": task.enabled,
|
||||
"interval_minutes": task.interval_minutes,
|
||||
**_schedule_fields(task),
|
||||
"max_notes_count": task.max_notes_count,
|
||||
"enable_comments": task.enable_comments,
|
||||
"max_comments_count": task.max_comments_count,
|
||||
@@ -644,6 +709,29 @@ async def list_tasks(
|
||||
]
|
||||
|
||||
|
||||
def _schedule_fields(task: MonitorTask) -> Dict[str, Any]:
|
||||
"""The schedule columns, plus the sentence the task list renders.
|
||||
|
||||
The label is composed here rather than in the frontend so the list and the
|
||||
editor cannot drift on what a given schedule means.
|
||||
"""
|
||||
hours = schedule.parse_hours(task.schedule_hours)
|
||||
days = schedule.parse_days(task.schedule_days)
|
||||
return {
|
||||
"schedule_mode": task.schedule_mode,
|
||||
"schedule_hours": hours,
|
||||
"schedule_days": days,
|
||||
"schedule_minute": task.schedule_minute,
|
||||
"schedule_label": schedule.describe(
|
||||
mode=task.schedule_mode,
|
||||
interval_minutes=task.interval_minutes,
|
||||
hours=hours,
|
||||
days=days,
|
||||
minute=task.schedule_minute,
|
||||
),
|
||||
}
|
||||
|
||||
|
||||
async def overview(session: AsyncSession, platform: Optional[str] = None) -> Dict[str, Any]:
|
||||
"""Headline numbers for the dashboard tiles, scoped to one platform."""
|
||||
now = get_current_timestamp()
|
||||
|
||||
Reference in New Issue
Block a user