import json import os import sys from pathlib import Path import redis from sqlalchemy.orm import Session # 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 TelegramSourceConfig, parse_source_config # noqa: E402 from ..models import ParseChannel, ParseJob, ParserProfile # 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 flatten_pair_config( channel: ParseChannel, profile: ParserProfile, *, limit: int = 100, ) -> dict: """Expand Profile + Channel into Redis/CP source_config.""" if not profile.heuristic_profile: raise ValueError("ParserProfile.heuristic_profile is empty") cfg = TelegramSourceConfig( channel=channel.channel.strip(), limit=limit, extract_mode="profile", heuristic_profile=profile.heuristic_profile, sample_post=profile.sample_post or None, ) return cfg.model_dump() def resolve_job_source_config(db: Session, job: ParseJob) -> dict: """ For pair jobs, rebuild flat source_config from current Profile/Channel. Legacy jobs keep stored source_config. """ if job.profile_id and job.channel_id: profile = db.query(ParserProfile).filter(ParserProfile.id == job.profile_id).first() channel = db.query(ParseChannel).filter(ParseChannel.id == job.channel_id).first() if not profile: raise ValueError(f"ParserProfile {job.profile_id} not found") if not channel: raise ValueError(f"ParseChannel {job.channel_id} not found") existing = job.source_config or {} limit = existing.get("limit", 100) if not isinstance(limit, int) or limit < 1: limit = 100 flattened = flatten_pair_config(channel, profile, limit=limit) job.source_config = flattened return flattened return dict(job.source_config or {}) 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)) def enqueue_parse_job(db: Session, job: ParseJob) -> None: """Resolve (flatten if pair) then push to Redis.""" source_config = resolve_job_source_config(db, job) db.commit() enqueue_job(job.id, job.source_type, source_config)