Add CP source adapter registry with multi-worker queues and LLM extract.

Replace legacy root backend/frontend with Telegram, Crawl4AI, and VIINA adapters routed by Redis job families.

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
2026-08-14 11:34:28 +03:00
co-authored by Cursor
parent 1492576fd9
commit 8fbabd3c11
54 changed files with 1625 additions and 2521 deletions
+83 -44
View File
@@ -1,29 +1,38 @@
"""CP worker: Redis jobs + real-time Telethon listener."""
"""CP worker: Redis jobs via adapter registry + optional Telethon listener."""
from __future__ import annotations
import asyncio
import json
import logging
import os
import sys
from pathlib import Path
# Monorepo local runs / Docker: ensure contracts/ is importable
_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
import httpx
import redis
from workers.converter import event_record_to_ingest
from workers.parsers.telegram_events import parse_event_posts
from workers.sources.telegram_client import (
TelegramAuthError,
TelegramConfigError,
fetch_channel_posts,
normalize_channel,
from contracts.queues import (
LEGACY_JOB_QUEUE_KEY,
family_for_source,
queue_key_for_family,
queue_key_for_source,
)
from workers.sources.telegram_listener import TelegramListener
from workers.sources.telegram_session import close_shared_client, get_shared_client
from workers.adapters import build_registry
from workers.adapters.base import WorkerContext
logging.basicConfig(level=logging.INFO, format="%(asctime)s %(levelname)s %(message)s")
logger = logging.getLogger("cp-worker")
REDIS_URL = os.getenv("REDIS_URL", "redis://redis:6379/0")
JOB_QUEUE_KEY = "cp:jobs"
CA_API_URL = os.getenv("CA_API_URL", "http://ca-api:8000")
INTERNAL_TOKEN = os.getenv("INTERNAL_TOKEN", "dev-internal-token")
POLL_TIMEOUT = int(os.getenv("WORKER_POLL_TIMEOUT", "5"))
@@ -33,32 +42,31 @@ LISTENER_ENABLED = os.getenv("TELEGRAM_LISTENER_ENABLED", "true").lower() not in
"no",
"off",
)
# Comma-separated families this process polls (default: derived from adapters)
WORKER_FAMILIES = os.getenv("WORKER_FAMILIES", "").strip()
REGISTRY = build_registry()
def get_redis() -> redis.Redis:
return redis.from_url(REDIS_URL, decode_responses=True)
async def process_telegram_job(
job_id: int,
source_config: dict,
*,
client=None,
) -> tuple[list[dict], str | None]:
channel = source_config.get("channel", "creamy_caprice")
limit = int(source_config.get("limit", 100))
try:
username = normalize_channel(channel)
posts = await fetch_channel_posts(username, limit=limit, client=client)
except (TelegramConfigError, TelegramAuthError, ValueError) as exc:
return [], str(exc)
except Exception as exc:
return [], f"Telegram: {exc}"
records = parse_event_posts(posts)
events = [event_record_to_ingest(r) for r in records]
return events, None
def _poll_keys() -> list[str]:
if WORKER_FAMILIES:
families = [f.strip() for f in WORKER_FAMILIES.split(",") if f.strip()]
else:
families = sorted(
{
family_for_source(source_type)
for source_type in REGISTRY.enabled_types()
}
)
keys = [queue_key_for_family(f) for f in families]
# Backward compatible: telegram workers also drain legacy cp:jobs
if "telegram" in families and LEGACY_JOB_QUEUE_KEY not in keys:
keys.append(LEGACY_JOB_QUEUE_KEY)
return keys
async def post_ingest(job_id: int, events: list[dict]) -> None:
@@ -95,18 +103,35 @@ async def patch_job_status(job_id: int, status: str, error: str | None = None) -
response.raise_for_status()
async def handle_job(payload: dict, *, tg_client=None) -> None:
async def handle_job(payload: dict, *, ctx: WorkerContext) -> None:
job_id = payload["job_id"]
source_type = payload["source_type"]
source_config = payload.get("source_config", {})
logger.info("Processing job %s (%s)", job_id, source_type)
await patch_job_status(job_id, "running")
if source_type == "telegram":
events, error = await process_telegram_job(job_id, source_config, client=tg_client)
else:
events, error = [], f"Unsupported source_type: {source_type}"
adapter = REGISTRY.get(source_type)
if adapter is None:
try:
target = queue_key_for_source(source_type)
except ValueError:
await patch_job_status(
job_id,
"failed",
error=f"Unsupported source_type: {source_type}",
)
return
get_redis().rpush(target, json.dumps(payload))
logger.warning(
"Job %s (%s) not enabled here; requeued to %s",
job_id,
source_type,
target,
)
return
await patch_job_status(job_id, "running")
events, error = await adapter.run(job_id, source_config, ctx=ctx)
if error:
logger.error("Job %s failed: %s", job_id, error)
@@ -123,18 +148,23 @@ async def handle_job(payload: dict, *, tg_client=None) -> None:
await patch_job_status(job_id, "failed", error=str(exc))
async def worker_loop(*, tg_client=None) -> None:
async def worker_loop(*, ctx: WorkerContext) -> None:
r = get_redis()
logger.info("CP worker started, polling %s", JOB_QUEUE_KEY)
keys = _poll_keys()
logger.info(
"CP worker started, polling %s (adapters: %s)",
keys,
", ".join(REGISTRY.enabled_types()) or "(none)",
)
while True:
try:
item = await asyncio.to_thread(r.blpop, JOB_QUEUE_KEY, POLL_TIMEOUT)
item = await asyncio.to_thread(r.blpop, keys, POLL_TIMEOUT)
if not item:
continue
_, raw = item
payload = json.loads(raw)
await handle_job(payload, tg_client=tg_client)
await handle_job(payload, ctx=ctx)
except redis.RedisError as exc:
logger.error("Redis error: %s", exc)
await asyncio.sleep(3)
@@ -144,9 +174,18 @@ async def worker_loop(*, tg_client=None) -> None:
async def run_with_listener() -> None:
if "telegram" not in REGISTRY:
logger.warning("Listener requested but telegram adapter is not enabled")
await worker_loop(ctx=WorkerContext())
return
from workers.sources.telegram_listener import TelegramListener
from workers.sources.telegram_session import close_shared_client, get_shared_client
client = await get_shared_client()
listener = TelegramListener(client)
worker_task = asyncio.create_task(worker_loop(tg_client=client))
ctx = WorkerContext(tg_client=client)
worker_task = asyncio.create_task(worker_loop(ctx=ctx))
listener_task = asyncio.create_task(listener.run())
logger.info("Telegram listener enabled (shared session with batch worker)")
@@ -160,12 +199,12 @@ async def run_with_listener() -> None:
async def run_batch_only() -> None:
logger.info("Telegram listener disabled")
await worker_loop(tg_client=None)
await worker_loop(ctx=WorkerContext(tg_client=None))
def main() -> None:
try:
if LISTENER_ENABLED:
if LISTENER_ENABLED and "telegram" in REGISTRY:
asyncio.run(run_with_listener())
else:
asyncio.run(run_batch_only())