# -*- coding: utf-8 -*- # Copyright (c) 2025 relakkes@gmail.com # # This file is part of MediaCrawler project. # Repository: https://github.com/NanmiCoder/MediaCrawler/blob/main/api/monitor/scheduler.py # GitHub: https://github.com/NanmiCoder # Licensed under NON-COMMERCIAL LEARNING LICENSE 1.1 # # 声明:本代码仅供学习和研究目的使用。使用者应遵守以下原则: # 1. 不得用于任何商业用途。 # 2. 使用时应遵守目标平台的使用条款和robots.txt规则。 # 3. 不得进行大规模爬取或对平台造成运营干扰。 # 4. 应合理控制请求频率,避免给目标平台带来不必要的负担。 # 5. 不得用于任何非法或不当的用途。 # # 详细许可条款请参阅项目根目录下的LICENSE文件。 # 使用本代码即表示您同意遵守上述原则和LICENSE中的所有条款。 """Background scheduler for monitor tasks. One asyncio loop polls for due tasks and hands them to the runner. A plain loop is enough here: there is exactly one process, one global crawler subprocess, and therefore no concurrency to coordinate -- a cron-style library would add a dependency without adding a capability. Two families of schedule, and the difference matters: * ``interval`` is **fixed-delay**, not fixed-rate -- ``next_run_at`` is measured from the moment a run starts, so a slow run cannot make its task fire back-to-back. * the clock modes (``daily``/``weekly``) are **fixed-time** -- recomputed from the calendar, so a run that starts late does not drag every later run with it. The arithmetic for both lives in schedule.py. """ import asyncio import random from datetime import datetime from typing import Optional from sqlalchemy import select from tools.time_util import get_current_timestamp from ..services import crawler_manager from . import app_settings, schedule from .db import get_session from .models import MonitorRun, MonitorTask, RUN_INTERRUPTED, RUN_RUNNING from .runner import execute_task from .settings import get_cookie POLL_INTERVAL_SECONDS = 20 # Spread tasks sharing an interval so they do not all come due on the same tick. # Applied to interval mode only -- see the advance step below. JITTER_SECONDS = 60 class MonitorScheduler: """Polls the task table and runs whatever is due.""" def __init__(self) -> None: self._loop_task: Optional[asyncio.Task] = None self._stopping = asyncio.Event() # Avoids logging "no cookie" on every single tick. self._warned_no_cookie = False async def start(self) -> None: if self._loop_task is not None and not self._loop_task.done(): return self._stopping.clear() self._loop_task = asyncio.create_task(self._run_loop()) async def stop(self) -> None: self._stopping.set() if self._loop_task is not None: self._loop_task.cancel() try: await self._loop_task except asyncio.CancelledError: pass self._loop_task = None async def _run_loop(self) -> None: try: await self.recover() except Exception as exc: # pragma: no cover - defensive print(f"[monitor.scheduler] recovery failed: {exc}") while not self._stopping.is_set(): try: await self.tick() except Exception as exc: # pragma: no cover - keep the loop alive print(f"[monitor.scheduler] tick failed: {exc}") await asyncio.sleep(POLL_INTERVAL_SECONDS) async def recover(self) -> None: """Clean up state left behind by a server restart. A run still marked ``running`` cannot be running -- its subprocess died with the previous process. Marking it interrupted stops it from blocking the UI as a phantom in-flight run. """ async with get_session() as session: stale = ( await session.scalars( select(MonitorRun).where(MonitorRun.status == RUN_RUNNING) ) ).all() for run in stale: run.status = RUN_INTERRUPTED run.finished_at = get_current_timestamp() if stale: print( f"[monitor.scheduler] marked {len(stale)} interrupted run(s) " f"left over from a previous process" ) async def tick(self) -> None: """Run one due task, if the crawler is free and we are in the active window.""" # The crawler subprocess is a global singleton, so a manual crawl and a # monitor run cannot overlap. Returning without advancing next_run_at # leaves the task due, and it is picked up on a later tick. if crawler_manager.is_busy(): return async with get_session() as session: if not await self._within_active_hours(session): # Deliberately does not advance next_run_at: the task simply runs # when the window next opens, rather than being skipped for a day. return await self._run_due_task() async def _within_active_hours(self, session) -> bool: """Whether scheduled runs are allowed right now (local time).""" start, end = await app_settings.active_hours(session) hour = datetime.now().hour if start <= end: return start <= hour <= end # Window wraps past midnight, e.g. 22 -> 6. return hour >= start or hour <= end async def _run_due_task(self) -> None: async with get_session() as session: task = await session.scalar( select(MonitorTask) .where( MonitorTask.enabled.is_(True), MonitorTask.next_run_at.is_not(None), MonitorTask.next_run_at <= get_current_timestamp(), ) .order_by(MonitorTask.next_run_at) .limit(1) ) if task is None: return # No cookie means every run would report an auth failure. Leave the # task due rather than advancing: it starts working the moment the # user pastes one. cookie = await get_cookie(session) if not cookie: if not self._warned_no_cookie: print( "[monitor.scheduler] no XHS cookie configured; " "scheduled tasks will not run until one is set" ) self._warned_no_cookie = True return self._warned_no_cookie = False # Advance before running so a crash mid-run cannot cause an immediate # re-fire, and so a long outage coalesces into a single run instead # of one run per missed interval. now = get_current_timestamp() following = schedule.next_occurrence( mode=task.schedule_mode, interval_minutes=task.interval_minutes, hours=schedule.parse_hours(task.schedule_hours), days=schedule.parse_days(task.schedule_days), minute=task.schedule_minute, after_ms=now, ) if following is None: # A clock schedule with no times can never fire. The API rejects # that shape, so this guards against a hand-edited row: park the # task with no next run rather than leaving it permanently due and # re-running it on every tick. task.next_run_at = None print( f"[monitor.scheduler] task {task.id} has no usable schedule " f"and will not run until one is set" ) elif task.schedule_mode == schedule.MODE_INTERVAL: # Jitter belongs to the interval mode only. Spreading identical # intervals apart is the point; nudging a time the operator # explicitly picked is not -- it just looks like a broken clock. task.next_run_at = following + random.randint(0, JITTER_SECONDS) * 1000 else: task.next_run_at = following task_id = task.id try: await execute_task(task_id, trigger="scheduled") except Exception as exc: print(f"[monitor.scheduler] task {task_id} failed: {exc}") # Global singleton, mirroring the crawler_manager pattern. monitor_scheduler = MonitorScheduler()