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 <cursoragent@cursor.com>
This commit is contained in:
@@ -7,6 +7,7 @@ import logging
|
||||
from contracts.sources import TelegramSourceConfig
|
||||
from workers.adapters.base import WorkerContext
|
||||
from workers.converter import event_record_to_ingest
|
||||
from workers.heuristic_profile import extract_with_profile
|
||||
from workers.llm_extract import (
|
||||
DEFAULT_EXTRACT_SCHEMA,
|
||||
DEFAULT_INSTRUCTION,
|
||||
@@ -57,11 +58,37 @@ class TelegramAdapter:
|
||||
return [], "extract_mode=llm requires DEEPSEEK_API_KEY in worker env"
|
||||
return await _extract_posts_with_llm(posts, cfg)
|
||||
|
||||
if cfg.extract_mode == "profile":
|
||||
return _extract_posts_with_profile(posts, cfg)
|
||||
|
||||
records = parse_event_posts(posts)
|
||||
events = [event_record_to_ingest(r) for r in records]
|
||||
return events, None
|
||||
|
||||
|
||||
def _extract_posts_with_profile(posts, cfg: TelegramSourceConfig) -> tuple[list[dict], str | None]:
|
||||
profile = cfg.heuristic_profile or {}
|
||||
events: list[dict] = []
|
||||
for post in posts:
|
||||
text = (post.text or "").strip()
|
||||
if not text:
|
||||
continue
|
||||
events.append(
|
||||
extract_with_profile(
|
||||
text,
|
||||
profile,
|
||||
source_url=post.url,
|
||||
source_type="telegram",
|
||||
extra_metadata={
|
||||
"channel": post.channel,
|
||||
"message_id": post.id,
|
||||
"post_date": post.date.isoformat() if post.date else None,
|
||||
},
|
||||
)
|
||||
)
|
||||
return events, None
|
||||
|
||||
|
||||
async def _extract_posts_with_llm(posts, cfg: TelegramSourceConfig) -> tuple[list[dict], str | None]:
|
||||
schema = cfg.extract_schema or DEFAULT_EXTRACT_SCHEMA
|
||||
instruction = cfg.instruction or DEFAULT_INSTRUCTION
|
||||
|
||||
@@ -0,0 +1,73 @@
|
||||
"""Apply HeuristicProfile to Telegram posts → IngestEventItem-shaped dicts."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from typing import Any
|
||||
|
||||
from contracts.heuristic_profile import HeuristicProfile, apply_profile
|
||||
from workers.llm_extract import parse_coords, parse_date
|
||||
|
||||
|
||||
def profile_fields_to_ingest(
|
||||
*,
|
||||
source_type: str,
|
||||
source_url: str,
|
||||
raw_text: str,
|
||||
fields: dict[str, str],
|
||||
extra_metadata: dict | None = None,
|
||||
) -> dict:
|
||||
lat, lng = parse_coords(str(fields.get("coords") or ""))
|
||||
description = str(fields.get("description") or raw_text)[:8000]
|
||||
locality = str(fields.get("locality") or "")
|
||||
region = str(fields.get("region") or locality or "") or None
|
||||
if fields.get("title"):
|
||||
title = str(fields["title"])
|
||||
elif locality:
|
||||
title = locality
|
||||
elif description:
|
||||
title = description.splitlines()[0][:120]
|
||||
else:
|
||||
title = ""
|
||||
topic = str(fields.get("topic") or "telegram_profile")
|
||||
event_date = parse_date(str(fields.get("event_date") or ""))
|
||||
|
||||
meta: dict[str, Any] = {
|
||||
"extract_mode": "profile",
|
||||
"extracted": dict(fields),
|
||||
}
|
||||
if extra_metadata:
|
||||
meta.update(extra_metadata)
|
||||
|
||||
return {
|
||||
"source_type": source_type,
|
||||
"source_url": source_url,
|
||||
"raw_text": raw_text[:20000],
|
||||
"title": title[:255],
|
||||
"description": description,
|
||||
"locality": locality,
|
||||
"latitude": lat,
|
||||
"longitude": lng,
|
||||
"event_date": event_date.isoformat() if event_date else None,
|
||||
"region": region,
|
||||
"topic": topic,
|
||||
"tags": [source_type, "profile"],
|
||||
"metadata": meta,
|
||||
}
|
||||
|
||||
|
||||
def extract_with_profile(
|
||||
text: str,
|
||||
profile: HeuristicProfile | dict[str, Any],
|
||||
*,
|
||||
source_url: str,
|
||||
source_type: str = "telegram",
|
||||
extra_metadata: dict | None = None,
|
||||
) -> dict:
|
||||
fields = apply_profile(text, profile)
|
||||
return profile_fields_to_ingest(
|
||||
source_type=source_type,
|
||||
source_url=source_url,
|
||||
raw_text=text,
|
||||
fields=fields,
|
||||
extra_metadata=extra_metadata,
|
||||
)
|
||||
@@ -9,6 +9,7 @@ import httpx
|
||||
from telethon import TelegramClient, events
|
||||
|
||||
from workers.converter import event_record_to_ingest
|
||||
from workers.heuristic_profile import extract_with_profile
|
||||
from workers.parsers.telegram_events import parse_event_post
|
||||
from workers.sources.telegram_client import message_to_post, normalize_channel
|
||||
|
||||
@@ -22,11 +23,12 @@ 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] = {}
|
||||
# channel -> {job_id, source_config}
|
||||
self._channels: dict[str, dict[str, Any]] = {}
|
||||
self._chat_ids: set[int] = set()
|
||||
self._handlers_registered = False
|
||||
|
||||
async def fetch_subscriptions(self) -> dict[str, int]:
|
||||
async def fetch_subscriptions(self) -> dict[str, dict[str, Any]]:
|
||||
async with httpx.AsyncClient(timeout=30.0) as http:
|
||||
response = await http.get(
|
||||
f"{CA_API_URL}/internal/listener/subscriptions",
|
||||
@@ -35,7 +37,7 @@ class TelegramListener:
|
||||
response.raise_for_status()
|
||||
data = response.json()
|
||||
|
||||
channels: dict[str, int] = {}
|
||||
channels: dict[str, dict[str, Any]] = {}
|
||||
for item in data:
|
||||
raw = item.get("channel")
|
||||
job_id = item.get("job_id")
|
||||
@@ -46,7 +48,13 @@ class TelegramListener:
|
||||
except ValueError:
|
||||
logger.warning("Skip invalid channel in subscription: %r", raw)
|
||||
continue
|
||||
channels.setdefault(key, int(job_id))
|
||||
channels.setdefault(
|
||||
key,
|
||||
{
|
||||
"job_id": int(job_id),
|
||||
"source_config": dict(item.get("source_config") or {}),
|
||||
},
|
||||
)
|
||||
return channels
|
||||
|
||||
async def refresh_subscriptions(self) -> None:
|
||||
@@ -73,10 +81,34 @@ class TelegramListener:
|
||||
", ".join(sorted(channels)) or "(none)",
|
||||
)
|
||||
|
||||
async def _ingest_post(self, channel: str, post) -> None:
|
||||
def _build_event(self, channel: str, post) -> dict:
|
||||
sub = self._channels.get(normalize_channel(channel)) or {}
|
||||
cfg = sub.get("source_config") or {}
|
||||
extract_mode = cfg.get("extract_mode") or "heuristic"
|
||||
text = (post.text or "").strip()
|
||||
|
||||
if extract_mode == "profile" and cfg.get("heuristic_profile"):
|
||||
return extract_with_profile(
|
||||
text,
|
||||
cfg["heuristic_profile"],
|
||||
source_url=post.url,
|
||||
source_type="telegram",
|
||||
extra_metadata={
|
||||
"channel": post.channel,
|
||||
"message_id": post.id,
|
||||
"post_date": post.date.isoformat() if post.date else None,
|
||||
"listener": True,
|
||||
},
|
||||
)
|
||||
|
||||
# llm jobs fall back to heuristic in listener (LLM is batch-only by design)
|
||||
record = parse_event_post(post)
|
||||
event = event_record_to_ingest(record)
|
||||
job_id = self._channels.get(normalize_channel(channel))
|
||||
return event_record_to_ingest(record)
|
||||
|
||||
async def _ingest_post(self, channel: str, post) -> None:
|
||||
event = self._build_event(channel, post)
|
||||
sub = self._channels.get(normalize_channel(channel)) or {}
|
||||
job_id = sub.get("job_id")
|
||||
|
||||
async with httpx.AsyncClient(timeout=60.0) as http:
|
||||
response = await http.post(
|
||||
|
||||
Reference in New Issue
Block a user