Add reusable LLM parser profiles with multi-event extract.
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>
This commit is contained in:
@@ -0,0 +1,36 @@
|
||||
"""Stale parse-job recovery helpers (running/queued left behind after worker crash)."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import os
|
||||
from datetime import datetime, timezone
|
||||
|
||||
from ..models import ParseJob
|
||||
|
||||
# Jobs stuck in running/queued longer than this are considered abandoned.
|
||||
STALE_JOB_SECONDS = int(os.getenv("STALE_JOB_SECONDS", "900"))
|
||||
|
||||
|
||||
def _aware(dt: datetime | None) -> datetime | None:
|
||||
if dt is None:
|
||||
return None
|
||||
if dt.tzinfo is None:
|
||||
return dt.replace(tzinfo=timezone.utc)
|
||||
return dt
|
||||
|
||||
|
||||
def job_anchor_time(job: ParseJob) -> datetime | None:
|
||||
"""Best available timestamp for staleness (prefer last_run_at)."""
|
||||
return _aware(job.last_run_at) or _aware(getattr(job, "created_at", None))
|
||||
|
||||
|
||||
def is_stale_job(job: ParseJob, now: datetime | None = None, *, ttl: int | None = None) -> bool:
|
||||
if job.status not in ("running", "queued"):
|
||||
return False
|
||||
now = now or datetime.now(timezone.utc)
|
||||
anchor = job_anchor_time(job)
|
||||
if anchor is None:
|
||||
# No timestamp — treat long-lived running as stale immediately for recovery
|
||||
return job.status == "running"
|
||||
limit = ttl if ttl is not None else STALE_JOB_SECONDS
|
||||
return (now - anchor).total_seconds() >= limit
|
||||
@@ -33,6 +33,25 @@ def flatten_pair_config(
|
||||
limit: int = 100,
|
||||
) -> dict:
|
||||
"""Expand Profile + Channel into Redis/CP source_config."""
|
||||
kind = (profile.kind or "heuristic").strip().lower()
|
||||
if kind == "llm":
|
||||
from contracts.llm_profile import LlmProfile
|
||||
|
||||
if not profile.llm_profile:
|
||||
raise ValueError("ParserProfile.llm_profile is empty")
|
||||
llm = LlmProfile.model_validate(profile.llm_profile)
|
||||
cfg = TelegramSourceConfig(
|
||||
channel=channel.channel.strip(),
|
||||
limit=limit,
|
||||
extract_mode="llm",
|
||||
extract_schema=llm.extract_schema,
|
||||
instruction=llm.instruction,
|
||||
required_fields=list(llm.required_fields),
|
||||
multi_event=bool(llm.multi_event),
|
||||
sample_post=profile.sample_post or None,
|
||||
)
|
||||
return cfg.model_dump()
|
||||
|
||||
if not profile.heuristic_profile:
|
||||
raise ValueError("ParserProfile.heuristic_profile is empty")
|
||||
cfg = TelegramSourceConfig(
|
||||
|
||||
@@ -23,11 +23,14 @@ def migrate_schema(engine: Engine) -> None:
|
||||
if "channel_id" not in columns:
|
||||
statements.append("ALTER TABLE parse_jobs ADD COLUMN channel_id INTEGER")
|
||||
|
||||
# create_all handles new tables; FKs on existing DBs may need indexes
|
||||
if "parse_jobs" in tables:
|
||||
# Re-inspect after potential adds is not needed for FK constraints here —
|
||||
# create_all + nullable FKs are enough for MVP; optional constraints below.
|
||||
pass
|
||||
if "parser_profiles" in tables:
|
||||
columns = {col["name"] for col in inspector.get_columns("parser_profiles")}
|
||||
if "kind" not in columns:
|
||||
statements.append(
|
||||
"ALTER TABLE parser_profiles ADD COLUMN kind VARCHAR(50) NOT NULL DEFAULT 'heuristic'"
|
||||
)
|
||||
if "llm_profile" not in columns:
|
||||
statements.append("ALTER TABLE parser_profiles ADD COLUMN llm_profile JSON")
|
||||
|
||||
if not statements:
|
||||
return
|
||||
|
||||
@@ -14,6 +14,7 @@ from contracts.heuristic_profile import (
|
||||
HeuristicProfile,
|
||||
TARGET_FIELDS,
|
||||
apply_profile,
|
||||
match_profile,
|
||||
target_field_specs,
|
||||
)
|
||||
|
||||
@@ -47,6 +48,14 @@ def preview_with_profile(sample_post: str, profile: dict[str, Any] | HeuristicPr
|
||||
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()]
|
||||
|
||||
@@ -191,3 +200,111 @@ async def generate_profile(
|
||||
|
||||
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)
|
||||
|
||||
|
||||
@@ -5,6 +5,7 @@ from datetime import datetime, timezone
|
||||
|
||||
from ..database import SessionLocal
|
||||
from ..models import ParseJob
|
||||
from .job_stale import STALE_JOB_SECONDS, is_stale_job
|
||||
from .jobs import enqueue_parse_job
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
@@ -13,10 +14,47 @@ TICK_SECONDS = 30
|
||||
RECURRING_STATUSES = ("completed", "failed")
|
||||
|
||||
|
||||
def recover_stale_jobs(db, now: datetime) -> None:
|
||||
"""Mark abandoned running/queued jobs as failed; re-queue if still active."""
|
||||
stuck = (
|
||||
db.query(ParseJob)
|
||||
.filter(ParseJob.status.in_(("running", "queued")))
|
||||
.all()
|
||||
)
|
||||
for job in stuck:
|
||||
if not is_stale_job(job, now):
|
||||
continue
|
||||
prev = job.status
|
||||
job.status = "failed"
|
||||
job.last_error = (
|
||||
f"Stale {prev} recovered after {STALE_JOB_SECONDS}s "
|
||||
"(worker likely restarted)"
|
||||
)
|
||||
job.last_run_at = now
|
||||
db.commit()
|
||||
logger.warning("Recovered stale job %s (was %s)", job.id, prev)
|
||||
|
||||
if not job.is_active or job.interval_seconds <= 0:
|
||||
continue
|
||||
job.status = "queued"
|
||||
job.last_error = None
|
||||
db.commit()
|
||||
try:
|
||||
enqueue_parse_job(db, job)
|
||||
logger.info("Re-queued recovered job %s", job.id)
|
||||
except ValueError as exc:
|
||||
job.status = "failed"
|
||||
job.last_error = str(exc)
|
||||
db.commit()
|
||||
logger.warning("Skip re-queue recovered job %s: %s", job.id, exc)
|
||||
|
||||
|
||||
def run_scheduler_tick() -> None:
|
||||
db = SessionLocal()
|
||||
try:
|
||||
now = datetime.now(timezone.utc)
|
||||
recover_stale_jobs(db, now)
|
||||
|
||||
jobs = (
|
||||
db.query(ParseJob)
|
||||
.filter(
|
||||
@@ -61,5 +99,9 @@ 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)
|
||||
logger.info(
|
||||
"Parse job scheduler started (tick every %ss, stale after %ss)",
|
||||
TICK_SECONDS,
|
||||
STALE_JOB_SECONDS,
|
||||
)
|
||||
return stop_event
|
||||
|
||||
Reference in New Issue
Block a user