new version

This commit is contained in:
UrloMythus
2026-04-15 19:23:14 +02:00
parent 5120b19d0b
commit 8134936d59
135 changed files with 3013 additions and 1589 deletions
+236 -38
View File
@@ -1,18 +1,23 @@
import asyncio
import logging
import math
import statistics
import time
from aiohttp import ClientTimeout
from fastapi import Request, Response, HTTPException
from mediaflow_proxy.drm.decrypter import decrypt_segment, process_drm_init_segment
from mediaflow_proxy.utils.crypto_utils import encryption_handler
from mediaflow_proxy.utils.http_client import create_aiohttp_session
from mediaflow_proxy.utils.http_utils import (
encode_mediaflow_proxy_url,
fetch_with_retry,
get_original_scheme,
ProxyRequestHeaders,
apply_header_manipulation,
)
from mediaflow_proxy.utils.mpd_utils import parse_sidx_fragments
from mediaflow_proxy.utils.dash_prebuffer import dash_prebuffer
from mediaflow_proxy.utils.cache_utils import get_cached_processed_init, set_cached_processed_init
from mediaflow_proxy.utils.m3u8_processor import SkipSegmentFilter
@@ -30,6 +35,83 @@ def _resolve_ts_mode(request: Request) -> bool:
return settings.remux_to_ts
def _resolve_nominal_duration_mpd_timescale(profile: dict, segments: list[dict]) -> int | None:
"""Resolve a stable nominal segment duration (MPD timescale units) for live sequence math."""
profile_duration = profile.get("nominal_duration_mpd_timescale")
if isinstance(profile_duration, (int, float)) and profile_duration > 0:
return int(profile_duration)
durations = []
for seg in segments:
seg_duration = seg.get("duration_mpd_timescale")
if isinstance(seg_duration, (int, float)) and seg_duration > 0:
durations.append(int(seg_duration))
if durations:
# Use median to avoid jumps when first/last live segments are shorter.
return int(statistics.median_low(durations))
return None
def _compute_live_media_sequence(first_segment: dict, profile: dict, segments: list[dict]) -> int:
"""
Compute a stable HLS media sequence for live playlists.
Strategy:
1) If MPD explicitly sets @startNumber, trust segment numbering.
2) Otherwise derive sequence from timeline time / nominal duration.
3) Fall back to segment number or template start number.
"""
segment_number = first_segment.get("number")
if profile.get("segment_template_start_number_explicit") and segment_number is not None:
return max(int(segment_number), 1)
timeline_time = first_segment.get("time")
nominal_duration = _resolve_nominal_duration_mpd_timescale(profile, segments)
if timeline_time is not None and nominal_duration and nominal_duration > 0:
return max(math.floor(int(timeline_time) / nominal_duration), 1)
if segment_number is not None:
return max(int(segment_number), 1)
template_start = profile.get("segment_template_start_number")
if isinstance(template_start, int) and template_start > 0:
return template_start
return 1
def _compute_live_playlist_depth(
is_ts_mode: bool,
effective_start_offset: float | None,
extinf_values: list[float],
) -> int:
"""
Compute a resilient live playlist depth to reduce segment expiry skips.
We keep a larger floor for fMP4 live (direct mode), and further expand
depth based on requested start_offset so players have enough headroom
during transient stalls.
"""
configured_depth = max(settings.mpd_live_playlist_depth, 1)
depth_floor = 20 if is_ts_mode else 15
depth = max(configured_depth, depth_floor)
if effective_start_offset is not None and effective_start_offset < 0:
if extinf_values:
segment_duration = statistics.median(extinf_values)
if segment_duration <= 0:
segment_duration = 4.0
else:
segment_duration = 4.0
segments_behind_live_edge = math.ceil(abs(effective_start_offset) / segment_duration)
safety_margin = 10 if is_ts_mode else 12
depth = max(depth, segments_behind_live_edge + safety_margin)
return max(depth, 1)
async def process_manifest(
request: Request,
mpd_dict: dict,
@@ -73,6 +155,68 @@ async def process_manifest(
return Response(content=hls_content, media_type="application/vnd.apple.mpegurl", headers=proxy_headers.response)
async def _expand_segment_base_segments(profiles: list, req_headers: dict) -> None:
"""
For SegmentBase profiles with a single segment that has an @indexRange, fetch the
SIDX box and expand the single segment entry into per-fragment segments.
This allows mpv/ffmpeg to seek by requesting only the relevant fragment from the
CDN instead of re-downloading the whole file from byte 938 every time.
Modifies *profiles* in-place; on any failure the original single-segment entry
is kept as a fallback.
"""
timeout = ClientTimeout(connect=10, sock_read=30, total=60)
for profile in profiles:
segments = profile.get("segments", [])
if len(segments) != 1:
continue
seg = segments[0]
index_range: str = seg.get("indexRange") or ""
media_url: str = seg.get("media") or ""
init_range: str = seg.get("initRange") or ""
if not index_range or not media_url:
continue # no SIDX info → keep single-segment fallback
try:
headers = dict(req_headers)
headers["Range"] = f"bytes={index_range}"
async with create_aiohttp_session(media_url, timeout=timeout) as (session, proxy_url):
response = await fetch_with_retry(session, "GET", media_url, headers, proxy=proxy_url)
sidx_bytes = await response.read()
index_range_start = int(index_range.split("-")[0])
fragments = parse_sidx_fragments(sidx_bytes, index_range_start)
if not fragments:
logger.warning(f"SIDX parse returned no fragments for {media_url}, keeping single-segment fallback")
continue
new_segments = [
{
"type": "segment",
"media": media_url,
"number": i + 1,
"extinf": frag["duration_timescale"] / frag["timescale"],
"initRange": init_range,
"mediaRange": f"{frag['start']}-{frag['end']}",
}
for i, frag in enumerate(fragments)
]
profile["segments"] = new_segments
logger.info(
f"SegmentBase SIDX expanded {media_url!r}{len(new_segments)} fragments "
f"({new_segments[0]['extinf']:.3f}s each)"
)
except Exception as exc:
logger.warning(f"SIDX expansion failed for {media_url}: {exc}; keeping single-segment fallback")
async def process_playlist(
request: Request,
mpd_dict: dict,
@@ -102,6 +246,10 @@ async def process_playlist(
if not matching_profiles:
raise HTTPException(status_code=404, detail="Profile not found")
# For SegmentBase profiles (single large file with SIDX), expand into per-fragment
# segments so that mpv/ffmpeg can seek by requesting only the relevant fragment.
await _expand_segment_base_segments(matching_profiles, dict(proxy_headers.request))
hls_content = build_hls_playlist(mpd_dict, matching_profiles, request, skip_segments, start_offset)
# Trigger prebuffering of upcoming segments for live streams
@@ -301,22 +449,79 @@ def build_hls(
max_height = max(p[0].get("height", 0) for p in video_profiles.values())
video_profiles = {k: v for k, v in video_profiles.items() if v[0].get("height", 0) >= max_height}
# Add audio streams
for i, (profile, playlist_url) in enumerate(audio_profiles.values()):
is_default = "YES" if i == 0 else "NO" # Set the first audio track as default
lang = profile.get("lang", "und")
bandwidth = profile.get("bandwidth", "128000")
name = f"Audio {lang} ({bandwidth})" if lang != "und" else f"Audio {i + 1} ({bandwidth})"
hls.append(
f'#EXT-X-MEDIA:TYPE=AUDIO,GROUP-ID="audio",NAME="{name}",DEFAULT={is_default},AUTOSELECT=YES,LANGUAGE="{lang}",URI="{playlist_url}"'
if not is_ts_mode and video_profiles:
# Sort by bandwidth descending; keep highest bandwidth per unique rep_id
sorted_by_bw = sorted(video_profiles.values(), key=lambda pv: pv[0].get("bandwidth", 0), reverse=True)
seen_rep_ids = set()
deduped = []
for profile, playlist_url in sorted_by_bw:
rep_id = profile.get("rep_id", profile["id"])
if rep_id not in seen_rep_ids:
seen_rep_ids.add(rep_id)
deduped.append((profile, playlist_url))
if resolution:
# Explicit resolution: single matching variant
deduped = deduped[:1]
else:
# Limit ABR ladder: one entry per unique height, capped at MAX_VIDEO_VARIANTS.
# libavformat (mpv/ffmpeg) fetches ALL #EXT-X-STREAM-INF playlists and their
# init segments before playback regardless of BANDWIDTH hints —
# N variants = N × ~2 s of init-segment probing at startup.
MAX_VIDEO_VARIANTS = 5
seen_heights: set = set()
height_deduped = []
for p, url in deduped:
h = p.get("height", 0)
if h not in seen_heights:
seen_heights.add(h)
height_deduped.append((p, url))
deduped = height_deduped[:MAX_VIDEO_VARIANTS]
video_profiles = {p["id"]: (p, url) for p, url in deduped}
# Determine the default audio (English preferred, else highest bandwidth).
default_audio_id = None
if audio_profiles:
all_audio = list(audio_profiles.values())
en_audio = [(p, u) for p, u in all_audio if (p.get("lang") or "").startswith("en")]
default_profile, _ = max(en_audio or all_audio, key=lambda pu: pu[0].get("bandwidth", 0))
default_audio_id = default_profile["id"]
# Audio tracks: one entry per unique language, capped at MAX_AUDIO_TRACKS.
# Sort default track first (always within cap), then by bandwidth descending
# so the highest-quality codec per language wins.
# libavformat probes every #EXT-X-MEDIA entry regardless of DEFAULT/AUTOSELECT;
# capping at 4 keeps language selection while bounding startup probing.
MAX_AUDIO_TRACKS = 4
first_audio_codec = None
if audio_profiles:
sorted_audio = sorted(
audio_profiles.values(),
key=lambda pv: (pv[0]["id"] != default_audio_id, -pv[0].get("bandwidth", 0)),
)
seen_langs = set()
count = 0
for profile, playlist_url in sorted_audio:
if count >= MAX_AUDIO_TRACKS:
break
lang = profile.get("lang", "und")
if lang in seen_langs:
continue
seen_langs.add(lang)
is_default = profile["id"] == default_audio_id
default_attr = "YES" if is_default else "NO"
autoselect_attr = "YES" if is_default else "NO"
bandwidth = profile.get("bandwidth", 128000)
name = f"Audio {lang} ({bandwidth})" if lang != "und" else f"Audio 1 ({bandwidth})"
hls.append(
f'#EXT-X-MEDIA:TYPE=AUDIO,GROUP-ID="audio",NAME="{name}",DEFAULT={default_attr},AUTOSELECT={autoselect_attr},LANGUAGE="{lang}",URI="{playlist_url}"'
)
if is_default:
first_audio_codec = profile.get("codecs", "")
count += 1
# Build combined codecs string (video + audio) for EXT-X-STREAM-INF
# ExoPlayer requires CODECS to list all codecs when AUDIO group is referenced
first_audio_codec = None
if audio_profiles:
first_audio_profile = next(iter(audio_profiles.values()))[0]
first_audio_codec = first_audio_profile.get("codecs", "")
# Add video streams
for profile, playlist_url in video_profiles.values():
@@ -440,9 +645,13 @@ def build_hls_playlist(
skip_filter = SkipSegmentFilter(skip_segments) if skip_segments else None
# In TS mode, we don't use EXT-X-MAP because TS segments are self-contained
# (PAT/PMT/VPS/SPS/PPS are embedded in each segment)
# Use EXT-X-MAP for live streams, but only for fMP4 (not TS)
use_map = is_live and not is_ts_mode
# (PAT/PMT/VPS/SPS/PPS are embedded in each segment).
# Use EXT-X-MAP for:
# - live fMP4 streams (init changes with discontinuities)
# - SegmentBase fMP4 (init and media are different byte ranges of the same file;
# without EXT-X-MAP every segment would redundantly include the moov box)
has_segment_base = not is_ts_mode and any(p.get("initRange") for p in profiles)
use_map = not is_ts_mode and (is_live or has_segment_base)
# Select appropriate endpoint based on remux mode
if is_ts_mode:
@@ -462,8 +671,8 @@ def build_hls_playlist(
continue
if is_live:
# TS mode uses deeper playlist for ExoPlayer buffering
depth = 20 if is_ts_mode else max(settings.mpd_live_playlist_depth, 1)
extinf_values_for_depth = [s["extinf"] for s in segments if "extinf" in s]
depth = _compute_live_playlist_depth(is_ts_mode, effective_start_offset, extinf_values_for_depth)
trimmed_segments = segments[-depth:]
else:
trimmed_segments = segments
@@ -479,31 +688,13 @@ def build_hls_playlist(
else:
target_duration = math.ceil(max(extinf_values)) if extinf_values else 3
# Align HLS media sequence with MPD-provided numbering when available
if is_ts_mode and is_live:
# For live TS, derive sequence from timeline first for stable continuity
time_val = first_segment.get("time")
duration_val = first_segment.get("duration_mpd_timescale")
if time_val is not None and duration_val and duration_val > 0:
sequence = math.floor(time_val / duration_val)
else:
sequence = first_segment.get("number") or profile.get("segment_template_start_number") or 1
if is_live:
sequence = _compute_live_media_sequence(first_segment, profile, trimmed_segments)
else:
mpd_start_number = profile.get("segment_template_start_number")
sequence = first_segment.get("number")
if sequence is None:
# Fallback to MPD template start number
if mpd_start_number is not None:
sequence = mpd_start_number
else:
# As a last resort, derive from timeline information
time_val = first_segment.get("time")
duration_val = first_segment.get("duration_mpd_timescale")
if time_val is not None and duration_val and duration_val > 0:
sequence = math.floor(time_val / duration_val)
else:
sequence = 1
sequence = mpd_start_number if mpd_start_number is not None else 1
hls.extend(
[
@@ -543,6 +734,10 @@ def build_hls_playlist(
if query_params.get("api_password"):
init_query_params["api_password"] = query_params["api_password"]
for k, v in query_params.items():
if k.startswith("rp_") or k.startswith("h_"):
init_query_params[k] = v
init_map_url = encode_mediaflow_proxy_url(
init_proxy_url,
query_params=init_query_params,
@@ -594,6 +789,9 @@ def build_hls_playlist(
# Segment may also have its own range (for SegmentBase)
if "initRange" in segment:
segment_query_params["init_range"] = segment["initRange"]
# Media byte range: bytes after the init segment (SegmentBase only)
if segment.get("mediaRange"):
segment_query_params["segment_range"] = segment["mediaRange"]
query_params.update(segment_query_params)
hls.append(