diff --git a/app/config.py b/app/config.py index b24464c..b9d6e1a 100644 --- a/app/config.py +++ b/app/config.py @@ -51,6 +51,14 @@ max_upload_mb = 500 # Сколько часов хранить результаты завершённых задач keep_results_hours = 72 +[webhook] +# Куда сообщать о готовых задачах. Пусто = не сообщать, забирайте опросом. +# Можно переопределить для отдельной задачи параметром webhook в запросе. +url = "" +# Секрет для подписи: сервис положит её в заголовок X-Talkscore-Signature, +# чтобы принимающая сторона убедилась, что запрос от вас. +secret = "" + [update] # Проверять обновления кода при каждом запуске. Обновляется только папка app, # это десятки килобайт: Python, библиотеки и модели остаются на месте. @@ -74,6 +82,8 @@ class Settings: speakers: int = 2 max_upload_mb: int = 500 keep_results_hours: float = 72.0 + webhook_url: str = "" + webhook_secret: str = "" update_enabled: bool = False update_server: str = "https://git.netranking.ru" update_repo: str = "bryzgalov/talkscore-asr" @@ -129,6 +139,7 @@ def load_settings(config_path: Path | None = None) -> Settings: security = data.get("security", {}) proc = data.get("processing", {}) upd = data.get("update", {}) + hook = data.get("webhook", {}) return Settings( host=server.get("host", "0.0.0.0"), @@ -141,6 +152,8 @@ def load_settings(config_path: Path | None = None) -> Settings: speakers=int(proc.get("speakers", 2)), max_upload_mb=int(proc.get("max_upload_mb", 500)), keep_results_hours=float(proc.get("keep_results_hours", 72)), + webhook_url=str(hook.get("url", "")), + webhook_secret=str(hook.get("secret", "")), update_enabled=bool(upd.get("enabled", False)), update_server=str(upd.get("server", "https://git.netranking.ru")), update_repo=str(upd.get("repo", "bryzgalov/talkscore-asr")), diff --git a/app/main.py b/app/main.py index b2f8dcf..51a9835 100644 --- a/app/main.py +++ b/app/main.py @@ -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: diff --git a/app/version.py b/app/version.py index 22049ab..49e0fc1 100644 --- a/app/version.py +++ b/app/version.py @@ -1 +1 @@ -__version__ = "0.6.2" +__version__ = "0.7.0" diff --git a/app/webhook.py b/app/webhook.py new file mode 100644 index 0000000..d49a072 --- /dev/null +++ b/app/webhook.py @@ -0,0 +1,59 @@ +"""Отправка готового результата на сторонний адрес. + +Опрос статуса работает, но заставляет принимающую сторону дёргать сервис +каждые несколько секунд. Вебхук снимает это: сервис сам постучится, когда +задача готова. +""" +import hashlib +import hmac +import json +import logging +import threading +import time + +__all__ = ["sign_payload", "deliver", "deliver_async"] + +log = logging.getLogger(__name__) + +# Задержки между попытками: сразу, через полминуты, через пять минут. +# Дольше ждать нет смысла - результат лежит в очереди и его можно забрать. +RETRY_DELAYS = (0, 30, 300) +TIMEOUT_SEC = 30 + + +def sign_payload(body: bytes, secret: str) -> str: + """Подпись тела запроса, чтобы принимающая сторона знала, что это мы.""" + return hmac.new(secret.encode("utf-8"), body, hashlib.sha256).hexdigest() + + +def deliver(url: str, payload: dict, secret: str = "", + delays: tuple = RETRY_DELAYS) -> bool: + """Отправляет результат, повторяя при неудаче. Возвращает признак успеха.""" + import requests + + body = json.dumps(payload, ensure_ascii=False).encode("utf-8") + headers = {"Content-Type": "application/json; charset=utf-8"} + if secret: + headers["X-Talkscore-Signature"] = sign_payload(body, secret) + + for attempt, delay in enumerate(delays, start=1): + if delay: + time.sleep(delay) + try: + response = requests.post(url, data=body, headers=headers, timeout=TIMEOUT_SEC) + if response.status_code < 300: + log.info("вебхук доставлен по задаче %s", payload.get("job_id")) + return True + log.warning("вебхук: попытка %d, ответ %s", attempt, response.status_code) + except Exception as exc: # noqa: BLE001 - причина неважна, важна повторная попытка + log.warning("вебхук: попытка %d не удалась (%s)", attempt, exc) + log.error("вебхук не доставлен по задаче %s, результат остаётся в очереди", + payload.get("job_id")) + return False + + +def deliver_async(url: str, payload: dict, secret: str = "") -> None: + """Отправляет в фоне: воркер не должен ждать чужой сервер.""" + thread = threading.Thread(target=deliver, args=(url, payload, secret), + name="webhook", daemon=True) + thread.start() diff --git a/original/original_02fef7a0-50a9-4af9-8e83-a88faf3c196f.mp3 b/original/original_02fef7a0-50a9-4af9-8e83-a88faf3c196f.mp3 new file mode 100644 index 0000000..ba21adc Binary files /dev/null and b/original/original_02fef7a0-50a9-4af9-8e83-a88faf3c196f.mp3 differ diff --git a/original/original_470e4858-1ef1-4ac3-b96c-0505ada2e0a6.mp3 b/original/original_470e4858-1ef1-4ac3-b96c-0505ada2e0a6.mp3 new file mode 100644 index 0000000..732fec5 Binary files /dev/null and b/original/original_470e4858-1ef1-4ac3-b96c-0505ada2e0a6.mp3 differ diff --git a/original/original_6886a557-cca2-4581-b4c4-214e29703f2c.mp3 b/original/original_6886a557-cca2-4581-b4c4-214e29703f2c.mp3 new file mode 100644 index 0000000..6eb1e92 Binary files /dev/null and b/original/original_6886a557-cca2-4581-b4c4-214e29703f2c.mp3 differ diff --git a/original/original_7444931c-6ce0-4f0a-a9cf-2b64905ecd02.mp3 b/original/original_7444931c-6ce0-4f0a-a9cf-2b64905ecd02.mp3 new file mode 100644 index 0000000..40b0610 Binary files /dev/null and b/original/original_7444931c-6ce0-4f0a-a9cf-2b64905ecd02.mp3 differ diff --git a/original/original_8dc00a27-5fe0-4d7e-8fc8-ef40820fcc36.mp3 b/original/original_8dc00a27-5fe0-4d7e-8fc8-ef40820fcc36.mp3 new file mode 100644 index 0000000..2483974 Binary files /dev/null and b/original/original_8dc00a27-5fe0-4d7e-8fc8-ef40820fcc36.mp3 differ diff --git a/original/original_9272254c-8a84-44a7-9c98-4ba64d7bd32e.mp3 b/original/original_9272254c-8a84-44a7-9c98-4ba64d7bd32e.mp3 new file mode 100644 index 0000000..19a1e7f Binary files /dev/null and b/original/original_9272254c-8a84-44a7-9c98-4ba64d7bd32e.mp3 differ diff --git a/original/original_ac20d666-0b46-44cd-b591-787b8e41e6af.mp3 b/original/original_ac20d666-0b46-44cd-b591-787b8e41e6af.mp3 new file mode 100644 index 0000000..32c96fc Binary files /dev/null and b/original/original_ac20d666-0b46-44cd-b591-787b8e41e6af.mp3 differ diff --git a/original/original_e287ae8b-3dc6-4b12-84c2-0d96387ef97d.mp3 b/original/original_e287ae8b-3dc6-4b12-84c2-0d96387ef97d.mp3 new file mode 100644 index 0000000..e499b7c Binary files /dev/null and b/original/original_e287ae8b-3dc6-4b12-84c2-0d96387ef97d.mp3 differ diff --git a/tests/test_webhook.py b/tests/test_webhook.py new file mode 100644 index 0000000..633491a --- /dev/null +++ b/tests/test_webhook.py @@ -0,0 +1,97 @@ +"""Тесты доставки результата на сторонний адрес.""" +import hashlib +import hmac +import json + +import pytest + +from app.webhook import deliver, sign_payload + + +class FakeResponse: + def __init__(self, status): + self.status_code = status + + +class TestSignature: + def test_signature_matches_hmac(self): + body = b'{"job_id":"x"}' + expected = hmac.new(b"secret", body, hashlib.sha256).hexdigest() + assert sign_payload(body, "secret") == expected + + def test_different_secret_gives_different_signature(self): + body = b'{"job_id":"x"}' + assert sign_payload(body, "a") != sign_payload(body, "b") + + def test_signature_changes_with_body(self): + assert sign_payload(b"one", "s") != sign_payload(b"two", "s") + + +class TestDelivery: + def test_successful_delivery_sends_once(self, monkeypatch): + calls = [] + + def fake_post(url, data=None, headers=None, timeout=None): + calls.append((url, data, headers)) + return FakeResponse(200) + + monkeypatch.setattr("requests.post", fake_post) + assert deliver("http://x/hook", {"job_id": "1"}, "s", delays=(0,)) is True + assert len(calls) == 1 + + def test_signature_header_present(self, monkeypatch): + seen = {} + + def fake_post(url, data=None, headers=None, timeout=None): + seen.update(headers) + return FakeResponse(200) + + monkeypatch.setattr("requests.post", fake_post) + deliver("http://x", {"job_id": "1"}, "секрет", delays=(0,)) + assert "X-Talkscore-Signature" in seen + + def test_no_signature_without_secret(self, monkeypatch): + seen = {} + monkeypatch.setattr("requests.post", + lambda url, data=None, headers=None, timeout=None: + (seen.update(headers), FakeResponse(200))[1]) + deliver("http://x", {"job_id": "1"}, "", delays=(0,)) + assert "X-Talkscore-Signature" not in seen + + def test_retries_on_server_error(self, monkeypatch): + attempts = [] + monkeypatch.setattr("requests.post", + lambda url, data=None, headers=None, timeout=None: + (attempts.append(1), FakeResponse(500))[1]) + assert deliver("http://x", {"job_id": "1"}, "", delays=(0, 0, 0)) is False + assert len(attempts) == 3 + + def test_retries_on_network_failure(self, monkeypatch): + attempts = [] + + def boom(url, data=None, headers=None, timeout=None): + attempts.append(1) + raise OSError("сеть недоступна") + + monkeypatch.setattr("requests.post", boom) + assert deliver("http://x", {"job_id": "1"}, "", delays=(0, 0)) is False + assert len(attempts) == 2 + + def test_stops_after_first_success(self, monkeypatch): + attempts = [] + + def flaky(url, data=None, headers=None, timeout=None): + attempts.append(1) + return FakeResponse(500 if len(attempts) == 1 else 200) + + monkeypatch.setattr("requests.post", flaky) + assert deliver("http://x", {"job_id": "1"}, "", delays=(0, 0, 0)) is True + assert len(attempts) == 2 + + def test_payload_is_valid_json_utf8(self, monkeypatch): + seen = {} + monkeypatch.setattr("requests.post", + lambda url, data=None, headers=None, timeout=None: + (seen.update({"body": data}), FakeResponse(200))[1]) + deliver("http://x", {"text": "русский текст"}, "", delays=(0,)) + assert json.loads(seen["body"].decode("utf-8"))["text"] == "русский текст"