import json import os import sys from pathlib import Path import redis # Allow importing shared contracts from monorepo root in local runs _HERE = Path(__file__).resolve() for _candidate in (_HERE.parent, *_HERE.parents): if (_candidate / "contracts").is_dir(): if str(_candidate) not in sys.path: sys.path.insert(0, str(_candidate)) break from contracts.queues import queue_key_for_source # noqa: E402 from contracts.sources import parse_source_config # noqa: E402 REDIS_URL = os.getenv("REDIS_URL", "redis://redis:6379/0") def get_redis() -> redis.Redis: return redis.from_url(REDIS_URL, decode_responses=True) def enqueue_job(job_id: int, source_type: str, source_config: dict) -> None: # Validate known configs early; unknown types still raise from queue_key_for_source try: parse_source_config(source_type, source_config or {}) except Exception: # Keep enqueue resilient for legacy rows; worker validates again pass payload = { "job_id": job_id, "source_type": source_type, "source_config": source_config or {}, } key = queue_key_for_source(source_type) get_redis().rpush(key, json.dumps(payload))