60 lines
1.9 KiB
Python
60 lines
1.9 KiB
Python
"""Small PostgreSQL-backed worker for post-call analysis."""
|
|
|
|
from __future__ import annotations
|
|
|
|
import asyncio
|
|
|
|
from db.models import ConversationSession
|
|
from db.session import SessionLocal
|
|
from loguru import logger
|
|
from services.post_call.analyzer import analyze_conversation
|
|
from sqlalchemy import select, update
|
|
|
|
|
|
POLL_INTERVAL_SECONDS = 2.0
|
|
|
|
|
|
async def recover_interrupted_analyses() -> None:
|
|
"""Return work interrupted by a previous process shutdown to the queue."""
|
|
async with SessionLocal() as session:
|
|
await session.execute(
|
|
update(ConversationSession)
|
|
.where(ConversationSession.analysis_status == "processing")
|
|
.values(analysis_status="pending", analysis_error="")
|
|
)
|
|
await session.commit()
|
|
|
|
|
|
async def claim_pending_analysis() -> str | None:
|
|
async with SessionLocal() as session:
|
|
row = (
|
|
await session.execute(
|
|
select(ConversationSession)
|
|
.where(ConversationSession.analysis_status == "pending")
|
|
.order_by(ConversationSession.ended_at)
|
|
.with_for_update(skip_locked=True)
|
|
.limit(1)
|
|
)
|
|
).scalar_one_or_none()
|
|
if not row:
|
|
return None
|
|
row.analysis_status = "processing"
|
|
row.analysis_error = ""
|
|
await session.commit()
|
|
return row.id
|
|
|
|
|
|
async def run_analysis_worker() -> None:
|
|
logger.info("通话后分析 worker 已启动")
|
|
while True:
|
|
try:
|
|
conversation_id = await claim_pending_analysis()
|
|
if conversation_id:
|
|
await analyze_conversation(conversation_id)
|
|
continue
|
|
except asyncio.CancelledError:
|
|
raise
|
|
except Exception as exc:
|
|
logger.error(f"通话后分析 worker 暂时不可用:{exc}")
|
|
await asyncio.sleep(POLL_INTERVAL_SECONDS)
|