Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 11 additions & 0 deletions src/dlna_stream_server/handlers/constants.py
Original file line number Diff line number Diff line change
Expand Up @@ -128,6 +128,17 @@
"aac": "audio/aac",
}

# Диагностика темпа подачи: длина окна замера и пороги скорости
# по кодеку, ниже которых подача считается медленнее реального
# времени (кбит/с). Пороги — с запасом ниже типичных битрейтов:
# FLAC ~1000-1400, MP3 320, радио-AAC ~128
STREAM_RATE_WINDOW = 30.0
STREAM_RATE_THRESHOLDS_KBPS = {
"flac": 700,
"mp3": 250,
"aac": 90,
}

# MIME-типы для protocolInfo в DIDL-метаданных привязки Ruark.
# Ресурс, объявленный как audio/aac, Ruark качает, но не воспроизводит —
# ADTS-поток радио играет только под вывеской audio/mpeg
Expand Down
1 change: 1 addition & 0 deletions src/dlna_stream_server/handlers/ffmpeg_supervisor.py
Original file line number Diff line number Diff line change
Expand Up @@ -342,6 +342,7 @@ async def _log_stderr(self, proc: asyncio.subprocess.Process) -> None:
"error",
"failed",
"connection",
"reconnect",
"broken",
"timeout",
"invalid data found",
Expand Down
52 changes: 52 additions & 0 deletions src/dlna_stream_server/handlers/stream_handler.py
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,8 @@
FFMPEG_LOCAL_MP3_PARAMS,
FFMPEG_MP3_PARAMS,
STREAM_MIME_TYPES,
STREAM_RATE_THRESHOLDS_KBPS,
STREAM_RATE_WINDOW,
)
from .ffmpeg_supervisor import FfmpegSupervisor

Expand Down Expand Up @@ -157,6 +159,33 @@ async def start_ffmpeg_stream(

await self._ffmpeg.start(yandex_url, params, radio)

def _log_stream_pacing(
self,
window_bytes: int,
elapsed: float,
source_wait: float,
client_wait: float,
) -> None:
"""Логирует темп подачи за окно; ниже реального времени — WARNING.

Доли ожиданий называют узкое место: источник — CDN/FFmpeg,
клиент — сеть до устройства (Wi-Fi Ruark).
"""
if elapsed <= 0:
return
kbps = window_bytes * 8 / 1000 / elapsed
message = (
f"📈 Темп подачи: {kbps:.0f} кбит/с за {elapsed:.0f}с "
f"({self._stream_codec}); ожидание источника "
f"{source_wait / elapsed * 100:.0f}%, "
f"ожидание клиента {client_wait / elapsed * 100:.0f}%"
)
threshold = STREAM_RATE_THRESHOLDS_KBPS.get(self._stream_codec, 0)
if kbps < threshold:
logger.warning(f"⚠️ Подача ниже реального времени! {message}")
else:
logger.debug(message)

async def stream_audio(self, radio: bool = False) -> StreamingResponse:
"""Отдаёт потоковый аудио-ответ клиенту.

Expand All @@ -180,6 +209,11 @@ async def generate():
total_bytes_sent = 0
serve_started = time.monotonic()
first_chunk_sent = False
# Окно замера темпа подачи и виновника ожиданий
window_started = serve_started
window_bytes = 0
window_source_wait = 0.0
window_client_wait = 0.0
while True:
if self._client_epoch != client_epoch:
logger.info(
Expand All @@ -195,6 +229,7 @@ async def generate():
)
break

read_started = time.monotonic()
try:
chunk = await asyncio.wait_for(
stdout.read(4096),
Expand Down Expand Up @@ -235,6 +270,7 @@ async def generate():
f"{time.monotonic() - serve_started:.2f}с "
f"после подключения"
)
window_source_wait += time.monotonic() - read_started
total_bytes_sent += len(chunk)
# Диагностика: логируем прогресс передачи данных
if total_bytes_sent % (1024 * 1024) == 0: # Каждый МБ
Expand All @@ -243,7 +279,23 @@ async def generate():
f"{total_bytes_sent // 1024 // 1024} МБ"
)

send_started = time.monotonic()
yield chunk
now = time.monotonic()
window_client_wait += now - send_started
window_bytes += len(chunk)

if now - window_started >= STREAM_RATE_WINDOW:
self._log_stream_pacing(
window_bytes,
now - window_started,
window_source_wait,
window_client_wait,
)
window_started = now
window_bytes = 0
window_source_wait = 0.0
window_client_wait = 0.0

# После выхода из цикла логируем завершение FFmpeg
if proc.returncode is not None:
Expand Down
35 changes: 35 additions & 0 deletions tests/test_stream_handler.py
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
import asyncio
import logging
import time
from unittest.mock import AsyncMock

Expand Down Expand Up @@ -193,6 +194,40 @@ async def test_new_client_displaces_previous_stream():
await handler.stop_ffmpeg()


def test_slow_flac_pacing_logged_as_warning(caplog):
"""FLAC на ~320 кбит/с — втрое ниже реального времени, WARNING."""
handler, _ = make_handler()
handler._stream_codec = "flac"

with caplog.at_level(logging.WARNING):
handler._log_stream_pacing(
window_bytes=1_200_000,
elapsed=30.0,
source_wait=25.0,
client_wait=1.0,
)

assert any(
"ниже реального времени" in record.message for record in caplog.records
)


def test_healthy_mp3_pacing_stays_quiet(caplog):
"""Тот же темп ~320 кбит/с для MP3 — норма, без WARNING."""
handler, _ = make_handler()
handler._stream_codec = "mp3"

with caplog.at_level(logging.WARNING):
handler._log_stream_pacing(
window_bytes=1_200_000,
elapsed=30.0,
source_wait=0.5,
client_wait=28.0,
)

assert not caplog.records


async def test_radio_declared_as_mpeg_in_didl(fast_sleep):
"""Радио: Content-Type остаётся aac, но в DIDL — audio/mpeg.

Expand Down
Loading