diff --git a/.github/workflows/interop.yml b/.github/workflows/interop.yml index 74a22e98bd..5a2fa6b717 100644 --- a/.github/workflows/interop.yml +++ b/.github/workflows/interop.yml @@ -138,6 +138,12 @@ jobs: env: MOQ_TEST_RUNS: ${{ runner.temp }}/moq-test-runs + # The T-STD model must pass a broadcast capture and fail the same capture + # delivered too slowly, too early, or in a burst. + - name: T-STD controls + run: nix develop --command just test ts-tstd + shell: bash -leo pipefail {0} + # A failing harness run keeps its directory, holding the process logs, the # relay config, and a Playwright trace of the failing page. Upload it for a # short while so a red cell is diagnosable; a passing run deletes its own. diff --git a/quest/m1/tstd/README.md b/quest/m1/tstd/README.md index ec8f9dc2cf..f681316236 100644 --- a/quest/m1/tstd/README.md +++ b/quest/m1/tstd/README.md @@ -16,6 +16,12 @@ Padding, pacing, and muxing all assume a fixed delay, so the export gets one first. t0ms is testing whether T-STD compliance is feasible at all; record the result here. If it isn't, re-plan this line. +The full-model check (`tstd` in `test/ts/compliance.py`) fails today's export +for two reasons, both delivery timing: 7-10 % of video access units finish +arriving after their DTS, and audio runs ~230 ms ahead in bursts that overflow +its TB and B. The media itself fits its buffers, so nothing yet says the bar +is out of reach. + This README owns the end-to-end proof: the #4613 netem rig (10% loss, a real ~10 Mb/s broadcast TS) passes the strict T-STD check, and the recipe runs nightly. @@ -23,5 +29,4 @@ nightly. ## Required - [Fixed-delay release](/quest/m1/tstd/delay.md) - frames go out at media time plus a fixed `--delay`, in one order under loss -- [T-STD check](/quest/m1/tstd/check.md) - the harness grades the full buffer model instead of the transport buffer alone - [TS byte schedule](/quest/m1/tstd/byte-schedule.md) - PCRs sit on the byte grid the mux rate implies, paced against the fixed delay diff --git a/quest/m1/tstd/check.md b/quest/m1/tstd/check.md deleted file mode 100644 index a7c42f511f..0000000000 --- a/quest/m1/tstd/check.md +++ /dev/null @@ -1,25 +0,0 @@ -# [M] The TS harness grades the full T-STD buffer model - -## Goal - -`test/ts/compliance.py` models every T-STD buffer an elementary stream passes -through, not only the transport buffer. It fails a stream where any buffer -overflows or where an access unit hasn't fully arrived by its DTS, and the TS -recipe runs it strictly against `moq export ts` output in CI. - -## Plan - -Today the `tstd` check models only the transport-buffer (TB) smoothing stage. -It's a shape check that warns unless `--strict` is set (`test/ts/README.md`). -Extend it to the multiplex buffer (MB, for video) and elementary buffer (EB) -with the ISO 13818-1 leak rates and sizes for each stream type and level, and -remove each access unit at its DTS. Timing stays on the stream's own PCR -clock, so the result is deterministic for a file. - -Validate the model before gating: a known-compliant reference TS passes, and -a deliberately bursty one fails. Where TSDuck or another maintained tool -already implements T-STD, prefer it over hand-rolled math. - -Gate: strict `tstd` in the TS recipe for export output. Until -[fixed-delay release](/quest/m1/tstd/delay.md) lands, the check will likely -fail, so land it report-only and flip the gate in that PR. diff --git a/quest/m1/tstd/delay.md b/quest/m1/tstd/delay.md index e07d26616e..df8b826567 100644 --- a/quest/m1/tstd/delay.md +++ b/quest/m1/tstd/delay.md @@ -62,6 +62,9 @@ compare with #4618's numbers. Update `doc/bin/cli.md` and the `moq export ts` examples. +Promote `tstd` in `test/ts/compliance.py` from shape to hard, so `just test +ts` fails a round-trip the T-STD model rejects; it reports only until then. + Public API: `ts::Export` takes the delay in place of its max age and loses the hold; breaking, on `dev`. Wire: none. diff --git a/test/justfile b/test/justfile index 4b46664051..523fd57f9c 100644 --- a/test/justfile +++ b/test/justfile @@ -77,8 +77,8 @@ drill-sensitivity *args: ./drill/sensitivity.sh "$@" # Round-trips a PCR-paced TS through a relay and checks the subscriber's `export -# ts` output with TSDuck plus a custom analyzer (PCR jitter/repetition, -# burstiness, instantaneous bitrate, T-STD). Flags: `--analyze-only file.ts` +# ts` output with TSDuck plus a custom analyzer (the T-STD buffer model, PCR +# spacing and byte schedule). Flags: `--analyze-only file.ts` # (skip the round-trip), `--strict` (fail on shape warnings), `--source file.ts` # (use a real capture), `--live` (grade PCR release timing and byte position off # the pipe instead, which a capture cannot carry; nightly runs this arm), @@ -96,6 +96,13 @@ ts *args: ts-eit *args: ./ts/eit-roundtrip.sh "$@" +# The T-STD model's controls: a broadcast capture must pass, and the same capture +# delivered too slowly, too early, or in a burst must fail. Needs only TSDuck. + +# Prove the T-STD model passes a compliant stream and fails broken ones. +ts-tstd: + ./ts/tstd-controls.py + # Runs the @moq/wasm bindings in headless Chromium against a real relay, one per # protocol flavour. The crate is `#![cfg(target_arch = "wasm32")]`, so nothing in # `just check` or `just rs wasm` gets past compiling it. See wasm/README.md. diff --git a/test/ts/README.md b/test/ts/README.md index 4fb5950b9d..b875227fd8 100644 --- a/test/ts/README.md +++ b/test/ts/README.md @@ -7,8 +7,8 @@ runs [TSDuck](https://tsduck.io) plus a custom analyzer over it. This is a diagnostic gate, not just a pass/fail: the exporter ([`rs/moq-mux/src/container/ts/export.rs`](../../rs/moq-mux/src/container/ts/export.rs)) -is VBR, inserts no null packets, and paces PCR once per media frame, so several -broadcast-shape checks are expected to flag. The report quantifies exactly where +pads to a constant rate only once the catalog carries the source's mux rate, so +several broadcast-shape checks are expected to flag. The report quantifies exactly where and by how much. Four instruments live here. `compliance.py` (via `run.sh`) grades a captured file @@ -39,10 +39,10 @@ That is the only arm that can see release timing at all, and it is what nightly runs (see [CI](#ci)). The default arm runs `pcr-timing.py` over its capture too, after `compliance.py`, -for the one thing the IRD model does not grade: whether the bytes between -consecutive PCRs are the ones the mux rate implies -([`pcr-schedule`](#byte-schedule)). It is a shape check, so it reports without -gating unless `--strict`. +for the PCR checks `compliance.py` leaves to it: the value interval (hard), and +whether the bytes between consecutive PCRs are the ones the mux rate implies +([`pcr-schedule`](#byte-schedule), a shape check that reports without gating +unless `--strict`). The live arm passes only when the grader's verdict *and* the publisher's exit status are clean. The grader can only speak for what reached it, and the sample floor @@ -57,7 +57,7 @@ moq --connect http://localhost:4443 --broadcast live.hang export ts > sub.ts ./run.sh --analyze-only sub.ts ``` -Requirements: `tsp` and `tsanalyze` (TSDuck) and `python3` for every mode; the +Requirements: `tsp`, `tsanalyze` and `tstables` (TSDuck) and `python3` for every mode; the round-trip modes also need `cargo`, `ffmpeg`, `curl`, and `timeout`. ## Checks @@ -80,14 +80,8 @@ Severities: **hard** checks fail the run by default; **shape** checks report as | `pcr-presence` | hard | a PCR PID is declared and carries PCR | | `pcr-monotonic` | hard | PCR strictly increases (one 33-bit wrap tolerated), except into a PCR that signals `discontinuity_indicator` | | `duration-fidelity` | hard | exported PCR span tracks the source's duration (round-trip only) | -| `pcr-repetition` | shape | consecutive PCRs within the limit (default 40 ms) | -| `pcr-jitter` | shape | per-interval PCR jitter vs the nominal bitrate (pcrverify model) | -| `null-ratio` | shape | null/stuffing fraction (flags only a pathological excess) | | `service-descriptors` | shape | an SDT naming the service is present | -| `bitrate-consistency` | shape | instantaneous-bitrate spread over 1 ms / 10 ms windows (CBR-ness) | -| `burstiness` | shape | peak/mean of windowed delivery | -| `inter-arrival` | shape | packet inter-arrival spread on the PCR clock (informational) | -| `tstd` | shape | transport-buffer smoothing (TB fills on arrival, leaks at Rx) | +| `tstd` | shape | the full T-STD buffer model: no TB, MB, EB or B overflow, no access unit incomplete at its decoding time (see [T-STD](#t-std)) | Every timing check reads the stream's own PCR, so a PCR emitted on the wrong clock rate stays internally consistent and passes them all. `duration-fidelity` @@ -96,10 +90,111 @@ independent duration, which pins the absolute rate. It runs only on a round-trip (where a source exists); `run.sh` passes the source automatically, and `--analyze-only` skips it. -Thresholds are CLI flags forwarded through `run.sh` (e.g. -`--pcr-repetition-ms`, `--pcr-jitter-us`, `--bitrate-cov-max`, `--burstiness-max`, -`--tb-size-bytes`, `--video-leak-bps`, `--audio-leak-bps`). `--report-json ` -writes the full machine-readable report. +`compliance.py` grades what TSDuck parses and the T-STD model, and leaves PCR +spacing, byte schedule and release timing to `pcr-timing.py`, which honours +signalled discontinuities. `--report-json ` writes the full +machine-readable report. + +## T-STD + +`tstd` runs every audio and video stream through the ISO 13818-1 system target +decoder (2.4.2; Rec. ITU-T H.222.0, whose 10/2014 edition is a free download), +fed on the stream's own PCR clock: + +```text +video (AVC 2.14.3.1, HEVC 2.17.2) TB --Rx--> MB --Rbx (leak)--> EB --DTS--> decoder +audio (2.4.2.3) TB --Rx--> B ----------------PTS--> decoder +``` + +It fails a stream where TB, MB, EB or B overflows, where TB stays occupied for a +second, where an access unit is not wholly in EB/B at its decoding time +(underflow), or where a byte waits longer than the STD delay bound (1 s, 10 s for +AVC/HEVC). The report gives each stream's peak fill per buffer, how late the worst +underflowed access unit finished arriving, and the longest any access unit waited. + +No maintained tool implements this. TSDuck has no T-STD analyzer, and nothing else +in nixpkgs does either, so the model is hand-rolled and its parameters are +transcribed from the specs: H.264 Table A-1 and H.265 Table A.8 for the level, the +ADTS and "other audio" rates and sizes in H.222.0 2.4.2.3, and ATSC A/52 and A/53 +Part 5 for AC-3 and E-AC-3. TSDuck does the parsing: `tstables` decodes the PMT +and `tsp -P pes --avc-access-unit` the SPS, AVC and HEVC alike. + +Video takes its buffers from the HRD the SPS declares, as H.222.0 2.14.3.1 +(AVC) and 2.17.2 (HEVC) specify. A NAL HRD sets Rx from its bit rate and EB to +its CPB size, and MB grows by whatever the level's CPB leaves over; Rbx stays the +level's. Without one, the level's limits stand in. A VCL HRD describes the VCL +alone, not the byte stream EB holds, so a stream declaring only that takes the +level defaults too. An `AVC_timing_and_HRD_descriptor` or +`HEVC_timing_and_HRD_descriptor` with `hrd_management_valid` switches MB-to-EB +transfer to the HRD's own schedule, which is not modelled, so such a stream is +refused. + +Opus is graded against ADTS's buffers for the same channel count: the Opus-in-TS +draft gives Rx (2 Mb/s for 1-2 channels, matching ADTS) but leaves the buffer size +unset, so that size is borrowed rather than specified. A stream with no parameters +at all is refused by name rather than skipped, which fails the check: MPEG-1/2 +video, DVB E-AC-3, and HEVC beyond Main/Main 10. Sections and private data +(SCTE-35, teletext) have no elementary-stream buffers and are listed as not +graded. A signalled PCR discontinuity starts fresh buffers, since the timestamps +on either side of it are on different clocks. + +Every packet on an elementary stream's PID enters TB, including adaptation-only +ones (a PCR, stuffing) and legal duplicates (2.4.3.3: the same counter and +payload twice); only PES bytes of a first copy go on to MB/B. Video access units +are read from the ES rather than taken from PES boundaries: one opens at the first +delimiter, parameter set or SEI after the previous picture's slices (H.264 +7.4.1.2.3, H.265 7.4.2.4.4), and H.222.0 requires a delimiter in each. A PES may +carry several, but only the first takes its timestamp, and deriving the rest from +the stream's own timing is not modelled, so that layout is refused. Video decode +times must strictly increase. + +One simplification, toward strictness: a packet's bytes reach MB/B when its last +byte leaves TB, up to one packet's drain time (0.75 ms for audio) later than +byte-by-byte, so underflow is judged that much stricter. + +### Controls + +`just test ts-tstd` (`tstd-controls.py`) proves the model can tell a compliant +stream from a broken one. The positive control is a real broadcast encoder's +output (`kyrion_dirtystart.ts` from the `moq-mux` test data: AVC High@4.0 with a +1.935 Mb/s CBR NAL HRD and a 755 kbit CPB, plus two MPEG-1 Layer II tracks), which +passes against its own declared buffer, filling EB to the brim as a CBR stream +should. The negatives restamp its PCRs with +`tsp -P pcradjust`, leaving every PES and timestamp alone, so only delivery +changes: + +| Case | Expected | +|---|---| +| as captured | pass | +| PCRs restamped at the capture's own rate | pass | +| delivered at 0.7x | EB and B underflow | +| delivered at 4x | TB and B overflow | +| delivered at 15x (a burst) | TB and B overflow | + +No restamp can overflow the video MB: it holds the level's whole CPB less the +declared one, about 3.6 MB, more than the 4 s capture carries. + +The packet layouts a capture cannot be edited into are built synthetically: a +10 Mb/s single-video stream carrying the Kyrion SPS, one access unit per PES, with +a PCR packet between them. Each case fails without the handling it names: + +| Case | Expected | +|---|---| +| as built | pass | +| two access units in one PES | refused | +| four adaptation-only packets after an access unit | TB overflow | +| an access unit's first packet sent twice | pass | + +The ffmpeg clip `run.sh` generates is not a positive control: its muxer sends +audio 0.7 s ahead by default (`-muxdelay`), which overflows the 3,584-byte ADTS +buffer, and even at 0.1 s a four-packet audio burst overflows TB. + +### Gate + +`tstd` is a shape check, so `just test ts` reports it without failing. `export +ts` output does not pass yet: some video access units arrive after their DTS, and +audio runs far enough ahead, in bursts, to overflow both TB and B. The change that +makes the exporter release on a fixed delay promotes `tstd` to a hard check. ## PCR timing (`pcr-timing.py`) @@ -112,11 +207,9 @@ cannot see the other two: | release | the bytes carrying a PCR were handed over when that PCR asserts | arrival stamps | | position | a PCR packet sits among the media bytes it describes | packet offsets | -`compliance.py` grades `value` from a file, deterministically and with no -wall-clock capture, which is the right basis for the model math it does. That -also means it cannot grade `release`: a change to *when* the exporter hands bytes -over is invisible to any harness that does not stamp arrivals. -`pcr-timing.py` reads a pipe and grades all three in one pass. +`pcr-timing.py` grades `value` and `position` from a file, and all three from a +pipe. A file carries no arrival stamps, so a change to *when* the exporter hands +bytes over is invisible to any harness that does not stamp them. A constant-rate stream makes a fourth claim, graded by `pcr-schedule`: that the bytes between consecutive PCRs are the bytes the mux rate implies for that @@ -143,13 +236,19 @@ grades only how evenly the bytes are laid over the PCRs. | Check | Severity | What it verifies | |---|---|---| -| `sync` | hard | no invalid sync bytes / transport-error packets | -| `continuity` | hard | no discontinuities, and a payload-less packet must not advance the counter (ISO 13818-1 2.4.3.3) | -| `pcr-single-pid` | hard | every PCR rides one PID | +| `sync` | hard | no invalid sync bytes / transport-error packets (`--live` only) | +| `continuity` | hard | no discontinuities, and a payload-less packet must not advance the counter (ISO 13818-1 2.4.3.3) (`--live` only) | | `pcr-value-interval` | hard | no interval above `--repetition-ms` (default 40, TR 101 290), within one time base | | `pcr-release-timing` | hard | no more than `--release-pct-max` of intervals arrive further than `--release-ms` from the interval their own values assert, and accumulated drift stays within `--drift-ms`, being the standing lag the sender is allowed to hold; a sample below `--live-min-pcr` PCRs or `--live-cover-pct` of the window is a failure, not a pass (`--live` only) | | `pcr-position` | shape | share of PCR packets within `--adjacent-packets` of the previous one | -| `pcr-schedule` | shape | share of PCR intervals, on the busiest PCR PID, whose bytes are within `--schedule-tolerance-pct` (default 1) or one packet of what `--mux-rate` implies (estimated from the capture if not given); hard, at that share, when `--schedule-pct-min` is given | +| `pcr-schedule` | shape | share of PCR intervals whose bytes are within `--schedule-tolerance-pct` (default 1) or one packet of what `--mux-rate` implies (estimated from the capture if not given); hard, at that share, when `--schedule-pct-min` is given | + +A stream carrying several PCR PIDs (one per program) is graded on the busiest, +since two correct grids offset from one another pool into one that neither keeps. +`sync` and `continuity` run only under `--live`: on a file, `compliance.py` +grades both through TSDuck's `tsanalyze`, which also catches a payload-less +packet advancing the counter. A pipe cannot go through TSDuck first without +rebuffering the arrivals `release` stamps. Accumulated drift has two shapes and only one is a defect, so the check bounds the total and reports the rate over the tail of the sample beside it. A sender @@ -444,8 +543,8 @@ exporter re-emits SI on its own repetition cadence rather than the source's. ## CI `.github/workflows/interop.yml` runs `just test ts`, `just test ts --open-gop`, -and `just test ts-eit` after the interop matrix (nightly, on demand, and on PRs -touching `test/ts/`). +`just test ts-eit`, and `just test ts-tstd` after the interop matrix (nightly, on +demand, and on PRs touching `test/ts/`). `ts-eit` is `eit-roundtrip.sh`: it builds the sparse-schedule and pending-version fixtures from a generated clip, round-trips them through a relay, and censuses the capture, so a break in the generators or in the SI @@ -469,6 +568,3 @@ that from a pipe running slow. - Wall-clock delivery jitter/burstiness is out of scope *for `compliance.py`*: all of its timing is derived from the stream's PCR, not from arrival times. `pcr-timing.py --live` covers that axis separately, by stamping a pipe. -- `tstd` models only the transport-buffer (TB) smoothing stage of the ISO 13818-1 - T-STD, not the full multiplex/elementary decode buffers. Its leak rates are - defaults, not level-derived, so treat overflow as a smell rather than proof. diff --git a/test/ts/compliance.py b/test/ts/compliance.py index 66aae98fe7..0609857cf7 100755 --- a/test/ts/compliance.py +++ b/test/ts/compliance.py @@ -6,19 +6,17 @@ summary, exiting non-zero on failure. Division of labour: TSDuck does the transport-stream parsing (we shell out to -`tsanalyze --json` for PSI/service/structure), and this script does the model -math TSDuck does not cover directly (PCR jitter/repetition, packet -inter-arrival, burstiness, instantaneous bitrate, and a transport-buffer model). -The PCR/PTS/DTS timeline the timing model needs comes from a 188/204-byte -packet-header scan done here, which also gives the per-packet PID that -`tsp -P pcrextract` does not expose. +`tsanalyze --json` for PSI/service/structure and `tstables` for the PMT), and +this script does what TSDuck does not cover: the ISO 13818-1 T-STD buffer +model, which no maintained tool implements, and the PCR checks that need the +time-base discontinuities a 188/204-byte header scan here recovers. PCR value +intervals, byte schedule and release timing live in `pcr-timing.py`. Checks split into two severities: - HARD (structural): fail the run by default. PAT/PMT, packet size, sync, continuity counters, PSI CRC, PCR presence, PCR monotonicity. - SHAPE (broadcast profile): reported as WARN and only fail the run under - `--strict`. PCR repetition interval, PCR jitter, null-packet ratio, bitrate - consistency / burstiness, service descriptors (SDT), transport-buffer model. + `--strict`. Service descriptors (SDT), T-STD buffer model. Timing basis is the stream's own PCR clock (an IRD locks to PCR), so the harness needs no wall-clock capture and results are deterministic for a given file. @@ -39,7 +37,6 @@ PTS_HZ = 90_000 # The PCR base field is 33 bits at 90 kHz, so the full 27 MHz PCR wraps here. PCR_WRAP = (1 << 33) * 300 -NULL_PID = 0x1FFF class Status(Enum): @@ -68,22 +65,6 @@ class Check: metrics: dict = field(default_factory=dict) -@dataclass -class Thresholds: - """IRD limits and T-STD model parameters, all overridable from the CLI.""" - - pcr_repetition_ms: float = 40.0 - pcr_jitter_us: float = 500.0 - null_ratio_max: float = 0.90 - bitrate_cov_max: float = 0.10 - burstiness_max: float = 3.0 - tb_size_bytes: int = 512 - video_leak_bps: float = 1.8e7 - audio_leak_bps: float = 2.0e6 - data_leak_bps: float = 1.0e6 - inst_windows_ms: tuple[float, ...] = (1.0, 10.0) - - # --------------------------------------------------------------------------- IO @@ -106,10 +87,7 @@ def run_tsanalyze(ts_path: str) -> dict: class Scan: """Per-packet facts recovered from a raw header scan of the TS file.""" - packet_size: int total_packets: int - # (ts_index, payload_bytes) per PID, in file order. - pid_packets: dict[int, list[tuple[int, int]]] # (ts_index, pcr_27mhz) samples for every PID that carries PCR. pcr_by_pid: dict[int, list[tuple[int, int]]] # ts_index of every PCR that starts a new time base (discontinuity_indicator). @@ -117,15 +95,14 @@ class Scan: def scan_packets(ts_path: str, packet_size: int) -> Scan: - """Walk the TS packet by packet, recovering PID, payload size, and PCR values. + """Walk the TS packet by packet, recovering the PCR timeline. - This is a header-only scan (no PES/PSI parsing): enough for the timing and - buffer models, which need per-packet PID and the PCR timeline. + A header-only scan: the PCR samples and which of them start a new time base, + which `tsp -P pcrextract` does not report. """ with open(ts_path, "rb") as handle: data = handle.read() - pid_packets: dict[int, list[tuple[int, int]]] = {} pcr_by_pid: dict[int, list[tuple[int, int]]] = {} pcr_new_base: set[int] = set() # PIDs whose discontinuity_indicator rode a packet without a PCR, so the new time @@ -147,8 +124,6 @@ def scan_packets(ts_path: str, packet_size: int) -> Scan: pid = ((b1 & 0x1F) << 8) | b2 afc = (b3 >> 4) & 0x3 - payload_len = 0 - af_len = 0 if afc in (2, 3): af_len = data[offset + 4] if af_len > 0: @@ -168,21 +143,12 @@ def scan_packets(ts_path: str, packet_size: int) -> Scan: if pid in pending_base: pending_base.discard(pid) pcr_new_base.add(index) - if afc in (1, 3): - # Payload = 184 minus the adaptation field (its length byte + body). - consumed = (1 + af_len) if afc == 3 else 0 - payload_len = 184 - consumed - - if pid != NULL_PID: - pid_packets.setdefault(pid, []).append((index, max(0, payload_len))) index += 1 offset += packet_size return Scan( - packet_size=packet_size, total_packets=index, - pid_packets=pid_packets, pcr_by_pid=pcr_by_pid, pcr_new_base=pcr_new_base, ) @@ -191,25 +157,6 @@ def scan_packets(ts_path: str, packet_size: int) -> Scan: # --------------------------------------------------------------- small helpers -def percentile(values: list[float], pct: float) -> float: - """Linear-interpolated percentile of an unsorted list (0..100). 0 for empty.""" - if not values: - return 0.0 - ordered = sorted(values) - if len(ordered) == 1: - return ordered[0] - rank = (pct / 100.0) * (len(ordered) - 1) - low = int(rank) - high = min(low + 1, len(ordered) - 1) - frac = rank - low - return ordered[low] * (1 - frac) + ordered[high] * frac - - -def mean(values: list[float]) -> float: - """Arithmetic mean, 0 for an empty list.""" - return sum(values) / len(values) if values else 0.0 - - def pcr_seconds(pcr_27mhz: int) -> float: """Convert a 27 MHz PCR value to seconds.""" return pcr_27mhz / PCR_HZ @@ -367,89 +314,6 @@ def check_pcr_monotonic(scan: Scan) -> Check: # ------------------------------------------------------------------- shape checks -def check_pcr_repetition(scan: Scan, th: Thresholds) -> Check: - """Consecutive PCRs on a PID should be no more than `pcr_repetition_ms` apart.""" - worst_ms = 0.0 - intervals: list[float] = [] - over = 0 - for samples in scan.pcr_by_pid.values(): - prev = None - for _i, pcr in samples: - if prev is not None: - delta = pcr - prev - if delta <= 0: - prev = pcr - continue - ms = pcr_seconds(delta) * 1000.0 - intervals.append(ms) - worst_ms = max(worst_ms, ms) - if ms > th.pcr_repetition_ms: - over += 1 - prev = pcr - metrics = { - "max_interval_ms": round(worst_ms, 3), - "mean_interval_ms": round(mean(intervals), 3), - "intervals_over_limit": over, - "limit_ms": th.pcr_repetition_ms, - } - detail = f"max {worst_ms:.1f} ms (limit {th.pcr_repetition_ms:.0f} ms), {over} over" - status = Status.PASS if over == 0 else Status.WARN - return Check("pcr-repetition", Severity.SHAPE, status, detail, metrics) - - -def check_pcr_jitter(scan: Scan, ts_bitrate: float, th: Thresholds) -> Check: - """Per-interval PCR jitter vs the nominal bitrate (pcrverify's model). - - For a true CBR mux the actual PCR delta matches the byte delta clocked at the - stream bitrate; the difference is the jitter. On a VBR stream it is large by - construction, which is exactly the IRD-relevant signal. - """ - if ts_bitrate <= 0: - return Check("pcr-jitter", Severity.SHAPE, Status.WARN, "unknown bitrate", {}) - bits_per_packet = scan.packet_size * 8 - jitters_us: list[float] = [] - for samples in scan.pcr_by_pid.values(): - prev = None - for i, pcr in samples: - if prev is not None: - pi, ppcr = prev - d_pcr = pcr - ppcr - if d_pcr <= 0: - prev = (i, pcr) - continue - expected_s = (i - pi) * bits_per_packet / ts_bitrate - actual_s = pcr_seconds(d_pcr) - jitters_us.append((actual_s - expected_s) * 1e6) - prev = (i, pcr) - if not jitters_us: - return Check("pcr-jitter", Severity.SHAPE, Status.WARN, "no PCR intervals", {}) - abs_jit = [abs(j) for j in jitters_us] - max_us = max(abs_jit) - p95 = percentile(abs_jit, 95) - metrics = { - "max_abs_us": round(max_us, 1), - "p95_abs_us": round(p95, 1), - "limit_us": th.pcr_jitter_us, - } - detail = f"max |jitter| {max_us:.0f} us, p95 {p95:.0f} us (limit {th.pcr_jitter_us:.0f} us)" - status = Status.PASS if max_us <= th.pcr_jitter_us else Status.WARN - return Check("pcr-jitter", Severity.SHAPE, status, detail, metrics) - - -def check_null_ratio(analysis: dict, th: Thresholds) -> Check: - """Report the null-packet (stuffing) fraction; flag only a pathological excess.""" - total = analysis["ts"]["packets"]["total"] or 1 - null = 0 - for pid in analysis.get("pids", []): - if pid["id"] == NULL_PID: - null = pid.get("packets", {}).get("total", 0) - ratio = null / total - metrics = {"null_ratio": round(ratio, 4), "null_packets": null, "limit": th.null_ratio_max} - detail = f"{ratio * 100:.2f}% null packets" - status = Status.PASS if ratio <= th.null_ratio_max else Status.WARN - return Check("null-ratio", Severity.SHAPE, status, detail, metrics) - - def check_service_descriptors(analysis: dict) -> Check: """An IRD expects an SDT naming the service; PAT/PMT-only streams get a WARN.""" tables = analysis.get("tables", []) @@ -468,152 +332,804 @@ def check_service_descriptors(analysis: dict) -> Check: ) -def windowed_bitrates(times: list[float], sizes: list[int], window_s: float) -> list[float]: - """Bytes-per-window converted to bit/s, over contiguous windows spanning the capture.""" - if not times: - return [] - start, end = times[0], times[-1] - if end <= start: - return [] - n_windows = max(1, int((end - start) / window_s) + 1) - buckets = [0] * n_windows - for t, size in zip(times, sizes): - idx = min(n_windows - 1, int((t - start) / window_s)) - buckets[idx] += size - # Drop the last (partial) window so a short tail doesn't skew the minimum. - if n_windows > 1: - buckets = buckets[:-1] - return [b * 8 / window_s for b in buckets] - - -def check_bitrate_and_burstiness(scan: Scan, clock: PcrClock, ts_bitrate: float, th: Thresholds) -> tuple[Check, Check]: - """Instantaneous bitrate spread (CBR-ness) and delivery burstiness, on the PCR clock.""" - # Every non-null packet, timed on the PCR clock, weighted by its full 188 bytes. - events: list[tuple[float, int]] = [] - for _pid, packets in scan.pid_packets.items(): - for i, _payload in packets: - events.append((clock.time_at(i), scan.packet_size)) - events.sort(key=lambda e: e[0]) - times = [t for t, _ in events] - sizes = [s for _, s in events] - - inst_metrics: dict = {"nominal_bps": round(ts_bitrate)} - worst_cov = 0.0 - worst_burst = 0.0 - for window_ms in th.inst_windows_ms: - rates = windowed_bitrates(times, sizes, window_ms / 1000.0) - if not rates: +# ---------------------------------------------------------------- T-STD model +# +# ISO/IEC 13818-1 2.4.2 (Rec. ITU-T H.222.0 10/2014, free from the ITU), fed each +# elementary stream's packets at the times its program's PCR assigns them: +# +# video (AVC 2.14.3.1, HEVC 2.17.2) TB -Rx-> MB -Rbx (leak)-> EB -DTS-> decoder +# audio (2.4.2.3) TB -Rx-> B ---------------PTS-> decoder +# +# and graded on 2.4.2.6 and its AVC/HEVC counterparts: no buffer overflows, TB +# empties at least once a second, every access unit is complete in EB/B at its +# decoding time, and no byte waits longer than the STD delay bound. +# +# No maintained tool implements this (TSDuck has no T-STD analyzer), so the +# parameters are transcribed from the specs named beside each table. + +TB_SIZE = 512 +# 2.4.2.6: TB must empty at least once a second. +TB_EMPTY_S = 1.0 +# PTS/DTS are 33 bits at 90 kHz. +STAMP_WRAP_S = (1 << 33) / PTS_HZ + +# H.264 Table A-1: level_idc -> (MaxBR, MaxCPB) as tabulated. H.222.0 2.14.3.1 +# scales both by 1200 bits for the buffers, and Rx by the profile's cpbBrNalFactor. +AVC_LEVELS = { + 10: (64, 175), + 11: (192, 500), + 12: (384, 1000), + 13: (768, 2000), + 20: (2000, 2000), + 21: (4000, 4000), + 22: (4000, 4000), + 30: (10000, 10000), + 31: (14000, 14000), + 32: (20000, 20000), + 40: (20000, 25000), + 41: (50000, 62500), + 42: (50000, 62500), + 50: (135000, 135000), + 51: (240000, 240000), + 52: (240000, 240000), + 60: (240000, 240000), + 61: (480000, 480000), + 62: (800000, 800000), +} +AVC_LEVEL_1B = (128, 350) +# H.264 Table A-2 cpbBrNalFactor, by profile_idc: the default BitRate that sets Rx. +AVC_NAL_FACTOR = {66: 1200, 77: 1200, 88: 1200, 100: 1500, 110: 3600, 122: 4800, 244: 4800, 44: 4800} + +# H.265 Table A.8: general_level_idc -> ((MaxBR, MaxCPB) Main tier, (MaxBR, MaxCPB) High tier). +HEVC_LEVELS = { + 30: ((128, 350), None), + 60: ((1500, 1500), None), + 63: ((3000, 3000), None), + 90: ((6000, 6000), None), + 93: ((10000, 10000), None), + 120: ((12000, 12000), (30000, 30000)), + 123: ((20000, 20000), (50000, 50000)), + 150: ((25000, 25000), (100000, 100000)), + 153: ((40000, 40000), (160000, 160000)), + 156: ((60000, 60000), (240000, 240000)), + 180: ((60000, 60000), (240000, 240000)), + 183: ((120000, 120000), (480000, 480000)), + 186: ((240000, 240000), (800000, 800000)), +} +# H.265 Table A.9 CpbNalFactor for Main, Main 10 and Main Still Picture (profile_idc 1-3). +HEVC_NAL_FACTOR = 1100 + +# H.222.0 2.4.2.3: ADTS (Rx, BSn) by channels, the LFE not counted. +ADTS_BUFFERS = ((2, 2_000_000, 3584), (8, 5_529_600, 8976), (12, 8_294_400, 12804), (48, 33_177_600, 51216)) +# channel_configuration -> full-bandwidth channels (5.1 and 7.1 drop the LFE). +ADTS_CHANNELS = {1: 1, 2: 2, 3: 3, 4: 4, 5: 5, 6: 5, 7: 7} +ADTS_RATES = (96000, 88200, 64000, 48000, 44100, 32000, 24000, 22050, 16000, 12000, 11025, 8000, 7350) +# 2.4.2.3 "other audio". +AUDIO_RX = 2_000_000 +MPEG_AUDIO_BS = 3584 +# ATSC A/53 Part 5 5.7 (AC-3) and A/52 Annex G 3.6.1 (E-AC-3: 736 + 64 + 12096). +AC3_ATSC_BS = 2592 +EAC3_ATSC_BS = 12896 +# ATSC A/52 Annex A 5.4: AC-3 carried as DVB private data. +AC3_DVB_BS = 5696 + +MPEG_AUDIO_KBPS = { + (1, 1): (0, 32, 64, 96, 128, 160, 192, 224, 256, 288, 320, 352, 384, 416, 448), + (1, 2): (0, 32, 48, 56, 64, 80, 96, 112, 128, 160, 192, 224, 256, 320, 384), + (1, 3): (0, 32, 40, 48, 56, 64, 80, 96, 112, 128, 160, 192, 224, 256, 320), + (2, 1): (0, 32, 48, 56, 64, 80, 96, 112, 128, 144, 160, 176, 192, 224, 256), + (2, 2): (0, 8, 16, 24, 32, 40, 48, 56, 64, 80, 96, 112, 128, 144, 160), + (2, 3): (0, 8, 16, 24, 32, 40, 48, 56, 64, 80, 96, 112, 128, 144, 160), +} +# stream_type -> model, for the types identified by stream_type alone. +STREAM_KINDS = { + 0x1B: "avc", + 0x24: "hevc", + 0x0F: "adts", + 0x03: "mpeg-audio", + 0x04: "mpeg-audio", + 0x81: "ac3-atsc", + 0x87: "eac3-atsc", +} +AC3_KBPS = (32, 40, 48, 56, 64, 80, 96, 112, 128, 160, 192, 224, 256, 320, 384, 448, 512, 576, 640) + + +class Refused(Exception): + """An elementary stream the model cannot grade: named, never skipped silently.""" + + +@dataclass +class Params: + """One elementary stream's T-STD buffers (bytes) and rates (bit/s).""" + + label: str + rx: float + # Video: MB and EB with the leak rate between them. Audio: the main buffer B alone. + mb: float = 0.0 + eb: float = 0.0 + rbx: float = 0.0 + b: float = 0.0 + # tdn(j) - t(i) bound: 10 s for AVC and HEVC, 1 s otherwise (2.4.2.6, 2.14.3.1, 2.17.2). + max_delay_s: float = 1.0 + + @property + def video(self) -> bool: + return self.rbx > 0 + + +@dataclass +class AccessUnit: + """ES byte span [start, end), where its first byte arrived, and its decode time.""" + + start: int + end: int + first_packet: int + # 90 kHz DTS (or PTS) when the PES stamps it, else derived from the one before. + stamp: int | None + duration_s: float = 0.0 + + +@dataclass +class Stream: + """One elementary stream's packets and PES timestamps, and its ES bytes where needed.""" + + pid: int + stream_type: int + # TSDuck's names for the PMT descriptors, plus "registration:" + # and, for one TSDuck has no name for, "tag:". + descriptors: set[str] + # (ts_index, PES header bytes, ES payload bytes, ES offset before this packet). + packets: list[tuple[int, int, int, int]] = field(default_factory=list) + # (ES offset, 90 kHz PTS or None, DTS or None, ts_index) per PES, from the first timestamped one. + pes: list[tuple[int, int | None, int | None, int]] = field(default_factory=list) + es: bytearray = field(default_factory=bytearray) + es_len: int = 0 + # Video: the first SPS as TSDuck decodes it. + sps: dict[str, str] = field(default_factory=dict) + + +def read_programs(ts_path: str) -> list[tuple[int, list[Stream]]]: + """(PCR PID, elementary streams) of each program's first PMT, as TSDuck decodes it.""" + proc = subprocess.run( + ["tstables", ts_path, "--psi-si", "--tid", "2", "--json-output", "-"], + capture_output=True, + text=True, + check=False, + ) + if proc.returncode != 0: + raise RuntimeError(f"tstables failed: {proc.stderr.strip()}") + programs: dict[int, tuple[int, list[Stream]]] = {} + for pmt in json.loads(proc.stdout or "[]"): + if pmt.get("service_id") in programs: continue - avg = mean(rates) - peak = max(rates) - low = min(rates) - var = mean([(r - avg) ** 2 for r in rates]) - cov = (var**0.5 / avg) if avg else 0.0 - burst = (peak / avg) if avg else 0.0 - worst_cov = max(worst_cov, cov) - worst_burst = max(worst_burst, burst) - inst_metrics[f"w{int(window_ms)}ms"] = { - "min_bps": round(low), - "mean_bps": round(avg), - "max_bps": round(peak), - "p95_bps": round(percentile(rates, 95)), - "cov": round(cov, 3), - "peak_over_mean": round(burst, 2), - } + streams = [] + for node in pmt.get("#nodes", []): + if node.get("#name") != "component": + continue + descriptors = [d for d in node.get("#nodes", []) if isinstance(d, dict)] + names = {d["#name"] for d in descriptors} + names |= {f"registration:{d['format_identifier']}" for d in descriptors if "format_identifier" in d} + names |= {"hrd_management_valid" for d in descriptors if d.get("hrd_management_valid") is True} + names |= { + f"tag{d['tag']}:" + "".join(d.get("#nodes", [])).replace(" ", "").lower() + for d in descriptors + if d["#name"] == "generic_descriptor" + } + streams.append(Stream(node["elementary_pid"], node["stream_type"], names)) + programs[pmt.get("service_id")] = (pmt["pcr_pid"], streams) + return list(programs.values()) + + +def _stamp(b: bytes) -> int: + """A 33-bit PTS/DTS from its 5-byte PES encoding.""" + return ((b[0] >> 1) & 0x07) << 30 | b[1] << 22 | (b[2] >> 1) << 15 | b[3] << 7 | b[4] >> 1 + + +def read_pes(data: bytes, packet_size: int, streams: dict[int, Stream]) -> dict[int, str]: + """Split each stream's packets into PES header and ES bytes, and note every PES start. + + Every packet on the PID enters TB (H.222.0 2.4.2.3), but only PES bytes go on to + MB/B. So an adaptation-only packet (a PCR, stuffing) costs TB and nothing after it, + as does a duplicate (2.4.3.3: the same counter and payload twice), which "is not + delivered" downstream. Packets before a PID's first timestamped PES likewise stop at + TB: their bytes belong to an access unit whose start was never captured. + Returns why each stream that could not be read was refused. + """ + refused: dict[int, str] = {} + last: dict[int, tuple[int, bytes]] = {} # PID -> (continuity_counter, payload) + for index, offset in enumerate(range(0, len(data) - packet_size + 1, packet_size)): + if data[offset] != 0x47: + continue + pid = ((data[offset + 1] & 0x1F) << 8) | data[offset + 2] + stream = streams.get(pid) + if stream is None or pid in refused: + continue + afc = (data[offset + 3] >> 4) & 0x3 + start = offset + 4 + (1 + data[offset + 4] if afc & 0x2 else 0) + payload = data[start : offset + 188] if afc & 0x1 else b"" + cc = data[offset + 3] & 0x0F + duplicate = bool(payload) and last.get(pid) == (cc, payload) + if payload: + last[pid] = (cc, payload) + if not payload or duplicate or not (stream.pes or data[offset + 1] & 0x40): + stream.packets.append((index, 0, 0, stream.es_len)) + continue + length = len(payload) + header = 0 + if data[offset + 1] & 0x40: + if payload[:3] != b"\x00\x00\x01": + refused[pid] = f"packet {index} starts a payload without a PES start code" + continue + if length < 9 or 9 + payload[8] > length: + refused[pid] = f"the PES header at packet {index} spans packets" + continue + header = 9 + payload[8] + flags = payload[7] + pts = _stamp(payload[9:14]) if flags & 0x80 else None + dts = _stamp(payload[14:19]) if flags & 0xC0 == 0xC0 else None + if pts is None and not stream.pes: + stream.packets.append((index, 0, 0, 0)) + continue + stream.pes.append((stream.es_len, pts, dts, index)) + body = length - header + stream.packets.append((index, header, body, stream.es_len)) + stream.es += payload[header:] + stream.es_len += body + return refused - bitrate_status = Status.PASS if worst_cov <= th.bitrate_cov_max else Status.WARN - bitrate = Check( - "bitrate-consistency", - Severity.SHAPE, - bitrate_status, - f"worst CoV {worst_cov:.2f} (limit {th.bitrate_cov_max:.2f})", - inst_metrics | {"worst_cov": round(worst_cov, 3)}, + +def read_sps(ts_path: str, pid: int, nal_type: int) -> dict[str, str]: + """The first SPS on `pid`, field by field, as TSDuck's `pes` plugin decodes it.""" + proc = subprocess.run( + ["tsp", "-I", "file", ts_path, "-P", "pes", "--pid", str(pid), "--avc-access-unit"] + + ["--nal-unit-type", str(nal_type), "--max-dump-count", "1", "-O", "drop"], + capture_output=True, + text=True, + check=False, ) - burst_status = Status.PASS if worst_burst <= th.burstiness_max else Status.WARN - burst = Check( - "burstiness", - Severity.SHAPE, - burst_status, - f"peak/mean {worst_burst:.2f} (limit {th.burstiness_max:.2f})", - {"worst_peak_over_mean": round(worst_burst, 2), "limit": th.burstiness_max}, + if proc.returncode != 0: + raise RuntimeError(f"tsp -P pes failed: {proc.stderr.strip()}") + fields = {} + for line in (proc.stdout + proc.stderr).splitlines(): + key, eq, value = line.strip().partition(" = ") + # A key repeats once per CPB schedule; the last is the one H.222.0 sizes EB from. + if eq: + fields[key] = value + if not fields: + raise Refused("the stream carries no SPS") + return fields + + +def _hrd(sps: dict[str, str], prefix: str, value: str) -> tuple[float, float] | None: + """(BitRate bit/s, CpbSize bits) of the highest-numbered CPB a NAL HRD declares. + + H.264 E.2.2 and H.265 E.3.3: BitRate = (bit_rate_value_minus1 + 1) * 2^(6 + bit_rate_scale) + and CpbSize = (cpb_size_value_minus1 + 1) * 2^(4 + cpb_size_scale). H.222.0 2.14.3.1 + sizes EB from CpbSize[cpb_cnt_minus1], so the last schedule is the one that counts. + """ + rates = [int(v) for k, v in sps.items() if k.startswith(prefix) and "bit_rate_value_minus1" in k] + sizes = [int(v) for k, v in sps.items() if k.startswith(prefix) and "cpb_size_value_minus1" in k] + if not rates or not sizes: + return None + rate = (rates[-1] + 1) << (6 + int(sps[f"{value}bit_rate_scale"])) + size = (sizes[-1] + 1) << (4 + int(sps[f"{value}cpb_size_scale"])) + return rate, size + + +def avc_params(sps: dict[str, str]) -> Params: + """H.222.0 2.14.3.1 buffers from the SPS: level limits, overridden by a declared NAL HRD. + + Only the NAL HRD counts: EB sizes the byte stream, which a VCL HRD does not describe, + so a stream declaring only a VCL HRD takes the level defaults (H.264 E.2.2). + """ + profile, level = int(sps["profile_idc"]), int(sps["level_idc"]) + if (level == 11 and sps.get("constraint_set3_flag") == "1" and profile in (66, 77, 88)) or level == 9: + max_br, max_cpb, name = *AVC_LEVEL_1B, "1b" + elif level in AVC_LEVELS: + max_br, max_cpb = AVC_LEVELS[level] + name = f"{level // 10}.{level % 10}" + else: + raise Refused(f"AVC level_idc {level} has no Table A-1 entry") + if profile not in AVC_NAL_FACTOR: + raise Refused(f"AVC profile_idc {profile} has no Table A-2 cpbBrNalFactor") + # Absent a NAL HRD, BitRate is cpbBrNalFactor * MaxBR (H.264 E.2.2) and cpb_size is + # 1200 * MaxCPB (H.222.0 2.14.3.1). + declared = ( + _hrd(sps, "vui.nal_hrd.", "vui.nal_hrd.") if sps.get("vui.nal_hrd_parameters_present_flag") == "1" else None ) - return bitrate, burst - - -def check_inter_arrival(scan: Scan, clock: PcrClock) -> Check: - """Report the packet inter-arrival spread on the PCR clock (informational).""" - indices = sorted(i for packets in scan.pid_packets.values() for i, _ in packets) - gaps_us: list[float] = [] - prev = None - for i in indices: - t = clock.time_at(i) - if prev is not None and t >= prev: - gaps_us.append((t - prev) * 1e6) - prev = t - if not gaps_us: - return Check("inter-arrival", Severity.SHAPE, Status.WARN, "no packets timed", {}) - metrics = { - "mean_us": round(mean(gaps_us), 2), - "p95_us": round(percentile(gaps_us, 95), 2), - "max_us": round(max(gaps_us), 2), - } - return Check( - "inter-arrival", - Severity.SHAPE, - Status.PASS, - f"mean {metrics['mean_us']:.1f} us, p95 {metrics['p95_us']:.1f} us, max {metrics['max_us']:.1f} us", - metrics, + bit_rate, cpb = declared or (AVC_NAL_FACTOR[profile] * max_br, 1200 * max_cpb) + overhead = (0.004 + 1 / 750) * max(1200 * max_br, 2_000_000) + return Params( + label=f"AVC profile {profile} level {name}, " + + (f"NAL HRD {bit_rate / 1e6:.3f} Mb/s cpb {cpb / 1e3:.0f} kbit" if declared else "level defaults"), + # H.222.0 2.14.3.1: Rx = 1.2 * BitRate[SchedSelIdx]; MBS = BSmux + BSoh + 1200 * + # MaxCPB - cpb_size; EBS = cpb_size; Rbx = 1200 * MaxBR whatever the HRD declares. + rx=1.2 * bit_rate, + mb=(overhead + 1200 * max_cpb - cpb) / 8, + eb=cpb / 8, + rbx=1200 * max_br, + max_delay_s=10.0, ) -def check_tstd(analysis: dict, scan: Scan, clock: PcrClock, th: Thresholds) -> Check: - """Approximate T-STD transport-buffer check: TB fills on arrival, leaks at Rx. +def hevc_params(sps: dict[str, str]) -> Params: + """H.222.0 2.17.2 buffers from the SPS: tier and level limits, overridden by a declared NAL HRD.""" + prefix = "profile_tier_level.general_" + tier, profile, level = (int(sps[f"{prefix}{k}"]) for k in ("tier_flag", "profile_idc", "level_idc")) + if profile not in (1, 2, 3): + raise Refused(f"HEVC general_profile_idc {profile} has no CpbNalFactor here (Main, Main 10, Still only)") + limits = HEVC_LEVELS.get(level, (None, None))[tier] + if limits is None: + raise Refused(f"HEVC level_idc {level} tier {tier} has no Table A.8 entry") + max_br, max_cpb = limits + declared = ( + _hrd(sps, "vui.hrd.nal_hrd_parameters", "vui.hrd.") + if sps.get("vui.hrd.nal_hrd_parameters_present_flag") == "1" + else None + ) + # Absent a NAL HRD, BitRate is CpbBrVclFactor * MaxBR (H.265 E.3.3), which H.222.0 then + # scales by CpbBrNalFactor / CpbBrVclFactor; cpb_size is CpbBrNalFactor * MaxCPB. + bit_rate, cpb = declared or (1000 * max_br, HEVC_NAL_FACTOR * max_cpb) + rate = HEVC_NAL_FACTOR * max_br + return Params( + label=f"HEVC profile {profile} {'high' if tier else 'main'} tier level {level / 30:g}, " + + (f"NAL HRD {bit_rate / 1e6:.3f} Mb/s cpb {cpb / 1e3:.0f} kbit" if declared else "level defaults"), + # H.222.0 2.17.2: Rx = CpbBrNalFactor / CpbBrVclFactor * BitRate[SchedSelIdx]; MBS = + # BSmux + BSoh + CpbBrNalFactor * MaxCPB - cpb_size; EBS = cpb_size; Rbx = + # CpbBrNalFactor * MaxBR. + rx=HEVC_NAL_FACTOR / 1000 * bit_rate, + mb=((0.004 + 1 / 750) * max(rate, 2_000_000) + HEVC_NAL_FACTOR * max_cpb - cpb) / 8, + eb=cpb / 8, + rbx=rate, + max_delay_s=10.0, + ) + - This models only the transport buffer (TB) smoothing stage of ISO 13818-1 - 2.4.2, not the full multiplex/elementary buffer decode model. TB size is - fixed at 512 bytes; the leak rate Rx defaults per stream type. An overflow - means the stream delivers a PID's bytes faster than a receiver drains them. +def adts_frame(es: bytes, pos: int) -> tuple[int, float]: + """(length, duration) of the ADTS frame at `pos`.""" + h = es[pos : pos + 7] + if len(h) < 7 or h[0] != 0xFF or h[1] & 0xF0 != 0xF0: + raise Refused(f"no ADTS sync at ES offset {pos}") + rate_index = (h[2] >> 2) & 0x0F + if rate_index >= len(ADTS_RATES): + raise Refused(f"ADTS sampling_frequency_index {rate_index} at ES offset {pos}") + length = ((h[3] & 0x03) << 11) | (h[4] << 3) | (h[5] >> 5) + return length, 1024 * ((h[6] & 0x03) + 1) / ADTS_RATES[rate_index] + + +def mpeg_audio_frame(es: bytes, pos: int) -> tuple[int, float]: + """(length, duration) of the MPEG-1/2 audio frame at `pos`.""" + h = es[pos : pos + 4] + if len(h) < 4 or h[0] != 0xFF or h[1] & 0xE0 != 0xE0: + raise Refused(f"no MPEG audio sync at ES offset {pos}") + version_bits = (h[1] >> 3) & 0x03 + layer = 4 - ((h[1] >> 1) & 0x03) + rate_index = (h[2] >> 2) & 0x03 + kbps_index = h[2] >> 4 + if version_bits == 1 or layer == 4 or rate_index == 3 or kbps_index in (0, 15): + raise Refused(f"MPEG audio header at ES offset {pos} is reserved or free-format") + # MPEG-1, MPEG-2 (half rate) and MPEG-2.5 (quarter rate). + version = 1 if version_bits == 3 else 2 + rate = (44100, 48000, 32000)[rate_index] >> {3: 0, 2: 1, 0: 2}[version_bits] + bits = MPEG_AUDIO_KBPS[(version, layer)][kbps_index] * 1000 + pad = (h[2] >> 1) & 0x01 + if layer == 1: + return (12 * bits // rate + pad) * 4, 384 / rate + samples = 576 if layer == 3 and version == 2 else 1152 + return samples // 8 * bits // rate + pad, samples / rate + + +def ac3_frame(es: bytes, pos: int) -> tuple[int, float]: + """(length, duration) of the AC-3 or E-AC-3 frame at `pos`; 0 duration for a dependent substream.""" + h = es[pos : pos + 6] + if len(h) < 6 or h[0] != 0x0B or h[1] != 0x77: + raise Refused(f"no AC-3 sync at ES offset {pos}") + if h[5] >> 3 > 10: # bsid 11-16: E-AC-3 (A/52 Annex E) + blocks = (1, 2, 3, 6)[(h[4] >> 4) & 0x03] if h[4] >> 6 != 3 else 6 + fscod = h[4] >> 6 + rate = (48000, 44100, 32000)[fscod] if fscod != 3 else (24000, 22050, 16000)[(h[4] >> 4) & 0x03] + dependent = h[2] >> 6 == 1 + return ((((h[2] & 0x07) << 8) | h[3]) + 1) * 2, 0.0 if dependent else 256 * blocks / rate + fscod, code = h[4] >> 6, h[4] & 0x3F + if fscod == 3 or code >> 1 >= len(AC3_KBPS): + raise Refused(f"AC-3 fscod/frmsizecod reserved at ES offset {pos}") + kbps = AC3_KBPS[code >> 1] + words = (2 * kbps, 320 * kbps // 147 + (code & 1), 3 * kbps)[fscod] + return words * 2, 1536 / (48000, 44100, 32000)[fscod] + + +def opus_frame(es: bytes, pos: int) -> tuple[int, float]: + """(length, duration) of the Opus-in-TS access unit at `pos`: control header, then one packet.""" + if es[pos] != 0x7F or es[pos + 1] & 0xE0 != 0xE0: + raise Refused(f"no Opus control header at ES offset {pos}") + flags = es[pos + 1] + at = pos + 2 + size = 0 + while True: + size += es[at] + at += 1 + if es[at - 1] != 0xFF: + break + at += 2 * bool(flags & 0x10) + 2 * bool(flags & 0x08) + if flags & 0x04: + at += 1 + es[at] + # RFC 6716 3.1: the TOC byte's config sets the frame size, its code the frame count. + config, code = es[at] >> 3, es[at] & 0x03 + if config < 12: + frame_ms = (10, 20, 40, 60)[config % 4] + elif config < 16: + frame_ms = (10, 20)[config % 2] + else: + frame_ms = (2.5, 5, 10, 20)[config % 4] + frames = (1, 2, 2, es[at + 1] & 0x3F)[code] + return at - pos + size, frame_ms * frames / 1000 + + +def stream_kind(stream: Stream) -> str | None: + """Which T-STD model the stream takes, or None for one with no elementary stream buffers (sections, data).""" + kind = stream.stream_type + tags = stream.descriptors + if kind in STREAM_KINDS: + return STREAM_KINDS[kind] + if kind == 0x06 and tags & {"DVB_AC3_descriptor", "AC3_descriptor"}: + return "ac3-dvb" + if kind == 0x06 and f"registration:{int.from_bytes(b'Opus', 'big')}" in tags: + return "opus" + if kind == 0x06 and tags & {"DVB_enhanced_AC3_descriptor", "enhanced_AC3_descriptor"}: + raise Refused("E-AC-3 as DVB private data takes its buffer from ETSI TS 101 154, not modelled") + if kind in (0x01, 0x02, 0x10, 0x11, 0x1C, 0x20, 0x21, 0x42, 0xD1, 0xEA): + raise Refused(f"stream_type 0x{kind:02X} is audio/video this model has no parameters for") + return None + + +def stream_params(kind: str, stream: Stream) -> tuple[Params, object]: + """The stream's T-STD parameters, and its audio frame parser (None for video).""" + es = stream.es + if kind == "avc": + return avc_params(stream.sps), None + if kind == "hevc": + return hevc_params(stream.sps), None + if kind == "adts": + adts_frame(es, 0) + config = ((es[2] & 0x01) << 2) | (es[3] >> 6) + if config not in ADTS_CHANNELS: + raise Refused("ADTS channel_configuration 0 defers the layout to a PCE, which this model does not read") + channels = ADTS_CHANNELS[config] + rx, bs = next((rx, bs) for top, rx, bs in ADTS_BUFFERS if channels <= top) + return Params(label=f"ADTS AAC {channels} ch", rx=rx, b=bs), adts_frame + if kind == "opus": + # Borrowed: the Opus-in-TS draft gives Rx (2 Mb/s for 1-2 channels, as here) but no + # buffer size, so Opus is graded against ADTS's buffers for the same channel count. + config = next((tag[len("tag127:80") :][:2] for tag in stream.descriptors if tag.startswith("tag127:80")), "") + if not config or not 0 <= int(config, 16) <= 8: + raise Refused("Opus without a channel_config_code of 0-8 in its extension descriptor") + channels = int(config, 16) or 2 # 0 is dual mono + rx, bs = next((rx, bs) for top, rx, bs in ADTS_BUFFERS if channels <= top) + return Params(label=f"Opus {channels} ch (ADTS buffers)", rx=rx, b=bs), opus_frame + label, bs, parse = { + "mpeg-audio": ("MPEG audio", MPEG_AUDIO_BS, mpeg_audio_frame), + "ac3-atsc": ("AC-3 (ATSC)", AC3_ATSC_BS, ac3_frame), + "eac3-atsc": ("E-AC-3 (ATSC)", EAC3_ATSC_BS, ac3_frame), + "ac3-dvb": ("AC-3 (DVB)", AC3_DVB_BS, ac3_frame), + }[kind] + return Params(label=label, rx=AUDIO_RX, b=bs), parse + + +# NAL unit types that open a new access unit when they follow the last VCL NAL unit +# of the previous one (H.264 7.4.1.2.3, H.265 7.4.2.4.4), and the VCL types. +AU_OPENERS = {"avc": {6, 7, 8, 9, 14, 15, 16, 17, 18}, "hevc": {32, 33, 34, 35, 39, 41, 42, 43, 44, *range(48, 56)}} +VCL = {"avc": set(range(1, 6)), "hevc": set(range(0, 32))} +DELIMITER = {"avc": 9, "hevc": 35} + + +def video_units(stream: Stream, kind: str) -> list[AccessUnit]: + """One access unit per coded picture, each decoded at its PES's DTS (or PTS). + + An access unit opens at the first delimiter, parameter set or SEI after the previous + picture's last slice (H.264 7.4.1.2.3, H.265 7.4.2.4.4); H.222.0 2.14.1 and 2.17.1 + also put a delimiter in every one, so a stream without them is refused. A PES may + carry several access units, but only the first takes its timestamp; the rest would + need decoding times derived from the stream's own timing (2.4.2.3), which this model + does not do, so such a layout is refused rather than graded as one unit. """ - kinds: dict[int, float] = {} - for pid in analysis.get("pids", []): - if pid.get("video"): - kinds[pid["id"]] = th.video_leak_bps - elif pid.get("audio"): - kinds[pid["id"]] = th.audio_leak_bps - - overflows = 0 - worst_pid = None - worst_occ = 0.0 - for pid, leak_bps in kinds.items(): - packets = scan.pid_packets.get(pid, []) - occ = 0.0 - last_t = None - pid_over = 0 - pid_peak = 0.0 - for i, payload in packets: - t = clock.time_at(i) - if last_t is not None and t > last_t: - occ = max(0.0, occ - leak_bps * (t - last_t) / 8.0) - occ += payload - pid_peak = max(pid_peak, occ) - if occ > th.tb_size_bytes: - pid_over += 1 - last_t = t - overflows += pid_over - if pid_peak > worst_occ: - worst_occ = pid_peak - worst_pid = pid - if not kinds: - return Check("tstd", Severity.SHAPE, Status.WARN, "no video/audio PID to model", {}) + es = bytes(stream.es) + starts = [] + delimited = False + after_vcl = True # the first opener starts the first whole access unit + pos = es.find(b"\x00\x00\x01") + while 0 <= pos < len(es) - 3: + nal = (es[pos + 3] >> 1) & 0x3F if kind == "hevc" else es[pos + 3] & 0x1F + delimited |= nal == DELIMITER[kind] + if nal in AU_OPENERS[kind] and after_vcl: + # The zero_byte ahead of the start code belongs to this NAL unit (H.264 B.1.1, + # H.265 B.2.1), so to this access unit rather than the last. + starts.append(pos - 1 if pos and es[pos - 1] == 0 else pos) + after_vcl = False + elif nal in VCL[kind]: + after_vcl = True + pos = es.find(b"\x00\x00\x01", pos + 3) + if not delimited: + raise Refused("no access unit delimiter, which H.222.0 2.14.1/2.17.1 requires in every access unit") + pes = stream.pes + pes_at = [offset for offset, *_ in pes] + packet_at = [offset for _i, _h, body, offset in stream.packets if body] + packets = [index for index, _h, body, _o in stream.packets if body] + units: list[AccessUnit] = [] + stamped: set[int] = set() + for n, start in enumerate(starts): + k = bisect.bisect_right(pes_at, start) - 1 + _offset, pts, dts, index = pes[k] + if pts is None: + raise Refused(f"an access unit starts in the untimestamped PES at packet {index}") + if k in stamped: + raise Refused(f"the PES at packet {index} carries several access units, whose decode times are not derived") + stamped.add(k) + end = starts[n + 1] if n + 1 < len(starts) else stream.es_len + first = packets[bisect.bisect_right(packet_at, start) - 1] + units.append(AccessUnit(start, end, first, dts if dts is not None else pts)) + return units + + +def access_units(stream: Stream, parse) -> list[AccessUnit]: + """Audio: one access unit per codec frame.""" + units: list[AccessUnit] = [] + pes = [entry for entry in stream.pes if entry[1] is not None] + # A PES timestamp belongs to the first frame that starts in that PES (2.4.3.7). + starts = [(offset, index) for index, _h, body, offset in stream.packets if body] + keys = [offset for offset, _ in starts] + stamps = iter(pes) + nxt = next(stamps, None) + pos = 0 + # A header cut short by the end of the capture is not a malformed frame. + while pos + 8 <= len(stream.es): + try: + length, duration = parse(stream.es, pos) + except IndexError: + break # a header that runs past the end of the capture + if length <= 0: + raise Refused(f"zero-length audio frame at ES offset {pos}") + stamp = None + while nxt is not None and nxt[0] <= pos: + stamp = nxt[1] + nxt = next(stamps, None) + if duration == 0.0 and units: + units[-1].end = pos + length + else: + first = starts[bisect.bisect_right(keys, pos) - 1][1] + units.append(AccessUnit(pos, pos + length, first, stamp, duration)) + pos += length + return units + + +@dataclass +class Grade: + """What one elementary stream did in the model.""" + + pid: int + label: str + violations: dict[str, int] = field(default_factory=dict) + peaks: dict[str, float] = field(default_factory=dict) + worst_late_ms: float = 0.0 + worst_delay_s: float = 0.0 + graded_units: int = 0 + + def flag(self, name: str) -> None: + self.violations[name] = self.violations.get(name, 0) + 1 + + def peak(self, name: str, fill: float, size: float) -> None: + self.peaks[name] = max(self.peaks.get(name, 0.0), fill / size) + + +def simulate( + stream: Stream, params: Params, units: list[AccessUnit], segments: list[tuple[PcrClock, int, int]] +) -> Grade: + """Run one elementary stream through its buffers, with fresh buffers for each time base. + + `segments` are (clock, first packet, end packet) per time base: the timestamps on + either side of a signalled discontinuity are on different clocks. + """ + grade = Grade(stream.pid, params.label) + eps = 1e-9 + for clock, lo, hi in segments: + # TB: packet i arrives over [t(i), t(i+1)] and leaves at Rx. `leave` is when its + # last byte does; the bytes still in TB as the packet finishes arriving are its + # peak, and a stretch where TB never empties may not exceed a second. + deliveries: list[tuple[float, int, int, int]] = [] + leave = float("-inf") + busy_since = None + for index, header, body, offset in stream.packets: + if not lo <= index < hi: + continue + start, end = clock.time_at(index), clock.time_at(index + 1) + if leave <= start: + busy_since = start + leave = max(end, max(leave, start) + 188 * 8 / params.rx) + fill = (leave - end) * params.rx / 8 + grade.peak("TB", fill, TB_SIZE) + if fill > TB_SIZE + 0.5: + grade.flag("TB overflow") + if busy_since is not None and leave - busy_since > TB_EMPTY_S: + grade.flag("TB not emptied within 1 s") + busy_since = None + if header or body: + deliveries.append((leave, header, body, offset)) + + horizon = clock.time_at(hi) + removals: list[tuple[float, AccessUnit]] = [] + td = None + for unit in units: + if not lo <= unit.first_packet < hi: + continue + if unit.stamp is not None: + base = unit.stamp / PTS_HZ + stamped = base + round((clock.time_at(unit.first_packet) - base) / STAMP_WRAP_S) * STAMP_WRAP_S + if removals and stamped <= removals[-1][0]: + grade.flag("decode time does not advance") + td = max(stamped, removals[-1][0]) if removals else stamped + elif td is None: + continue # an audio frame ahead of the first timestamp has no decoding time + if td > horizon: + break + removals.append((td, unit)) + td += unit.duration_s + + # MB -> EB (video, leak method) or B (audio): walk deliveries and removals in time + # order. EB/B is tracked by ES offset: everything below `into` has reached it and + # everything below `out` has been removed, so an access unit whose bytes arrive + # after its decoding time passes through as underflow rather than lingering as fill. + into = out = deliveries[0][3] if deliveries else 0 + mb: list[list[int]] = [] # [header, payload] per delivered packet, FIFO + mb_header = mb_payload = 0 + headers: list[tuple[int, int]] = [] # audio: (ES offset, bytes) of PES headers held in B + b_header = 0 + now = float("-inf") + late: list[tuple[int, float]] = [] + + def settle(t: float, moved_from: int, rate: float) -> None: + # Record how late each underflowed unit finished arriving. + while late and into >= late[0][0]: + end, td = late.pop(0) + done = t if rate <= 0 else now + (end - moved_from) * 8 / rate + grade.worst_late_ms = max(grade.worst_late_ms, (done - td) * 1000) + + def leak(until: float) -> None: + nonlocal into, mb_header, mb_payload, now + if not params.video or until <= now or not mb_payload: + now = max(now, until) + return + room = params.eb - max(0, into - out) + max(0, out - into) + amount = min(params.rbx * (until - now) / 8, mb_payload, room) + before = into + remaining = amount + while remaining > eps and mb: + head = mb[0] + mb_header -= head[0] + head[0] = 0 + take = min(head[1], remaining) + head[1] -= take + remaining -= take + if head[1] <= eps: + mb.pop(0) + mb_payload -= amount + into += amount + settle(until, before, params.rbx) + grade.peak("EB", max(0, into - out), params.eb) + now = until + + events = sorted( + [(t, 1, n) for n, (t, *_rest) in enumerate(deliveries)] + [(t, 0, n) for n, (t, _u) in enumerate(removals)] + ) + for t, kind, n in events: + leak(t) + if kind == 1: + _t, header, body, offset = deliveries[n] + if params.video: + mb.append([header, body]) + mb_header += header + mb_payload += body + grade.peak("MB", mb_header + mb_payload, params.mb) + if mb_header + mb_payload > params.mb + 0.5: + grade.flag("MB overflow") + else: + if header: + headers.append((offset, header)) + b_header += header + into += body + settle(t, into, 0) + fill = max(0, into - out) + b_header + grade.peak("B", fill, params.b) + if fill > params.b + 0.5: + grade.flag("B overflow") + continue + _t, unit = removals[n] + grade.graded_units += 1 + grade.worst_delay_s = max(grade.worst_delay_s, t - clock.time_at(unit.first_packet)) + if t - clock.time_at(unit.first_packet) > params.max_delay_s: + grade.flag(f"held over {params.max_delay_s:g} s") + if into + eps < unit.end: + grade.flag("EB underflow" if params.video else "B underflow") + late.append((unit.end, t)) + out = max(out, unit.end) + while headers and headers[0][0] < unit.end: + b_header -= headers.pop(0)[1] + return grade + + +def check_tstd(ts_path: str, packet_size: int, scan: Scan) -> Check: + """Full T-STD buffer model (ISO 13818-1 2.4.2) for every audio and video stream. + + Each program's streams run on that program's PCR, one time base at a time: a + signalled discontinuity starts fresh buffers, since the timestamps before and + after it are on different clocks. + """ + with open(ts_path, "rb") as handle: + data = handle.read() + programs = read_programs(ts_path) + + grades: list[Grade] = [] + refused: dict[int, str] = {} + skipped: list[int] = [] + for pcr_pid, streams in programs: + kinds: dict[int, str] = {} + for stream in streams: + try: + kind = stream_kind(stream) + except Refused as err: + refused[stream.pid] = str(err) + continue + if kind is None: + skipped.append(stream.pid) + else: + kinds[stream.pid] = kind + modelled = {s.pid: s for s in streams if s.pid in kinds} + for pid, why in read_pes(data, packet_size, modelled).items(): + refused[pid] = why + del modelled[pid] + samples = scan.pcr_by_pid.get(pcr_pid, []) + bounds = [0, *sorted(i for i, _ in samples if i in scan.pcr_new_base), scan.total_packets] + segments = [(PcrClock([s for s in samples if lo <= s[0] < hi]), lo, hi) for lo, hi in zip(bounds, bounds[1:])] + segments = [segment for segment in segments if segment[0].ok()] + for stream in modelled.values(): + try: + if kinds[stream.pid] in ("avc", "hevc"): + if "hrd_management_valid" in stream.descriptors: + raise Refused("the HRD-scheduled MB to EB transfer (H.222.0 2.14.3.1) is not modelled") + stream.sps = read_sps(ts_path, stream.pid, 7 if kinds[stream.pid] == "avc" else 33) + params, parse = stream_params(kinds[stream.pid], stream) + if parse is None: + units = video_units(stream, kinds[stream.pid]) + else: + units = access_units(stream, parse) + except Refused as err: + refused[stream.pid] = str(err) + continue + grades.append(simulate(stream, params, units, segments)) + metrics = { - "tb_overflows": overflows, - "worst_pid": worst_pid, - "worst_peak_bytes": round(worst_occ, 1), - "tb_size_bytes": th.tb_size_bytes, + "streams": { + str(g.pid): { + "type": g.label, + "access_units": g.graded_units, + "violations": g.violations, + "peak_fill_pct": {k: round(v * 100, 1) for k, v in g.peaks.items()}, + "worst_underflow_late_ms": round(g.worst_late_ms, 1), + "worst_delay_s": round(g.worst_delay_s, 3), + } + for g in grades + }, + "refused": {str(pid): why for pid, why in refused.items()}, + "not_elementary": skipped, } - detail = f"TB overflows {overflows}, worst peak {worst_occ:.0f}B on PID {worst_pid} (TB {th.tb_size_bytes}B)" - status = Status.PASS if overflows == 0 else Status.WARN - return Check("tstd", Severity.SHAPE, status, detail, metrics) + faults = [f"PID {g.pid} {name} x{count}" for g in grades for name, count in sorted(g.violations.items())] + faults += [f"PID {pid} refused: {why}" for pid, why in refused.items()] + if not grades and not refused: + return Check("tstd", Severity.SHAPE, Status.WARN, "no audio/video stream to model", metrics) + if not faults and not any(g.graded_units for g in grades): + return Check("tstd", Severity.SHAPE, Status.WARN, "no access unit fell inside a modelled time base", metrics) + if faults: + return Check("tstd", Severity.SHAPE, Status.WARN, "; ".join(faults), metrics) + peaks = ", ".join(f"PID {g.pid} " + "/".join(f"{k} {v * 100:.0f}%" for k, v in g.peaks.items()) for g in grades) + return Check("tstd", Severity.SHAPE, Status.PASS, f"no overflow or underflow; peak fill {peaks}", metrics) def detect_packet_size(analysis: dict) -> int: @@ -660,30 +1176,17 @@ def check_duration_fidelity(captured_s: float, reference_s: float) -> Check: # ------------------------------------------------------------------------ driver -def analyze(ts_path: str, th: Thresholds, reference_seconds: float | None = None) -> list[Check]: +def analyze(ts_path: str, reference_seconds: float | None = None) -> list[Check]: """Run every check against `ts_path` and return the ordered results. `reference_seconds` (the source's PCR span, round-trip only) enables the duration-fidelity check that pins the exported stream's absolute rate. """ analysis = run_tsanalyze(ts_path) - ts = analysis["ts"] packet_size = detect_packet_size(analysis) scan = scan_packets(ts_path, packet_size) clock_by_pid = {pid: PcrClock(samples) for pid, samples in scan.pcr_by_pid.items()} - # The reference clock is the PID with the most PCR samples (the PCR PID). - main_clock = max(clock_by_pid.values(), key=lambda c: len(c.idx), default=PcrClock([])) - - # Nominal bitrate: total bytes clocked over the PCR span (the rate an IRD would - # play the stream at). Self-consistent with the PCR clock used everywhere else, - # and far more stable than tsanalyze's instantaneous PCR bitrate on a bursty - # capture. Fall back to tsanalyze only when there is no usable PCR clock. - span = main_clock.sec[-1] - main_clock.sec[0] if main_clock.ok() else 0.0 - if span > 0: - ts_bitrate = scan.total_packets * scan.packet_size * 8 / span - else: - ts_bitrate = float(ts.get("bitrate") or ts.get("pcr-bitrate") or 0) checks: list[Check] = [] checks.append(check_packet_size(analysis)) @@ -696,22 +1199,10 @@ def analyze(ts_path: str, th: Thresholds, reference_seconds: float | None = None checks.append(check_pcr_presence(analysis, clock_by_pid)) checks.append(check_pcr_monotonic(scan)) if reference_seconds is not None: - checks.append(check_duration_fidelity(span, reference_seconds)) + checks.append(check_duration_fidelity(pcr_span_seconds(scan), reference_seconds)) checks.append(check_service_descriptors(analysis)) - checks.append(check_pcr_repetition(scan, th)) - checks.append(check_pcr_jitter(scan, ts_bitrate, th)) - checks.append(check_null_ratio(analysis, th)) - - if main_clock.ok(): - bitrate, burst = check_bitrate_and_burstiness(scan, main_clock, ts_bitrate, th) - checks.append(bitrate) - checks.append(burst) - checks.append(check_inter_arrival(scan, main_clock)) - checks.append(check_tstd(analysis, scan, main_clock, th)) - else: - for name in ("bitrate-consistency", "burstiness", "inter-arrival", "tstd"): - checks.append(Check(name, Severity.SHAPE, Status.WARN, "not enough PCRs to build a clock", {})) + checks.append(check_tstd(ts_path, packet_size, scan)) return checks @@ -741,20 +1232,6 @@ def verdict(checks: list[Check], strict: bool) -> int: return 0 -def build_thresholds(args: argparse.Namespace) -> Thresholds: - """Assemble a Thresholds from parsed CLI arguments.""" - return Thresholds( - pcr_repetition_ms=args.pcr_repetition_ms, - pcr_jitter_us=args.pcr_jitter_us, - null_ratio_max=args.null_ratio_max, - bitrate_cov_max=args.bitrate_cov_max, - burstiness_max=args.burstiness_max, - tb_size_bytes=args.tb_size_bytes, - video_leak_bps=args.video_leak_bps, - audio_leak_bps=args.audio_leak_bps, - ) - - def main() -> int: """CLI entry point: parse args, analyze the TS, print the report, return the exit code.""" parser = argparse.ArgumentParser(description="MPEG-TS / IRD compliance analyzer") @@ -765,20 +1242,11 @@ def main() -> int: "--reference", help="source TS the capture was muxed from; enables the duration-fidelity check", ) - parser.add_argument("--pcr-repetition-ms", type=float, default=Thresholds.pcr_repetition_ms) - parser.add_argument("--pcr-jitter-us", type=float, default=Thresholds.pcr_jitter_us) - parser.add_argument("--null-ratio-max", type=float, default=Thresholds.null_ratio_max) - parser.add_argument("--bitrate-cov-max", type=float, default=Thresholds.bitrate_cov_max) - parser.add_argument("--burstiness-max", type=float, default=Thresholds.burstiness_max) - parser.add_argument("--tb-size-bytes", type=int, default=Thresholds.tb_size_bytes) - parser.add_argument("--video-leak-bps", type=float, default=Thresholds.video_leak_bps) - parser.add_argument("--audio-leak-bps", type=float, default=Thresholds.audio_leak_bps) args = parser.parse_args() - th = build_thresholds(args) try: reference_seconds = source_duration(args.reference) if args.reference else None - checks = analyze(args.ts, th, reference_seconds) + checks = analyze(args.ts, reference_seconds) except (RuntimeError, FileNotFoundError, json.JSONDecodeError) as err: print(f"error: {err}", file=sys.stderr) return 2 diff --git a/test/ts/pcr-timing.py b/test/ts/pcr-timing.py index b318fe3904..643613d997 100755 --- a/test/ts/pcr-timing.py +++ b/test/ts/pcr-timing.py @@ -80,6 +80,8 @@ def __init__(self): self.cc_on_empty = [] self.wraps = 0 self.new_base = set() # positions in self.pcr that state a new time base + self.pcr_pid = None # the PID graded, once keep_busiest_pcr_pid has run + self.other_pcr_pids = [] self._cc = {} self._dup = {} self._new_base = set() # PIDs whose next PCR states a new time base @@ -157,6 +159,21 @@ def feed(self, p, arrival=None): elif cc != prev: self.cc_on_empty.append((index, pid, prev, cc)) + def keep_busiest_pcr_pid(self): + """Grade one clock: the PID carrying the most PCRs. + + Each program may carry its own PCR, and two correct grids offset from one another + pool into one grid of half the interval that neither keeps. + """ + pids = collections.Counter(pid for _, _, _, pid in self.pcr) + self.pcr_pid = pids.most_common(1)[0][0] if pids else None + self.other_pcr_pids = sorted(p for p in pids if p != self.pcr_pid) + if not self.other_pcr_pids: + return + kept = [(k, e) for k, e in enumerate(self.pcr) if e[3] == self.pcr_pid] + self.new_base = {n for n, (k, _) in enumerate(kept) if k in self.new_base} + self.pcr = [e for _, e in kept] + def scan_file(path): scan = Scan() @@ -255,18 +272,6 @@ def check_continuity(scan, args): ) -def check_pcr_single_pid(scan, args): - """Every PCR must ride the one PID the PMT declares.""" - pids = collections.Counter(pid for _, _, _, pid in scan.pcr) - return ( - "pcr-single-pid", - HARD, - len(pids) <= 1, - f"PCR carried on {len(pids)} PID(s): " + ", ".join(f"{k} ({v})" for k, v in pids.most_common()), - {"pids": {str(k): v for k, v in pids.items()}}, - ) - - def check_value_interval(scan, args): """PCR values must be spaced within the repetition limit. @@ -493,16 +498,9 @@ def check_schedule(scan, args): hard = args.schedule_pct_min is not None severity = HARD if hard else SHAPE required = args.schedule_pct_min if hard else 99.0 - # One PID only. Two PIDs each on a correct grid, offset from one another, pool into a - # grid of half the interval with half the bytes in each slot, which is a schedule - # neither of them keeps. pcr-single-pid already fails a stream carrying two; grading - # the busier one keeps this check's answer about one clock rather than about both. - pids = collections.Counter(pid for _, _, _, pid in scan.pcr) - pid = pids.most_common(1)[0][0] if pids else None - on_pid = [(k, e) for k, e in enumerate(scan.pcr) if e[3] == pid] graded = [] skipped = 0 - for (_, a), (kb, b) in zip(on_pid, on_pid[1:]): + for kb, (a, b) in enumerate(zip(scan.pcr, scan.pcr[1:]), start=1): seconds = (b[1] - a[1]) / TICKS_PER_MS / 1000.0 # An interval spanning a signalled new time base measures nothing, as in the value # check. A non-positive one is a duplicate packet (legal) or a backwards clock (the @@ -513,7 +511,7 @@ def check_schedule(scan, args): graded.append(((b[0] - a[0]) * PKT, seconds)) if len(graded) < 3: # Asked to gate, an ungradable stream is a failure rather than a clean one. - return ("pcr-schedule", severity, not hard, "not measured (too few PCR intervals on one PID)", {}) + return ("pcr-schedule", severity, not hard, "not measured (too few PCR intervals)", {}) aggregate = 8.0 * sum(n for n, _ in graded) / sum(s for _, s in graded) rate = args.mux_rate if args.mux_rate else aggregate @@ -529,8 +527,8 @@ def check_schedule(scan, args): within = sum(1 for (n, _), w, s in zip(graded, want, slack) if abs(n - w) <= s) median_s = statistics.median(s for _, s in graded) detail = { - "pid": pid, - "other_pcr_pids": sorted(p for p in pids if p != pid), + "pid": scan.pcr_pid, + "other_pcr_pids": scan.other_pcr_pids, "count": len(graded), "skipped": skipped, "rate_bps": round(rate), @@ -565,15 +563,10 @@ def check_schedule(scan, args): ) -CHECKS = [ - check_sync, - check_continuity, - check_pcr_single_pid, - check_value_interval, - check_release, - check_position, - check_schedule, -] +# A file is graded for sync and continuity by compliance.py, through TSDuck's tsanalyze. +# A pipe cannot be: tsp in front of it would rebuffer the very arrivals `release` stamps. +LIVE_CHECKS = [check_sync, check_continuity] +CHECKS = [check_value_interval, check_release, check_position, check_schedule] def coincidence(scan, args): @@ -598,15 +591,11 @@ def coincidence(scan, args): def main(): - ap = argparse.ArgumentParser( - description=__doc__, formatter_class=argparse.RawDescriptionHelpFormatter - ) + ap = argparse.ArgumentParser(description=__doc__, formatter_class=argparse.RawDescriptionHelpFormatter) ap.add_argument("path", nargs="?", help="TS file; omit with --live to read stdin") ap.add_argument("--live", action="store_true", help="read stdin and stamp arrivals, grading release timing") ap.add_argument("--seconds", type=float, default=45.0, help="live capture window (default 45)") - ap.add_argument( - "--repetition-ms", type=float, default=40.0, help="max PCR value interval, TR 101 290 (default 40)" - ) + ap.add_argument("--repetition-ms", type=float, default=40.0, help="max PCR value interval, TR 101 290 (default 40)") ap.add_argument( "--release-ms", type=float, @@ -630,8 +619,7 @@ def main(): "--live-cover-pct", type=float, default=50.0, - help="share of the --seconds window a --live sample must span, or it is treated as " - "truncated (default 50)", + help="share of the --seconds window a --live sample must span, or it is treated as truncated (default 50)", ) ap.add_argument( "--drift-ms", @@ -684,7 +672,8 @@ def main(): else: ap.error("give a path or --live") - results = [check(scan, args) for check in CHECKS] + scan.keep_busiest_pcr_pid() + results = [check(scan, args) for check in (LIVE_CHECKS if args.live else []) + CHECKS] width = max(len(r[0]) for r in results) print(f"### PCR timing report - {scan.packets} packets, {len(scan.pcr)} PCR") @@ -721,8 +710,7 @@ def main(): "pcr_count": len(scan.pcr), "live": args.live, "checks": [ - {"name": n, "severity": s, "ok": ok, "headline": h, "detail": d} - for n, s, ok, h, d in results + {"name": n, "severity": s, "ok": ok, "headline": h, "detail": d} for n, s, ok, h, d in results ], }, f, diff --git a/test/ts/run.sh b/test/ts/run.sh index 0858b89dad..3ceb480b7f 100755 --- a/test/ts/run.sh +++ b/test/ts/run.sh @@ -7,7 +7,7 @@ # TSDuck + custom analyzer in compliance.py against the capture. The point is to # tell whether what the subscriber emits is something an Integrated # Receiver/Decoder would accept, and to quantify where it diverges (the exporter -# is VBR, emits no null packets, and paces PCR per frame). +# pads to a constant rate only once the catalog carries the source's mux rate). # # Modes: # ./run.sh # generate a clip, round-trip it, analyze @@ -160,7 +160,7 @@ have() { command -v "$1" >/dev/null 2>&1; } require_tools() { local missing=() t - for t in tsp tsanalyze python3; do + for t in tsp tsanalyze tstables python3; do have "$t" || missing+=("$t") done # ffmpeg + cargo are only needed for the round-trip, not for --analyze-only. @@ -174,7 +174,7 @@ require_tools() { fi if [[ ${#missing[@]} -gt 0 ]]; then echo "error: missing required tools: ${missing[*]}" >&2 - echo " TSDuck (tsp, tsanalyze) is required; install from https://tsduck.io" >&2 + echo " TSDuck (tsp, tsanalyze, tstables) is required; install from https://tsduck.io" >&2 exit 1 fi } diff --git a/test/ts/tstd-controls.py b/test/ts/tstd-controls.py new file mode 100755 index 0000000000..ccebb6f8f0 --- /dev/null +++ b/test/ts/tstd-controls.py @@ -0,0 +1,220 @@ +#!/usr/bin/env python3 +"""Positive and negative controls for the T-STD model in compliance.py. + +The positive control is a real broadcast encoder's output (a Kyrion contribution +feed: AVC High@4.0 plus two MPEG-1 Layer II tracks), which a broadcast chain only +carries because it is T-STD compliant. Each negative restamps that capture's PCRs +at a different constant rate with TSDuck's `pcradjust`, leaving every PES, PTS and +DTS untouched, so the only thing that changes is when the bytes arrive: + + as captured the broadcast itself must pass + restamped 1x PCRs rewritten at the capture's own rate must pass (the rewrite is not the fault) + 0.7x delivered too slowly EB and B underflow + 4x delivered too early TB and B overflow + 15x delivered in a burst TB and B overflow + +The video TB drains at 1.2x the bit rate the SPS's NAL HRD declares (1.935 Mb/s), +so any delivery well above real time overflows it. MB holds the level's whole CPB +less the declared one, ~3.6 MB, more than this 4 s capture carries, so no restamp +of it can overflow MB. + +A capture cannot be edited into the packet layouts the model has to get right, so +the rest are built here: a 10 Mb/s single-video stream carrying the Kyrion SPS, +with each AU in its own PES, a PCR packet between AUs and nulls elsewhere. + + synthetic as built must pass + two AUs, one PES an AU without its own timestamp refused + adaptation burst four payload-less packets after an AU TB overflow (they cost TB) + duplicate packet one AU's first packet sent twice (2.4.3.3) must pass (TB only, not MB) + +A model that passes everything, or fails everything, cannot tell these apart. +""" + +import os +import subprocess +import sys +import tempfile + +DIR = os.path.dirname(os.path.abspath(__file__)) +sys.path.insert(0, DIR) + +import compliance # noqa: E402 + +FIXTURE = os.path.join(DIR, "../../rs/moq-mux/src/container/ts/test_data/scte35/kyrion_dirtystart.ts") + +# The Kyrion SPS (High@4.0, NAL HRD 1.935 Mb/s CBR, 755 kbit CPB), so the synthetic +# streams get the same declared buffers. +SPS = bytes.fromhex("6764 0028 acd1 0078 044f de03 6a02 0202 8000 01f4 8000 7530 7500 0762 0002 e11f af7f 072a 1629 92") +VIDEO_PID, PMT_PID = 0x100, 0x1000 +RATE = 10_000_000 +SLOT_S = 188 * 8 / RATE + + +def crc32_mpeg(data: bytes) -> int: + """The MPEG-2 section CRC: polynomial 0x04C11DB7, not reflected, initial all ones.""" + crc = 0xFFFFFFFF + for byte in data: + crc ^= byte << 24 + for _ in range(8): + crc = (crc << 1) ^ 0x04C11DB7 if crc & 0x80000000 else crc << 1 + crc &= 0xFFFFFFFF + return crc + + +def section(table_id: int, extension: int, body: bytes) -> bytes: + """A long-form PSI section, CRC included.""" + head = bytes([table_id, 0xB0 | ((len(body) + 9) >> 8), (len(body) + 9) & 0xFF]) + head += extension.to_bytes(2, "big") + bytes([0xC1, 0, 0]) + return head + body + crc32_mpeg(head + body).to_bytes(4, "big") + + +def packet(pid: int, cc: int, payload: bytes = b"", pusi: bool = False, adaptation: bytes | None = None) -> bytes: + """One 188-byte packet; any room the payload leaves is adaptation-field stuffing.""" + room = 184 - len(payload) + if adaptation is None and room > 0: + adaptation = b"\x00" if room > 1 else b"" + if adaptation is not None: + adaptation = bytes([room - 1]) + adaptation.ljust(room - 1, b"\xff") if room > 1 else b"\x00" + afc = (2 if adaptation is not None else 0) | (1 if payload else 0) + head = bytes([0x47, (0x40 if pusi else 0) | pid >> 8, pid & 0xFF, afc << 4 | cc]) + return head + (adaptation or b"") + payload + + +def stamp(prefix: int, ticks: int) -> bytes: + """A 5-byte PES PTS/DTS field.""" + return bytes( + [ + prefix << 4 | (ticks >> 29) & 0x0E | 1, + (ticks >> 22) & 0xFF, + (ticks >> 14) & 0xFE | 1, + (ticks >> 7) & 0xFF, + (ticks << 1) & 0xFE | 1, + ] + ) + + +def synthetic(path: str, layout: str) -> None: + """Write a compliant 2 s stream, bent into `layout` at one access unit.""" + pat = section(0x00, 1, (1).to_bytes(2, "big") + (0xE000 | PMT_PID).to_bytes(2, "big")) + pmt = section(0x02, 1, bytes([0xE1, 0x00, 0xF0, 0x00, 0x1B, 0xE1, 0x00, 0xF0, 0x00])) + cc: dict[int, int] = {} + out: list[bytes] = [] + + def emit(pid: int, payload: bytes = b"", pusi: bool = False, adaptation: bytes | None = None) -> bytes: + if payload: + cc[pid] = (cc.get(pid, -1) + 1) & 0x0F + made = packet(pid, cc.get(pid, 0), payload, pusi, adaptation) + out.append(made) + return made + + def pes(units: list[bytes], dts_s: float) -> None: + ticks = round((1.0 + dts_s) * compliance.PTS_HZ) + data = bytes.fromhex("000001e0 0000 80c0 0a") + stamp(3, ticks) + stamp(1, ticks) + b"".join(units) + first = True + for at in range(0, len(data), 184): + made = emit(VIDEO_PID, data[at : at + 184], pusi=first) + if first and layout == "duplicate packet" and dts_s == 0.5 + 30 * 0.04: + out.append(made) + first = False + + frames = 50 + pending = None + for k in range(frames): + begin = len(out) + emit(0, b"\x00" + pat, pusi=True) + emit(PMT_PID, b"\x00" + pmt, pusi=True) + unit = b"\x00\x00\x00\x01\x09\xf0" + (b"\x00\x00\x00\x01" + SPS if k == 0 else b"") + unit += b"\x00\x00\x01\x65" + b"\xaa" * 250 + # Each AU decodes 0.5 s after its first byte arrives. + if layout == "two AUs, one PES" and k == 10: + pending = unit + else: + pes([pending, unit] if pending else [unit], 0.5 + (k - bool(pending)) * 0.04) + pending = None + if layout == "adaptation burst" and k == 20: + for _ in range(4): + emit(VIDEO_PID, adaptation=b"\x00") + slots = round(0.04 / SLOT_S) + while len(out) - begin < slots: + if len(out) - begin == slots // 2: + pcr = round((1.0 + len(out) * SLOT_S) * compliance.PCR_HZ) + field = bytes([0x10]) + ((pcr // 300) << 15 | 0x7E00 | pcr % 300).to_bytes(6, "big") + emit(VIDEO_PID, adaptation=field) + else: + out.append(packet(0x1FFF, 0, b"\xff" * 184)) + with open(path, "wb") as handle: + handle.write(b"".join(out)) + + +def restamp(scale: float): + """A builder that restamps the Kyrion capture's PCRs at `scale` times its own rate.""" + + def build(path: str) -> str: + rate = compliance.run_tsanalyze(FIXTURE)["ts"]["bitrate"] + subprocess.run( + ["tsp", "-I", "file", FIXTURE, "-P", "pcradjust", "--bitrate", f"{rate * scale:.0f}"] + + ["--ignore-pts", "--ignore-dts", "-O", "file", path], + check=True, + ) + return path + + return build + + +def built(layout: str): + """A builder for one synthetic layout.""" + + def build(path: str) -> str: + synthetic(path, layout) + return path + + return build + + +# (name, builder, expected): violations that must all appear, or "pass", or a refusal's text. +CASES = [ + ("as captured", lambda _path: FIXTURE, "pass"), + ("restamped 1x", restamp(1.0), "pass"), + ("0.7x", restamp(0.7), {"EB underflow", "B underflow"}), + ("4x", restamp(4.0), {"TB overflow", "B overflow"}), + ("15x", restamp(15.0), {"TB overflow", "B overflow"}), + ("synthetic", built("synthetic"), "pass"), + ("two AUs, one PES", built("two AUs, one PES"), "several access units"), + ("adaptation burst", built("adaptation burst"), {"TB overflow"}), + ("duplicate packet", built("duplicate packet"), "pass"), +] + + +def grade(path: str) -> compliance.Check: + """compliance.py's tstd verdict on one file.""" + size = compliance.detect_packet_size(compliance.run_tsanalyze(path)) + return compliance.check_tstd(path, size, compliance.scan_packets(path, size)) + + +def main() -> int: + """Run every case and exit non-zero if any verdict is not the expected one.""" + failed = 0 + with tempfile.TemporaryDirectory() as tmp: + for n, (name, build, expected) in enumerate(CASES): + check = grade(build(os.path.join(tmp, f"{n}.ts"))) + streams = check.metrics.get("streams", {}) + seen = {v for s in streams.values() for v in s["violations"]} + if expected == "pass": + # Not vacuous: every stream must have had access units to grade. + ok = check.status == compliance.Status.PASS and all(s["access_units"] for s in streams.values()) + elif isinstance(expected, str): + ok = any(expected in why for why in check.metrics.get("refused", {}).values()) + else: + ok = expected <= seen + failed += not ok + want = expected if isinstance(expected, str) else ", ".join(sorted(expected)) + print(f" {'ok ' if ok else 'FAIL'} {name:<18} want {want:<40} got {check.detail}") + if failed: + print(f"tstd controls: {failed} of {len(CASES)} cases wrong", file=sys.stderr) + return 1 + print("tstd controls: PASS") + return 0 + + +if __name__ == "__main__": + sys.exit(main())