Files
gitrusprusandCursor 3f9dc6643b 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>
2026-09-13 20:22:38 +03:00

108 lines
3.2 KiB
Python

import logging
import threading
import time
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__)
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(
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()
try:
enqueue_parse_job(db, job)
except ValueError as exc:
job.status = "failed"
job.last_error = str(exc)
db.commit()
logger.warning("Skip re-queue job %s: %s", job.id, exc)
continue
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, stale after %ss)",
TICK_SECONDS,
STALE_JOB_SECONDS,
)
return stop_event