mirror of
https://git.vladimir.cc/vladimir/ewsdr.git
synced 2026-08-25 20:37:33 +00:00
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
This commit is contained in:
Executable
+329
@@ -0,0 +1,329 @@
|
||||
#!/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())
|
||||
Reference in New Issue
Block a user