Support kind=llm profiles (instruction/schema), optional multi-event posts via #eN URLs, and recover stale running/queued parse jobs after worker crashes. Co-authored-by: Cursor <cursoragent@cursor.com>
311 lines
11 KiB
Python
311 lines
11 KiB
Python
"""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,
|
|
match_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. "
|
|
"Every rule must actually extract a non-empty value from the given sample when applied."
|
|
)
|
|
|
|
|
|
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 match_preview(
|
|
sample_post: str,
|
|
profile: dict[str, Any] | HeuristicProfile,
|
|
) -> tuple[dict[str, str], bool, list[str]]:
|
|
matched, fields, missing = match_profile(sample_post, profile)
|
|
return fields, matched, missing
|
|
|
|
|
|
def empty_preview_fields(preview: dict[str, str]) -> list[str]:
|
|
return [name for name, value in preview.items() if not (value or "").strip()]
|
|
|
|
|
|
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
|
|
|
|
|
|
def _base_user_prompt(sample: str) -> str:
|
|
fields_help = "\n".join(
|
|
f"- {spec['name']}: {spec['description']}" for spec in target_field_specs()
|
|
)
|
|
return (
|
|
"Build a HeuristicProfile JSON for this sample Telegram post.\n\n"
|
|
"Schema:\n"
|
|
'{"version": 1, "notes": "...", "fields": {'
|
|
'"<field>": {"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/region tags.\n"
|
|
"- title: usually the first non-empty line (strategy line, line_index 0).\n"
|
|
"- description: the multi-line body AFTER the title and BEFORE coords/hashtags. "
|
|
"Prefer strategy between or regex with flags including \"s\" (DOTALL). "
|
|
"Do NOT leave description empty if the sample has a body paragraph.\n"
|
|
"- locality: settlement name from the title line when present.\n"
|
|
"- coords should capture 'lat, lon' when present.\n"
|
|
"- event_date should capture DD.MM.YYYY / DD.MM.YY or similar.\n"
|
|
"- topic/region: hashtags like #ru → topic \"ru\"; avoid stuffing unrelated tags into region.\n"
|
|
"- Do not invent Python code; only declarative rules.\n"
|
|
"- Mentally apply each rule to the sample: every included field must yield a non-empty string.\n\n"
|
|
f"Target fields:\n{fields_help}\n\n"
|
|
f"Sample post:\n{sample[:12000]}"
|
|
)
|
|
|
|
|
|
def _parse_profile_response(content: str) -> HeuristicProfile:
|
|
raw = _extract_json_object(content)
|
|
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)
|
|
|
|
|
|
async def generate_profile(
|
|
sample_post: str,
|
|
*,
|
|
hint: str | None = None,
|
|
current_profile: dict[str, Any] | HeuristicProfile | None = None,
|
|
) -> 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")
|
|
|
|
hint_text = (hint or "").strip()
|
|
current: HeuristicProfile | None = None
|
|
if current_profile is not None:
|
|
current = (
|
|
current_profile
|
|
if isinstance(current_profile, HeuristicProfile)
|
|
else HeuristicProfile.model_validate(current_profile)
|
|
)
|
|
|
|
if current is not None or hint_text:
|
|
preview = preview_with_profile(sample, current) if current is not None else {}
|
|
empty = empty_preview_fields(preview) if preview else []
|
|
parts = [
|
|
"Revise the HeuristicProfile JSON for this sample Telegram post.",
|
|
"Keep strategies that already work; fix only what the manager asks for "
|
|
"and any fields that preview as empty.",
|
|
"Output the full updated profile JSON (same schema), not a patch.",
|
|
"",
|
|
_base_user_prompt(sample),
|
|
]
|
|
if current is not None:
|
|
parts.extend(
|
|
[
|
|
"",
|
|
"Current profile JSON:",
|
|
json.dumps(current.model_dump(), ensure_ascii=False, indent=2)[:12000],
|
|
]
|
|
)
|
|
if empty:
|
|
parts.append(
|
|
"Fields currently empty in preview (must become non-empty if present in sample): "
|
|
+ ", ".join(empty)
|
|
)
|
|
else:
|
|
parts.append("Current preview fields: " + json.dumps(preview, ensure_ascii=False))
|
|
if hint_text:
|
|
parts.extend(["", "Manager hint (follow this):", hint_text[:4000]])
|
|
user_prompt = "\n".join(parts)
|
|
else:
|
|
user_prompt = _base_user_prompt(sample)
|
|
|
|
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"]
|
|
return _parse_profile_response(content)
|
|
|
|
|
|
async def preview_llm_extract(
|
|
sample_post: str,
|
|
llm_profile: dict[str, Any],
|
|
) -> tuple[list[dict[str, str]], dict[str, str], bool, list[str], bool, int]:
|
|
"""One-shot DeepSeek extract for CA preview (does not persist).
|
|
|
|
Returns (events, fields, matched, missing_required, is_event, matched_count).
|
|
``fields`` is the first event (or empty) for legacy UI compatibility.
|
|
"""
|
|
from contracts.llm_profile import DEFAULT_INSTRUCTION, LlmProfile, match_llm_required
|
|
|
|
settings = deepseek_settings()
|
|
if not settings["api_key"]:
|
|
raise RuntimeError(
|
|
"DEEPSEEK_API_KEY is not set on ca-api. Add it to .env for LLM preview."
|
|
)
|
|
|
|
sample = (sample_post or "").strip()
|
|
if not sample:
|
|
raise ValueError("sample_post is required")
|
|
|
|
profile = LlmProfile.model_validate(llm_profile)
|
|
schema = profile.extract_schema
|
|
instr = profile.instruction or DEFAULT_INSTRUCTION
|
|
schema_lines = "\n".join(f"- {k}: {v}" for k, v in schema.items())
|
|
multi = bool(profile.multi_event)
|
|
|
|
if multi:
|
|
user_prompt = (
|
|
f"{instr}\n\n"
|
|
"If the text describes multiple distinct events (different places, "
|
|
"coords, or dates), return one object per event in \"events\".\n"
|
|
f"Fields per event:\n{schema_lines}\n\n"
|
|
'Return JSON: {"is_event": true|false, "events": [{<field>: <string>}, ...]}\n'
|
|
"If there is no event, return is_event=false and events=[].\n\n"
|
|
f"Text:\n{sample[:12000]}"
|
|
)
|
|
else:
|
|
user_prompt = (
|
|
f"{instr}\n\n"
|
|
f"Fields to extract:\n{schema_lines}\n\n"
|
|
'Return JSON: {"is_event": true|false, "fields": {<field>: <string>}}\n\n'
|
|
f"Text:\n{sample[:12000]}"
|
|
)
|
|
|
|
payload = {
|
|
"model": settings["model"],
|
|
"messages": [
|
|
{
|
|
"role": "system",
|
|
"content": (
|
|
"You extract structured event data for a geoint map. "
|
|
"Output valid JSON only, no markdown."
|
|
),
|
|
},
|
|
{"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"]
|
|
parsed = _extract_json_object(content)
|
|
is_event = bool(parsed.get("is_event", True))
|
|
empty_fields = {key: "" for key in schema}
|
|
|
|
if not is_event:
|
|
return [], empty_fields, False, [], False, 0
|
|
|
|
events: list[dict[str, str]] = []
|
|
if multi and isinstance(parsed.get("events"), list):
|
|
for item in parsed["events"]:
|
|
if isinstance(item, dict):
|
|
events.append({key: str(item.get(key) or "").strip() for key in schema})
|
|
else:
|
|
fields_raw = parsed.get("fields") if isinstance(parsed.get("fields"), dict) else parsed
|
|
if not isinstance(fields_raw, dict):
|
|
fields_raw = {}
|
|
events.append({key: str(fields_raw.get(key) or "").strip() for key in schema})
|
|
|
|
if not events:
|
|
return [], empty_fields, False, [], False, 0
|
|
|
|
matched_events = [ev for ev in events if match_llm_required(ev, list(profile.required_fields))]
|
|
fields = events[0]
|
|
missing = [
|
|
name
|
|
for name in profile.required_fields
|
|
if not str(fields.get(name) or "").strip()
|
|
]
|
|
matched = match_llm_required(fields, list(profile.required_fields))
|
|
return events, fields, matched, missing, True, len(matched_events)
|
|
|