diff --git a/README.md b/README.md index 370bee2..72db3f6 100644 --- a/README.md +++ b/README.md @@ -417,8 +417,12 @@ systemd 定时备份示例为 `deploy/fluxip-backup.service` 和 ## 健康检查与排错 - `GET /healthz`:进程存活状态,返回 `status` 和应用版本;不探测 SQLite 或云端账号。 -- `GET /readyz`:检查 SQLite 是否可用;就绪时返回 HTTP 200,数据库异常时返回 - HTTP 503 和 `degraded`。该接口不会调用 AWS 或 Cloudflare。 +- `GET /readyz`:检查 SQLite 和已启用调度器的运行状态;就绪时返回 HTTP 200,数据库 + 异常或已启用调度器意外退出时返回 HTTP 503 和 `degraded`。响应中的 `scheduler` + 会显示开关、运行状态、最近扫描/派发/异常时间;该接口不会调用 AWS 或 Cloudflare。 +- **到期后没有自动轮换**:先检查 `/readyz` 的 `scheduler.enabled` 和 `running`。系统环境 + 变量优先于 `.env`;请删除意外继承的 `ROTATOR_SCHEDULER_ENABLED=false`,或明确设为 + `true` 后重启。已逾期的实例组会在调度器启动后的首轮扫描立即执行。 - **没有默认管理员**:首次打开页面自行创建;远程初始化需 Bootstrap Token。 - **账号连接测试失败**:确认实例选中了正确类型的 AWS 与 Cloudflare 账号,账号密钥 完整,并且 Region、Zone ID 和记录名属于对应账号。 diff --git a/app/container.py b/app/container.py index 95fb3e4..4a2034c 100644 --- a/app/container.py +++ b/app/container.py @@ -1,5 +1,6 @@ from __future__ import annotations +import logging from pathlib import Path from app.accounts.repository import AccountRepository @@ -18,6 +19,8 @@ from app.rotation.repository import RotationRepository from app.rotation.schedule_service import SchedulerLoop from app.rotation.service import RotationService +logger = logging.getLogger(__name__) + class AppContainer: def __init__(self, settings: AppSettings) -> None: @@ -72,6 +75,8 @@ class AppContainer: self.rotation_service.recover() if self.settings.scheduler_enabled: self.scheduler_loop.start(self.settings.scheduler_tick_seconds) + else: + logger.warning("scheduler_disabled") async def shutdown(self) -> None: await self.scheduler_loop.stop() diff --git a/app/main.py b/app/main.py index dc26f1f..44efa88 100644 --- a/app/main.py +++ b/app/main.py @@ -194,10 +194,17 @@ def create_app(settings: AppSettings | None = None) -> FastAPI: @application.get("/readyz", include_in_schema=False) async def ready(request: Request) -> JSONResponse: - healthy = request.app.state.container.database.health_check() + container = request.app.state.container + database_healthy = container.database.health_check() + scheduler = container.scheduler_loop.status(enabled=container.settings.scheduler_enabled) + healthy = database_healthy and bool(scheduler["healthy"]) return JSONResponse( status_code=200 if healthy else 503, - content={"status": "ok" if healthy else "degraded", "database": healthy}, + content={ + "status": "ok" if healthy else "degraded", + "database": database_healthy, + "scheduler": scheduler, + }, ) if STATIC_DIR.exists(): diff --git a/app/rotation/schedule_service.py b/app/rotation/schedule_service.py index 00fb5f4..ec18c23 100644 --- a/app/rotation/schedule_service.py +++ b/app/rotation/schedule_service.py @@ -2,8 +2,11 @@ from __future__ import annotations import asyncio import logging +import time +from typing import Any from app.core.errors import ConflictError +from app.core.time import to_iso from app.rotation.repository import RotationRepository from app.rotation.service import RotationService @@ -20,11 +23,27 @@ class SchedulerLoop: self.rotation_service = rotation_service self._task: asyncio.Task[None] | None = None self._stopping = asyncio.Event() + self._tick_seconds: float | None = None + self._started_at: str | None = None + self._last_tick_at: str | None = None + self._last_success_at: str | None = None + self._last_dispatch_at: str | None = None + self._last_error_at: str | None = None + self._consecutive_errors = 0 + self._started_monotonic: float | None = None + self._last_tick_monotonic: float | None = None - def start(self, tick_seconds: int) -> None: + def start(self, tick_seconds: float) -> None: if self._task and not self._task.done(): return + self._stopping = asyncio.Event() + self._tick_seconds = tick_seconds + self._started_at = to_iso() + self._started_monotonic = time.monotonic() + self._last_tick_monotonic = None self._task = asyncio.create_task(self._run(tick_seconds), name="fleet-schedule-loop") + self._task.add_done_callback(self._task_done) + logger.info("scheduler_started", extra={"tick_seconds": tick_seconds}) async def stop(self) -> None: self._stopping.set() @@ -32,27 +51,114 @@ class SchedulerLoop: self._task.cancel() await asyncio.gather(self._task, return_exceptions=True) - async def _run(self, tick_seconds: int) -> None: + def status(self, *, enabled: bool) -> dict[str, Any]: + running = bool(self._task and not self._task.done()) + state = "disabled" if not enabled else "running" if running else "stopped" + healthy = not enabled or ( + running and self._consecutive_errors < 3 and self._tick_is_fresh() + ) + return { + "enabled": enabled, + "running": running, + "healthy": healthy, + "state": state, + "tick_seconds": self._tick_seconds, + "started_at": self._started_at, + "last_tick_at": self._last_tick_at, + "last_success_at": self._last_success_at, + "last_dispatch_at": self._last_dispatch_at, + "last_error_at": self._last_error_at, + "consecutive_errors": self._consecutive_errors, + } + + async def _run(self, tick_seconds: float) -> None: while not self._stopping.is_set(): - for group_id in self.repository.due_group_ids(): - try: - self.rotation_service.start_group(group_id, trigger="scheduled") - break - except ConflictError: - break - except Exception: - logger.exception( - "Scheduled group rotation dispatch failed", - extra={"group_id": group_id}, - ) - try: - self.repository.defer_group(group_id) - except Exception: - logger.exception( - "Failed to defer group schedule", - extra={"group_id": group_id}, - ) + self._last_tick_at = to_iso() + self._last_tick_monotonic = time.monotonic() + try: + tick_succeeded = self._dispatch_due_group() + except asyncio.CancelledError: + raise + except Exception: + self._record_error() + self._consecutive_errors += 1 + logger.exception("scheduler_tick_failed") + else: + if tick_succeeded: + self._last_success_at = to_iso() + self._consecutive_errors = 0 + else: + self._consecutive_errors += 1 try: await asyncio.wait_for(self._stopping.wait(), timeout=tick_seconds) except TimeoutError: continue + + def _dispatch_due_group(self) -> bool: + succeeded = True + due_group_ids = self.repository.due_group_ids() + if due_group_ids: + logger.info( + "scheduler_due_groups_found", + extra={"due_group_count": len(due_group_ids)}, + ) + for group_id in due_group_ids: + try: + run = self.rotation_service.start_group(group_id, trigger="scheduled") + except ConflictError as exc: + logger.warning( + "scheduler_dispatch_conflict", + extra={"group_id": group_id, "code": exc.code}, + ) + if exc.code == "ROTATION_IN_PROGRESS": + break + except Exception: + self._record_error() + succeeded = False + logger.exception("scheduler_dispatch_failed", extra={"group_id": group_id}) + try: + self.repository.defer_group(group_id) + except Exception: + self._record_error() + logger.exception("scheduler_defer_failed", extra={"group_id": group_id}) + else: + self._last_dispatch_at = to_iso() + logger.info( + "scheduler_group_dispatched", + extra={"group_id": group_id, "run_id": run.id}, + ) + break + return succeeded + + def _record_error(self) -> None: + self._last_error_at = to_iso() + + def _tick_is_fresh(self) -> bool: + reference = self._last_tick_monotonic or self._started_monotonic + if reference is None: + return False + allowed_silence = max((self._tick_seconds or 10) * 3, 10) + return time.monotonic() - reference <= allowed_silence + + def _task_done(self, task: asyncio.Task[None]) -> None: + if task.cancelled(): + if self._stopping.is_set(): + logger.info("scheduler_stopped") + else: + self._record_error() + self._consecutive_errors += 1 + logger.error("scheduler_cancelled_unexpectedly") + return + try: + task.result() + except Exception: + self._record_error() + self._consecutive_errors += 1 + logger.exception("scheduler_exited_unexpectedly") + return + if not self._stopping.is_set(): + self._record_error() + self._consecutive_errors += 1 + logger.error("scheduler_exited_unexpectedly") + else: + logger.info("scheduler_stopped") diff --git a/tests/test_schedule_service.py b/tests/test_schedule_service.py new file mode 100644 index 0000000..bab32db --- /dev/null +++ b/tests/test_schedule_service.py @@ -0,0 +1,215 @@ +from __future__ import annotations + +import asyncio +from collections.abc import Callable +from datetime import timedelta +from types import SimpleNamespace +from typing import cast + +import pytest +from fastapi.testclient import TestClient + +from app.core.errors import ConflictError +from app.core.time import to_iso, utc_now +from app.rotation.repository import RotationRepository +from app.rotation.schedule_service import SchedulerLoop +from app.rotation.service import RotationService + + +class RecordingRotationService: + def __init__(self, outcomes: dict[str, Exception] | None = None) -> None: + self.outcomes = outcomes or {} + self.calls: list[tuple[str, str]] = [] + self.dispatched = asyncio.Event() + + def start_group(self, group_id: str, trigger: str = "manual") -> SimpleNamespace: + self.calls.append((group_id, trigger)) + outcome = self.outcomes.get(group_id) + if outcome is not None: + raise outcome + self.dispatched.set() + return SimpleNamespace(id=f"run-{group_id}") + + +class ScriptedRepository: + def __init__(self, scans: list[list[str] | Exception]) -> None: + self.scans = scans + self.scan_count = 0 + self.deferred: list[str] = [] + + def due_group_ids(self) -> list[str]: + self.scan_count += 1 + result = self.scans.pop(0) if self.scans else [] + if isinstance(result, Exception): + raise result + return result + + def defer_group(self, group_id: str) -> None: + self.deferred.append(group_id) + + +def scheduler_for( + repository: RotationRepository | ScriptedRepository, + service: RecordingRotationService, +) -> SchedulerLoop: + return SchedulerLoop( + cast(RotationRepository, repository), + cast(RotationService, service), + ) + + +async def wait_until(predicate: Callable[[], bool], timeout: float = 0.5) -> None: + deadline = asyncio.get_running_loop().time() + timeout + while not predicate(): + if asyncio.get_running_loop().time() >= deadline: + pytest.fail("等待调度器状态变化超时") + await asyncio.sleep(0.005) + + +@pytest.mark.asyncio +async def test_scheduler_dispatches_due_group_from_sqlite( + client: TestClient, + managed_instance_payload: dict[str, object], + group_payload: dict[str, object], +) -> None: + container = client.app.state.container + container.fleet_repository.create_instance(managed_instance_payload) + group = container.fleet_repository.create_group(group_payload) + with container.database.connect() as connection, connection: + connection.execute( + "UPDATE instance_groups SET next_run_at = ? WHERE id = ?", + (to_iso(utc_now() - timedelta(seconds=1)), group.id), + ) + service = RecordingRotationService() + scheduler = scheduler_for(container.rotation_repository, service) + + scheduler.start(0.01) + try: + await wait_until(service.dispatched.is_set) + status = scheduler.status(enabled=True) + finally: + await scheduler.stop() + + assert service.calls[0] == (group.id, "scheduled") + assert status["state"] == "running" + assert status["last_tick_at"] is not None + assert status["last_success_at"] is not None + assert status["last_dispatch_at"] is not None + + +@pytest.mark.asyncio +async def test_scheduler_recovers_after_scan_error(caplog: pytest.LogCaptureFixture) -> None: + repository = ScriptedRepository([RuntimeError("database is temporarily locked"), ["group"]]) + service = RecordingRotationService() + scheduler = scheduler_for(repository, service) + + scheduler.start(0.01) + try: + await wait_until(service.dispatched.is_set) + status = scheduler.status(enabled=True) + finally: + await scheduler.stop() + + assert repository.scan_count >= 2 + assert service.calls == [("group", "scheduled")] + assert status["running"] is True + assert status["last_error_at"] is not None + assert status["consecutive_errors"] == 0 + assert "scheduler_tick_failed" in caplog.messages + + +@pytest.mark.asyncio +async def test_scheduler_reports_unhealthy_after_repeated_scan_errors() -> None: + repository = ScriptedRepository([RuntimeError("database locked") for _ in range(100)]) + service = RecordingRotationService() + scheduler = scheduler_for(repository, service) + + scheduler.start(0.01) + try: + await wait_until(lambda: repository.scan_count >= 3) + status = scheduler.status(enabled=True) + finally: + await scheduler.stop() + + assert status["running"] is True + assert status["healthy"] is False + assert int(status["consecutive_errors"]) >= 3 + + +@pytest.mark.asyncio +async def test_non_active_conflict_does_not_block_next_due_group() -> None: + repository = ScriptedRepository([["disabled-group", "ready-group"]]) + service = RecordingRotationService( + { + "disabled-group": ConflictError( + "实例组计划未启用", + code="GROUP_SCHEDULE_DISABLED", + ) + } + ) + scheduler = scheduler_for(repository, service) + + scheduler.start(0.01) + try: + await wait_until(service.dispatched.is_set) + finally: + await scheduler.stop() + + assert service.calls == [ + ("disabled-group", "scheduled"), + ("ready-group", "scheduled"), + ] + assert repository.deferred == [] + + +@pytest.mark.asyncio +async def test_unexpected_scheduler_cancellation_is_reported( + caplog: pytest.LogCaptureFixture, +) -> None: + repository = ScriptedRepository([[]]) + service = RecordingRotationService() + scheduler = scheduler_for(repository, service) + scheduler.start(1) + await wait_until(lambda: repository.scan_count >= 1) + task = scheduler._task + assert task is not None + + task.cancel() + await asyncio.gather(task, return_exceptions=True) + await asyncio.sleep(0) + status = scheduler.status(enabled=True) + + assert status["state"] == "stopped" + assert status["healthy"] is False + assert status["last_error_at"] is not None + assert "scheduler_cancelled_unexpectedly" in caplog.messages + + +def test_readyz_reports_scheduler_disabled(client: TestClient) -> None: + response = client.get("/readyz") + + assert response.status_code == 200 + assert response.json()["scheduler"] == { + "enabled": False, + "running": False, + "healthy": True, + "state": "disabled", + "tick_seconds": None, + "started_at": None, + "last_tick_at": None, + "last_success_at": None, + "last_dispatch_at": None, + "last_error_at": None, + "consecutive_errors": 0, + } + + +def test_readyz_is_degraded_when_enabled_scheduler_is_not_running(client: TestClient) -> None: + client.app.state.container.settings.scheduler_enabled = True + + response = client.get("/readyz") + + assert response.status_code == 503 + assert response.json()["status"] == "degraded" + assert response.json()["scheduler"]["healthy"] is False + assert response.json()["scheduler"]["state"] == "stopped"