import asyncio
import json

from app.models.events import FileChangedEvent
from app.services import redis_consumer_service


class _FakeRedis:
    def __init__(self):
        self.acked = []
        self.dlq = []

    async def xack(self, stream, group, message_id):
        self.acked.append((stream, group, message_id))

    async def xpending_range(self, stream, group, start, end, count):
        return [{"times_delivered": 1}]

    async def xadd(self, stream, fields):
        self.dlq.append((stream, fields))


def _event_payload():
    return {
        "event_id": "evt_1",
        "event_version": "v1",
        "created_at": "2026-05-27T00:00:00Z",
        "repo_id": "ad/repo",
        "snapshot_id": "snap_1",
        "commit_sha": "abc123",
        "file_path": "README.md",
        "language": "markdown",
        "file_kind": "docs",
        "change_type": "added",
    }


async def _skipped_handler(event):
    return "skipped"


async def _processed_handler(event):
    return "processed"


async def _retry_handler(event):
    raise RuntimeError("transient failure")


def test_process_stream_message_acks_skipped_event(monkeypatch):
    redis = _FakeRedis()
    increments = []

    async def _record_increment(repo_id, snapshot_id):
        increments.append((repo_id, snapshot_id))

    monkeypatch.setattr(
        redis_consumer_service,
        "process_file_changed_event",
        _skipped_handler,
    )
    monkeypatch.setattr(
        redis_consumer_service.indexing_run_repo,
        "increment_processed_docs",
        _record_increment,
    )

    result = asyncio.run(
        redis_consumer_service.process_stream_message(
            redis,
            "files.changed",
            "1-0",
            {"payload": json.dumps(_event_payload())},
        )
    )

    assert result == "acked"
    assert redis.acked == [("files.changed", "docs-ingestion-service", "1-0")]
    assert increments == [("ad/repo", "snap_1")]


def test_process_stream_message_retries_transient_failure(monkeypatch):
    redis = _FakeRedis()
    monkeypatch.setattr(
        redis_consumer_service,
        "process_file_changed_event",
        _retry_handler,
    )
    monkeypatch.setenv("REDIS_STREAM_MAX_RETRIES", "5")

    result = asyncio.run(
        redis_consumer_service.process_stream_message(
            redis,
            "files.changed",
            "2-0",
            {"payload": json.dumps(_event_payload())},
        )
    )

    assert result == "retry"
    assert redis.acked == []


def test_build_dlq_stream_key_uses_prefix(monkeypatch):
    monkeypatch.setenv("REDIS_STREAM_PREFIX", "dev")
    assert redis_consumer_service.build_dlq_stream_key() == "dev.files.changed.dlq"
