Files
ewsdr/test/tci/tx_chrono_bench.py
ew8bakandClaude Opus 5 d2096976df fix(tci): всплески на передаче — TX-аудио просилось общим тиком 20 мс
Посреди передачи из MSHV на водопаде появлялись всплески своего сигнала.
Цепочка: очередь DUC пустеет дольше подушки отправителя (DUC_FIFO_THROTTLE
= 2000 отсчётов = 10.4 мс) → FIFO радио сохнет → модуляция обрывается → в
эфире остаётся голая несущая на частоте гетеродина DUC, в стороне от тона
ровно на звуковой сдвиг. Доказано pcap-съёмом: шесть всплесков в дампе —
ровно столько, сколько видел оператор, и каждый стоит за паузой 10.3-23.5 мс,
а паузы 8 мс и короче не дали ни одного.

Виноват не клиент и не блокировка UI, а зернистость НАШЕГО запроса. Слой
первый: PushTxChrono жил на общем тике сервера 20 мс, а просил блок клиента
целиком (2048 отсчётов = 42.7 мс) — маркер выходил через два или три тика,
то есть через 40 или 60 мс. Слой второй: MSHV отвечает пачками по 4-5 блоков
раз в ~44 мс (STREAM_C = 4096 при 96 кГц), и мелкий квант этого не лечит —
нужен запас не меньше пачки.

Сделано:
* квант запроса = один блок TXA (512 отсчётов движка), а не блок клиента;
* свой поток-планировщик TxTickLoop с АБСОЛЮТНЫМИ дедлайнами (опоздание
  одного пробуждения не сдвигает сетку); общий тик маркеров больше не шлёт;
* бухгалтерия Owed/InFlight в кадрах на канал, гасится по k до интерполятора;
  потолок долга обязан быть выше окна в полёте (TCI_TX_OWED_HEADROOM_Q),
  иначе связывающим становится он и подача падает до 58% реального времени
  при полностью исправном клиенте;
* SendBinNow: маркеры пишутся в сокет напрямую под FWriteLock, минуя очередь
  (та выпускается лишь на пробуждении потока клиента, TCI_POLL_MS = 20 мс —
  вдвое больше кванта, и подача снова рвалась);
* FReapLock: планировщик TX — новый поток, а правило «клиента освобождает
  только тик-поток» держалось на том, что им же он и пользуется;
* аванс под зернистость клиента: измеряется по ПЕРИОДУ между пачками (размер
  пачки зависит от того, сколько мы запросили ⇒ положительная обратная связь),
  переживает конец передачи, умеет уменьшаться по выдержке TCI_TX_LEAD_DOWN_MS,
  потолок TCI_TX_LEAD_MAX_MS;
* старт передачи: KickTxTick будит планировщика на фронте PTT, TxPreWarm шлёт
  один маркер ДО SetMOX (41 мс раздумий клиента накладываются на нашу же
  подготовку тракта) под гейтом «передатчик свободен и чужого источника нет»,
  PrimeDUCIQ для источника TCI растянут до Max(6096, аванс×4) — путь микрофона
  радио, CW и web не затронут;
* посев аванса TCI_TX_LEAD_DEF_MS = 50 мс, пока про клиента ничего не известно:
  обучение к первому осушению физически не успевает.

Монотонные часы одного источника для всех потоков — PlatformUtils.MonotonicUs
(абсолютные дедлайны не терпят часов, способных прыгнуть от NTP).

На железе: опасных осушений посреди передачи НОЛЬ (было 12 за 11 с), четыре
передачи из пяти вообще без единого, включая старт; прогон 15:21 чист везде,
в том числе на первой передаче после подключения. Всплесков оператор больше
не видит.

Приборы: TxTrace.pas (EWSDR_TXTRACE=1) и стенд test/hpsdr — кольцевой tcpdump
capture.sh, разбор дампа pcap_tx_scan.py, разбор трассы txtrace_scan.py.
Стенд test/tci: часть F «Пейсинг TX», 260 проверок, провалов нет; живой клиент
с рампой, RTT и потерями — test/tci/tx_chrono_bench.py.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_014gmVQnna1i4EbSZ2VGm6KD
2026-08-23 22:06:10 +03:00

330 lines
16 KiB
Python
Executable File
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
#!/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())