import uuid from datetime import UTC, datetime import pytest from dirty_equals import IsDatetime from freezegun import freeze_time from starlette.testclient import TestClient from resonance_engine.dsp.enums import DSPClientName, DSPId from resonance_engine.dsp.models import DSPClient from resonance_engine.tasks.enums import FanoutSource, TaskFinishedReason, TaskStatus from resonance_engine.tasks.models import CollectTask, FanoutTask from tests.unit.helpers import create_model @pytest.mark.db def test_list_fanout_tasks_returns_paginated_response(client: TestClient) -> None: task = create_model( FanoutTask, started_at=datetime(2026, 5, 1, 12, 0, 0, tzinfo=UTC), finished_at=datetime(2026, 5, 1, 12, 5, 0, tzinfo=UTC), source=FanoutSource.scheduled, status=TaskStatus.done, finished_reason=TaskFinishedReason.done, fans_dispatched=10, messages_sent=2, ) create_model( CollectTask, fanout_task_id=task.id, dsp_client_name=DSPClientName.spotify_songwhip, status=TaskStatus.done, started_at=datetime(2026, 5, 1, 12, 0, 0, tzinfo=UTC), finished_at=datetime(2026, 5, 1, 12, 4, 0, tzinfo=UTC), fans_processed=8, fans_errors=2, requests=10, requests_rate_limited=1, ) create_model( CollectTask, fanout_task_id=task.id, dsp_client_name=DSPClientName.spotify_songwhip, status=TaskStatus.running, started_at=datetime(2026, 5, 1, 12, 0, 0, tzinfo=UTC), finished_at=None, ) response = client.get("/tasks/fanout") assert response.status_code == 200 assert response.json() == { "items": [ { "id": str(task.id), "started_at": IsDatetime(iso_string=True), "finished_at": IsDatetime(iso_string=True), "source": FanoutSource.scheduled, "status": TaskStatus.done, "finished_reason": TaskFinishedReason.done, "fans_dispatched": 10, "messages_sent": 2, "fans_collected": 8, "fans_errors": 2, "fans_stale_tokens": 0, "fans_skipped": 0, "requests": 10, "requests_rate_limited": 1, "collect_tasks_done": 1, "filters": { "dsp_client_id": None, "status": "active", "not_collected_since": None, "not_last_dispatched_since": None, "max_consecutive_failures": None, "limit": None, }, } ], "next_cursor": None, "total": 1, } @pytest.mark.db def test_get_fanout_task_returns_task(client: TestClient) -> None: task = create_model( FanoutTask, started_at=datetime(2026, 5, 1, 12, 0, 0, tzinfo=UTC), source=FanoutSource.manual, status=TaskStatus.running, ) response = client.get(f"/tasks/fanout/{task.id}") assert response.status_code == 200 assert response.json() == { "id": str(task.id), "started_at": IsDatetime(iso_string=True), "finished_at": None, "source": FanoutSource.manual, "status": TaskStatus.running, "finished_reason": None, "fans_dispatched": 0, "messages_sent": 0, "fans_collected": 0, "fans_errors": 0, "fans_stale_tokens": 0, "fans_skipped": 0, "requests": 0, "requests_rate_limited": 0, "collect_tasks_done": 0, "filters": { "dsp_client_id": None, "status": "active", "not_collected_since": None, "not_last_dispatched_since": None, "max_consecutive_failures": None, "limit": None, }, } @pytest.mark.db def test_list_fanout_tasks_filters_by_status(client: TestClient) -> None: running = create_model( FanoutTask, started_at=datetime(2026, 5, 1, 12, 0, 0, tzinfo=UTC), status=TaskStatus.running, ) create_model( FanoutTask, started_at=datetime(2026, 5, 1, 11, 0, 0, tzinfo=UTC), status=TaskStatus.done, finished_at=datetime(2026, 5, 1, 11, 5, 0, tzinfo=UTC), finished_reason=TaskFinishedReason.done, ) response = client.get("/tasks/fanout", params={"status": "running"}) assert response.status_code == 200 assert response.json() == { "items": [ { "id": str(running.id), "started_at": IsDatetime(iso_string=True), "finished_at": None, "source": FanoutSource.scheduled, "status": TaskStatus.running, "finished_reason": None, "fans_dispatched": 0, "messages_sent": 0, "fans_collected": 0, "fans_errors": 0, "fans_stale_tokens": 0, "fans_skipped": 0, "requests": 0, "requests_rate_limited": 0, "collect_tasks_done": 0, "filters": { "dsp_client_id": None, "status": "active", "not_collected_since": None, "not_last_dispatched_since": None, "max_consecutive_failures": None, "limit": None, }, } ], "next_cursor": None, "total": 1, } @pytest.mark.db def test_list_fanout_tasks_filters_by_dsp_client(client: TestClient) -> None: match = create_model( FanoutTask, started_at=datetime(2026, 5, 1, 12, 0, 0, tzinfo=UTC), status=TaskStatus.running, ) create_model( CollectTask, fanout_task_id=match.id, dsp_client_name=DSPClientName.spotify_songwhip, status=TaskStatus.done, fans_total=5, started_at=datetime(2026, 5, 1, 12, 0, 0, tzinfo=UTC), ) other = create_model( FanoutTask, started_at=datetime(2026, 5, 1, 11, 0, 0, tzinfo=UTC), status=TaskStatus.running, ) create_model( CollectTask, fanout_task_id=other.id, dsp_client_name=DSPClientName.spotify_smf_sme, status=TaskStatus.done, fans_total=5, started_at=datetime(2026, 5, 1, 11, 0, 0, tzinfo=UTC), ) response = client.get( "/tasks/fanout", params={"dsp_client_name": "spotify_songwhip"} ) assert response.status_code == 200 assert response.json() == { "items": [ { "id": str(match.id), "started_at": IsDatetime(iso_string=True), "finished_at": None, "source": FanoutSource.scheduled, "status": TaskStatus.running, "finished_reason": None, "fans_dispatched": 0, "messages_sent": 0, "fans_collected": 0, "fans_errors": 0, "fans_stale_tokens": 0, "fans_skipped": 0, "requests": 0, "requests_rate_limited": 0, "collect_tasks_done": 0, "filters": { "dsp_client_id": None, "status": "active", "not_collected_since": None, "not_last_dispatched_since": None, "max_consecutive_failures": None, "limit": None, }, } ], "next_cursor": None, "total": 1, } @pytest.mark.db def test_get_fanout_task_returns_404_for_unknown(client: TestClient) -> None: response = client.get(f"/tasks/fanout/{uuid.uuid4()}") assert response.status_code == 404 @pytest.mark.db def test_list_collect_tasks_returns_tasks_for_fanout(client: TestClient) -> None: fanout = create_model( FanoutTask, started_at=datetime(2026, 5, 1, 12, 0, 0, tzinfo=UTC), finished_at=datetime(2026, 5, 1, 12, 5, 0, tzinfo=UTC), source=FanoutSource.scheduled, status=TaskStatus.done, finished_reason=TaskFinishedReason.done, ) collect = create_model( CollectTask, fanout_task_id=fanout.id, started_at=datetime(2026, 5, 1, 12, 0, 0, tzinfo=UTC), finished_at=datetime(2026, 5, 1, 12, 5, 0, tzinfo=UTC), dsp_client_name=DSPClientName.spotify_songwhip, status=TaskStatus.done, finished_reason=TaskFinishedReason.done, fans_total=50, fans_skipped=0, requests_skipped=0, fans_processed=48, fans_errors=2, fans_stale_tokens=1, requests_rate_limited=0, requests=50, ) response = client.get(f"/tasks/fanout/{fanout.id}/collect") assert response.status_code == 200 assert response.json() == [ { "id": str(collect.id), "fanout_task_id": str(fanout.id), "dsp_client_name": DSPClientName.spotify_songwhip, "started_at": IsDatetime(iso_string=True), "finished_at": IsDatetime(iso_string=True), "status": TaskStatus.done, "finished_reason": TaskFinishedReason.done, "fans_total": 50, "fans_skipped": 0, "requests_skipped": 0, "fans_processed": 48, "fans_errors": 2, "fans_stale_tokens": 1, "requests_rate_limited": 0, "requests": 50, } ] @pytest.mark.db @freeze_time("2026-05-25T12:00:00Z") def test_activity_returns_bucketed_collect_stats(client: TestClient) -> None: create_model( DSPClient, dsp_id=DSPId.spotify, name=DSPClientName.spotify_songwhip, display_name="Spotify Songwhip", ) fanout = create_model( FanoutTask, started_at=datetime(2026, 5, 24, 10, 0, 0, tzinfo=UTC), status=TaskStatus.done, finished_at=datetime(2026, 5, 24, 10, 5, 0, tzinfo=UTC), finished_reason=TaskFinishedReason.done, ) create_model( CollectTask, fanout_task_id=fanout.id, started_at=datetime(2026, 5, 24, 10, 0, 0, tzinfo=UTC), finished_at=datetime(2026, 5, 24, 10, 4, 0, tzinfo=UTC), status=TaskStatus.done, finished_reason=TaskFinishedReason.done, fans_total=100, fans_processed=90, fans_errors=5, fans_stale_tokens=2, requests_rate_limited=3, requests=95, ) response = client.get("/tasks/activity", params={"days": 7, "granularity": "daily"}) assert response.status_code == 200 assert response.json() == [ { "bucket": "2026-05-24", "dsp_client_name": "spotify_songwhip", "dsp_client_display_name": "Spotify Songwhip", "fans_processed": 90, "fans_errors": 5, "fans_stale_tokens": 2, "requests_rate_limited": 3, "requests": 95, } ] def test_activity_rejects_unknown_granularity(client: TestClient) -> None: response = client.get("/tasks/activity", params={"granularity": "weekly"}) assert response.status_code == 422