Keep AI workers running after failed batch
This commit is contained in:
@@ -39,7 +39,7 @@ def parse_recipients(value: Any) -> list[int]:
|
||||
return recipients
|
||||
|
||||
|
||||
async def send_ai_worker_disabled_alert(worker_name: str, model: str, post_ids: list[int], error: str) -> None:
|
||||
async def send_ai_worker_error_alert(worker_name: str, model: str, post_ids: list[int], error: str) -> None:
|
||||
token = (
|
||||
str(await fetch_setting("daily_report_bot_token", "") or "").strip()
|
||||
or str(await fetch_setting("tg_poster_bot_token", "") or "").strip()
|
||||
@@ -52,7 +52,7 @@ async def send_ai_worker_disabled_alert(worker_name: str, model: str, post_ids:
|
||||
|
||||
post_part = ", ".join(str(post_id) for post_id in post_ids) if post_ids else "-"
|
||||
text = (
|
||||
"AI worker auto-disabled\n"
|
||||
"AI worker batch failed\n"
|
||||
f"worker: {worker_name}\n"
|
||||
f"model: {model or '-'}\n"
|
||||
f"posts: {post_part}\n"
|
||||
|
||||
@@ -15,7 +15,7 @@ from ..constants import WORKER_AI_QUALIFIER
|
||||
from ..db import fetch_bool_setting, fetch_float_setting, fetch_int_setting, fetch_setting, get_pool
|
||||
from ..heartbeat import HeartbeatReporter
|
||||
from ..jobs import is_worker_enabled
|
||||
from .ai_alerts import send_ai_worker_disabled_alert
|
||||
from .ai_alerts import send_ai_worker_error_alert
|
||||
|
||||
|
||||
def now_utc() -> datetime:
|
||||
@@ -292,19 +292,11 @@ class AIQualifierWorker:
|
||||
error[:1000],
|
||||
)
|
||||
|
||||
async def disable_after_error(self, model: str, post_ids: list[int], error: str) -> None:
|
||||
await self.pool.execute(
|
||||
"""
|
||||
UPDATE app_settings
|
||||
SET value_json='false'::jsonb,
|
||||
updated_at=NOW()
|
||||
WHERE key='ai_qualifier_enabled'
|
||||
"""
|
||||
)
|
||||
async def report_failed_batch(self, model: str, post_ids: list[int], error: str) -> None:
|
||||
meta = {"error": error[:300], "post_ids": post_ids, "model": model}
|
||||
await self.heartbeat.beat(self.pool, status="disabled_after_error", meta=meta, force=True)
|
||||
await self.heartbeat.beat(self.pool, status="error", meta=meta, force=True)
|
||||
try:
|
||||
await send_ai_worker_disabled_alert(WORKER_AI_QUALIFIER, model, post_ids, error)
|
||||
await send_ai_worker_error_alert(WORKER_AI_QUALIFIER, model, post_ids, error)
|
||||
except Exception as alert_exc:
|
||||
logger.warning("AI qualifier alert failed: {}", alert_exc)
|
||||
|
||||
@@ -429,7 +421,7 @@ class AIQualifierWorker:
|
||||
except Exception as exc:
|
||||
await self.mark_posts_failed(post_ids, str(exc))
|
||||
await self.finish_batch(batch_id, "failed", error=str(exc))
|
||||
await self.disable_after_error(normalize_model(provider, model), post_ids, str(exc))
|
||||
await self.report_failed_batch(normalize_model(provider, model), post_ids, str(exc))
|
||||
logger.exception("AI qualifier batch failed: id={} error={}", batch_id, exc)
|
||||
return True
|
||||
|
||||
|
||||
@@ -17,7 +17,7 @@ from ..db import fetch_bool_setting, fetch_float_setting, fetch_int_setting, fet
|
||||
from ..heartbeat import HeartbeatReporter
|
||||
from ..jobs import is_worker_enabled
|
||||
from ..text_utils import build_publication_text, normalize_hash_tag, parse_categories
|
||||
from .ai_alerts import send_ai_worker_disabled_alert
|
||||
from .ai_alerts import send_ai_worker_error_alert
|
||||
|
||||
|
||||
NON_TARGET_CATEGORY_ID = 18
|
||||
@@ -373,19 +373,11 @@ class AIWriterWorker:
|
||||
error[:1000],
|
||||
)
|
||||
|
||||
async def disable_after_error(self, model: str, post_ids: list[int], error: str) -> None:
|
||||
await self.pool.execute(
|
||||
"""
|
||||
UPDATE app_settings
|
||||
SET value_json='false'::jsonb,
|
||||
updated_at=NOW()
|
||||
WHERE key='ai_writer_enabled'
|
||||
"""
|
||||
)
|
||||
async def report_failed_batch(self, model: str, post_ids: list[int], error: str) -> None:
|
||||
meta = {"error": error[:300], "post_ids": post_ids, "model": model}
|
||||
await self.heartbeat.beat(self.pool, status="disabled_after_error", meta=meta, force=True)
|
||||
await self.heartbeat.beat(self.pool, status="error", meta=meta, force=True)
|
||||
try:
|
||||
await send_ai_worker_disabled_alert(WORKER_AI_WRITER, model, post_ids, error)
|
||||
await send_ai_worker_error_alert(WORKER_AI_WRITER, model, post_ids, error)
|
||||
except Exception as alert_exc:
|
||||
logger.warning("AI writer alert failed: {}", alert_exc)
|
||||
|
||||
@@ -535,7 +527,7 @@ class AIWriterWorker:
|
||||
except Exception as exc:
|
||||
await self.mark_posts_failed(post_ids, str(exc))
|
||||
await self.finish_batch(batch_id, "failed", response=response, error=str(exc), usage=usage)
|
||||
await self.disable_after_error(normalize_model(provider, model), post_ids, str(exc))
|
||||
await self.report_failed_batch(normalize_model(provider, model), post_ids, str(exc))
|
||||
logger.exception("AI writer batch failed: id={} error={}", batch_id, exc)
|
||||
return True
|
||||
|
||||
|
||||
Reference in New Issue
Block a user