From 699e9be503899571da4d8419b0354e46c260d49f Mon Sep 17 00:00:00 2001 From: gitrusprus Date: Sun, 16 Aug 2026 14:14:18 +0300 Subject: [PATCH] Add parser builder: one-shot DeepSeek profile for Telegram extract_mode=profile. Generate static HeuristicProfile in CA admin, preview and run without LLM on each post via shared interpreter in CP batch and listener. Co-authored-by: Cursor --- .env.example | 4 +- centers/analytics/ARCHITECTURE.md | 20 +- centers/analytics/api/app/main.py | 3 +- centers/analytics/api/app/routers/internal.py | 8 +- .../api/app/routers/parser_builder.py | 108 ++++++ centers/analytics/api/app/schemas.py | 1 + .../api/app/services/parser_builder.py | 130 +++++++ centers/analytics/api/requirements.txt | 1 + .../frontend/src/api/parserBuilder.ts | 52 +++ .../frontend/src/layouts/AppLayout.vue | 1 + .../analytics/frontend/src/router/index.ts | 2 + .../frontend/src/views/ParserBuilderView.vue | 320 ++++++++++++++++++ .../frontend/src/views/ParsersView.vue | 61 +++- centers/parsing/ARCHITECTURE.md | 9 + .../workers/workers/adapters/telegram.py | 27 ++ .../workers/workers/heuristic_profile.py | 73 ++++ .../workers/sources/telegram_listener.py | 46 ++- contracts/heuristic_profile.py | 165 +++++++++ contracts/sources.py | 17 +- docker-compose.yml | 2 + docs/architecture-overview.md | 5 +- docs/contracts.md | 4 +- docs/local-dev.md | 3 +- 23 files changed, 1042 insertions(+), 20 deletions(-) create mode 100644 centers/analytics/api/app/routers/parser_builder.py create mode 100644 centers/analytics/api/app/services/parser_builder.py create mode 100644 centers/analytics/frontend/src/api/parserBuilder.ts create mode 100644 centers/analytics/frontend/src/views/ParserBuilderView.vue create mode 100644 centers/parsing/workers/workers/heuristic_profile.py create mode 100644 contracts/heuristic_profile.py diff --git a/.env.example b/.env.example index e49efc1..15de841 100644 --- a/.env.example +++ b/.env.example @@ -12,7 +12,9 @@ TELEGRAM_SESSION_PATH=/data/telegram.session INTERNAL_TOKEN=dev-internal-token TEST_PI_API_KEY=test-pi-api-key-change-me -# DeepSeek LLM for extract_mode=llm (Telegram unstructured + Crawl4AI LLM) +# DeepSeek LLM: +# - extract_mode=llm on CP workers (Telegram unstructured + Crawl4AI) +# - parser-builder Generate on ca-api (one-shot profile; runtime stays rule-based) # DEEPSEEK_API_KEY=sk-... # DEEPSEEK_BASE_URL=https://api.deepseek.com # DEEPSEEK_MODEL=deepseek-chat diff --git a/centers/analytics/ARCHITECTURE.md b/centers/analytics/ARCHITECTURE.md index 84b0391..73220d2 100644 --- a/centers/analytics/ARCHITECTURE.md +++ b/centers/analytics/ARCHITECTURE.md @@ -59,12 +59,14 @@ centers/analytics/ │ │ ├── objects.py # /api/health, /api/objects, media │ │ ├── map.py # /api/map/* │ │ ├── admin.py # /admin/* +│ │ ├── parser_builder.py # /admin/parser-builder/* │ │ ├── internal.py # /internal/* (только ЦП) │ │ └── v1.py # /api/v1/* (ПИ) │ └── services/ │ ├── jobs.py # Redis RPUSH │ ├── scheduler.py # периодический re-queue │ ├── ingest.py # дедуп + map sync +│ ├── parser_builder.py # DeepSeek → HeuristicProfile (один раз) │ ├── filtering.py │ ├── map_query.py │ ├── events_query.py @@ -74,7 +76,7 @@ centers/analytics/ ├── Dockerfile # Vite build + nginx ├── nginx.conf # proxy /api /admin /internal → ca-api └── src/ - ├── views/ # Map, Parsers, Events, Analytics, Consumers + ├── views/ # Map, Parsers, ParserBuilder, Events, … ├── components/ # карта, CRUD объектов ├── api/ # HTTP-клиенты └── router/index.ts @@ -96,12 +98,23 @@ centers/analytics/ | Prefix | Кто вызывает | Содержание | |--------|--------------|------------| | `/api/*` | UI, публичный health | Карта, объекты, медиа | -| `/admin/*` | UI admin | Jobs, events, analytics, consumers | +| `/admin/*` | UI admin | Jobs, events, analytics, consumers, parser-builder | | `/internal/*` | Только ЦП | ingest, job status, listener subscriptions | | `/api/v1/*` | Внешние клиенты | Events с Bearer-ключом | Internal защищён заголовком `X-Internal-Token` (`INTERNAL_TOKEN`). +### Конструктор парсера + +Раздел UI `/parser-builder` → `POST /admin/parser-builder/generate|preview|jobs`: + +1. Менеджер вставляет образец поста. +2. DeepSeek (`DEEPSEEK_API_KEY` на **ca-api**) один раз возвращает `HeuristicProfile` (regex/line/marker). +3. Preview и сохранение job с `extract_mode=profile` + JSON профиля в `source_config`. +4. ЦП применяет профиль статически (batch + listener); LLM на ingest не вызывается. + +Roadmap: кастомные пользовательские таблицы подменяют только список `target-fields` при генерации. + ### Jobs и scheduler 1. `POST /admin/jobs` / retry → запись `ParseJob` + `enqueue_job` (`services/jobs.py`). @@ -123,6 +136,7 @@ Internal защищён заголовком `X-Internal-Token` (`INTERNAL_TOKEN |------|------| | `/` | `MapViewPage.vue` | | `/parsers` | `ParsersView.vue` | +| `/parser-builder` | `ParserBuilderView.vue` | | `/events` | `EventsView.vue` | | `/analytics` | `AnalyticsView.vue` | | `/consumers` | `ConsumersView.vue` | @@ -142,6 +156,8 @@ Internal защищён заголовком `X-Internal-Token` (`INTERNAL_TOKEN | `INTERNAL_TOKEN` | Auth ЦП ↔ ЦА | | `TEST_PI_API_KEY` | Seed consumer `test-pi` | | `UPLOAD_DIR` | Медиа (по умолчанию `/data/uploads`) | +| `DEEPSEEK_API_KEY` | Конструктор парсера (Generate); не нужен для preview/runtime profile | +| `DEEPSEEK_BASE_URL` / `DEEPSEEK_MODEL` | Опционально | ### Типовые точки входа в код diff --git a/centers/analytics/api/app/main.py b/centers/analytics/api/app/main.py index 68c80f3..9d237ec 100644 --- a/centers/analytics/api/app/main.py +++ b/centers/analytics/api/app/main.py @@ -13,7 +13,7 @@ for _candidate in (_HERE.parent, *_HERE.parents): break from .database import Base, engine, get_db -from .routers import admin, internal, map, objects, v1 +from .routers import admin, internal, map, objects, parser_builder, v1 from .seed import seed_objects, seed_test_consumer from .services.migrations import migrate_schema from .services.scheduler import start_scheduler @@ -50,4 +50,5 @@ app.include_router(objects.router) app.include_router(map.router) app.include_router(internal.router) app.include_router(admin.router) +app.include_router(parser_builder.router) app.include_router(v1.router) diff --git a/centers/analytics/api/app/routers/internal.py b/centers/analytics/api/app/routers/internal.py index aed78b2..b5effa1 100644 --- a/centers/analytics/api/app/routers/internal.py +++ b/centers/analytics/api/app/routers/internal.py @@ -71,5 +71,11 @@ def listener_subscriptions( 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))) + result.append( + ListenerSubscription( + job_id=job.id, + channel=str(channel), + source_config=dict(job.source_config or {}), + ) + ) return result diff --git a/centers/analytics/api/app/routers/parser_builder.py b/centers/analytics/api/app/routers/parser_builder.py new file mode 100644 index 0000000..af60556 --- /dev/null +++ b/centers/analytics/api/app/routers/parser_builder.py @@ -0,0 +1,108 @@ +"""Admin API: construct static telegram parsers from a sample post.""" + +from __future__ import annotations + +from typing import Any + +from fastapi import APIRouter, Depends, HTTPException +from pydantic import BaseModel, Field +from sqlalchemy.orm import Session + +from contracts.heuristic_profile import HeuristicProfile +from contracts.sources import TelegramSourceConfig + +from ..database import get_db +from ..models import ParseJob +from ..schemas import ParseJobRead +from ..services.jobs import enqueue_job +from ..services import parser_builder as builder + +router = APIRouter(prefix="/admin/parser-builder", tags=["parser-builder"]) + + +class GenerateRequest(BaseModel): + sample_post: str = Field(min_length=1) + + +class GenerateResponse(BaseModel): + profile: dict[str, Any] + + +class PreviewRequest(BaseModel): + sample_post: str = Field(min_length=1) + profile: dict[str, Any] + + +class PreviewResponse(BaseModel): + fields: dict[str, str] + + +class CreateProfileJobRequest(BaseModel): + channel: str = Field(min_length=1) + profile: dict[str, Any] + sample_post: str | None = None + limit: int = Field(default=100, ge=1, le=1000) + interval_seconds: int = Field(default=3600, ge=60, le=604800) + is_active: bool = True + + +@router.get("/target-fields") +def target_fields(): + return {"fields": builder.get_target_fields()} + + +@router.post("/generate", response_model=GenerateResponse) +async def generate_profile(payload: GenerateRequest): + if not builder.deepseek_enabled(): + raise HTTPException( + status_code=503, + detail="DEEPSEEK_API_KEY is not set on ca-api. Add it to .env for parser generation.", + ) + try: + profile = await builder.generate_profile(payload.sample_post) + except ValueError as exc: + raise HTTPException(status_code=400, detail=str(exc)) from exc + except Exception as exc: + raise HTTPException( + status_code=502, + detail=f"DeepSeek generate failed: {exc}", + ) from exc + return GenerateResponse(profile=profile.model_dump()) + + +@router.post("/preview", response_model=PreviewResponse) +def preview_profile(payload: PreviewRequest): + try: + HeuristicProfile.model_validate(payload.profile) + fields = builder.preview_with_profile(payload.sample_post, payload.profile) + except Exception as exc: + raise HTTPException(status_code=400, detail=str(exc)) from exc + return PreviewResponse(fields=fields) + + +@router.post("/jobs", response_model=ParseJobRead, status_code=201) +def create_profile_job(payload: CreateProfileJobRequest, db: Session = Depends(get_db)): + try: + cfg = TelegramSourceConfig( + channel=payload.channel.strip(), + limit=payload.limit, + extract_mode="profile", + heuristic_profile=payload.profile, + sample_post=payload.sample_post, + ) + except Exception as exc: + raise HTTPException(status_code=400, detail=str(exc)) from exc + + job = ParseJob( + source_type="telegram", + source_config=cfg.model_dump(), + interval_seconds=payload.interval_seconds, + is_active=payload.is_active, + status="queued", + ) + db.add(job) + db.commit() + db.refresh(job) + + enqueue_job(job.id, job.source_type, job.source_config) + return job diff --git a/centers/analytics/api/app/schemas.py b/centers/analytics/api/app/schemas.py index 6d04282..670878c 100644 --- a/centers/analytics/api/app/schemas.py +++ b/centers/analytics/api/app/schemas.py @@ -117,6 +117,7 @@ class IngestRequest(BaseModel): class ListenerSubscription(BaseModel): job_id: int channel: str + source_config: dict[str, Any] = Field(default_factory=dict) class IngestResponse(BaseModel): diff --git a/centers/analytics/api/app/services/parser_builder.py b/centers/analytics/api/app/services/parser_builder.py new file mode 100644 index 0000000..b9a4396 --- /dev/null +++ b/centers/analytics/api/app/services/parser_builder.py @@ -0,0 +1,130 @@ +"""Parser builder: one-shot DeepSeek profile generation + static preview.""" + +from __future__ import annotations + +import json +import logging +import os +import re +from typing import Any + +import httpx + +from contracts.heuristic_profile import ( + HeuristicProfile, + TARGET_FIELDS, + apply_profile, + target_field_specs, +) + +logger = logging.getLogger("ca.parser_builder") + +GENERATE_SYSTEM = ( + "You design a static heuristic parser profile for structured Telegram posts. " + "Output valid JSON only, no markdown, no Python code. " + "The profile is applied with regex/line/marker rules at runtime — never with an LLM." +) + + +def deepseek_enabled() -> bool: + return bool(os.getenv("DEEPSEEK_API_KEY", "").strip()) + + +def deepseek_settings() -> dict[str, str]: + return { + "api_key": os.getenv("DEEPSEEK_API_KEY", "").strip(), + "base_url": os.getenv("DEEPSEEK_BASE_URL", "https://api.deepseek.com").rstrip("/"), + "model": os.getenv("DEEPSEEK_MODEL", "deepseek-chat"), + } + + +def get_target_fields() -> list[dict[str, str]]: + return target_field_specs() + + +def preview_with_profile(sample_post: str, profile: dict[str, Any] | HeuristicProfile) -> dict[str, str]: + return apply_profile(sample_post, profile) + + +def _extract_json_object(content: str) -> dict[str, Any]: + content = (content or "").strip() + if content.startswith("```"): + content = re.sub(r"^```(?:json)?\s*", "", content) + content = re.sub(r"\s*```$", "", content) + try: + parsed = json.loads(content) + except json.JSONDecodeError: + match = re.search(r"\{[\s\S]*\}", content) + if not match: + raise + parsed = json.loads(match.group(0)) + if not isinstance(parsed, dict): + raise ValueError("LLM response must be a JSON object") + return parsed + + +async def generate_profile(sample_post: str) -> HeuristicProfile: + settings = deepseek_settings() + if not settings["api_key"]: + raise RuntimeError( + "DEEPSEEK_API_KEY is not set on ca-api. Add it to .env for parser generation." + ) + + sample = (sample_post or "").strip() + if not sample: + raise ValueError("sample_post is required") + + fields_help = "\n".join( + f"- {spec['name']}: {spec['description']}" for spec in target_field_specs() + ) + user_prompt = ( + "Build a HeuristicProfile JSON for this sample Telegram post.\n\n" + "Schema:\n" + '{"version": 1, "notes": "...", "fields": {' + '"": {"strategy": "regex|line|after_marker|between|full_text|literal", ' + '"pattern": "...", "group": 1, "line_index": 0, "marker": "...", ' + '"end_marker": "...", "value": "...", "flags": "im", "strip": true}' + "}}\n\n" + "Rules:\n" + "- Only include fields you can extract reliably from the sample.\n" + f"- Allowed field names: {', '.join(TARGET_FIELDS)}.\n" + "- Prefer regex / line / after_marker / between over literal.\n" + "- Use literal only for constant topic tags.\n" + "- coords should capture 'lat, lon' when present.\n" + "- event_date should capture DD.MM.YYYY or similar.\n" + "- Do not invent Python code; only declarative rules.\n\n" + f"Target fields:\n{fields_help}\n\n" + f"Sample post:\n{sample[:12000]}" + ) + + payload = { + "model": settings["model"], + "messages": [ + {"role": "system", "content": GENERATE_SYSTEM}, + {"role": "user", "content": user_prompt}, + ], + "temperature": 0.1, + "response_format": {"type": "json_object"}, + } + + url = f"{settings['base_url']}/chat/completions" + async with httpx.AsyncClient(timeout=90.0) as client: + response = await client.post( + url, + headers={ + "Authorization": f"Bearer {settings['api_key']}", + "Content-Type": "application/json", + }, + json=payload, + ) + response.raise_for_status() + data = response.json() + + content = data["choices"][0]["message"]["content"] + raw = _extract_json_object(content) + # Accept either top-level profile or {"profile": {...}} + if "fields" not in raw and isinstance(raw.get("profile"), dict): + raw = raw["profile"] + if "version" not in raw: + raw["version"] = 1 + return HeuristicProfile.model_validate(raw) diff --git a/centers/analytics/api/requirements.txt b/centers/analytics/api/requirements.txt index d44b3a0..fb6c4c0 100644 --- a/centers/analytics/api/requirements.txt +++ b/centers/analytics/api/requirements.txt @@ -5,3 +5,4 @@ pydantic==2.10.3 python-multipart==0.0.20 psycopg2-binary==2.9.10 redis==5.2.1 +httpx==0.28.1 diff --git a/centers/analytics/frontend/src/api/parserBuilder.ts b/centers/analytics/frontend/src/api/parserBuilder.ts new file mode 100644 index 0000000..3d52560 --- /dev/null +++ b/centers/analytics/frontend/src/api/parserBuilder.ts @@ -0,0 +1,52 @@ +import { ADMIN_BASE, request } from "./client"; +import type { ParseJob } from "../types/admin"; + +export type TargetField = { + name: string; + type: string; + description: string; +}; + +export type HeuristicProfile = { + version: number; + fields: Record>; + notes?: string; +}; + +export function fetchTargetFields(): Promise<{ fields: TargetField[] }> { + return request<{ fields: TargetField[] }>("/parser-builder/target-fields", undefined, ADMIN_BASE); +} + +export function generateParserProfile(sample_post: string): Promise<{ profile: HeuristicProfile }> { + return request<{ profile: HeuristicProfile }>( + "/parser-builder/generate", + { method: "POST", body: JSON.stringify({ sample_post }) }, + ADMIN_BASE, + ); +} + +export function previewParserProfile( + sample_post: string, + profile: HeuristicProfile | Record, +): Promise<{ fields: Record }> { + return request<{ fields: Record }>( + "/parser-builder/preview", + { method: "POST", body: JSON.stringify({ sample_post, profile }) }, + ADMIN_BASE, + ); +} + +export function createProfileJob(payload: { + channel: string; + profile: HeuristicProfile | Record; + sample_post?: string; + limit?: number; + interval_seconds?: number; + is_active?: boolean; +}): Promise { + return request( + "/parser-builder/jobs", + { method: "POST", body: JSON.stringify(payload) }, + ADMIN_BASE, + ); +} diff --git a/centers/analytics/frontend/src/layouts/AppLayout.vue b/centers/analytics/frontend/src/layouts/AppLayout.vue index 0f7d200..f607455 100644 --- a/centers/analytics/frontend/src/layouts/AppLayout.vue +++ b/centers/analytics/frontend/src/layouts/AppLayout.vue @@ -7,6 +7,7 @@ const route = useRoute(); const navItems = [ { to: "/", label: "Карта", exact: true }, { to: "/parsers", label: "Парсеры" }, + { to: "/parser-builder", label: "Конструктор" }, { to: "/events", label: "События" }, { to: "/analytics", label: "Аналитика" }, { to: "/consumers", label: "ПИ" }, diff --git a/centers/analytics/frontend/src/router/index.ts b/centers/analytics/frontend/src/router/index.ts index 9f856ed..43dd613 100644 --- a/centers/analytics/frontend/src/router/index.ts +++ b/centers/analytics/frontend/src/router/index.ts @@ -5,6 +5,7 @@ import AnalyticsView from "../views/AnalyticsView.vue"; import ConsumersView from "../views/ConsumersView.vue"; import EventsView from "../views/EventsView.vue"; import MapViewPage from "../views/MapViewPage.vue"; +import ParserBuilderView from "../views/ParserBuilderView.vue"; import ParsersView from "../views/ParsersView.vue"; const router = createRouter({ @@ -16,6 +17,7 @@ const router = createRouter({ children: [ { path: "", name: "map", component: MapViewPage }, { path: "parsers", name: "parsers", component: ParsersView }, + { path: "parser-builder", name: "parser-builder", component: ParserBuilderView }, { path: "events", name: "events", component: EventsView }, { path: "analytics", name: "analytics", component: AnalyticsView }, { path: "consumers", name: "consumers", component: ConsumersView }, diff --git a/centers/analytics/frontend/src/views/ParserBuilderView.vue b/centers/analytics/frontend/src/views/ParserBuilderView.vue new file mode 100644 index 0000000..6549ec1 --- /dev/null +++ b/centers/analytics/frontend/src/views/ParserBuilderView.vue @@ -0,0 +1,320 @@ + + +