Source.sync(ctx) -> SyncResult пишет только raw_* и возвращает курсор; локи, журнал, ошибки и продвижение курсора берёт на себя worker/runner. ZenMoney читается единственным доступным способом — POST /v8/diff/ по serverTimestamp; токен живёт сутки, поэтому worker ротирует refresh_token через source_credential. Маппер всегда пересобирает core из полных raw_*, так что удаление в ZenMoney исчезает и у нас. ЦБ ходит мимо прокси (trust_env=False) и отдаёт cp1251 с делением на Nominal. Курсы только по рабочим дням — протяжку по календарю делает аналитика. Планировщик — APScheduler в отдельном процессе, на источник advisory-лок sync:<name>, чтобы ручной запуск не пересёкся с плановым.
116 lines
4.1 KiB
Python
116 lines
4.1 KiB
Python
import pytest
|
|
|
|
from fintracker.sources import registry
|
|
from fintracker.sources.base import SyncContext, SyncResult
|
|
|
|
|
|
class FakeSource:
|
|
name = "fake"
|
|
|
|
def __init__(self) -> None:
|
|
self.calls = 0
|
|
|
|
async def sync(self, ctx: SyncContext) -> SyncResult:
|
|
self.calls += 1
|
|
before = int(ctx.cursor_before or 0)
|
|
return SyncResult(cursor_after=str(before + 1), counts={"rows": 3})
|
|
|
|
|
|
class FailingSource:
|
|
name = "boom"
|
|
|
|
async def sync(self, ctx: SyncContext) -> SyncResult:
|
|
raise RuntimeError("upstream down")
|
|
|
|
|
|
@pytest.fixture
|
|
def fake_source():
|
|
src = registry.register(FakeSource())
|
|
registry.register(FailingSource())
|
|
yield src
|
|
registry.unregister("fake")
|
|
registry.unregister("boom")
|
|
|
|
|
|
async def test_trigger_dedupes_queued_jobs(client, auth_headers, fake_source):
|
|
r1 = await client.post("/api/v1/sync/fake", headers=auth_headers)
|
|
r2 = await client.post("/api/v1/sync/fake", headers=auth_headers)
|
|
assert r1.status_code == 202 and r2.status_code == 202
|
|
assert r1.json()["id"] == r2.json()["id"]
|
|
|
|
r = await client.post("/api/v1/sync/unknown", headers=auth_headers)
|
|
assert r.status_code == 404
|
|
|
|
status = await client.get("/api/v1/sync/status", headers=auth_headers)
|
|
assert status.status_code == 200
|
|
by_name = {s["source"]: s for s in status.json()}
|
|
assert by_name["fake"]["queued"] is True
|
|
assert by_name["fake"]["cursor"] is None
|
|
|
|
|
|
async def test_worker_runs_queued_job_and_advances_cursor(client, auth_headers, fake_source):
|
|
from fintracker.worker.scheduler import _process_queued_jobs
|
|
|
|
await client.post("/api/v1/sync/fake", headers=auth_headers)
|
|
await _process_queued_jobs()
|
|
assert fake_source.calls == 1
|
|
|
|
status = await client.get("/api/v1/sync/status", headers=auth_headers)
|
|
fake = next(s for s in status.json() if s["source"] == "fake")
|
|
assert fake["cursor"] == "1"
|
|
assert fake["queued"] is False
|
|
assert fake["last_run_status"] == "ok"
|
|
|
|
runs = await client.get("/api/v1/sync/runs", params={"source": "fake"}, headers=auth_headers)
|
|
assert runs.json()[0]["counts"] == {"rows": 3}
|
|
assert runs.json()[0]["triggered_by"] == "manual"
|
|
|
|
|
|
async def test_failed_sync_is_recorded_not_raised(client, auth_headers, fake_source):
|
|
from fintracker.worker.runner import run_source
|
|
|
|
run = await run_source("boom", triggered_by="cli")
|
|
assert run.status.value == "error"
|
|
assert run.error is not None and "upstream down" in run.error
|
|
|
|
status = await client.get("/api/v1/sync/status", headers=auth_headers)
|
|
boom = next(s for s in status.json() if s["source"] == "boom")
|
|
assert boom["last_run_status"] == "error"
|
|
assert boom["last_success_at"] is None
|
|
|
|
|
|
async def test_metrics_failure_after_a_sync_marks_the_run_as_error(app, fake_source):
|
|
"""The download succeeded, so the cursor stays advanced; the run is still an error,
|
|
because the metrics the UI reads are stale."""
|
|
from fintracker.db import get_sessionmaker
|
|
from fintracker.metrics import refresh as refresh_module
|
|
from fintracker.models import SyncState
|
|
from fintracker.worker.runner import run_source
|
|
|
|
async def failing_step(session) -> None:
|
|
raise RuntimeError("step exploded")
|
|
|
|
refresh_module.STEPS.append(("test_failing_step", failing_step))
|
|
try:
|
|
run = await run_source("fake", triggered_by="cli")
|
|
finally:
|
|
refresh_module.STEPS[:] = [s for s in refresh_module.STEPS if s[0] != "test_failing_step"]
|
|
|
|
assert run.status.value == "error"
|
|
assert run.error is not None
|
|
assert run.error.startswith("metrics refresh failed:")
|
|
assert "step exploded" in run.error
|
|
# the sync itself succeeded: the cursor must not replay
|
|
assert run.cursor_after == "1"
|
|
async with get_sessionmaker()() as session:
|
|
state = await session.get(SyncState, "fake")
|
|
assert state is not None and state.cursor == "1"
|
|
|
|
|
|
async def test_a_later_refresh_is_clean_again(app, fake_source):
|
|
"""The temporary failing step is gone, so nothing leaks into other tests."""
|
|
from fintracker.worker.runner import run_source
|
|
|
|
run = await run_source("fake", triggered_by="cli")
|
|
assert run.status.value == "ok", run.error
|