diff --git a/.gitignore b/.gitignore index 037b4ae..b124d43 100644 --- a/.gitignore +++ b/.gitignore @@ -1,6 +1,5 @@ .env -data/telegram.session -data/*.session +data/ __pycache__/ *.pyc node_modules/ diff --git a/README.md b/README.md index 0abede5..b8099de 100644 --- a/README.md +++ b/README.md @@ -18,7 +18,7 @@ flowchart LR | Центр | Контейнеры | Назначение | |-------|------------|------------| | **ЦА** | `ca-db`, `ca-api`, `ca-frontend` | PostgreSQL, ingest API, карта, distribution API | -| **ЦП** | `cp-workers` | Парсинг Telegram (`centers/parsing/workers`) | +| **ЦП** | `cp-workers` | Парсинг Telegram: real-time listener (Telethon) + batch-задания из Redis | | **Общее** | `redis` | Очередь заданий ЦА → ЦП | ## Структура monorepo @@ -67,7 +67,7 @@ docker compose up --build | Раздел | Путь | Описание | |--------|------|----------| | **Карта** | `/` | Интерактивная карта событий: навигация по датам (flatpickr), пресеты периода, фильтры региона/темы/источника, подложки Яндекс/OSM/Topo/ESRI, линейка, полноэкранный режим, центрирование по координатам и городам, поиск населённых пунктов (Nominatim). CRUD для ручных объектов (ПКМ). Поддерживает `?eventId=` | -| **Парсеры** | `/parsers` | Создание заданий Telegram-парсинга, таблица статусов с автообновлением (5 с), повтор failed-заданий | +| **Парсеры** | `/parsers` | Telegram-парсеры с периодическим запуском, настройка интервала, редактирование и удаление; дубликаты по `source_url` не записываются | | **События** | `/events` | Фильтрация, пагинация, просмотр деталей, ссылка «На карте» для событий с координатами | | **Аналитика** | `/analytics` | KPI-карточки, график динамики ingest за 30 дней, топ населённых пунктов и регионов | | **ПИ** | `/consumers` | CRUD подписчиков distribution API, ротация ключей, тест среза через `/api/v1/events` | @@ -108,8 +108,10 @@ docker compose up --build | Метод | Путь | Описание | |-------|------|----------| -| POST | `/admin/jobs` | Создать задание парсинга (ставится в Redis) | -| GET | `/admin/jobs` | Список заданий | +| POST | `/admin/jobs` | Создать парсер (сразу в очередь + периодический запуск) | +| GET | `/admin/jobs` | Список парсеров | +| PATCH | `/admin/jobs/{id}` | Изменить канал/лимит, `interval_seconds`, `is_active` | +| DELETE | `/admin/jobs/{id}` | Удалить парсер | | POST | `/admin/jobs/{id}/retry` | Повторить failed-задание | | GET | `/admin/events` | События с фильтрами (`{ items, total }`) | | GET | `/admin/analytics/summary` | KPI-сводка | @@ -148,14 +150,35 @@ curl http://localhost:8080/api/v1/events \ -H 'Authorization: Bearer test-pi-api-key-change-me' ``` +## Парсинг Telegram + +`cp-workers` объединяет два режима в **одном процессе** (общая сессия `telegram.session`): + +| Режим | Как работает | +|-------|----------------| +| **Listener (real-time)** | Telethon `NewMessage` / `Album` на активных парсерах (`is_active=true`); новый пост сразу уходит в ingest | +| **Batch (по расписанию)** | Планировщик в `ca-api` ставит задание в Redis; воркер забирает последние N постов (`iter_messages`) | + +Дубликаты по `source_url` при ingest пропускаются. Listener не меняет `status` парсера (флаг `listener: true` в ingest). + +Переменные `cp-workers`: + +| Переменная | По умолчанию | Описание | +|------------|--------------|----------| +| `TELEGRAM_LISTENER_ENABLED` | `true` | Включить real-time listener | +| `TELEGRAM_LISTENER_REFRESH_SECONDS` | `60` | Как часто обновлять список каналов из БД | + +Internal API: `GET /internal/listener/subscriptions` — список активных каналов для listener. + ## Поток данных -1. Аналитик создаёт задание: `POST /admin/jobs` -2. `ca-api` ставит задание в Redis (`cp:jobs`) -3. `cp-workers` забирает задание, парсит Telegram через существующую сессию -4. Результаты отправляются в `POST /internal/ingest` -5. События с координатами автоматически появляются на карте как `MapObject` -6. Внешние ПИ получают отфильтрованный срез через `/api/v1/events` +1. Аналитик создаёт парсер: `POST /admin/jobs` (канал, лимит, интервал в секундах) +2. **Listener** сразу подписывается на канал и ingest-ит новые посты +3. **Планировщик** `ca-api` периодически ставит batch-задание в Redis (`cp:jobs`) +4. `cp-workers` забирает batch-задание, парсит последние N постов +5. Результаты → `POST /internal/ingest` (дубликаты по `source_url` пропускаются) +6. Новые события с координатами появляются на карте как `MapObject` +7. Внешние ПИ получают срез через `/api/v1/events` ## Миграция EventRecord → Event diff --git a/centers/analytics/api/app/main.py b/centers/analytics/api/app/main.py index 3119f17..606d912 100644 --- a/centers/analytics/api/app/main.py +++ b/centers/analytics/api/app/main.py @@ -6,6 +6,8 @@ from fastapi.middleware.cors import CORSMiddleware from .database import Base, engine, get_db from .routers import admin, internal, map, objects, v1 from .seed import seed_objects, seed_test_consumer +from .services.migrations import migrate_schema +from .services.scheduler import start_scheduler from .storage import ensure_upload_dir @@ -13,6 +15,8 @@ from .storage import ensure_upload_dir async def lifespan(_: FastAPI): ensure_upload_dir() Base.metadata.create_all(bind=engine) + migrate_schema(engine) + scheduler_stop = start_scheduler() db = next(get_db()) try: seed_objects(db) @@ -20,6 +24,7 @@ async def lifespan(_: FastAPI): finally: db.close() yield + scheduler_stop.set() app = FastAPI(title="CA API (Analytics Center)", lifespan=lifespan) diff --git a/centers/analytics/api/app/models.py b/centers/analytics/api/app/models.py index a9466cf..bbd69d1 100644 --- a/centers/analytics/api/app/models.py +++ b/centers/analytics/api/app/models.py @@ -100,6 +100,8 @@ class ParseJob(Base): source_type: Mapped[str] = mapped_column(String(50), nullable=False) source_config: Mapped[dict] = mapped_column(JSON, nullable=False, default=dict) schedule: Mapped[str | None] = mapped_column(String(100), nullable=True) + interval_seconds: Mapped[int] = mapped_column(Integer, default=3600, nullable=False) + is_active: Mapped[bool] = mapped_column(Boolean, default=True, nullable=False) status: Mapped[str] = mapped_column(String(50), default="pending", index=True) last_run_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True) last_error: Mapped[str | None] = mapped_column(Text, nullable=True) diff --git a/centers/analytics/api/app/routers/admin.py b/centers/analytics/api/app/routers/admin.py index 1aeaa2d..e38bb70 100644 --- a/centers/analytics/api/app/routers/admin.py +++ b/centers/analytics/api/app/routers/admin.py @@ -14,6 +14,7 @@ from ..schemas import ( EventRead, ParseJobCreate, ParseJobRead, + ParseJobUpdate, TimelinePoint, TopItem, ) @@ -42,6 +43,8 @@ def create_parse_job(payload: ParseJobCreate, db: Session = Depends(get_db)): source_type=payload.source_type, source_config=payload.source_config, schedule=payload.schedule, + interval_seconds=payload.interval_seconds, + is_active=payload.is_active, status="queued", ) db.add(job) @@ -70,6 +73,8 @@ def retry_parse_job(job_id: int, db: Session = Depends(get_db)): job = db.query(ParseJob).filter(ParseJob.id == job_id).first() if not job: raise HTTPException(status_code=404, detail="Job not found") + if job.status in ("queued", "running"): + raise HTTPException(status_code=409, detail="Job is already running or queued") job.status = "queued" job.last_error = None @@ -80,6 +85,40 @@ def retry_parse_job(job_id: int, db: Session = Depends(get_db)): return job +@router.patch("/jobs/{job_id}", response_model=ParseJobRead) +def update_parse_job( + job_id: int, + payload: ParseJobUpdate, + db: Session = Depends(get_db), +): + job = db.query(ParseJob).filter(ParseJob.id == job_id).first() + if not job: + raise HTTPException(status_code=404, detail="Job not found") + + if payload.source_config is not None: + job.source_config = payload.source_config + if payload.interval_seconds is not None: + job.interval_seconds = payload.interval_seconds + if payload.is_active is not None: + job.is_active = payload.is_active + + db.commit() + db.refresh(job) + return job + + +@router.delete("/jobs/{job_id}", status_code=204) +def delete_parse_job(job_id: int, db: Session = Depends(get_db)): + job = db.query(ParseJob).filter(ParseJob.id == job_id).first() + if not job: + raise HTTPException(status_code=404, detail="Job not found") + if job.status == "running": + raise HTTPException(status_code=409, detail="Cannot delete a running job") + + db.delete(job) + db.commit() + + @router.get("/events", response_model=EventListResponse) def list_events( limit: int = Query(default=50, ge=1, le=1000), diff --git a/centers/analytics/api/app/routers/internal.py b/centers/analytics/api/app/routers/internal.py index 5e0c116..aed78b2 100644 --- a/centers/analytics/api/app/routers/internal.py +++ b/centers/analytics/api/app/routers/internal.py @@ -6,7 +6,7 @@ from sqlalchemy.orm import Session from ..database import get_db from ..deps import verify_internal_token from ..models import ParseJob -from ..schemas import IngestRequest, IngestResponse +from ..schemas import IngestRequest, IngestResponse, ListenerSubscription from ..services.ingest import ingest_events router = APIRouter(prefix="/internal", tags=["internal"]) @@ -18,10 +18,16 @@ def internal_ingest( _: None = Depends(verify_internal_token), db: Session = Depends(get_db), ): - ingested, updated, map_synced = ingest_events(db, payload.events, payload.job_id) + ingested, updated, skipped, map_synced = ingest_events( + db, + payload.events, + payload.job_id, + touch_job=not payload.listener, + ) return IngestResponse( ingested=ingested, updated=updated, + skipped=skipped, map_objects_synced=map_synced, ) @@ -39,7 +45,31 @@ def update_job_status( raise HTTPException(status_code=404, detail="Job not found") job.status = status - job.last_run_at = datetime.now(timezone.utc) + if status in ("completed", "failed"): + job.last_run_at = datetime.now(timezone.utc) job.last_error = error db.commit() return {"id": job.id, "status": job.status} + + +@router.get("/listener/subscriptions", response_model=list[ListenerSubscription]) +def listener_subscriptions( + _: None = Depends(verify_internal_token), + db: Session = Depends(get_db), +): + jobs = ( + db.query(ParseJob) + .filter( + ParseJob.is_active.is_(True), + ParseJob.source_type == "telegram", + ) + .order_by(ParseJob.id.asc()) + .all() + ) + result: list[ListenerSubscription] = [] + for job in jobs: + channel = job.source_config.get("channel") if job.source_config else None + if not channel: + continue + result.append(ListenerSubscription(job_id=job.id, channel=str(channel))) + return result diff --git a/centers/analytics/api/app/schemas.py b/centers/analytics/api/app/schemas.py index 27468fd..6d04282 100644 --- a/centers/analytics/api/app/schemas.py +++ b/centers/analytics/api/app/schemas.py @@ -111,11 +111,18 @@ class IngestEventItem(BaseModel): class IngestRequest(BaseModel): job_id: int | None = None events: list[IngestEventItem] = Field(default_factory=list) + listener: bool = False + + +class ListenerSubscription(BaseModel): + job_id: int + channel: str class IngestResponse(BaseModel): ingested: int updated: int + skipped: int = 0 map_objects_synced: int @@ -123,6 +130,14 @@ class ParseJobCreate(BaseModel): source_type: str = "telegram" source_config: dict[str, Any] = Field(default_factory=dict) schedule: str | None = None + interval_seconds: int = Field(default=3600, ge=60, le=604800) + is_active: bool = True + + +class ParseJobUpdate(BaseModel): + source_config: dict[str, Any] | None = None + interval_seconds: int | None = Field(default=None, ge=60, le=604800) + is_active: bool | None = None class ParseJobRead(BaseModel): @@ -132,6 +147,8 @@ class ParseJobRead(BaseModel): source_type: str source_config: dict[str, Any] schedule: str | None + interval_seconds: int + is_active: bool status: str last_run_at: datetime | None last_error: str | None diff --git a/centers/analytics/api/app/services/ingest.py b/centers/analytics/api/app/services/ingest.py index 63e1576..a220b9b 100644 --- a/centers/analytics/api/app/services/ingest.py +++ b/centers/analytics/api/app/services/ingest.py @@ -9,9 +9,12 @@ def ingest_events( db: Session, items: list[IngestEventItem], job_id: int | None = None, -) -> tuple[int, int, int]: + *, + touch_job: bool = True, +) -> tuple[int, int, int, int]: ingested = 0 updated = 0 + skipped = 0 map_synced = 0 for item in items: @@ -22,44 +25,32 @@ def ingest_events( ) if existing: - existing.source_type = item.source_type - existing.raw_text = item.raw_text - existing.title = item.title - existing.description = item.description - existing.locality = item.locality - existing.latitude = item.latitude - existing.longitude = item.longitude - existing.event_date = item.event_date - existing.region = item.region - existing.topic = item.topic - existing.tags = item.tags - existing.metadata_ = item.metadata - event = existing - updated += 1 - else: - event = Event( - source_type=item.source_type, - source_url=item.source_url, - raw_text=item.raw_text, - title=item.title, - description=item.description, - locality=item.locality, - latitude=item.latitude, - longitude=item.longitude, - event_date=item.event_date, - region=item.region, - topic=item.topic, - tags=item.tags, - metadata_=item.metadata, - ) - db.add(event) - ingested += 1 + skipped += 1 + continue + + event = Event( + source_type=item.source_type, + source_url=item.source_url, + raw_text=item.raw_text, + title=item.title, + description=item.description, + locality=item.locality, + latitude=item.latitude, + longitude=item.longitude, + event_date=item.event_date, + region=item.region, + topic=item.topic, + tags=item.tags, + metadata_=item.metadata, + ) + db.add(event) + ingested += 1 db.flush() if sync_event_to_map_object(db, event): map_synced += 1 - if job_id is not None: + if touch_job and job_id is not None: job = db.query(ParseJob).filter(ParseJob.id == job_id).first() if job: from datetime import datetime, timezone @@ -69,4 +60,4 @@ def ingest_events( job.last_error = None db.commit() - return ingested, updated, map_synced + return ingested, updated, skipped, map_synced diff --git a/centers/analytics/api/app/services/migrations.py b/centers/analytics/api/app/services/migrations.py new file mode 100644 index 0000000..b0707f5 --- /dev/null +++ b/centers/analytics/api/app/services/migrations.py @@ -0,0 +1,28 @@ +from sqlalchemy import inspect, text +from sqlalchemy.engine import Engine + + +def migrate_schema(engine: Engine) -> None: + """Apply lightweight schema updates for existing deployments.""" + inspector = inspect(engine) + if "parse_jobs" not in inspector.get_table_names(): + return + + columns = {col["name"] for col in inspector.get_columns("parse_jobs")} + statements: list[str] = [] + + if "interval_seconds" not in columns: + statements.append( + "ALTER TABLE parse_jobs ADD COLUMN interval_seconds INTEGER NOT NULL DEFAULT 3600" + ) + if "is_active" not in columns: + statements.append( + "ALTER TABLE parse_jobs ADD COLUMN is_active BOOLEAN NOT NULL DEFAULT TRUE" + ) + + if not statements: + return + + with engine.begin() as conn: + for stmt in statements: + conn.execute(text(stmt)) diff --git a/centers/analytics/api/app/services/scheduler.py b/centers/analytics/api/app/services/scheduler.py new file mode 100644 index 0000000..16ab4e4 --- /dev/null +++ b/centers/analytics/api/app/services/scheduler.py @@ -0,0 +1,58 @@ +import logging +import threading +import time +from datetime import datetime, timezone + +from ..database import SessionLocal +from ..models import ParseJob +from .jobs import enqueue_job + +logger = logging.getLogger(__name__) + +TICK_SECONDS = 30 +RECURRING_STATUSES = ("completed", "failed") + + +def run_scheduler_tick() -> None: + db = SessionLocal() + try: + now = datetime.now(timezone.utc) + jobs = ( + db.query(ParseJob) + .filter( + ParseJob.is_active.is_(True), + ParseJob.interval_seconds > 0, + ParseJob.status.in_(RECURRING_STATUSES), + ) + .all() + ) + for job in jobs: + if job.last_run_at is None: + continue + elapsed = (now - job.last_run_at).total_seconds() + if elapsed < job.interval_seconds: + continue + + job.status = "queued" + job.last_error = None + db.commit() + enqueue_job(job.id, job.source_type, job.source_config) + logger.info("Re-queued recurring job %s (interval %ss)", job.id, job.interval_seconds) + except Exception: + logger.exception("Scheduler tick failed") + db.rollback() + finally: + db.close() + + +def _scheduler_loop(stop_event: threading.Event) -> None: + while not stop_event.wait(TICK_SECONDS): + run_scheduler_tick() + + +def start_scheduler() -> threading.Event: + stop_event = threading.Event() + thread = threading.Thread(target=_scheduler_loop, args=(stop_event,), daemon=True) + thread.start() + logger.info("Parse job scheduler started (tick every %ss)", TICK_SECONDS) + return stop_event diff --git a/centers/analytics/frontend/src/api/admin.ts b/centers/analytics/frontend/src/api/admin.ts index 30a403a..8af399a 100644 --- a/centers/analytics/frontend/src/api/admin.ts +++ b/centers/analytics/frontend/src/api/admin.ts @@ -9,6 +9,7 @@ import type { EventRecord, ParseJob, ParseJobCreate, + ParseJobUpdate, TimelinePoint, TopItem, } from "../types/admin"; @@ -44,6 +45,22 @@ export function retryJob(jobId: number): Promise { ); } +export function updateJob(jobId: number, payload: ParseJobUpdate): Promise { + return request( + `/jobs/${jobId}`, + { method: "PATCH", body: JSON.stringify(payload) }, + ADMIN_BASE, + ); +} + +export function deleteJob(jobId: number): Promise { + return request( + `/jobs/${jobId}`, + { method: "DELETE" }, + ADMIN_BASE, + ); +} + export function fetchEvents(filters: EventFilters = {}): Promise { return request( `/events${buildQuery(filters as Record)}`, diff --git a/centers/analytics/frontend/src/components/map/CoordsTools.vue b/centers/analytics/frontend/src/components/map/CoordsTools.vue index f4a9f2d..2e16915 100644 --- a/centers/analytics/frontend/src/components/map/CoordsTools.vue +++ b/centers/analytics/frontend/src/components/map/CoordsTools.vue @@ -1,5 +1,5 @@ diff --git a/centers/analytics/frontend/src/components/map/MapToolbar.vue b/centers/analytics/frontend/src/components/map/MapToolbar.vue index 1370880..8a770c9 100644 --- a/centers/analytics/frontend/src/components/map/MapToolbar.vue +++ b/centers/analytics/frontend/src/components/map/MapToolbar.vue @@ -21,6 +21,8 @@ const props = defineProps<{ topic: string; sourceType: string; loading?: boolean; + statusText?: string; + statusError?: boolean; }>(); const emit = defineEmits<{ @@ -33,8 +35,10 @@ const emit = defineEmits<{ }>(); const dateInput = ref(null); +const rangeBtnRef = ref(null); const rangeOpen = ref(false); const filtersOpen = ref(false); +const dropdownPos = ref({ top: 0, left: 0 }); let picker: flatpickr.Instance | null = null; const availableDates = computed(() => props.filters?.available_dates ?? []); @@ -74,6 +78,31 @@ function selectRange(preset: DateRangePreset) { emit("apply"); } +function toggleRangeOpen(event: MouseEvent) { + event.stopPropagation(); + if (!rangeOpen.value && rangeBtnRef.value) { + const rect = rangeBtnRef.value.getBoundingClientRect(); + dropdownPos.value = { top: rect.bottom + 4, left: rect.left }; + } + rangeOpen.value = !rangeOpen.value; + if (rangeOpen.value) { + filtersOpen.value = false; + } +} + +function toggleFiltersOpen(event: MouseEvent) { + event.stopPropagation(); + filtersOpen.value = !filtersOpen.value; + if (filtersOpen.value) { + rangeOpen.value = false; + } +} + +function closeDropdowns() { + rangeOpen.value = false; + filtersOpen.value = false; +} + function resetFilters() { emit("update:region", ""); emit("update:topic", ""); @@ -82,6 +111,8 @@ function resetFilters() { } onMounted(() => { + document.addEventListener("click", closeDropdowns); + if (!dateInput.value) return; picker = flatpickr(dateInput.value, { locale: Russian, @@ -101,6 +132,7 @@ onMounted(() => { }); onUnmounted(() => { + document.removeEventListener("click", closeDropdowns); picker?.destroy(); }); @@ -115,7 +147,8 @@ watch( diff --git a/centers/analytics/frontend/src/components/map/PlaceSearch.vue b/centers/analytics/frontend/src/components/map/PlaceSearch.vue index 418bf1d..3a85c7d 100644 --- a/centers/analytics/frontend/src/components/map/PlaceSearch.vue +++ b/centers/analytics/frontend/src/components/map/PlaceSearch.vue @@ -30,12 +30,13 @@ function selectPlace(place: PlaceSearchResult) { + + diff --git a/centers/parsing/workers/worker.py b/centers/parsing/workers/worker.py index 3628242..346de76 100644 --- a/centers/parsing/workers/worker.py +++ b/centers/parsing/workers/worker.py @@ -1,10 +1,9 @@ -"""CP worker: poll Redis jobs, parse Telegram, ingest to CA.""" +"""CP worker: Redis jobs + real-time Telethon listener.""" import asyncio import json import logging import os -import time import httpx import redis @@ -17,6 +16,8 @@ from workers.sources.telegram_client import ( fetch_channel_posts, normalize_channel, ) +from workers.sources.telegram_listener import TelegramListener +from workers.sources.telegram_session import close_shared_client, get_shared_client logging.basicConfig(level=logging.INFO, format="%(asctime)s %(levelname)s %(message)s") logger = logging.getLogger("cp-worker") @@ -26,19 +27,30 @@ 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")) +LISTENER_ENABLED = os.getenv("TELEGRAM_LISTENER_ENABLED", "true").lower() not in ( + "0", + "false", + "no", + "off", +) 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) -> tuple[list[dict], str | None]: +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) + posts = await fetch_channel_posts(username, limit=limit, client=client) except (TelegramConfigError, TelegramAuthError, ValueError) as exc: return [], str(exc) except Exception as exc: @@ -60,7 +72,14 @@ async def post_ingest(job_id: int, events: list[dict]) -> None: headers={"X-Internal-Token": INTERNAL_TOKEN}, ) response.raise_for_status() - logger.info("Ingested %s events for job %s: %s", len(events), job_id, response.json()) + result = response.json() + logger.info( + "Ingested job %s: new=%s skipped=%s map=%s", + job_id, + result.get("ingested", 0), + result.get("skipped", 0), + result.get("map_objects_synced", 0), + ) async def patch_job_status(job_id: int, status: str, error: str | None = None) -> None: @@ -76,7 +95,7 @@ async def patch_job_status(job_id: int, status: str, error: str | None = None) - response.raise_for_status() -async def handle_job(payload: dict) -> None: +async def handle_job(payload: dict, *, tg_client=None) -> None: job_id = payload["job_id"] source_type = payload["source_type"] source_config = payload.get("source_config", {}) @@ -85,7 +104,7 @@ async def handle_job(payload: dict) -> None: await patch_job_status(job_id, "running") if source_type == "telegram": - events, error = await process_telegram_job(job_id, source_config) + events, error = await process_telegram_job(job_id, source_config, client=tg_client) else: events, error = [], f"Unsupported source_type: {source_type}" @@ -104,28 +123,54 @@ async def handle_job(payload: dict) -> None: await patch_job_status(job_id, "failed", error=str(exc)) -async def worker_loop() -> None: +async def worker_loop(*, tg_client=None) -> None: r = get_redis() logger.info("CP worker started, polling %s", JOB_QUEUE_KEY) while True: try: - item = r.blpop(JOB_QUEUE_KEY, timeout=POLL_TIMEOUT) + item = await asyncio.to_thread(r.blpop, JOB_QUEUE_KEY, POLL_TIMEOUT) if not item: continue _, raw = item payload = json.loads(raw) - await handle_job(payload) + await handle_job(payload, tg_client=tg_client) except redis.RedisError as exc: logger.error("Redis error: %s", exc) - time.sleep(3) + await asyncio.sleep(3) except Exception: logger.exception("Unexpected worker error") - time.sleep(1) + await asyncio.sleep(1) + + +async def run_with_listener() -> None: + client = await get_shared_client() + listener = TelegramListener(client) + worker_task = asyncio.create_task(worker_loop(tg_client=client)) + listener_task = asyncio.create_task(listener.run()) + + logger.info("Telegram listener enabled (shared session with batch worker)") + try: + await client.run_until_disconnected() + finally: + worker_task.cancel() + listener_task.cancel() + await close_shared_client() + + +async def run_batch_only() -> None: + logger.info("Telegram listener disabled") + await worker_loop(tg_client=None) def main() -> None: - asyncio.run(worker_loop()) + try: + if LISTENER_ENABLED: + asyncio.run(run_with_listener()) + else: + asyncio.run(run_batch_only()) + except KeyboardInterrupt: + pass if __name__ == "__main__": diff --git a/centers/parsing/workers/workers/sources/telegram_client.py b/centers/parsing/workers/workers/sources/telegram_client.py index cd0aac1..78b4b67 100644 --- a/centers/parsing/workers/workers/sources/telegram_client.py +++ b/centers/parsing/workers/workers/sources/telegram_client.py @@ -4,6 +4,7 @@ import os import re from pathlib import Path +from telethon import TelegramClient from telethon.errors import AuthKeyUnregisteredError, SessionPasswordNeededError from workers.models import TelegramPost @@ -46,54 +47,76 @@ def _build_post_url(channel: str, message_id: int) -> str: return f"https://t.me/{username}/{message_id}" -async def fetch_channel_posts(channel: str, limit: int = 100) -> list[TelegramPost]: - limit = max(1, min(limit, MAX_POSTS)) - try: - api_id, api_hash, session_path = _get_config() - except ValueError as exc: - raise TelegramConfigError(str(exc)) from exc - +def message_to_post(message, channel: str) -> TelegramPost | None: + text = (message.text or message.message or "").strip() + if not text: + return None username = normalize_channel(channel) + return TelegramPost( + id=message.id, + text=text, + date=message.date, + url=_build_post_url(username, message.id), + channel=username, + ) - if not Path(session_path).exists(): - raise TelegramAuthError( - f"Файл сессии не найден: {session_path}. " - "Выполните: python scripts/telegram_auth.py" - ) - client = create_client(session_path, api_id, api_hash) - posts: list[TelegramPost] = [] +async def fetch_channel_posts( + channel: str, + limit: int = 100, + *, + client: TelegramClient | None = None, +) -> list[TelegramPost]: + limit = max(1, min(limit, MAX_POSTS)) + username = normalize_channel(channel) + own_client = client is None - try: - await client.connect() - if not await client.is_user_authorized(): + if own_client: + try: + api_id, api_hash, session_path = _get_config() + except ValueError as exc: + raise TelegramConfigError(str(exc)) from exc + + if not Path(session_path).exists(): raise TelegramAuthError( - "Telegram-сессия не авторизована. Выполните: python scripts/telegram_auth.py" + f"Файл сессии не найден: {session_path}. " + "Выполните: python scripts/telegram_auth.py" ) - entity = await client.get_entity(username) - async for message in client.iter_messages(entity, limit=limit): - if not message.text: - continue - posts.append( - TelegramPost( - id=message.id, - text=message.text.strip(), - date=message.date, - url=_build_post_url(username, message.id), - channel=username, + client = create_client(session_path, api_id, api_hash) + posts: list[TelegramPost] = [] + + try: + await client.connect() + if not await client.is_user_authorized(): + raise TelegramAuthError( + "Telegram-сессия не авторизована. Выполните: python scripts/telegram_auth.py" ) - ) - except AuthKeyUnregisteredError as exc: - raise TelegramAuthError( - "Сессия Telegram недействительна. Переавторизуйтесь: python scripts/telegram_auth.py" - ) from exc - except SessionPasswordNeededError as exc: - raise TelegramAuthError( - "Для аккаунта включена двухфакторная аутентификация. " - "Авторизуйтесь через scripts/telegram_auth.py с паролем 2FA." - ) from exc - finally: - await client.disconnect() + entity = await client.get_entity(username) + async for message in client.iter_messages(entity, limit=limit): + post = message_to_post(message, username) + if post: + posts.append(post) + except AuthKeyUnregisteredError as exc: + raise TelegramAuthError( + "Сессия Telegram недействительна. Переавторизуйтесь: python scripts/telegram_auth.py" + ) from exc + except SessionPasswordNeededError as exc: + raise TelegramAuthError( + "Для аккаунта включена двухфакторная аутентификация. " + "Авторизуйтесь через scripts/telegram_auth.py с паролем 2FA." + ) from exc + finally: + await client.disconnect() + + return posts + + assert client is not None + posts = [] + entity = await client.get_entity(username) + async for message in client.iter_messages(entity, limit=limit): + post = message_to_post(message, username) + if post: + posts.append(post) return posts diff --git a/centers/parsing/workers/workers/sources/telegram_listener.py b/centers/parsing/workers/workers/sources/telegram_listener.py new file mode 100644 index 0000000..df65c4f --- /dev/null +++ b/centers/parsing/workers/workers/sources/telegram_listener.py @@ -0,0 +1,150 @@ +"""Real-time Telethon listener for active Telegram parse jobs.""" + +import asyncio +import logging +import os +from typing import Any + +import httpx +from telethon import TelegramClient, events + +from workers.converter import event_record_to_ingest +from workers.parsers.telegram_events import parse_event_post +from workers.sources.telegram_client import message_to_post, normalize_channel + +logger = logging.getLogger("cp-listener") + +CA_API_URL = os.getenv("CA_API_URL", "http://ca-api:8000") +INTERNAL_TOKEN = os.getenv("INTERNAL_TOKEN", "dev-internal-token") +REFRESH_SECONDS = int(os.getenv("TELEGRAM_LISTENER_REFRESH_SECONDS", "60")) + + +class TelegramListener: + def __init__(self, client: TelegramClient) -> None: + self.client = client + self._channels: dict[str, int] = {} + self._chat_ids: set[int] = set() + self._handlers_registered = False + + async def fetch_subscriptions(self) -> dict[str, int]: + async with httpx.AsyncClient(timeout=30.0) as http: + response = await http.get( + f"{CA_API_URL}/internal/listener/subscriptions", + headers={"X-Internal-Token": INTERNAL_TOKEN}, + ) + response.raise_for_status() + data = response.json() + + channels: dict[str, int] = {} + for item in data: + raw = item.get("channel") + job_id = item.get("job_id") + if not raw or job_id is None: + continue + try: + key = normalize_channel(str(raw)) + except ValueError: + logger.warning("Skip invalid channel in subscription: %r", raw) + continue + channels.setdefault(key, int(job_id)) + return channels + + async def refresh_subscriptions(self) -> None: + try: + channels = await self.fetch_subscriptions() + except Exception: + logger.exception("Failed to load listener subscriptions") + return + + if channels == self._channels: + return + + self._channels = channels + self._chat_ids = set() + for username in channels: + try: + entity = await self.client.get_entity(username) + self._chat_ids.add(entity.id) + except Exception: + logger.exception("Cannot resolve channel entity: %s", username) + + logger.info( + "Listener subscriptions updated: %s", + ", ".join(sorted(channels)) or "(none)", + ) + + async def _ingest_post(self, channel: str, post) -> None: + record = parse_event_post(post) + event = event_record_to_ingest(record) + job_id = self._channels.get(normalize_channel(channel)) + + async with httpx.AsyncClient(timeout=60.0) as http: + response = await http.post( + f"{CA_API_URL}/internal/ingest", + json={ + "job_id": job_id, + "events": [event], + "listener": True, + }, + headers={"X-Internal-Token": INTERNAL_TOKEN}, + ) + response.raise_for_status() + result = response.json() + logger.info( + "Listener ingest %s: new=%s skipped=%s", + post.url, + result.get("ingested", 0), + result.get("skipped", 0), + ) + + def _channel_from_event(self, event: events.common.EventCommon) -> str | None: + chat = event.chat + if chat is None: + return None + username = getattr(chat, "username", None) + if username: + return normalize_channel(username) + return None + + async def _process_messages(self, event: events.common.EventCommon, messages: list[Any]) -> None: + channel = self._channel_from_event(event) + if not channel or channel not in self._channels: + return + if event.chat_id not in self._chat_ids: + return + + for message in messages: + post = message_to_post(message, channel) + if not post: + continue + try: + await self._ingest_post(channel, post) + except Exception: + logger.exception("Listener ingest failed for %s", post.url) + + def register_handlers(self) -> None: + if self._handlers_registered: + return + + @self.client.on(events.NewMessage(incoming=True)) + async def on_new_message(event: events.NewMessage.Event) -> None: + if event.grouped_id: + return + await self._process_messages(event, [event.message]) + + @self.client.on(events.Album) + async def on_album(event: events.Album.Event) -> None: + await self._process_messages(event, list(event.messages)) + + self._handlers_registered = True + logger.info("Telegram NewMessage/Album handlers registered") + + async def refresh_loop(self) -> None: + while True: + await self.refresh_subscriptions() + await asyncio.sleep(REFRESH_SECONDS) + + async def run(self) -> None: + self.register_handlers() + await self.refresh_subscriptions() + await self.refresh_loop() diff --git a/centers/parsing/workers/workers/sources/telegram_session.py b/centers/parsing/workers/workers/sources/telegram_session.py new file mode 100644 index 0000000..1632144 --- /dev/null +++ b/centers/parsing/workers/workers/sources/telegram_session.py @@ -0,0 +1,55 @@ +"""Единое подключение Telethon для listener и batch-заданий.""" + +from pathlib import Path + +from telethon import TelegramClient +from telethon.errors import AuthKeyUnregisteredError, SessionPasswordNeededError + +from workers.sources.telegram_settings import create_client, get_api_credentials, get_session_path +from workers.sources.telegram_client import TelegramAuthError, TelegramConfigError + +_shared_client: TelegramClient | None = None + + +async def get_shared_client() -> TelegramClient: + global _shared_client + if _shared_client is not None and _shared_client.is_connected(): + return _shared_client + + try: + api_id, api_hash = get_api_credentials() + except ValueError as exc: + raise TelegramConfigError(str(exc)) from exc + + session_path = get_session_path() + if not Path(session_path).exists(): + raise TelegramAuthError( + f"Файл сессии не найден: {session_path}. " + "Выполните: python scripts/telegram_auth.py" + ) + + client = create_client(session_path, api_id, api_hash) + try: + await client.connect() + if not await client.is_user_authorized(): + raise TelegramAuthError( + "Telegram-сессия не авторизована. Выполните: python scripts/telegram_auth.py" + ) + except AuthKeyUnregisteredError as exc: + raise TelegramAuthError( + "Сессия Telegram недействительна. Переавторизуйтесь: python scripts/telegram_auth.py" + ) from exc + except SessionPasswordNeededError as exc: + raise TelegramAuthError( + "Для аккаунта включена 2FA. Авторизуйтесь через scripts/telegram_auth.py." + ) from exc + + _shared_client = client + return client + + +async def close_shared_client() -> None: + global _shared_client + if _shared_client is not None: + await _shared_client.disconnect() + _shared_client = None diff --git a/docker-compose.yml b/docker-compose.yml index 883ca4f..b4442c2 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -48,6 +48,8 @@ services: - .env environment: TELEGRAM_SESSION_PATH: /data/telegram.session + TELEGRAM_LISTENER_ENABLED: "true" + TELEGRAM_LISTENER_REFRESH_SECONDS: "60" CA_API_URL: http://ca-api:8000 REDIS_URL: redis://redis:6379/0 INTERNAL_TOKEN: dev-internal-token