refactor(gateway): every Ogg/Opus voice transcode routes through base.transcode_to_ogg_opus

Matrix, WhatsApp Cloud and the TTS tool each ran their own ffmpeg argv for the same
speech-tuned libopus encode; they predate the shared helper and never migrated, so the
codec flags, timeout handling and error reporting drifted (matrix 48k/30s, whatsapp_cloud
async subprocess with no timeout and no `-ac 1`, tts an in-place sidecar repair).

`transcode_to_ogg_opus` gains `timeout=` and `output_path=` (sibling-file and in-place
writes go through a `.tmp.ogg` sidecar so a failed encode never truncates the source).
Deleted: matrix `_matrix_transcode_voice_to_ogg`, tts `_ffmpeg_transcode_to_opus`;
whatsapp_cloud `_convert_to_opus` keeps only its warn-once ffmpeg install hint and calls
the helper via `asyncio.to_thread`. Matrix and tts keep their 48k bitrate.

Behavior change: whatsapp_cloud transcodes now use mono (`-ac 1`), `-compression_level 10`
and a 60s timeout like every other voice bubble; a failed encode logs at WARNING for all
three sites (matrix previously DEBUG).
This commit is contained in:
teknium1
2026-09-13 05:32:38 -07:00
committed by Teknium
parent 8185834bc9
commit 3a179fe524
6 changed files with 111 additions and 86 deletions
+22 -10
View File
@@ -49,27 +49,39 @@ _TELEGRAM_AUDIO_ATTACHMENT_EXTS = frozenset({'.mp3', '.m4a'})
_TELEGRAM_VOICE_EXTS = frozenset({'.ogg', '.opus'})
def transcode_to_ogg_opus(path: str, *, bitrate: str = "32k") -> "str | None":
"""Best-effort ffmpeg transcode to Ogg/Opus (voip-tuned) for native voice bubbles: a NEW temp
``.ogg`` path (caller cleans up), or None when ffmpeg is missing/fails. Blocking (to_thread)."""
def transcode_to_ogg_opus(path: str, *, bitrate: str = "32k", timeout: int = 60,
output_path: "str | None" = None) -> "str | None":
"""Best-effort ffmpeg transcode to Ogg/Opus (voip-tuned) for native voice bubbles: the written
``.ogg`` path (a NEW temp file unless ``output_path`` is given; caller cleans up), or None when
ffmpeg is missing/fails. ``output_path`` may equal ``path`` (in-place container repair) — the
encode goes through a sidecar so a failed run never truncates the source. Blocking (to_thread)."""
import shutil as _shutil
ffmpeg = _shutil.which("ffmpeg")
if not ffmpeg:
return None
fd, ogg_path = tempfile.mkstemp(prefix="voice_transcode_", suffix=".ogg")
os.close(fd)
if output_path is None:
fd, ogg_path = tempfile.mkstemp(prefix="voice_transcode_", suffix=".ogg")
os.close(fd)
else:
ogg_path = output_path
in_place = os.path.abspath(str(path)) == os.path.abspath(ogg_path)
work_path = ogg_path + ".tmp.ogg" if in_place else ogg_path
try:
result = subprocess.run(
[ffmpeg, "-v", "error", "-y", "-i", str(path),
"-acodec", "libopus", "-ac", "1", "-b:a", bitrate, "-vbr", "on",
"-application", "voip", "-compression_level", "10", ogg_path],
capture_output=True, timeout=60, stdin=subprocess.DEVNULL)
if result.returncode == 0 and os.path.getsize(ogg_path) > 0:
"-application", "voip", "-compression_level", "10", "-f", "ogg", work_path],
capture_output=True, timeout=timeout, stdin=subprocess.DEVNULL)
if result.returncode == 0 and os.path.getsize(work_path) > 0:
if in_place:
os.replace(work_path, ogg_path)
return ogg_path
logger.warning("ffmpeg Ogg/Opus transcode of %s failed (returncode=%s): %s", path, result.returncode,
(result.stderr or b"").decode("utf-8", errors="replace")[:500])
except Exception:
logger.debug("voice transcode to Ogg/Opus failed for %s", path, exc_info=True)
logger.warning("voice transcode to Ogg/Opus failed for %s", path, exc_info=True)
with contextlib.suppress(OSError):
os.unlink(ogg_path)
os.unlink(work_path)
return None
_POST_DELIVERY_CALLBACK_TIMEOUT_SECONDS = 30.0
# History dedup is best-effort: stay well below the Discord heartbeat watchdog and fail open.
+4 -20
View File
@@ -39,7 +39,7 @@ except ImportError:
httpx = None # type: ignore[assignment]
from gateway.config import Platform, PlatformConfig
from gateway.platforms.base import BasePlatformAdapter, ExecApprovalPrompt, SendResult
from gateway.platforms.base import BasePlatformAdapter, ExecApprovalPrompt, SendResult, transcode_to_ogg_opus
from gateway.platforms.event import MessageEvent, MessageType
from gateway.platforms.whatsapp_common import WhatsAppBehaviorMixin, _get_wsecret
from gateway.platforms.access_policy_mixin import OPTIN_TRUTHY as _OPTIN_TRUTHY
@@ -604,8 +604,8 @@ class WhatsAppCloudAdapter(WhatsAppBehaviorMixin, BasePlatformAdapter):
return await self._send_media_from_path_or_link(chat_id, audio_path, "audio", caption=caption, reply_to=reply_to, mime_type=mime_type)
async def _convert_to_opus(self, mp3_path: str) -> Optional[str]:
"""MP3 → ``audio/ogg; codecs=opus``; None if ffmpeg is missing or fails. ``-application voip``
tunes for speech; ``-b:a 32k -vbr on`` matches WhatsApp's native voice-note bitrate."""
"""MP3 → ``audio/ogg; codecs=opus`` sibling file; None if ffmpeg is missing or fails. The
missing-ffmpeg warning fires once per adapter: it is an install hint, not a per-message error."""
if not _FFMPEG_PATH:
if not self._warned_no_ffmpeg:
self._warned_no_ffmpeg = True
@@ -615,23 +615,7 @@ class WhatsAppCloudAdapter(WhatsAppBehaviorMixin, BasePlatformAdapter):
"Windows `winget install Gyan.FFmpeg`, macOS `brew install ffmpeg`, Linux package manager."
)
return None
out_path = mp3_path.rsplit(".", 1)[0] + ".ogg"
try:
proc = await asyncio.create_subprocess_exec(
_FFMPEG_PATH, "-y", "-i", mp3_path, "-c:a", "libopus", "-b:a", "32k", "-vbr", "on", "-application", "voip",
out_path, stdout=asyncio.subprocess.DEVNULL, stderr=asyncio.subprocess.PIPE,
)
_, stderr = await proc.communicate()
except Exception:
logger.exception("[whatsapp_cloud] ffmpeg subprocess raised")
return None
if proc.returncode == 0 and Path(out_path).exists():
return out_path
logger.error(
"[whatsapp_cloud] ffmpeg opus conversion failed (returncode=%s): %s",
proc.returncode, (stderr or b"").decode("utf-8", errors="replace")[:500],
)
return None
return await asyncio.to_thread(transcode_to_ogg_opus, mp3_path, output_path=mp3_path.rsplit(".", 1)[0] + ".ogg")
# ------------------------------------------------------------------ inbound media
async def _graph_get(self, url: str, headers: Dict[str, str], what: str, media_id: str) -> Any:
+3 -24
View File
@@ -62,6 +62,7 @@ from gateway.platforms.base import (
gateway_trust_env, BasePlatformAdapter, ExecApprovalPrompt,
SendResult, resolve_proxy_url, proxy_kwargs_for_aiohttp, _ssrf_redirect_guard,
)
from gateway.platforms.base import transcode_to_ogg_opus
from gateway.platforms.event import MessageEvent, MessageType, ProcessingOutcome
from gateway.platforms.helpers import ThreadParticipationTracker
@@ -112,29 +113,6 @@ def _matrix_voice_metadata_for_file(path: Path) -> Dict[str, Any]:
logger.debug("Matrix: failed to build voice waveform for %s", path, exc_info=True)
return metadata
def _matrix_transcode_voice_to_ogg(path: str) -> Optional[str]:
"""Transcode to a NEW temp .ogg (caller owns cleanup); None if ffmpeg is missing/fails.
Blocking subprocess work — call via ``asyncio.to_thread`` from async code."""
ffmpeg = shutil.which("ffmpeg")
if not ffmpeg:
return None
import tempfile
fd, ogg_path = tempfile.mkstemp(prefix="matrix_voice_", suffix=".ogg")
os.close(fd)
try:
result = _run_media_tool(
[ffmpeg, "-v", "error", "-y", "-i", str(path), "-acodec", "libopus", "-ac", "1", "-b:a", "48k",
"-vbr", "on", "-application", "voip", "-compression_level", "10", ogg_path],
timeout=30)
if result.returncode == 0 and os.path.getsize(ogg_path) > 0:
return ogg_path
except Exception:
logger.debug("Matrix: voice transcode to Ogg/Opus failed for %s", path, exc_info=True)
with suppress(OSError):
os.unlink(ogg_path)
return None
_MATRIX_BANG_COMMAND_RE = re.compile(r"^!([A-Za-z][A-Za-z0-9_-]*)(?=$|\s)(.*)$", re.DOTALL)
@@ -1578,7 +1556,8 @@ class MatrixAdapter(BasePlatformAdapter):
format (e.g. TTS output), so transcode here — best-effort: without ffmpeg the original is sent."""
converted_path: Optional[str] = None
if not str(audio_path).lower().endswith((".ogg", ".oga", ".opus")):
converted_path = await asyncio.to_thread(_matrix_transcode_voice_to_ogg, audio_path)
# 48k (not the 32k default): Element renders voice bubbles at a higher quality tier.
converted_path = await asyncio.to_thread(transcode_to_ogg_opus, audio_path, bitrate="48k", timeout=30)
try:
return await self._send_local_file(
chat_id, converted_path or audio_path, "m.audio", caption, reply_to,
+2 -2
View File
@@ -236,7 +236,7 @@ class TestMatrixSendVoiceMSC3245:
self.adapter._client.send_message_event = mock_send_message_event
with patch(
"plugins.platforms.matrix.adapter._matrix_transcode_voice_to_ogg",
"plugins.platforms.matrix.adapter.transcode_to_ogg_opus",
return_value=converted_path,
) as mock_transcode, patch(
"plugins.platforms.matrix.adapter._matrix_voice_metadata_for_file",
@@ -248,7 +248,7 @@ class TestMatrixSendVoiceMSC3245:
caption="Test voice",
)
mock_transcode.assert_called_once_with(temp_path)
mock_transcode.assert_called_once_with(temp_path, bitrate="48k", timeout=30)
assert sent_content is not None, "No message was sent"
assert "org.matrix.msc3245.voice" in sent_content
assert sent_content["info"]["mimetype"] == "audio/ogg"
@@ -0,0 +1,76 @@
"""Every voice-bubble transcode routes through ``gateway.platforms.base.transcode_to_ogg_opus``.
Matrix, WhatsApp Cloud and the TTS tool each used to run their own ffmpeg argv; the shared
helper is the only place the codec flags, the timeout and the in-place-safe write live.
"""
from __future__ import annotations
import asyncio
import subprocess
from types import SimpleNamespace
def _capture_ffmpeg(monkeypatch):
calls = []
def fake_run(argv, **kwargs):
calls.append((argv, kwargs))
out = argv[-1]
with open(out, "wb") as fh:
fh.write(b"OggS")
return SimpleNamespace(returncode=0, stderr=b"")
monkeypatch.setattr("gateway.platforms.base.subprocess.run", fake_run)
monkeypatch.setattr("shutil.which", lambda name: "/usr/bin/ffmpeg")
return calls
def _bitrate(argv):
return argv[argv.index("-b:a") + 1]
def test_all_sites_share_one_ffmpeg_invocation(monkeypatch, tmp_path):
from gateway.platforms import whatsapp_cloud
from plugins.platforms.matrix import adapter as matrix
from tools import tts_tool_delivery as tts
calls = _capture_ffmpeg(monkeypatch)
src = tmp_path / "speech.mp3"
src.write_bytes(b"ID3")
# Matrix: temp output at 48k.
matrix_out = asyncio.run(asyncio.to_thread(matrix.transcode_to_ogg_opus, str(src), bitrate="48k", timeout=30))
assert matrix_out and matrix_out.endswith(".ogg") and _bitrate(calls[-1][0]) == "48k"
assert calls[-1][1]["timeout"] == 30
# WhatsApp Cloud: sibling .ogg at the 32k default, warn-once on missing ffmpeg untouched.
monkeypatch.setattr(whatsapp_cloud, "_FFMPEG_PATH", "/usr/bin/ffmpeg")
wa = object.__new__(whatsapp_cloud.WhatsAppCloudAdapter)
wa._warned_no_ffmpeg = False
wa_out = asyncio.run(wa._convert_to_opus(str(src)))
assert wa_out == str(tmp_path / "speech.ogg") and _bitrate(calls[-1][0]) == "32k"
# TTS tool: sibling .ogg at 48k, and in-place repair never writes straight over the source.
assert tts._convert_to_opus(str(src)) == str(tmp_path / "speech.ogg") and _bitrate(calls[-1][0]) == "48k"
bad = tmp_path / "bad.ogg"
bad.write_bytes(b"ID3")
assert tts._repair_ogg_container(str(bad)) == str(bad)
assert calls[-1][0][-1] == str(bad) + ".tmp.ogg" and bad.read_bytes() == b"OggS"
# Every argv carries the voice-tuned flag set exactly once.
for argv, kwargs in calls:
assert argv[argv.index("-application") + 1] == "voip" and "-compression_level" in argv
assert kwargs["stdin"] is subprocess.DEVNULL
def test_failed_in_place_repair_keeps_the_source(monkeypatch, tmp_path):
from gateway.platforms.base import transcode_to_ogg_opus
monkeypatch.setattr("shutil.which", lambda name: "/usr/bin/ffmpeg")
monkeypatch.setattr("gateway.platforms.base.subprocess.run",
lambda argv, **kw: SimpleNamespace(returncode=1, stderr=b"boom"))
bad = tmp_path / "bad.ogg"
bad.write_bytes(b"ID3")
assert transcode_to_ogg_opus(str(bad), output_path=str(bad)) is None
assert bad.read_bytes() == b"ID3" and not (tmp_path / "bad.ogg.tmp.ogg").exists()
+4 -30
View File
@@ -288,35 +288,8 @@ def _write_wav_bytes_as(wav_bytes: bytes, output_path: str) -> str:
def _convert_to_opus(mp3_path: str) -> Optional[str]:
"""Convert any ffmpeg-readable audio file to OGG Opus next to it; None on failure."""
return _ffmpeg_transcode_to_opus(mp3_path, mp3_path.rsplit(".", 1)[0] + ".ogg")
def _ffmpeg_transcode_to_opus(input_path: str, ogg_path: str) -> Optional[str]:
"""Transcode *input_path* to real Ogg/Opus at *ogg_path* (in-place safe via temp file); None on failure."""
if shutil.which("ffmpeg") is None:
return None
in_place = os.path.abspath(input_path) == os.path.abspath(ogg_path)
work_path = ogg_path + ".tmp.ogg" if in_place else ogg_path
try:
result = _ffmpeg_run("ffmpeg", ["-i", input_path, *_OPUS_VOICE_ARGS, "-f", "ogg", work_path, "-y"])
if result.returncode != 0:
logger.warning("ffmpeg conversion failed with return code %d: %s",
result.returncode, result.stderr.decode('utf-8', errors='ignore')[:200])
return None
if os.path.exists(work_path) and os.path.getsize(work_path) > 0:
if in_place:
os.replace(work_path, ogg_path)
return ogg_path
except subprocess.TimeoutExpired:
logger.warning("ffmpeg OGG conversion timed out after 30s")
except FileNotFoundError:
logger.warning("ffmpeg not found in PATH")
except Exception as e:
logger.warning("ffmpeg OGG conversion failed: %s", e, exc_info=True)
finally:
if in_place and os.path.exists(work_path):
_remove_quietly(work_path)
return None
from gateway.platforms.base import transcode_to_ogg_opus
return transcode_to_ogg_opus(mp3_path, bitrate="48k", timeout=30, output_path=mp3_path.rsplit(".", 1)[0] + ".ogg")
# --- Container sniffing / repair ---
@@ -341,7 +314,8 @@ def _repair_ogg_container(file_str: str) -> str:
if container in ("ogg", "unknown"):
return file_str
logger.info("TTS wrote %s bytes into a .ogg path (%s) — transcoding to real Ogg/Opus", container, file_str)
repaired = _ffmpeg_transcode_to_opus(file_str, file_str)
from gateway.platforms.base import transcode_to_ogg_opus
repaired = transcode_to_ogg_opus(file_str, bitrate="48k", timeout=30, output_path=file_str)
if repaired:
return repaired
honest = f"{file_str[:-4]}.{container}"