fix: fail open on watermark dispatch deadlines

This commit is contained in:
debpalash
2026-08-20 10:32:02 +05:30
parent b7f14ce4ad
commit aa7c2f5801
2 changed files with 42 additions and 1 deletions
+11 -1
View File
@@ -298,7 +298,12 @@ async def mark_synthetic_async(
import asyncio
import functools
from services.model_manager import get_watermark_pool, run_on_gpu_pool_guarded
from services.model_manager import (
GpuJobTimeoutError,
GpuPoolBusyError,
get_watermark_pool,
run_on_gpu_pool_guarded,
)
try:
pool = get_watermark_pool()
@@ -315,6 +320,11 @@ async def mark_synthetic_async(
job, what="Audio watermark", timeout=timeout, executor=pool
)
return await asyncio.get_running_loop().run_in_executor(pool, job)
except (GpuJobTimeoutError, GpuPoolBusyError):
# Watermarking is provenance best-effort: a typed execution overrun or
# queue saturation must not discard synthesis that already completed.
logger.warning("Watermark skipped after its bounded dispatch expired")
return waveform
except asyncio.CancelledError:
# A queued future is cancelled during pool teardown. Caller-driven
# cancellation while the pool is live must retain normal semantics.
@@ -314,6 +314,37 @@ def test_watermark_preserves_caller_cancellation(monkeypatch, watermark):
asyncio.run(_cancel())
@pytest.mark.parametrize("error_name", ["GpuJobTimeoutError", "GpuPoolBusyError"])
def test_timed_watermark_deadline_returns_finished_audio(
monkeypatch, watermark, error_name
):
"""Both guarded deadline phases are watermark-only fail-open outcomes."""
import asyncio
import torch
model_manager = importlib.import_module("services.model_manager")
error_type = getattr(model_manager, error_name)
async def _expired(*_args, **_kwargs):
raise error_type("watermark deadline expired")
pool = model_manager.get_watermark_pool()
monkeypatch.setattr(model_manager, "get_watermark_pool", lambda: pool)
monkeypatch.setattr(model_manager, "run_on_gpu_pool_guarded", _expired)
audio = torch.zeros(1, 240)
marked = asyncio.run(
watermark.mark_synthetic_async(
audio,
24000,
context="test.watermark_typed_deadline",
timeout=0.01,
)
)
assert marked is audio
def test_watermark_pool_shutdown_waits_for_active_worker():
"""Lifespan teardown cannot finish while AudioSeal is still loading."""
from services.model_manager import (