app.bettersight.io/backend/jobs.py

128 lines
5.0 KiB
Python

import os
import time
import logging
from datetime import datetime, timezone
from redis import Redis
from rq import Queue
from repositories.scrape_repository import scrape_repository
from repositories.competitor_repository import competitor_repository
from services import analysis_service, alert_service
logger = logging.getLogger(__name__)
REDIS_URL = os.getenv('REDIS_URL', 'redis://bettersight-redis:6379')
redis_conn = Redis.from_url(REDIS_URL)
# Three priority queues — Intelligence tier (not yet built) will use
# `high`; Analyse-tier on-demand and n8n cron jobs use `normal`;
# background/cleanup work uses `low`. Worker command:
# rq worker high normal low -u $REDIS_URL
high_q = Queue('high', connection=redis_conn)
normal_q = Queue('normal', connection=redis_conn)
low_q = Queue('low', connection=redis_conn)
JOB_SLOW_SECONDS = 300 # 5 minutes — §29 THRESHOLDS['job_slow_mins']
JOB_STUCK_SECONDS = 600 # 10 minutes — §29 THRESHOLDS['job_stuck_mins']
def enqueue_research_job(job_id, data):
"""
Enqueues a research job to the `normal` priority queue. Never runs
synchronously — every /research-style route calls this instead of
invoking run_research_job() directly.
Data flow:
job_id (the scrape_runs record id, reused as the RQ job id) + data →
normal_q.enqueue(run_research_job, job_id, data) →
RQ job id returned (== job_id, since we pass it explicitly)
"""
normal_q.enqueue(run_research_job, job_id, data, job_id=job_id, job_timeout=900)
return job_id
def run_research_job(job_id, data):
"""
The RQ job function. Loops over every requested competitor, running
the fast/refresh path decision per CLAUDE.md §7. One competitor's
failure is logged and skipped — it never aborts the rest of the
batch (§18 resilience rule). Records duration and fires Gotify
alerts on slow/stuck/failed completion (§29 Component 2).
Data flow:
job_id + { tenant_id, competitors, destination, duration,
travel_style, tab_type } →
scrape_runs status → running →
per competitor: analysis_service.run_competitor_analysis() →
success: append to results, competitors_done += 1 →
failure: append to results.failed, continue →
scrape_runs status → complete (or failed if every competitor failed) →
duration computed → scrape_runs.duration_seconds updated →
slow/stuck duration → Gotify alert fired
"""
start = time.time()
scrape_repository.update(job_id, {'status': 'running'})
tenant_id = data.get('tenant_id')
competitors_input = data.get('competitors', [])
context = {
'destination': data.get('destination'),
'duration': data.get('duration'),
'travel_style': data.get('travel_style'),
'tab_type': data.get('tab_type'),
'product_name': data.get('productName') or data.get('product_name'),
}
success, failed = [], []
done = 0
for entry in competitors_input:
competitor = competitor_repository.get_by_id(entry['id']) if isinstance(entry, dict) else None
if not competitor:
failed.append({'name': entry.get('name', 'Unknown') if isinstance(entry, dict) else str(entry),
'error': 'Competitor not found'})
continue
try:
outcome = analysis_service.run_competitor_analysis(tenant_id, competitor, context)
success.append({'competitor_id': competitor['id'], 'name': competitor['name'], **outcome})
except Exception as e:
logger.error(f'{competitor.get("name", entry)} failed in job {job_id}: {str(e)}')
failed.append({'competitor_id': competitor['id'], 'name': competitor['name'], 'error': str(e)})
finally:
done += 1
scrape_repository.update(job_id, {'competitors_done': done})
status = 'complete' if success or not competitors_input else 'failed'
scrape_repository.update(job_id, {
'status': status,
'results': {'success': success, 'failed': failed},
'error_log': '; '.join(f['error'] for f in failed) if failed else '',
'completed_at': datetime.now(timezone.utc).isoformat(),
})
if status == 'failed':
alert_service.send_gotify(
title='🔴 Analysis job failed',
message=f'Job failed for tenant {tenant_id}.\nErrors: {"; ".join(f["error"] for f in failed)[:200]}',
priority=7
)
duration = round(time.time() - start)
scrape_repository.update(job_id, {'duration_seconds': duration})
if duration > JOB_STUCK_SECONDS:
alert_service.send_gotify(
title='🔴 Job likely stuck',
message=f'Job took {round(duration / 60, 1)} mins. Tenant: {tenant_id}. '
f'Possible proxy or Playwright issue.',
priority=10
)
elif duration > JOB_SLOW_SECONDS:
alert_service.send_gotify(
title='🟡 Slow job detected',
message=f'Job took {round(duration / 60, 1)} mins. Tenant: {tenant_id}.',
priority=5
)