From 3b3ee4d682d841e34c6e6633df881805625bbed2 Mon Sep 17 00:00:00 2001 From: Dmitry Date: Sat, 19 Sep 2026 21:54:30 +0300 Subject: [PATCH] =?UTF-8?q?feat(worker):=20=D1=80=D1=83=D1=87=D0=BD=D0=BE?= =?UTF-8?q?=D0=B9=20=D0=BF=D0=B5=D1=80=D0=B5=D1=81=D1=87=D1=91=D1=82=20?= =?UTF-8?q?=D0=BC=D0=B5=D1=82=D1=80=D0=B8=D0=BA=20=D1=87=D0=B5=D1=80=D0=B5?= =?UTF-8?q?=D0=B7=20=D0=BE=D1=87=D0=B5=D1=80=D0=B5=D0=B4=D1=8C,=20=D0=B4?= =?UTF-8?q?=D0=B8=D0=B0=D0=B3=D0=BD=D0=BE=D1=81=D1=82=D0=B8=D0=BA=D0=B0=20?= =?UTF-8?q?=D1=88=D0=B0=D0=B3=D0=BE=D0=B2=20=D0=B8=20heartbeat?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit POST /metrics/refresh ставит задачу в sync_job (source=METRICS_JOB) и отвечает 202, пересчёт делает worker; GET /metrics/status отдаёт refreshing и consistent. metric_refresh_log хранит failed_step и step_timings. Источники с needs="tinvest_token" не попадают в расписание без токена, tinvest/moex добавлены в default_schedule. Воркер трогает heartbeat-файл для healthcheck. --- ...b05_metric_refresh_log_step_diagnostics.py | 34 ++++++ backend/src/fintracker/api/routers/metrics.py | 68 +++++++++--- backend/src/fintracker/api/schemas/metrics.py | 14 +++ backend/src/fintracker/metrics/refresh.py | 16 ++- backend/src/fintracker/models/__init__.py | 2 + backend/src/fintracker/models/metrics.py | 5 + backend/src/fintracker/models/sync.py | 8 +- backend/src/fintracker/worker/jobs.py | 27 ++++- backend/src/fintracker/worker/scheduler.py | 45 +++++++- backend/tests/api/test_metrics_api.py | 103 ++++++++++++++++-- backend/tests/test_schedule.py | 49 +++++++++ 11 files changed, 341 insertions(+), 30 deletions(-) create mode 100644 backend/alembic/versions/c41e7a9d2b05_metric_refresh_log_step_diagnostics.py create mode 100644 backend/tests/test_schedule.py diff --git a/backend/alembic/versions/c41e7a9d2b05_metric_refresh_log_step_diagnostics.py b/backend/alembic/versions/c41e7a9d2b05_metric_refresh_log_step_diagnostics.py new file mode 100644 index 0000000..98a5299 --- /dev/null +++ b/backend/alembic/versions/c41e7a9d2b05_metric_refresh_log_step_diagnostics.py @@ -0,0 +1,34 @@ +"""metric_refresh_log: имя упавшего шага и тайминги шагов + +Revision ID: c41e7a9d2b05 +Revises: 2bf84b07fd5e +Create Date: 2026-09-19 17:10:00.000000 +""" + +from __future__ import annotations + +from collections.abc import Sequence + +import sqlalchemy as sa +from alembic import op +from sqlalchemy.dialects import postgresql + +revision: str = "c41e7a9d2b05" +down_revision: str | None = "2bf84b07fd5e" +branch_labels: str | Sequence[str] | None = None +depends_on: str | Sequence[str] | None = None + + +def upgrade() -> None: + op.add_column( + "metric_refresh_log", sa.Column("failed_step", sa.String(length=64), nullable=True) + ) + op.add_column( + "metric_refresh_log", + sa.Column("step_timings", postgresql.JSONB(astext_type=sa.Text()), nullable=True), + ) + + +def downgrade() -> None: + op.drop_column("metric_refresh_log", "step_timings") + op.drop_column("metric_refresh_log", "failed_step") diff --git a/backend/src/fintracker/api/routers/metrics.py b/backend/src/fintracker/api/routers/metrics.py index 76613ca..55bb0ce 100644 --- a/backend/src/fintracker/api/routers/metrics.py +++ b/backend/src/fintracker/api/routers/metrics.py @@ -4,9 +4,9 @@ from fastapi import APIRouter, status from sqlalchemy import select from fintracker.api.deps import CurrentUser, SessionDep -from fintracker.api.schemas.metrics import DataQualityRow, RefreshLogOut -from fintracker.metrics.refresh import refresh_all -from fintracker.models import MetricDataQuality, MetricRefreshLog +from fintracker.api.schemas.metrics import DataQualityRow, MetricsStatusOut, RefreshLogOut +from fintracker.api.schemas.sync import SyncJobOut +from fintracker.models import METRICS_JOB, JobStatus, MetricDataQuality, MetricRefreshLog, SyncJob router = APIRouter(tags=["metrics"]) @@ -22,9 +22,10 @@ async def data_quality(session: SessionDep, _: CurrentUser) -> list[DataQualityR @router.get("/metrics/status", name="status") -async def metrics_status(session: SessionDep, _: CurrentUser) -> RefreshLogOut | None: - """When the metric tables were last rebuilt, and whether it failed.""" - row = ( +async def metrics_status(session: SessionDep, _: CurrentUser) -> MetricsStatusOut: + """When the metric tables were last rebuilt, whether they are one consistent snapshot, + and whether another rebuild is on its way.""" + latest = ( ( await session.execute( select(MetricRefreshLog).order_by(MetricRefreshLog.started_at.desc()).limit(1) @@ -33,13 +34,54 @@ async def metrics_status(session: SessionDep, _: CurrentUser) -> RefreshLogOut | .scalars() .first() ) - if row is None: - return None - return RefreshLogOut.model_validate(row, from_attributes=True) + last_finished = ( + ( + await session.execute( + select(MetricRefreshLog) + .where(MetricRefreshLog.finished_at.is_not(None)) + .order_by(MetricRefreshLog.started_at.desc()) + .limit(1) + ) + ) + .scalars() + .first() + ) + queued = ( + await session.execute( + select(SyncJob.id) + .where( + SyncJob.source == METRICS_JOB, + SyncJob.status.in_([JobStatus.queued, JobStatus.running]), + ) + .limit(1) + ) + ).first() + return MetricsStatusOut( + last_refresh=RefreshLogOut.model_validate(latest, from_attributes=True) if latest else None, + consistent=last_finished is None or last_finished.error is None, + refreshing=queued is not None or (latest is not None and latest.finished_at is None), + ) @router.post("/metrics/refresh", name="refresh", status_code=status.HTTP_202_ACCEPTED) -async def metrics_refresh(session: SessionDep, _: CurrentUser) -> RefreshLogOut: - """Rebuild every metric_* table inline (seconds at personal volumes).""" - entry = await refresh_all(session, trigger="manual") - return RefreshLogOut.model_validate(entry, from_attributes=True) +async def metrics_refresh(session: SessionDep, _: CurrentUser) -> SyncJobOut: + """Queue a rebuild of every metric_* table; the worker runs it. Poll `/metrics/status` + until `refreshing` is false. + + A request while one is already waiting shares it. One that arrives while a rebuild is + RUNNING queues another, because the running one may have read the data before the change + that prompted this call.""" + waiting = ( + await session.execute( + select(SyncJob) + .where(SyncJob.source == METRICS_JOB, SyncJob.status == JobStatus.queued) + .limit(1) + ) + ).scalar_one_or_none() + if waiting is not None: + return SyncJobOut.model_validate(waiting, from_attributes=True) + job = SyncJob(source=METRICS_JOB, status=JobStatus.queued) + session.add(job) + await session.commit() + await session.refresh(job) + return SyncJobOut.model_validate(job, from_attributes=True) diff --git a/backend/src/fintracker/api/schemas/metrics.py b/backend/src/fintracker/api/schemas/metrics.py index f6b9e3c..1d084da 100644 --- a/backend/src/fintracker/api/schemas/metrics.py +++ b/backend/src/fintracker/api/schemas/metrics.py @@ -23,3 +23,17 @@ class RefreshLogOut(BaseModel): finished_at: datetime | None trigger: str error: str | None + failed_step: str | None + """The step that raised; earlier steps committed, later ones did not run.""" + step_timings: dict[str, float] | None + """Seconds per step that completed, in run order.""" + + +class MetricsStatusOut(BaseModel): + last_refresh: RefreshLogOut | None + """Latest run, finished or not; null before the first one.""" + consistent: bool + """False when the latest FINISHED run failed part-way: some metric tables are from it and + the rest from an earlier run, so numbers on different screens may disagree.""" + refreshing: bool + """A rebuild is queued or running, so `last_refresh` is about to be superseded.""" diff --git a/backend/src/fintracker/metrics/refresh.py b/backend/src/fintracker/metrics/refresh.py index fc4cf3a..88ce11a 100644 --- a/backend/src/fintracker/metrics/refresh.py +++ b/backend/src/fintracker/metrics/refresh.py @@ -1,8 +1,8 @@ """Rebuild every metric_* table after data changed (plan §3). -Order (each step is a function taking an AsyncSession and committing its own tables): - fx -> classify -> net worth -> cash flow -> spending -> runway -> data quality -Phase 2 inserts prices/lots/valuation/holdings/returns between fx and net worth. +Each step is a function taking an AsyncSession and committing its own tables. The steps and +their order live in `analytics.register_steps`, with the reason for each dependency — this +module only runs whatever is registered, in registration order. `refresh_all` is what the worker calls after a sync that reported `changed=True`, what the CLI `fintracker metrics refresh` runs, and what `POST /metrics/refresh` queues. @@ -16,6 +16,7 @@ connection serialises them: a refresh that arrives during a sync WAITS and then from __future__ import annotations import logging +import time import traceback from collections.abc import Awaitable, Callable from datetime import UTC, datetime @@ -52,15 +53,22 @@ async def refresh_all(session: AsyncSession, trigger: str) -> MetricRefreshLog: entry = MetricRefreshLog(trigger=trigger) session.add(entry) await session.commit() + timings: dict[str, float] = {} + current: str | None = None try: for name, step in STEPS: + current = name log.info("metrics: %s", name) + started = time.monotonic() await step(session) await session.commit() + timings[name] = round(time.monotonic() - started, 3) except Exception as exc: await session.rollback() entry.error = f"{type(exc).__name__}: {exc}\n{traceback.format_exc()[-4000:]}" - log.exception("metrics refresh failed at step") + entry.failed_step = current + log.exception("metrics refresh failed at step %s", current) + entry.step_timings = timings entry.finished_at = datetime.now(UTC) await session.commit() return entry diff --git a/backend/src/fintracker/models/__init__.py b/backend/src/fintracker/models/__init__.py index 4ee695a..23af6c2 100644 --- a/backend/src/fintracker/models/__init__.py +++ b/backend/src/fintracker/models/__init__.py @@ -72,6 +72,7 @@ from fintracker.models.reports import ( ReportParseStatus, ) from fintracker.models.sync import ( + METRICS_JOB, JobStatus, RunStatus, SourceCredential, @@ -101,6 +102,7 @@ from fintracker.models.zenmoney import ( __all__ = [ "EXTERNAL_FLOW_KINDS", + "METRICS_JOB", "POSITION_KINDS", "Account", "AccountKind", diff --git a/backend/src/fintracker/models/metrics.py b/backend/src/fintracker/models/metrics.py index 2dc5fc3..0d32d35 100644 --- a/backend/src/fintracker/models/metrics.py +++ b/backend/src/fintracker/models/metrics.py @@ -135,6 +135,11 @@ class MetricRefreshLog(Base): trigger: Mapped[str] = mapped_column(String(32)) """sync: | manual | cli""" error: Mapped[str | None] = mapped_column(Text) + failed_step: Mapped[str | None] = mapped_column(String(64)) + """Name of the step that raised. Steps before it committed and the ones after it did not + run, so the metric tables then come from two different runs.""" + step_timings: Mapped[dict[str, Any] | None] + """Seconds per step that completed, in run order.""" class MetricPortfolioValueDaily(Base): diff --git a/backend/src/fintracker/models/sync.py b/backend/src/fintracker/models/sync.py index addb6c0..73a0dc1 100644 --- a/backend/src/fintracker/models/sync.py +++ b/backend/src/fintracker/models/sync.py @@ -19,6 +19,11 @@ class RunStatus(enum.StrEnum): error = "error" +METRICS_JOB = "metrics" +"""`sync_job.source` of a queued metrics rebuild. Not a registered source: the worker +routes it to `refresh_all` instead of `run_source`, and `POST /sync/{source}` answers 404.""" + + class JobStatus(enum.StrEnum): queued = "queued" running = "running" @@ -55,7 +60,8 @@ class SyncRun(Base): class SyncJob(Base): - """Manual trigger queue: API inserts, worker picks up (plan §2.4).""" + """Manual trigger queue: API inserts, worker picks up (plan §2.4). Also carries the metrics + rebuild requested by `POST /metrics/refresh`, as `source == METRICS_JOB`.""" __tablename__ = "sync_job" diff --git a/backend/src/fintracker/worker/jobs.py b/backend/src/fintracker/worker/jobs.py index 3559e07..133ee8a 100644 --- a/backend/src/fintracker/worker/jobs.py +++ b/backend/src/fintracker/worker/jobs.py @@ -9,6 +9,8 @@ from apscheduler.triggers.base import BaseTrigger from apscheduler.triggers.combining import OrTrigger from apscheduler.triggers.cron import CronTrigger +from fintracker.config import Settings + @dataclass(frozen=True) class JobSpec: @@ -16,10 +18,13 @@ class JobSpec: trigger: BaseTrigger """A source appears at most once: `worker/scheduler.py` keys jobs by `sync:`, so several times of day are one `OrTrigger`, not several specs.""" + needs: str | None = None + """Name of a `Settings` field that must be set for the source to run at all. Without it + the job is left off the timetable: a run would only fail, hourly, and bury real errors.""" MSK = ZoneInfo("Europe/Moscow") -"""CBR publishes on Moscow time, so those jobs pin the zone instead of following settings.""" +"""CBR and MOEX publish on Moscow time, so those jobs pin the zone instead of following settings.""" def default_schedule() -> list[JobSpec]: @@ -36,4 +41,24 @@ def default_schedule() -> list[JobSpec]: ] ), ), + # Broker operations and balances, every three hours. The MOEX run below is placed + # after a T-Invest one, because it prices whatever the ledger says is held. + JobSpec( + "tinvest", CronTrigger(hour="8-23/3", minute=10, timezone=MSK), needs="tinvest_token" + ), + # Prices: after the open, mid-session, after the close (main session ends 18:50 MSK), + # and once more for the evening session. + JobSpec("moex", CronTrigger(hour="10,14,19,23", minute=20, timezone=MSK)), + # Announced dividends and coupons change rarely: once a day, before the market opens. + JobSpec( + "tinvest_events", CronTrigger(hour=6, minute=30, timezone=MSK), needs="tinvest_token" + ), + JobSpec("moex_payouts", CronTrigger(hour=6, minute=45, timezone=MSK)), ] + + +def scheduled(settings: Settings) -> tuple[list[JobSpec], list[JobSpec]]: + """(to run, left off because a required setting is missing).""" + specs = default_schedule() + ready = [s for s in specs if s.needs is None or getattr(settings, s.needs)] + return ready, [s for s in specs if s not in ready] diff --git a/backend/src/fintracker/worker/scheduler.py b/backend/src/fintracker/worker/scheduler.py index 9e5c03f..6f7e4ab 100644 --- a/backend/src/fintracker/worker/scheduler.py +++ b/backend/src/fintracker/worker/scheduler.py @@ -6,20 +6,30 @@ import asyncio import logging import signal from datetime import UTC, datetime +from pathlib import Path from apscheduler.schedulers.asyncio import AsyncIOScheduler from sqlalchemy import select from fintracker.config import get_settings from fintracker.db import get_sessionmaker, reset_engine -from fintracker.models import JobStatus, SyncJob +from fintracker.metrics.refresh import refresh_all +from fintracker.models import METRICS_JOB, JobStatus, SyncJob from fintracker.sources import registry -from fintracker.worker.jobs import default_schedule +from fintracker.worker.jobs import scheduled from fintracker.worker.runner import Skipped, run_source log = logging.getLogger(__name__) POLL_SECONDS = 5 +HEARTBEAT_SECONDS = 10 +HEARTBEAT_FILE = Path("/tmp/fintracker-worker-heartbeat") +"""Touched by its own job, so a long sync (which blocks `_process_queued_jobs`, not the event +loop) does not look like a dead worker. The compose healthcheck reads its mtime.""" + + +async def _heartbeat() -> None: + HEARTBEAT_FILE.touch() async def _scheduled(source: str) -> None: @@ -43,6 +53,9 @@ async def _process_queued_jobs() -> None: .all() ) for job in jobs: + if job.source == METRICS_JOB: + await _run_metrics_job(job.id) + continue if job.source not in registry.names(): await _finish_job(job.id, JobStatus.error, error=f"unknown source {job.source}") continue @@ -56,6 +69,20 @@ async def _process_queued_jobs() -> None: await _finish_job(job.id, status, run_id=run.id, error=run.error) +async def _run_metrics_job(job_id) -> None: + await _mark_running(job_id) + try: + async with get_sessionmaker()() as session: + entry = await refresh_all(session, trigger="manual") + except Exception as exc: # the lock connection or the log insert, not a step (those are caught) + log.exception("metrics job failed") + await _finish_job(job_id, JobStatus.error, error=f"{type(exc).__name__}: {exc}") + return + await _finish_job( + job_id, JobStatus.done if entry.error is None else JobStatus.error, error=entry.error + ) + + async def _mark_running(job_id) -> None: async with get_sessionmaker()() as session: job = await session.get(SyncJob, job_id) @@ -79,7 +106,10 @@ async def _finish_job(job_id, status: JobStatus, *, run_id=None, error: str | No def build_scheduler() -> AsyncIOScheduler: settings = get_settings() scheduler = AsyncIOScheduler(timezone=settings.timezone) - for spec in default_schedule(): + specs, skipped = scheduled(settings) + for spec in skipped: + log.warning("not scheduling %s: %s is not set", spec.source, spec.needs) + for spec in specs: scheduler.add_job( _scheduled, spec.trigger, @@ -89,6 +119,15 @@ def build_scheduler() -> AsyncIOScheduler: coalesce=True, misfire_grace_time=600, ) + scheduler.add_job( + _heartbeat, + "interval", + seconds=HEARTBEAT_SECONDS, + id="heartbeat", + max_instances=1, + coalesce=True, + next_run_time=datetime.now(UTC), + ) scheduler.add_job( _process_queued_jobs, "interval", diff --git a/backend/tests/api/test_metrics_api.py b/backend/tests/api/test_metrics_api.py index 77a0c9e..21cba65 100644 --- a/backend/tests/api/test_metrics_api.py +++ b/backend/tests/api/test_metrics_api.py @@ -1,8 +1,19 @@ from datetime import timedelta +from sqlalchemy import select + from factories import make_account, make_category, make_txn, month_back from fintracker.analytics import today_local -from fintracker.models import AccountRole +from fintracker.models import METRICS_JOB, AccountRole, JobStatus, SyncJob +from fintracker.worker.scheduler import _process_queued_jobs + + +async def _queue_and_run(client, auth_headers): + """POST /metrics/refresh, then let the worker take the job, as it does every 5 s.""" + r = await client.post("/api/v1/metrics/refresh", headers=auth_headers) + assert r.status_code == 202, r.text + await _process_queued_jobs() + return r.json() async def test_refresh_populates_every_metric_endpoint(client, auth_headers): @@ -20,15 +31,20 @@ async def test_refresh_populates_every_metric_endpoint(client, auth_headers): ) await make_txn(today_local() - timedelta(days=1), outcome="500", outcome_account_id=card) - assert (await client.get("/api/v1/metrics/status", headers=auth_headers)).json() is None + before = (await client.get("/api/v1/metrics/status", headers=auth_headers)).json() + assert before == {"last_refresh": None, "consistent": True, "refreshing": False} - r = await client.post("/api/v1/metrics/refresh", headers=auth_headers) - assert r.status_code == 202 - assert r.json()["error"] is None - assert r.json()["finished_at"] is not None + job = await _queue_and_run(client, auth_headers) + assert job["status"] == "queued" status = (await client.get("/api/v1/metrics/status", headers=auth_headers)).json() - assert status["trigger"] == "manual" + assert status["last_refresh"]["trigger"] == "manual" + assert status["last_refresh"]["error"] is None + assert status["last_refresh"]["finished_at"] is not None + assert status["last_refresh"]["failed_step"] is None + assert status["last_refresh"]["step_timings"] # every step is timed + assert status["consistent"] is True + assert status["refreshing"] is False series = (await client.get("/api/v1/networth/series", headers=auth_headers)).json() assert series @@ -81,7 +97,7 @@ async def test_spending_rejects_a_bad_month(client, auth_headers): async def test_empty_database_reports_no_transactions(client, auth_headers): - await client.post("/api/v1/metrics/refresh", headers=auth_headers) + await _queue_and_run(client, auth_headers) quality = (await client.get("/api/v1/data-quality", headers=auth_headers)).json() assert [row["check_name"] for row in quality] == ["no_transactions"] assert (await client.get("/api/v1/runway", headers=auth_headers)).json()[ @@ -89,3 +105,74 @@ async def test_empty_database_reports_no_transactions(client, auth_headers): ] is None assert (await client.get("/api/v1/networth/series", headers=auth_headers)).json() == [] assert (await client.get("/api/v1/spending/categories", headers=auth_headers)).json() == [] + + +async def test_refresh_returns_before_the_work_is_done(client, auth_headers): + """The 202 is honest: nothing has been rebuilt until the worker picks the job up.""" + r = await client.post("/api/v1/metrics/refresh", headers=auth_headers) + assert r.status_code == 202 + assert r.json()["source"] == METRICS_JOB + + status = (await client.get("/api/v1/metrics/status", headers=auth_headers)).json() + assert status["last_refresh"] is None + assert status["refreshing"] is True + + await _process_queued_jobs() + + status = (await client.get("/api/v1/metrics/status", headers=auth_headers)).json() + assert status["refreshing"] is False + assert status["last_refresh"] is not None + + +async def test_a_waiting_refresh_is_shared_not_duplicated(client, auth_headers): + first = (await client.post("/api/v1/metrics/refresh", headers=auth_headers)).json() + second = (await client.post("/api/v1/metrics/refresh", headers=auth_headers)).json() + assert first["id"] == second["id"] + + from fintracker.db import get_sessionmaker + + async with get_sessionmaker()() as session: + jobs = ( + (await session.execute(select(SyncJob).where(SyncJob.source == METRICS_JOB))) + .scalars() + .all() + ) + assert len(jobs) == 1 + + +async def test_a_refresh_that_fails_part_way_is_reported_as_inconsistent( + client, auth_headers, monkeypatch +): + from fintracker.db import get_sessionmaker + from fintracker.metrics import refresh as refresh_module + + async def fine(session): + pass + + async def boom(session): + raise RuntimeError("step blew up") + + monkeypatch.setattr("fintracker.analytics.register_steps", lambda: None) + monkeypatch.setattr(refresh_module, "STEPS", [("fine", fine), ("boom", boom), ("never", fine)]) + + await _queue_and_run(client, auth_headers) + + status = (await client.get("/api/v1/metrics/status", headers=auth_headers)).json() + last = status["last_refresh"] + assert last["failed_step"] == "boom" + assert "step blew up" in last["error"] + assert list(last["step_timings"]) == ["fine"] # 'never' did not run, 'boom' did not finish + assert status["consistent"] is False + + async with get_sessionmaker()() as session: + job = ( + await session.execute(select(SyncJob).where(SyncJob.source == METRICS_JOB)) + ).scalar_one() + assert job.status == JobStatus.error + + # a later clean run makes the tables one snapshot again + monkeypatch.setattr(refresh_module, "STEPS", [("fine", fine)]) + await _queue_and_run(client, auth_headers) + status = (await client.get("/api/v1/metrics/status", headers=auth_headers)).json() + assert status["consistent"] is True + assert status["last_refresh"]["failed_step"] is None diff --git a/backend/tests/test_schedule.py b/backend/tests/test_schedule.py new file mode 100644 index 0000000..42d482c --- /dev/null +++ b/backend/tests/test_schedule.py @@ -0,0 +1,49 @@ +"""The timetable only names sources that exist, and leaves off the ones that cannot run.""" + +from __future__ import annotations + +from fintracker.config import Settings +from fintracker.sources import registry +from fintracker.worker.jobs import default_schedule, scheduled + + +def test_every_scheduled_source_is_registered_once(): + names = [s.source for s in default_schedule()] + assert len(names) == len(set(names)) + assert set(names) <= set(registry.names()) + + +def test_every_live_source_is_on_the_timetable(): + assert {s.source for s in default_schedule()} == set(registry.names()) + + +def test_sources_needing_a_token_are_left_off_without_one(): + ready, skipped = scheduled(Settings.model_construct(tinvest_token=None)) + assert {s.source for s in skipped} == {"tinvest", "tinvest_events"} + assert not {s.source for s in ready} & {"tinvest", "tinvest_events"} + + ready, skipped = scheduled(Settings.model_construct(tinvest_token="t.x")) + assert skipped == [] + assert {"tinvest", "moex", "tinvest_events", "moex_payouts"} <= {s.source for s in ready} + + +def test_needs_names_a_real_setting(): + for spec in default_schedule(): + if spec.needs is not None: + assert spec.needs in Settings.model_fields + + +def test_worker_timetable_has_the_heartbeat_and_the_queue_poller(): + from fintracker.worker.scheduler import build_scheduler + + ids = {j.id for j in build_scheduler().get_jobs()} + assert {"heartbeat", "poll-sync-jobs", "sync:zenmoney", "sync:moex", "sync:moex_payouts"} <= ids + + +async def test_heartbeat_touches_its_file(tmp_path, monkeypatch): + from fintracker.worker import scheduler + + beat = tmp_path / "beat" + monkeypatch.setattr(scheduler, "HEARTBEAT_FILE", beat) + await scheduler._heartbeat() + assert beat.exists()