Вебхук: сервис сам сообщает о готовой задаче
Опрос статуса заставлял принимающую сторону дёргать сервис каждые несколько секунд. Теперь при завершении задачи результат уходит POST-ом на заданный адрес, с подписью HMAC-SHA256 в заголовке и тремя попытками при неудаче. Адрес задаётся в настройках или параметром запроса. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 5
parent
710fc8b90b
commit
e6e10dddd1
+17
-2
@@ -21,6 +21,7 @@ from app.logbuffer import install as install_log_buffer
|
||||
from app.pipeline import ModelsMissing, Pipeline, to_wav16k
|
||||
from app.security import check_token, ip_allowed, parse_allowlist
|
||||
from app.store import JobStatus, JobStore
|
||||
from app.webhook import deliver_async
|
||||
from app.version import __version__
|
||||
|
||||
log = logging.getLogger("talkscore-asr")
|
||||
@@ -72,13 +73,24 @@ def _process(job_id: str, worker: Pipeline, pool) -> None:
|
||||
store.mark_done(job_id, result)
|
||||
log.info("задача %s готова: %.1f с аудио, x%s", job_id,
|
||||
result["duration_sec"], result["timing"]["realtime_factor"])
|
||||
_notify(job, {"job_id": job_id, "status": JobStatus.DONE, **result})
|
||||
except Exception as exc: # noqa: BLE001 - в статус задачи должна попасть любая причина
|
||||
log.exception("задача %s провалилась", job_id)
|
||||
store.mark_failed(job_id, error=f"{type(exc).__name__}: {exc}")
|
||||
error = f"{type(exc).__name__}: {exc}"
|
||||
store.mark_failed(job_id, error=error)
|
||||
_notify(job, {"job_id": job_id, "status": JobStatus.FAILED,
|
||||
"filename": job["filename"], "error": error})
|
||||
finally:
|
||||
upload.unlink(missing_ok=True)
|
||||
|
||||
|
||||
def _notify(job: dict, payload: dict) -> None:
|
||||
"""Сообщает о результате, если для задачи задан адрес."""
|
||||
url = job["options"].get("webhook") or settings.webhook_url
|
||||
if url:
|
||||
deliver_async(url, payload, settings.webhook_secret)
|
||||
|
||||
|
||||
def _worker_loop(index: int, worker: Pipeline, pool=None) -> None:
|
||||
"""Разбирает очередь. Задача захватывается атомарно, поэтому воркеров может быть много."""
|
||||
last_cleanup = 0.0
|
||||
@@ -263,13 +275,16 @@ async def create_job(
|
||||
file: UploadFile = File(...),
|
||||
speakers: int | None = Query(None, ge=0, le=10,
|
||||
description="число говорящих, 0 = определить автоматически"),
|
||||
webhook: str | None = Query(None,
|
||||
description="куда сообщить о готовности; заменяет адрес из настроек"),
|
||||
) -> dict:
|
||||
if not _state["ready"]:
|
||||
raise HTTPException(status_code=503, detail=_state["error"] or "сервис ещё не готов")
|
||||
|
||||
job_id = store.create(filename=file.filename or "audio",
|
||||
duration_sec=0.0,
|
||||
options={"speakers": settings.speakers if speakers is None else speakers})
|
||||
options={"speakers": settings.speakers if speakers is None else speakers,
|
||||
"webhook": webhook or ""})
|
||||
target = settings.data_dir / "uploads" / job_id
|
||||
size = 0
|
||||
try:
|
||||
|
||||
Reference in New Issue
Block a user