from __future__ import annotations

import asyncio
import signal

from app.core.config import settings
from app.core.logger import logger
from app.core.redis_client import get_redis
from app.redis.stream_consumer import StreamConsumer


async def _run() -> None:
    settings.validate_runtime()
    stop_event = asyncio.Event()
    loop = asyncio.get_running_loop()

    for sig in (signal.SIGINT, signal.SIGTERM):
        loop.add_signal_handler(sig, stop_event.set)

    redis = get_redis()
    consumer = StreamConsumer(redis)
    logger.info(
        "Starting Redis stream consumer group={group} consumer={consumer}",
        group=settings.CONSUMER_GROUP,
        consumer=consumer.consumer_name,
    )

    try:
        await consumer.run(stop_event)
    finally:
        await redis.aclose()
        logger.info("Redis stream consumer stopped")


def main() -> None:
    asyncio.run(_run())


if __name__ == "__main__":
    main()
