#!/usr/bin/env python3 """Стенд подачи TX-аудио: проверяет планировщик TX_CHRONO под потерями и задержкой. Зачем: всплески своего сигнала на передаче оказались не сбоем клиента и не блокировкой UI, а зернистостью НАШЕГО запроса — маркеры уходили через 40 или 60 мс при звуке на 42.7 мс в каждом, и на трёхтактном интервале очередь DUC пересыхала. Здесь проверяется новый планировщик: квант в один блок TXA, абсолютные дедлайны, бухгалтерия Owed/InFlight с прощением зависших кредитов. Клиент отвечает на маркеры НУМЕРОВАННОЙ РАМПОЙ, а не тишиной: осушение очереди — не единственный способ испортить звук. Чтение неготовых данных, потерянный или задвоенный кусок уходят в эфир молча, и видно их только сверкой с эталоном. ★Период рампы простой и не связан с квантом (1021): будь он равен 512 или кратен ему, потеря ровно одного кванта дала бы последовательность, неотличимую от правильной — тест прошёл бы при том самом дефекте, который ищет. Сверка содержимого — по сырому отводу приложения (моно-кадры ДО интерполятора): EWSDR_TXAUDIO_DUMP=/tmp/txaudio.f32 ./bin/x86_64-linux/ewsdr Сценарии (--scenario): clean — здоровый клиент, ответ сразу rtt — ответ с задержкой (--rtt-ms), окно обязано вырасти без прощений drop — потерять N ответов вразнобой (--drops), система обязана нагнать freeze — заморозить N слотов навсегда, остальные обслуживать: конвейер продолжает отдавать звук медленнее реального времени — ровно тот случай, который сторож по одному лишь молчанию не ловит latch — потерять четыре ответа подряд: выход из полной защёлки Запуск (приложение работает, TCI-сервер поднят, передатчик не в эфире): test/tci/tx_chrono_bench.py --scenario freeze --freeze 2 --rtt-ms 30 """ import argparse import math import os import struct import sys import time sys.path.insert(0, os.path.dirname(os.path.abspath(__file__))) from live_rx_test import WS, split_commands, parse_command # noqa: E402 HDR = struct.Struct("<16I") ST_TX_AUDIO = 3 ST_TX_CHRONO = 4 RAMP_PERIOD = 1021 # простое, взаимно простое с квантом 512 RAMP_AMPL = 0.25 # ★внутри ±1: дальше по тракту ScaleIQ24 зажимает SYNC = [0.9, -0.9, 0.8, -0.8, 0.7, -0.7] # уникальная синхропоследовательность def ramp_value(index): return RAMP_AMPL * ((index % RAMP_PERIOD) / RAMP_PERIOD * 2.0 - 1.0) class Bench: def __init__(self, ws, rx, args): self.ws = ws self.rx = str(rx) self.args = args self.vfo_hz = None self.markers = [] # (время, запрошено кадров) self.answers = [] # (время, отдано кадров) self.sent_index = 0 # позиция в эталонном сигнале self.sync_done = False self.pending = [] # отложенные ответы: (когда_слать, кадров) self.dropped = 0 self.frozen = 0 self.chan = 2 self.rate = 48000 self.fmt = 3 self.failures = [] self.notes = [] # ── транспорт ──────────────────────────────────────────────────────── def pump(self, seconds): deadline = time.monotonic() + seconds while True: now = time.monotonic() self.flush_pending(now) if now >= deadline: return frame = self.ws.recv_frame(min(0.005, max(0.0, deadline - now))) if frame is None: continue opcode, payload = frame if opcode == 1: for command in split_commands(payload.decode("utf-8", "replace")): self.on_text(command) elif opcode == 2: self.on_binary(payload) def on_text(self, command): name, args = parse_command(command) if name == "vfo" and len(args) >= 3 and args[0] == self.rx and args[1] == "0": try: self.vfo_hz = int(args[2]) except ValueError: pass # ── ядро: ответ на маркер ──────────────────────────────────────────── def on_binary(self, payload): if len(payload) < HDR.size: return head = HDR.unpack(payload[:HDR.size]) receiver, rate, fmt, _codec, _crc, length, stype, chan = head[:8] if stype != ST_TX_CHRONO or receiver != int(self.rx): return self.rate, self.fmt, self.chan = rate, fmt, max(1, chan) frames = length // self.chan # length — значения всего блока self.markers.append((time.monotonic(), frames)) n = len(self.markers) if self.args.scenario == "drop" and n in self.args.drop_set: self.dropped += 1 return if self.args.scenario == "latch" and self.args.latch_at <= n < self.args.latch_at + 4: self.dropped += 1 return if self.args.scenario == "freeze" and self.frozen < self.args.freeze: self.frozen += 1 return # слот зависает навсегда due = time.monotonic() + self.args.rtt_ms / 1000.0 self.pending.append((due, frames)) def flush_pending(self, now): while self.pending and self.pending[0][0] <= now: _due, frames = self.pending.pop(0) self.send_audio(frames) def send_audio(self, frames): values = [] for _ in range(frames): if not self.sync_done and self.sent_index < len(SYNC): v = SYNC[self.sent_index] else: v = ramp_value(self.sent_index - len(SYNC)) self.sent_index += 1 if self.sent_index >= len(SYNC): self.sync_done = True # ★При двух каналах MSHV шлёт СОСЕДНИЕ отсчёты своего буфера 96 кГц, # а приёмная сторона усредняет пару — это децимация, не стерео. # Повторяем это же: два одинаковых значения на кадр дают после # усреднения ровно его. values.extend([v] * self.chan) body = struct.pack(f"<{len(values)}f", *values) head = [0] * 16 head[0] = int(self.rx) head[1] = self.rate head[2] = self.fmt head[5] = len(values) head[6] = ST_TX_AUDIO head[7] = self.chan self.ws.send_frame(2, HDR.pack(*head) + body) self.answers.append((time.monotonic(), frames)) # ── сценарий ───────────────────────────────────────────────────────── def run(self): for _ in range(5): self.ws.send_text(f"vfo:{self.rx},0;") self.pump(0.4) if self.vfo_hz is not None: break if self.vfo_hz is None: self.failures.append("сервер не ответил vfo — приёмника нет") return self.ws.send_text("audio_samplerate:48000;") self.ws.send_text("audio_stream_sample_type:float32;") self.ws.send_text(f"audio_stream_channels:{self.args.channels};") self.ws.send_text(f"audio_stream_samples:{self.args.block};") self.ws.send_text("tx_stream_audio_buffering:50;") self.pump(0.3) self.ws.send_text(f"trx:{self.rx},true,tci;") self.pump(self.args.seconds) self.ws.send_text(f"trx:{self.rx},false;") self.pump(0.5) # ── разбор ─────────────────────────────────────────────────────────── def report(self): print(f"\n=== сценарий {self.args.scenario}, RTT {self.args.rtt_ms} мс, " f"каналов {self.args.channels}, блок клиента {self.args.block}") if not self.markers: self.failures.append("маркеров TX_CHRONO не пришло вовсе") return quanta = sorted({m[1] for m in self.markers}) periods = [1000.0 * (b[0] - a[0]) for a, b in zip(self.markers, self.markers[1:])] periods.sort() q = self.markers[0][1] nominal = 1000.0 * q / self.rate span = self.markers[-1][0] - self.markers[0][0] asked = sum(m[1] for m in self.markers) given = sum(a[1] for a in self.answers) print(f" квант: {quanta} кадров (номинальный период {nominal:.2f} мс)") print(f" маркеров {len(self.markers)} за {span:.2f} с, " f"ответов {len(self.answers)}, потеряно намеренно " f"{self.dropped + self.frozen}") if periods: print(f" период маркеров: медиана {periods[len(periods)//2]:.2f} " f"макс {periods[-1]:.2f} мс") print(f" запрошено {asked} кадров, отдано {given} " f"({given / max(span, 1e-9):.0f} кадров/с при номинале {self.rate})") # 1. Зернистость: систематических 2P/3P быть не должно. if periods: long_share = sum(1 for p in periods if p > 1.6 * nominal) / len(periods) if long_share > 0.10: self.failures.append( f"{100*long_share:.0f}% интервалов маркеров длиннее 1.6 периода " f"— зернистость запроса осталась") else: print(f" ok доля длинных интервалов {100*long_share:.1f}%") # 2. Темп подачи: после любых потерь система обязана нагнать. if span > 2.0: rate_got = given / span if rate_got < 0.97 * self.rate: self.failures.append( f"подача {rate_got:.0f} кадров/с ниже номинала {self.rate} " f"— планировщик не нагнал потери") else: print(f" ok темп подачи {rate_got:.0f} кадров/с") # 3. Защёлка: маркеры не должны прекратиться до конца сценария. gap = self.markers[-1][0] last_gap = 1000.0 * (self.answers[-1][0] - self.markers[-1][0]) if self.answers else 0 tail = span - (self.markers[-1][0] - self.markers[0][0]) if periods and periods[-1] > 20 * nominal: self.failures.append( f"пауза в выдаче маркеров {periods[-1]:.0f} мс " f"— похоже на защёлку по InFlight") _ = gap, last_gap, tail def verify_audio(self, path): if not path or not os.path.exists(path): self.notes.append(f"отвод содержимого не включён (нет {path}) — " f"проверены только темп и зернистость") return import array data = array.array("f") with open(path, "rb") as fh: raw = fh.read() data.frombytes(raw[:len(raw) - len(raw) % 4]) got = list(data) if len(got) < len(SYNC) + 100: self.failures.append(f"в отводе всего {len(got)} отсчётов") return # ★Выравнивание — по ВСЕЙ синхропоследовательности и ОДИН раз. # По первому ненулевому отсчёту нельзя: так замаскируется ровно то, # что ищем — потерянный или лишний начальный участок. offset = -1 for i in range(0, len(got) - len(SYNC)): if all(abs(got[i + j] - SYNC[j]) < 0.02 for j in range(len(SYNC))): offset = i + len(SYNC) break if offset < 0: self.failures.append("синхропоследовательность в отводе не найдена") return zeros = sum(1 for v in got[:offset - len(SYNC)] if abs(v) < 1e-6) print(f" синхронизация на позиции {offset - len(SYNC)}, " f"перед ней {zeros} нулевых отсчётов") bad = 0 first_bad = None n = min(len(got) - offset, self.sent_index - len(SYNC)) for i in range(n): if abs(got[offset + i] - ramp_value(i)) > 0.02: bad += 1 if first_bad is None: first_bad = i if bad: self.failures.append( f"содержимое разошлось с эталоном на {bad} из {n} отсчётов, " f"первое расхождение на {first_bad} — вставка, пропуск или " f"чтение неготовых данных") else: print(f" ok содержимое совпало с эталоном на {n} отсчётах") def main(): ap = argparse.ArgumentParser(description=__doc__, formatter_class=argparse.RawDescriptionHelpFormatter) ap.add_argument("--host", default="127.0.0.1") ap.add_argument("--port", type=int, default=40001) ap.add_argument("--rx", type=int, default=0) ap.add_argument("--seconds", type=float, default=10.0) ap.add_argument("--scenario", default="clean", choices=["clean", "rtt", "drop", "freeze", "latch"]) ap.add_argument("--rtt-ms", type=float, default=0.0) ap.add_argument("--drops", default="20,55,90") ap.add_argument("--freeze", type=int, default=2) ap.add_argument("--latch-at", type=int, default=40) ap.add_argument("--channels", type=int, default=2, choices=[1, 2]) ap.add_argument("--block", type=int, default=2048) ap.add_argument("--audio-dump", default="", help="файл EWSDR_TXAUDIO_DUMP для сверки содержимого") args = ap.parse_args() args.drop_set = {int(x) for x in args.drops.split(",") if x.strip()} if args.scenario == "rtt" and args.rtt_ms == 0: args.rtt_ms = 30.0 ws = WS(args.host, args.port) ws.connect() bench = Bench(ws, args.rx, args) try: bench.run() finally: ws.close() bench.report() bench.verify_audio(args.audio_dump) for note in bench.notes: print(f" .. {note}") if bench.failures: print("\nПРОВАЛЕНО:") for f in bench.failures: print(f" FAIL {f}") return 1 print("\nвсё сошлось") return 0 if __name__ == "__main__": sys.exit(main())