blitz_api/app/celery_app.py
fusion44 4e038e8fc7 feat: fetch app status via a celery task
This is a feature that allows to fetch the app status via a celery
task which stores the result in the Redis database and notifies the API
of the change. The API sends a notification to connected clients via the
SSE mechanism.

Using a background task allows to avoid crashing the whole API if the
script call fails.

The by default the cache is refreshed every 30 minutes. This can be
changed by setting the `BAPI_APP_STATUS_UPDATE_INTERVAL_MIN` environment
variable.

refs #123
2025-05-06 09:16:46 +02:00

44 lines
1.2 KiB
Python

from celery import Celery
from celery.schedules import crontab
from loguru import logger
from app.api.config import config
BAPI_REDIS_URL = config("BAPI_REDIS_URL", "redis://127.0.0.1:6379/0")
APP_STATUS_UPDATE_INTERVAL_MIN = config(
"BAPI_APP_STATUS_UPDATE_INTERVAL_MIN", cast=int, default=30
)
if not isinstance(APP_STATUS_UPDATE_INTERVAL_MIN, int):
raise TypeError("BAPI_APP_STATUS_UPDATE_INTERVAL_MIN must be an integer")
logger.info(f"Celery app started with interval: {APP_STATUS_UPDATE_INTERVAL_MIN} mins")
celery_app = Celery(
"worker",
broker=BAPI_REDIS_URL,
backend=BAPI_REDIS_URL,
include=["app.apps.tasks"],
# cache results for 1 hour, will be printed to logs anyway
result_expires=3600,
)
celery_app.conf.update(
task_serializer="json",
accept_content=["json"],
result_serializer="json",
timezone="UTC",
enable_utc=True,
)
celery_app.conf.beat_schedule = {
f"update-app-status-every-{APP_STATUS_UPDATE_INTERVAL_MIN}-mins": {
"task": "app.apps.tasks.update_app_state_task",
"schedule": crontab(minute=f"*/{APP_STATUS_UPDATE_INTERVAL_MIN}"),
},
}
if __name__ == "__main__":
celery_app.start()