Files
ewsdr/TCIServer.pas
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

1660 lines
71 KiB
ObjectPascal

unit TCIServer;
{
TCIServer.pas — транспорт TCI: WebSocket-сервер (роль сервера играем мы,
как ExpertSDR3; клиенты — логгеры, скиммеры, программы цифровых видов).
Что делает:
• слушает TCP-порт (умолчание 40001), принимает HTTP-Upgrade на WebSocket
по любому пути (клиенты ходят на ws://host:40001/);
• режет входящие текстовые фреймы на команды и отдаёт их наверх
(OnCommand) прямо в потоке клиента — маршалинг в поток контроллера
делает TCIAdapter, как это устроено у CAT;
• рассылает строки всем клиентам (Broadcast) — сервер TCI обязан
синхронизировать всех подключённых (§3.5);
• тикает OnTick (умолчание 20 мс) — по нему адаптер шлёт показания
измерителей с индивидуальным для каждого клиента периодом.
Отправка НИКОГДА не блокирует того, кто зовёт Send/Broadcast: строка кладётся
в очередь клиента, а в сокет её пишет ЕГО СОБСТВЕННЫЙ поток (Flush в цикле
HandleClient, recv просыпается каждые TCI_POLL_MS). Иначе медленный клиент
останавливал бы UI-поток на секунду за раз — уведомления рождаются в OnState,
то есть внутри Changed() контроллера, — а общий поток отправки задерживал бы
на его таймаут ещё и всех остальных клиентов.
Владение объектом клиента: создаёт accept-поток, освобождает ТОЛЬКО тик-поток
(ReapClients) и только после того, как клиентский поток честно вышел. Никто
больше клиентов не освобождает — поэтому указатель, взятый под FClientLock,
остаётся валидным, пока тик-поток не сделает следующий проход. На остановке
освобождает Stop, но лишь дождавшись выхода ВСЕХ клиентских потоков.
Два счёта соединений. Слот из TCI_MAX_CLIENTS занимает только клиент,
прошедший handshake (Up); сокет до handshake живёт в общем массиве
(TCI_MAX_SOCKETS) и убивается по таймауту TCI_HANDSHAKE_MS. Иначе восемь
молчащих TCP-соединений навсегда закрывали дверь настоящим клиентам.
Бинарные фреймы (потоки IQ/аудио, §3.4) ходят в обе стороны: блоки наружу
кладутся в отдельное кольцо клиента (SendBin) и уходят его же потоком вместе
с командами, входящие собираются из фрагментов и отдаются наверх (OnBinary)
— там их разбирает TCIAdapter. У двух очередей разная политика переполнения:
команду терять нельзя (клиент выбрасывается), блок потока — можно и нужно
(теряется самый старый), иначе отставший скиммер рвал бы себе управление.
Сокеты и WS-фреймы переиспользованы из веб-подсистемы (WebUtils/WsClient):
тот же код handshake и та же схема «поток на клиента + accept-поток», что в
WebServer.
Авторизации у TCI нет by design. Порт слушается там, где сказано в
настройках; умолчание — 127.0.0.1, чтобы наружу он не торчал без спроса.
Отсюда же отказ браузерным клиентам (заголовок Origin): страница, открытая
в браузере, иначе дотянулась бы до петлевого порта и до передатчика.
}
{$IFDEF FPC}
{$MODE Delphi}
{$LONGSTRINGS ON}
{$ENDIF}
interface
uses
Classes, SysUtils,
WebUtils, WsClient, TCIProtocol, PlatformUtils
{$IFDEF WINDOWS}, Windows, WinSock2{$ELSE}, Sockets{$ENDIF},
SyncObjs; // ← после платформенных юнитов (конфликт идентификатора Create)
const
TCI_MAX_CLIENTS = 8; // прошедших handshake (слоты протокола)
TCI_MAX_SOCKETS = 32; // всего сокетов, включая ещё не поднявшиеся
TCI_TICK_MS = 20; // период OnTick (сенсоры троттлятся адаптером)
// ★TX-аудио тикает ОТДЕЛЬНО от сенсоров. Общий тик в 20 мс не может выдавать
// запросы с шагом в один TXA-блок (10.667 мс): маркеры выходили через два или
// три тика, то есть через 40 или 60 мс, и на каждом трёхтактном интервале
// очередь DUC пересыхала — это и были всплески на передаче. Планировщик ниже
// ведёт АБСОЛЮТНЫЕ дедлайны, период ему задаёт адаптер (квант / частота
// клиента), а опоздание одного пробуждения не сдвигает всю сетку.
TCI_TX_TICK_MIN_MS = 1; // ниже не опускаемся: спать точнее всё равно не выйдет
TCI_TX_TICK_IDLE_MS = 20; // не передаём — будим планировщик редко
TCI_WS_GUID = '258EAFA5-E914-47DA-95CA-C5AB0DC85B11';
TCI_SEND_TIMEOUT = 300; // мс на SockSend, иначе клиент считается мёртвым
TCI_HANDSHAKE_MS = 5000; // мс на HTTP-запрос от подключившегося
TCI_POLL_MS = 20; // на столько recv клиента засыпает между кадрами
TCI_OUT_MAX = 4000; // потолок очереди отправки на клиента (строк)
TCI_OUT_CHUNK = 3800; // склейка очереди в один фрейм, символов
TCI_MSG_MAX = 65536; // потолок собираемого из фрагментов сообщения
TCI_STOP_KILL_MS = 500; // как часто добиваем клиентов, ожидая их выхода
// Очередь бинарных блоков (§3.4) на клиента. Переполнение здесь НЕ повод
// рвать соединение, в отличие от очереди команд: поток — это данные
// реального времени, и клиент, не успевший забрать блок, должен потерять
// именно блок. Выбрасываем самый старый: свежий звук полезнее протухшего.
TCI_BIN_QUEUE = 48;
type
TTCIServer = class;
{ Один подключённый клиент: WS-сокет, его личные подписки и очередь
отправки. Подписки на сенсоры в TCI индивидуальны (RX_SENSORS_ENABLE
«отправляется только клиентом»), поэтому живут здесь, а не в адаптере.
Параметры потоков (§4.3) — тоже клиентские.
Подписки и параметры потоков пишет поток клиента, а читает тик-поток,
поэтому и те и другие ходят через FStateLock: набор «включено + период +
последняя отправка» обязан меняться и читаться целиком. }
TTCIClient = class
private
FWs: TWsClient;
FUp: Boolean; // handshake прошёл: клиент занимает слот
FReady: Boolean; // пачка инициализации отправлена
FStateLock: TCriticalSection;
FRxSensors: Boolean;
FRxSensorsMs: Integer;
FRxSensorsAt: QWord; // тик последней отправки
FTxSensors: Boolean;
FTxSensorsMs: Integer;
FTxSensorsAt: QWord;
// Параметры бинарных потоков (§3.4): по спецификации это настройки
// КЛИЕНТА, а не устройства — свои у каждого подключения.
FIQRate: Integer;
FAudioRate: Integer;
FAudioSamples: Integer;
FAudioChannels: Integer;
FAudioSampleType: string;
FTxBuffering: Integer;
// Очередь отправки: пишут любые потоки, читает тик-поток.
FOutLock: TCriticalSection;
FOut: array of string;
FOutCount: Integer;
// Очередь бинарных блоков потоков. Отдельная от командной: у них разная
// политика переполнения (команду терять нельзя, блок потока — можно) и
// разные производители (блоки кладёт DSP-поток).
FBinOut: array[0..TCI_BIN_QUEUE-1] of TBytes;
FBinHead: Integer; // куда класть
FBinTail: Integer; // откуда брать
FBinDropped: LongInt; // сколько блоков выброшено (диагностика)
// ★Запись в сокет клиента. Обычно пишет только его собственный поток, но
// маркеры TX_CHRONO шлёт планировщик TX — им очередь не годится: она
// выпускается лишь на пробуждении потока клиента (recv с таймаутом
// TCI_POLL_MS), то есть запрос на модуляцию задерживался бы на все 20 мс.
// При кванте в 10.7 мс это вдвое больше самого кванта, и подача клиента
// становилась рваной ровно так же, как от старого тика 20 мс.
FWriteLock: TCriticalSection;
FDead: Boolean; // сокет уже не пишется — гасим соединение
FKilled: Boolean; // shutdown сокета уже сделан
FClosed: Boolean; // клиентский поток вышел (можно освобождать)
function GetReady: Boolean;
procedure SetReady(V: Boolean);
function GetIQRate: Integer; procedure SetIQRate(V: Integer);
function GetAudioRate: Integer; procedure SetAudioRate(V: Integer);
function GetAudioSamples: Integer; procedure SetAudioSamples(V: Integer);
function GetAudioChannels: Integer; procedure SetAudioChannels(V: Integer);
function GetAudioSampleType: string; procedure SetAudioSampleType(const V: string);
function GetTxBuffering: Integer; procedure SetTxBuffering(V: Integer);
public
constructor Create(AWs: TWsClient);
destructor Destroy; override;
{ Строку в очередь клиенту. False — соединение уже мертво. Не блокирует. }
function Send(const S: string): Boolean;
{ Блок бинарного потока (заголовок + сэмплы) в очередь. Зовётся из
DSP-потока, поэтому только копирование под коротким локом: сеть тут не
трогается. False — клиент мёртв (блок никуда не пошёл). }
{ Немедленная отправка блока, минуя очередь: для TX_CHRONO, где 20 мс
задержки очереди сопоставимы с самим квантом. Пишет из чужого потока,
поэтому под FWriteLock. }
function SendBinNow(const Hdr: TTCIStreamHeader): Boolean;
function SendBin(const Hdr: TTCIStreamHeader; Data: Pointer;
Bytes: Integer): Boolean;
{ Сколько блоков потока выброшено из-за отставания клиента. }
function BinDropped: LongInt;
{ Слить очередь в сокет. Зовёт ТОЛЬКО собственный поток клиента: запись
может ждать до TCI_SEND_TIMEOUT, и общий поток на этом задерживал бы
всех остальных. False — клиент умер. }
function Flush: Boolean;
{ Пометить мёртвым и разбудить его поток (shutdown сокета). }
procedure Kill;
{ Подписки на измерители — целиком под локом. }
procedure SetRxSensors(On_: Boolean);
procedure SetRxSensorsMs(Ms: Integer);
procedure SetTxSensors(On_: Boolean);
procedure SetTxSensorsMs(Ms: Integer);
{ Пора ли слать измеритель: проверка периода и отметка отправки — один
атомарный шаг, иначе тик-поток и клиентский расходятся в наборе. }
function DueRxSensors(Now_: QWord): Boolean;
function DueTxSensors(Now_: QWord): Boolean;
property Ws: TWsClient read FWs;
property Up: Boolean read FUp;
property Ready: Boolean read GetReady write SetReady;
property Dead: Boolean read FDead;
property IQRate: Integer read GetIQRate write SetIQRate;
property AudioRate: Integer read GetAudioRate write SetAudioRate;
property AudioSamples: Integer read GetAudioSamples write SetAudioSamples;
property AudioChannels: Integer read GetAudioChannels write SetAudioChannels;
property AudioSampleType: string read GetAudioSampleType write SetAudioSampleType;
property TxBuffering: Integer read GetTxBuffering write SetTxBuffering;
end;
TTCIClientEvent = procedure(Client: TTCIClient) of object;
TTCICommandEvent = procedure(Client: TTCIClient; const Cmd: string) of object;
{ Собранное бинарное сообщение от клиента (TX-аудио, §3.4). Данные живут
только на время вызова — обработчик обязан их скопировать. }
TTCIBinaryEvent = procedure(Client: TTCIClient; Data: PByte;
Len: Integer) of object;
TTCIServer = class
private
FListenSock: TSocket;
FClients: array[0..TCI_MAX_SOCKETS-1] of TTCIClient;
FClientCount: Integer; // всего сокетов в массиве (с не поднявшимися)
FUpCount: LongInt; // прошедших handshake (Interlocked*)
FClientLock: TCriticalSection;
FAcceptThread: TThread;
FTickThread: TThread;
FTxTickThread: TThread;
// ★Освобождение клиентов и планировщик TX_CHRONO обязаны быть взаимно
// исключены. Правило «указатель, взятый под FClientLock, живёт до
// следующего прохода тика» держалось на том, что клиентов освобождает
// ТОЛЬКО тик-поток и он же ими пользуется. У планировщика TX свой поток, и
// без этого лока он может писать в клиента, которого ReapClients уже
// освободил. Одного FClientLock мало: Reap освобождает память, уже выйдя
// из него.
FReapLock: TCriticalSection;
// ★Будилка планировщика TX. Без неё первый маркер уходил в среднем на 10, а
// в худшем на 20 мс позже фронта PTT (поток дремал с шагом
// TCI_TX_TICK_IDLE_MS), и ровно эти миллисекунды не хватало нулевому
// pre-roll, чтобы дожить до первого ответа клиента.
FTxWake: PRTLEvent;
FTxKick: Boolean; // «проснись и тикай сейчас», ставит KickTxTick
FThreadCount: LongInt; // живых клиентских потоков (Interlocked*)
FRunning: Boolean;
FStopping: Boolean;
FPort: Word;
FBindIP: string;
FOnCommand: TTCICommandEvent;
FOnBinary: TTCIBinaryEvent;
FOnConnect: TTCIClientEvent;
FOnDisconnect: TTCIClientEvent;
FOnTick: TThreadMethod;
FOnTxTick: TThreadMethod; // планировщик TX_CHRONO (свой поток)
FTxTickPeriodMs: Integer; // период планировщика, мс (ставит адаптер)
function InitListen: Boolean;
procedure ReapClients; // освободить клиентов, чьи потоки вышли
procedure KillAll;
procedure Disconnected(Client: TTCIClient); // OnDisconnect, единая точка
public
constructor Create;
destructor Destroy; override;
{ Проверка настроек без побочных эффектов: можно ли вообще открыть такой
слушатель. Зовётся ДО остановки работающего сервера. }
class function ValidSettings(APort: Word; const ABindIP: string): Boolean;
{ Настройка слушателя. Применяется при следующем Start.
False — адрес не разобран (порт не откроется). }
function Configure(APort: Word; const ABindIP: string): Boolean;
function Start: Boolean;
procedure Stop;
function Running: Boolean;
{ Всем клиентам, прошедшим инициализацию. Skip — кого пропустить
(обычно автора изменения не пропускаем: сервер отвечает и ему тоже,
это и есть подтверждение установки). Только кладёт в очереди. }
procedure Broadcast(const S: string; Skip: TTCIClient = nil);
{ Обход клиентов под локом — для рассылки с индивидуальным периодом.
Proc обязана быть быстрой: она держит FClientLock. }
procedure EnumClients(Proc: TTCIClientEvent);
function ClientCount: Integer;
{ Внутреннее (зовётся потоками сервера). }
procedure AcceptLoop;
procedure TickLoop;
procedure TxTickLoop;
procedure BeginTxClientUse;
procedure EndTxClientUse;
{ Немедленно разбудить планировщик TX_CHRONO — зовётся на фронте передачи. }
procedure KickTxTick;
procedure HandleClient(Client: TTCIClient);
function Promote(Client: TTCIClient): Boolean; // handshake прошёл
procedure ThreadDone; // клиентский поток отработал
property Port: Word read FPort;
property BindIP: string read FBindIP;
{ Идёт остановка: адаптер не должен начинать новых вызовов в поток
контроллера — тот, кто нас останавливает, обычно и есть поток
контроллера, и Synchronize из клиентского потока в него не вернётся. }
property Stopping: Boolean read FStopping;
property OnCommand: TTCICommandEvent read FOnCommand write FOnCommand;
property OnBinary: TTCIBinaryEvent read FOnBinary write FOnBinary;
property OnConnect: TTCIClientEvent read FOnConnect write FOnConnect;
property OnDisconnect: TTCIClientEvent read FOnDisconnect write FOnDisconnect;
property OnTick: TThreadMethod read FOnTick write FOnTick;
property OnTxTick: TThreadMethod read FOnTxTick write FOnTxTick;
// Период планировщика TX_CHRONO. Адаптер ставит его равным длительности
// кванта запроса; 0 — передачи нет, планировщик дремлет.
property TxTickPeriodMs: Integer read FTxTickPeriodMs write FTxTickPeriodMs;
end;
{ IPv4 из строки в сетевом порядке. Строгий: ровно четыре десятичных октета
0..255. '' и '0.0.0.0' — это INADDR_ANY (слушать везде), и только они:
«ошибка разбора = слушаем всё» в протоколе без авторизации недопустима. }
function TCIParseIPv4(const S: string; out Addr: LongWord): Boolean;
{ Значение HTTP-заголовка (Name — в нижнем регистре, без ':'). Разбор
построчный: точное сравнение подстроки «upgrade: websocket» отвергало
валидные запросы с табуляцией или без пробела после двоеточия. }
function TCIHttpHeader(const Header, Name: string): string;
{ Проверка UTF-8: текстовые кадры WebSocket обязаны быть корректным UTF-8
(RFC 6455 §5.6). Отвергает и оборванные последовательности, и избыточно
длинные формы, и суррогаты, и всё выше U+10FFFF. }
function TCIValidUTF8(const S: string): Boolean;
implementation
type
TTCIAcceptThread = class(TThread)
private FServer: TTCIServer;
protected procedure Execute; override;
public constructor Create(AServer: TTCIServer);
end;
TTCITxTickThread = class(TThread)
private
FServer: TTCIServer;
protected
procedure Execute; override;
public
constructor Create(AServer: TTCIServer);
end;
TTCITickThread = class(TThread)
private FServer: TTCIServer;
protected procedure Execute; override;
public constructor Create(AServer: TTCIServer);
end;
TTCIClientThread = class(TThread)
private FServer: TTCIServer; FClient: TTCIClient;
protected procedure Execute; override;
public constructor Create(AServer: TTCIServer; AClient: TTCIClient);
end;
constructor TTCIAcceptThread.Create(AServer: TTCIServer);
begin
inherited Create(True);
FServer := AServer;
FreeOnTerminate := False;
end;
procedure TTCIAcceptThread.Execute;
begin
FServer.AcceptLoop;
end;
constructor TTCITickThread.Create(AServer: TTCIServer);
begin
inherited Create(True);
FServer := AServer;
FreeOnTerminate := False;
end;
procedure TTCITickThread.Execute;
begin
FServer.TickLoop;
end;
constructor TTCITxTickThread.Create(AServer: TTCIServer);
begin
inherited Create(True);
FServer := AServer;
FreeOnTerminate := False;
end;
procedure TTCITxTickThread.Execute;
begin
FServer.TxTickLoop;
end;
constructor TTCIClientThread.Create(AServer: TTCIServer; AClient: TTCIClient);
begin
inherited Create(True);
FServer := AServer;
FClient := AClient;
FreeOnTerminate := True;
end;
procedure TTCIClientThread.Execute;
begin
try
try
FServer.HandleClient(FClient);
except
// Разбор фрейма рухнул — соединение всё равно закрываем штатно, иначе
// клиент остался бы висеть в массиве до остановки сервера.
end;
finally
// Освобождать себя нельзя: объект переиспользуется рассылкой из чужих
// потоков. Помечаем «поток вышел» — освободит тик-поток (ReapClients)
// или Stop, который ждёт именно этого.
FClient.FClosed := True;
FServer.ThreadDone;
end;
end;
function TCIParseIPv4(const S: string; out Addr: LongWord): Boolean;
var
Oct: array[0..3] of LongWord;
N, i, Start, V: Integer;
Part: string;
c: Char;
begin
Addr := 0; // INADDR_ANY
Result := False;
if (Trim(S) = '') or (Trim(S) = '0.0.0.0') then Exit(True);
N := 0;
Start := 1;
for i := 1 to Length(S) + 1 do
if (i > Length(S)) or (S[i] = '.') then
begin
if N > 3 then Exit(False);
Part := Copy(S, Start, i - Start);
if (Part = '') or (Length(Part) > 3) then Exit(False);
for c in Part do
if (c < '0') or (c > '9') then Exit(False);
V := StrToIntDef(Part, -1);
if (V < 0) or (V > 255) then Exit(False);
Oct[N] := LongWord(V);
Inc(N);
Start := i + 1;
end;
if N <> 4 then Exit(False);
// Сетевой порядок байт: первый октет — младший байт in_addr.
Addr := Oct[0] or (Oct[1] shl 8) or (Oct[2] shl 16) or (Oct[3] shl 24);
Result := True;
end;
function TCIHttpHeader(const Header, Name: string): string;
var
i, Start, P: Integer;
Line, LName: string;
begin
Result := '';
Start := 1;
for i := 1 to Length(Header) + 1 do
if (i > Length(Header)) or (Header[i] = #10) then
begin
Line := Trim(Copy(Header, Start, i - Start)); // Trim снимет и #13
Start := i + 1;
P := Pos(':', Line);
if P <= 1 then Continue;
LName := LowerCase(Trim(Copy(Line, 1, P - 1)));
if LName = Name then
Exit(Trim(Copy(Line, P + 1, MaxInt)));
end;
end;
function TCIValidUTF8(const S: string): Boolean;
var
i, k, N, Len: Integer;
B: Byte;
Cp: LongWord;
begin
i := 1;
Len := Length(S);
while i <= Len do
begin
B := Byte(S[i]);
if B < $80 then begin Inc(i); Continue; end
else if (B >= $C2) and (B <= $DF) then begin N := 1; Cp := B and $1F; end
else if (B >= $E0) and (B <= $EF) then begin N := 2; Cp := B and $0F; end
else if (B >= $F0) and (B <= $F4) then begin N := 3; Cp := B and $07; end
else Exit(False); // $80..$C1 и $F5.. началом последовательности не бывают
if i + N > Len then Exit(False); // оборвано на середине символа
for k := 1 to N do
begin
if (Byte(S[i + k]) and $C0) <> $80 then Exit(False);
Cp := (Cp shl 6) or (Byte(S[i + k]) and $3F);
end;
// Избыточно длинная форма, суррогатная пара и выход за U+10FFFF.
if ((N = 2) and (Cp < $800)) or ((N = 3) and (Cp < $10000)) or
((Cp >= $D800) and (Cp <= $DFFF)) or (Cp > $10FFFF) then Exit(False);
Inc(i, N + 1);
end;
Result := True;
end;
{ ═══════════════════════════════════════════════════════════════════════════
TTCIClient
═══════════════════════════════════════════════════════════════════════════ }
constructor TTCIClient.Create(AWs: TWsClient);
begin
inherited Create;
FWs := AWs;
FUp := False;
FReady := False;
FStateLock := TCriticalSection.Create;
FRxSensors := False;
FRxSensorsMs := 200;
FTxSensors := False;
FTxSensorsMs := 200;
FOutLock := TCriticalSection.Create;
FWriteLock := TCriticalSection.Create;
FOutCount := 0;
SetLength(FOut, 64);
FBinHead := 0;
FBinTail := 0;
FBinDropped := 0;
// Умолчания параметров потоков — как в §4.3 (клиент их обычно переопределяет).
FIQRate := TCI_IQ_RATE_DEF;
FAudioRate := TCI_AUDIO_RATE_DEF;
FAudioSamples := TCIDefaultAudioSamples(TCI_AUDIO_RATE_DEF);
FAudioChannels := TCI_AUDIO_CHAN_DEF;
FAudioSampleType := 'float32';
FTxBuffering := TCI_TX_BUFFERING_DEF;
end;
destructor TTCIClient.Destroy;
var i: Integer;
begin
for i := 0 to TCI_BIN_QUEUE - 1 do FBinOut[i] := nil;
FOutLock.Free;
FWriteLock.Free;
FStateLock.Free;
inherited;
end;
function TTCIClient.GetReady: Boolean;
begin
FStateLock.Enter;
try Result := FReady; finally FStateLock.Leave; end;
end;
procedure TTCIClient.SetReady(V: Boolean);
begin
FStateLock.Enter;
try FReady := V; finally FStateLock.Leave; end;
end;
function TTCIClient.GetIQRate: Integer;
begin
FStateLock.Enter;
try Result := FIQRate; finally FStateLock.Leave; end;
end;
procedure TTCIClient.SetIQRate(V: Integer);
begin
FStateLock.Enter;
try FIQRate := V; finally FStateLock.Leave; end;
end;
function TTCIClient.GetAudioRate: Integer;
begin
FStateLock.Enter;
try Result := FAudioRate; finally FStateLock.Leave; end;
end;
procedure TTCIClient.SetAudioRate(V: Integer);
begin
FStateLock.Enter;
try FAudioRate := V; finally FStateLock.Leave; end;
end;
function TTCIClient.GetAudioSamples: Integer;
begin
FStateLock.Enter;
try Result := FAudioSamples; finally FStateLock.Leave; end;
end;
procedure TTCIClient.SetAudioSamples(V: Integer);
begin
FStateLock.Enter;
try FAudioSamples := V; finally FStateLock.Leave; end;
end;
function TTCIClient.GetAudioChannels: Integer;
begin
FStateLock.Enter;
try Result := FAudioChannels; finally FStateLock.Leave; end;
end;
procedure TTCIClient.SetAudioChannels(V: Integer);
begin
FStateLock.Enter;
try FAudioChannels := V; finally FStateLock.Leave; end;
end;
function TTCIClient.GetAudioSampleType: string;
begin
FStateLock.Enter;
try Result := FAudioSampleType; finally FStateLock.Leave; end;
end;
procedure TTCIClient.SetAudioSampleType(const V: string);
begin
FStateLock.Enter;
try FAudioSampleType := V; finally FStateLock.Leave; end;
end;
function TTCIClient.GetTxBuffering: Integer;
begin
FStateLock.Enter;
try Result := FTxBuffering; finally FStateLock.Leave; end;
end;
procedure TTCIClient.SetTxBuffering(V: Integer);
begin
FStateLock.Enter;
try FTxBuffering := V; finally FStateLock.Leave; end;
end;
procedure TTCIClient.SetRxSensors(On_: Boolean);
begin
FStateLock.Enter;
try
FRxSensors := On_;
if On_ then FRxSensorsAt := 0; // первую посылку не ждём период
finally
FStateLock.Leave;
end;
end;
procedure TTCIClient.SetRxSensorsMs(Ms: Integer);
begin
FStateLock.Enter;
try FRxSensorsMs := Ms; finally FStateLock.Leave; end;
end;
procedure TTCIClient.SetTxSensors(On_: Boolean);
begin
FStateLock.Enter;
try
FTxSensors := On_;
if On_ then FTxSensorsAt := 0;
finally
FStateLock.Leave;
end;
end;
procedure TTCIClient.SetTxSensorsMs(Ms: Integer);
begin
FStateLock.Enter;
try FTxSensorsMs := Ms; finally FStateLock.Leave; end;
end;
function TTCIClient.DueRxSensors(Now_: QWord): Boolean;
begin
FStateLock.Enter;
try
Result := FRxSensors and (Now_ - FRxSensorsAt >= QWord(FRxSensorsMs));
if Result then FRxSensorsAt := Now_;
finally
FStateLock.Leave;
end;
end;
function TTCIClient.DueTxSensors(Now_: QWord): Boolean;
begin
FStateLock.Enter;
try
Result := FTxSensors and (Now_ - FTxSensorsAt >= QWord(FTxSensorsMs));
if Result then FTxSensorsAt := Now_;
finally
FStateLock.Leave;
end;
end;
function TTCIClient.Send(const S: string): Boolean;
begin
Result := False;
if (S = '') or FDead or (FWs = nil) then Exit;
FOutLock.Enter;
try
if FDead then Exit;
// Очередь переполнилась: клиент не читает сокет быстрее, чем мы пишем.
// Копить дальше нечестно (память + отставшее состояние), рвём соединение.
if FOutCount >= TCI_OUT_MAX then
begin
FDead := True;
Exit;
end;
if FOutCount >= Length(FOut) then SetLength(FOut, Length(FOut) * 2);
FOut[FOutCount] := S;
Inc(FOutCount);
Result := True;
finally
FOutLock.Leave;
end;
end;
function TTCIClient.SendBinNow(const Hdr: TTCIStreamHeader): Boolean;
var
Blk: TBytes;
begin
Result := False;
if FDead or (FWs = nil) then Exit;
SetLength(Blk, SizeOf(Hdr));
Move(Hdr, Blk[0], SizeOf(Hdr));
FWriteLock.Enter;
try
if FDead or (FWs = nil) then Exit;
Result := FWs.SendBinary(Blk[0], Length(Blk));
finally
FWriteLock.Leave;
end;
if not Result then FDead := True; // гасит сокет поток клиента (см. Flush)
end;
function TTCIClient.SendBin(const Hdr: TTCIStreamHeader; Data: Pointer;
Bytes: Integer): Boolean;
// Кладёт готовый блок в кольцо. Зовётся из DSP-потока: единственное, что тут
// разрешено — копирование под коротким локом. Кольцо полное — выбрасываем
// САМЫЙ СТАРЫЙ блок: рвать соединение из-за отставания в потоке нельзя
// (команды при этом продолжают ходить), а протухший звук клиенту не нужен.
var
Blk: TBytes;
NewH: Integer;
begin
Result := False;
if FDead or (FWs = nil) or (Bytes < 0) then Exit;
if Bytes > TCI_STREAM_DATA_MAX then Exit; // блок не по протоколу
SetLength(Blk, SizeOf(Hdr) + Bytes);
Move(Hdr, Blk[0], SizeOf(Hdr));
if Bytes > 0 then Move(Data^, Blk[SizeOf(Hdr)], Bytes);
FOutLock.Enter;
try
if FDead then Exit;
NewH := (FBinHead + 1) mod TCI_BIN_QUEUE;
if NewH = FBinTail then
begin
FBinOut[FBinTail] := nil;
FBinTail := (FBinTail + 1) mod TCI_BIN_QUEUE;
Inc(FBinDropped);
end;
FBinOut[FBinHead] := Blk;
FBinHead := NewH;
Result := True;
finally
FOutLock.Leave;
end;
end;
function TTCIClient.BinDropped: LongInt;
begin
FOutLock.Enter;
try Result := FBinDropped; finally FOutLock.Leave; end;
end;
function TTCIClient.Flush: Boolean;
var
Batch: array of string;
Bins: array[0..TCI_BIN_QUEUE-1] of TBytes;
N, i, NB: Integer;
Chunk: string;
begin
Result := not FDead;
// Мёртвым клиента могла пометить и очередь (переполнилась в чужом потоке —
// там гасить сокет нельзя, лок чужой). Добиваем здесь: без shutdown его
// поток так и висел бы в recv, а объект никогда бы не освободился.
if FDead then begin Kill; Exit; end;
NB := 0;
FOutLock.Enter;
try
N := FOutCount;
if N > 0 then
begin
SetLength(Batch, N);
for i := 0 to N - 1 do
begin
Batch[i] := FOut[i];
FOut[i] := '';
end;
FOutCount := 0;
end;
// Бинарные блоки забираем тем же заходом: лишний Enter/Leave на каждый
// блок потока — это тысячи лишних локов в секунду.
while FBinTail <> FBinHead do
begin
Bins[NB] := FBinOut[FBinTail];
FBinOut[FBinTail] := nil;
FBinTail := (FBinTail + 1) mod TCI_BIN_QUEUE;
Inc(NB);
end;
finally
FOutLock.Leave;
end;
if (N = 0) and (NB = 0) then Exit;
// Склейка: несколько команд в одном фрейме протокол разрешает (§3.1), а
// syscall'ов и заголовков становится в разы меньше.
// ★Под FWriteLock: в тот же сокет пишет планировщик TX (SendBinNow).
FWriteLock.Enter;
try
Chunk := '';
for i := 0 to N - 1 do
begin
if (Chunk <> '') and (Length(Chunk) + Length(Batch[i]) > TCI_OUT_CHUNK) then
begin
if not FWs.SendText(Chunk) then begin Kill; Exit(False); end;
Chunk := '';
end;
Chunk := Chunk + Batch[i];
end;
if Chunk <> '' then
if not FWs.SendText(Chunk) then begin Kill; Exit(False); end;
// Блоки потоков — каждый отдельным binary-фреймом: клиент читает их по
// одному заголовку на кадр, склейка тут запрещена протоколом.
for i := 0 to NB - 1 do
begin
if not FWs.SendBinary(Bins[i][0], Length(Bins[i])) then
begin
Kill;
Exit(False);
end;
Bins[i] := nil;
end;
finally
FWriteLock.Leave;
end;
end;
procedure TTCIClient.Kill;
var i: Integer;
begin
FOutLock.Enter;
try
FDead := True;
FOutCount := 0;
for i := 0 to TCI_BIN_QUEUE - 1 do FBinOut[i] := nil;
FBinHead := 0;
FBinTail := 0;
if FKilled then Exit; // shutdown уже был — второй раз незачем
FKilled := True;
finally
FOutLock.Leave;
end;
if FWs <> nil then
begin
FWs.State := wsClosed;
SockShutdown(FWs.Socket); // будим поток клиента, висящий в recv
end;
end;
{ ═══════════════════════════════════════════════════════════════════════════
TTCIServer — жизненный цикл
═══════════════════════════════════════════════════════════════════════════ }
constructor TTCIServer.Create;
{$IFDEF WINDOWS}
var WSAData: TWSAData;
{$ENDIF}
begin
inherited Create;
{$IFDEF WINDOWS}
WSAStartup($0202, WSAData); // refcounted: свой вызов на каждый сервер
{$ENDIF}
FListenSock := SOCK_INVALID;
FClientCount := 0;
FUpCount := 0;
FClientLock := TCriticalSection.Create;
FReapLock := TCriticalSection.Create;
FTxWake := RTLEventCreate;
FTxKick := False;
FPort := TCI_DEFAULT_PORT;
FBindIP := '127.0.0.1';
FTxTickPeriodMs := 0; // передачи нет — планировщик TX_CHRONO дремлет
end;
destructor TTCIServer.Destroy;
begin
Stop;
FClientLock.Free;
FReapLock.Free;
RTLEventDestroy(FTxWake);
{$IFDEF WINDOWS}
WSACleanup;
{$ENDIF}
inherited;
end;
class function TTCIServer.ValidSettings(APort: Word; const ABindIP: string): Boolean;
var Dummy: LongWord;
begin
Result := (APort <> 0) and TCIParseIPv4(ABindIP, Dummy);
end;
function TTCIServer.Configure(APort: Word; const ABindIP: string): Boolean;
begin
FPort := APort;
FBindIP := ABindIP;
Result := ValidSettings(APort, ABindIP);
end;
function TTCIServer.Running: Boolean;
begin
Result := FRunning;
end;
function TTCIServer.InitListen: Boolean;
var
Addr: {$IFDEF WINDOWS}TSockAddrIn{$ELSE}TInetSockAddr{$ENDIF};
One: Integer;
IP: LongWord;
begin
Result := False;
// Кривой адрес — отказ. Молча свалиться в INADDR_ANY нельзя: в TCI нет
// авторизации, и открытый наружу порт отдаёт управление передатчиком.
if not TCIParseIPv4(FBindIP, IP) then Exit;
if FPort = 0 then Exit;
{$IFDEF WINDOWS}
FListenSock := socket(AF_INET, SOCK_STREAM, IPPROTO_TCP);
{$ELSE}
FListenSock := fpSocket(AF_INET, SOCK_STREAM, IPPROTO_TCP);
{$ENDIF}
if FListenSock = SOCK_INVALID then Exit;
One := 1;
{$IFDEF WINDOWS}
setsockopt(FListenSock, SOL_SOCKET, SO_REUSEADDR, @One, SizeOf(One));
FillChar(Addr, SizeOf(Addr), 0);
Addr.sin_family := AF_INET;
Addr.sin_port := htons(FPort);
Addr.sin_addr.S_addr := IP;
if bind(FListenSock, @Addr, SizeOf(Addr)) = SOCKET_ERROR then Exit;
if listen(FListenSock, 5) = SOCKET_ERROR then Exit;
{$ELSE}
fpSetSockOpt(FListenSock, SOL_SOCKET, SO_REUSEADDR, @One, SizeOf(One));
FillChar(Addr, SizeOf(Addr), 0);
Addr.sin_family := AF_INET;
Addr.sin_port := htons(FPort);
Addr.sin_addr.s_addr := IP;
if fpBind(FListenSock, @Addr, SizeOf(Addr)) <> 0 then Exit;
if fpListen(FListenSock, 5) <> 0 then Exit;
{$ENDIF}
Result := True;
end;
function TTCIServer.Start: Boolean;
begin
Result := False;
if FRunning then Exit;
if not InitListen then
begin
if FListenSock <> SOCK_INVALID then
begin
SockClose(FListenSock);
FListenSock := SOCK_INVALID;
end;
Exit;
end;
FStopping := False;
FRunning := True;
FAcceptThread := TTCIAcceptThread.Create(Self);
TTCIAcceptThread(FAcceptThread).Start;
FTickThread := TTCITickThread.Create(Self);
TTCITickThread(FTickThread).Start;
FTxTickThread := TTCITxTickThread.Create(Self);
TTCITxTickThread(FTxTickThread).Start;
Result := True;
end;
procedure JoinPumped(var T: TThread);
// Ожидание выхода потока с прокачкой очереди Synchronize — то же, что делает
// шаг 4 в Stop, но для accept- и тик-потока. Глухой WaitFor тут — взаимный
// клин: тик-поток из ReapClients зовёт Disconnected → адаптер → Invoke, а
// Invoke это TThread.Synchronize к потоку контроллера, то есть ровно к тому,
// кто сейчас ждёт в WaitFor. Проверка CanInvoke у адаптера не спасает: она
// читает Stopping ДО входа в Synchronize, и клиент, отвалившийся в момент
// остановки сервера, успевает проскочить в это окно.
begin
if T = nil then Exit;
while not T.Finished do
if GetCurrentThreadId = MainThreadID then CheckSynchronize(5) else Sleep(5);
T.WaitFor; // Finished взводится уже после DoTerminate — не блокирует
FreeAndNil(T);
end;
procedure TTCIServer.KillAll;
var i: Integer;
begin
FClientLock.Enter;
try
for i := 0 to FClientCount - 1 do
if FClients[i] <> nil then FClients[i].Kill;
finally
FClientLock.Leave;
end;
end;
procedure TTCIServer.Stop;
var i, Waited: Integer;
begin
if not FRunning then Exit;
FStopping := True; // адаптер перестаёт звать Invoke (см. property Stopping)
FRunning := False;
// Шаг 1: гасим listen-сокет. SockShutdown обязателен до close — иначе
// fpAccept в accept-потоке не разблокируется (см. WebServer.Stop).
if FListenSock <> SOCK_INVALID then
begin
SockShutdown(FListenSock);
SockClose(FListenSock);
FListenSock := SOCK_INVALID;
end;
// Шаг 2: дожидаемся accept-потока — после него новых клиентов не появится.
JoinPumped(FAcceptThread);
// Шаг 3: будим клиентские потоки, висящие в recv, и останавливаем тик.
KillAll;
JoinPumped(FTickThread);
JoinPumped(FTxTickThread);
// Шаг 4: ждём выхода клиентских потоков — БЕЗ таймаута. Прокачивая очередь
// Synchronize: Stop зовёт поток контроллера (UI), а клиентский поток может
// как раз в нём висеть на FController.Invoke; без прокачки это взаимный
// клин. Выйти отсюда по таймауту нельзя: следом освобождаются и клиенты, и
// сам сервер с адаптером, а живой поток вернулся бы в эту память.
Waited := 0;
while FThreadCount > 0 do
begin
if GetCurrentThreadId = MainThreadID then CheckSynchronize(5) else Sleep(5);
Inc(Waited, 5);
// Повторный shutdown: клиент мог быть принят между шагом 2 и шагом 3
// (accept уже вернул сокет, поток стартовал позже) и Kill его не застал.
if (Waited mod TCI_STOP_KILL_MS) = 0 then KillAll;
end;
// Шаг 5: зачистка. Потоков больше нет — освобождать безопасно.
FClientLock.Enter;
try
for i := 0 to FClientCount - 1 do
if FClients[i] <> nil then
begin
Disconnected(FClients[i]); // и на остановке тоже: захваты снимаются
FClients[i].Ws.Free;
FreeAndNil(FClients[i]);
end;
FClientCount := 0;
FUpCount := 0;
finally
FClientLock.Leave;
end;
end;
{ ═══════════════════════════════════════════════════════════════════════════
Приём соединений
═══════════════════════════════════════════════════════════════════════════ }
procedure TTCIServer.AcceptLoop;
var
CSock: TSocket;
Addr: {$IFDEF WINDOWS}TSockAddrIn{$ELSE}TInetSockAddr{$ENDIF};
ALen: {$IFDEF WINDOWS}Integer{$ELSE}TSockLen{$ENDIF};
Client: TTCIClient;
T: TTCIClientThread;
Full: Boolean;
begin
while FRunning do
begin
ALen := SizeOf(Addr);
{$IFDEF WINDOWS}
CSock := accept(FListenSock, @Addr, @ALen);
{$ELSE}
CSock := fpAccept(FListenSock, @Addr, @ALen);
{$ENDIF}
if CSock = SOCK_INVALID then
begin
if FRunning then Sleep(10);
Continue;
end;
if not FRunning then
begin
SockClose(CSock);
Break;
end;
SockSetSndTimeout(CSock, TCI_SEND_TIMEOUT);
// До конца handshake сокет не должен молчать вечно: иначе горстка пустых
// соединений держала бы место, ничего не сказав.
SockSetRcvTimeout(CSock, TCI_HANDSHAKE_MS);
Client := TTCIClient.Create(TWsClient.Create(CSock));
Full := False;
FClientLock.Enter;
try
// Место в массиве освобождает тик-поток; слот протокола (Up) клиент
// получит позже — после успешного Upgrade (см. Promote).
if FClientCount >= TCI_MAX_SOCKETS then Full := True
else
begin
FClients[FClientCount] := Client;
Inc(FClientCount);
end;
finally
FClientLock.Leave;
end;
if Full then
begin
Client.Ws.Free; // закрывает сокет
Client.Free;
Continue;
end;
// Счётчик — ДО старта: иначе Stop успел бы проскочить шаг 4 между Start и
// первой строкой потока. ★А не родился поток (память, лимит потоков) —
// сразу возвращаем счётчик назад: ждут его на шаге 4 без таймаута, и
// лишняя единица подвешивала бы и выход из программы, и любую смену
// настроек TCI. Объект потока уносит себя сам: упавший конструктор — сразу,
// а поднявшийся — по FreeOnTerminate.
InterLockedIncrement(FThreadCount);
try
T := TTCIClientThread.Create(Self, Client);
T.Start;
except
InterLockedDecrement(FThreadCount);
// Обслуживать клиента больше некому: гасим сокет и метим на освобождение
// — заберёт тик-поток (ReapClients), как и обычное отключение.
Client.Kill;
Client.FClosed := True;
end;
end;
end;
function TTCIServer.Promote(Client: TTCIClient): Boolean;
// Слот протокола выдаётся ТОЛЬКО тут — после разбора HTTP-запроса и до ответа
// 101. Считаем поднявшихся: молчащие сокеты слотов не занимают.
var i, N: Integer;
begin
Result := False;
if not FRunning then Exit;
FClientLock.Enter;
try
N := 0;
for i := 0 to FClientCount - 1 do
if (FClients[i] <> nil) and FClients[i].FUp then Inc(N);
if N >= TCI_MAX_CLIENTS then Exit;
Client.FUp := True;
InterLockedIncrement(FUpCount);
Result := True;
finally
FClientLock.Leave;
end;
end;
procedure TTCIServer.KickTxTick;
begin
FTxKick := True;
RTLEventSetEvent(FTxWake);
end;
procedure TTCIServer.BeginTxClientUse;
begin
FReapLock.Enter;
end;
procedure TTCIServer.EndTxClientUse;
begin
FReapLock.Leave;
end;
procedure TTCIServer.ReapClients;
// Освобождение клиентов — единственное место, кроме Stop. Зовёт только
// тик-поток, поэтому указатель, взятый кем угодно под FClientLock, живёт до
// следующего прохода тика (а вне лока указателей никто не держит).
var
i, j, N: Integer;
Doomed: array[0..TCI_MAX_SOCKETS-1] of TTCIClient;
begin
N := 0;
FReapLock.Enter;
try
FClientLock.Enter;
try
i := 0;
while i < FClientCount do
if (FClients[i] <> nil) and FClients[i].FClosed then
begin
if FClients[i].FUp then InterLockedDecrement(FUpCount);
Doomed[N] := FClients[i];
Inc(N);
for j := i to FClientCount - 2 do FClients[j] := FClients[j + 1];
FClients[FClientCount - 1] := nil;
Dec(FClientCount);
end
else
Inc(i);
finally
FClientLock.Leave;
end;
for i := 0 to N - 1 do
begin
Disconnected(Doomed[i]);
Doomed[i].Ws.Free; // закрывает сокет
Doomed[i].Free;
end;
finally
FReapLock.Leave;
end;
end;
procedure TTCIServer.Disconnected(Client: TTCIClient);
// Единственное место, где наверх уходит «клиент ушёл»: и обычное отключение
// (ReapClients), и остановка сервера. Иначе после Stop у адаптера оставались
// висеть захваты параметров ушедших клиентов (§3.5).
begin
if (Client <> nil) and Assigned(FOnDisconnect) then FOnDisconnect(Client);
end;
{ ═══════════════════════════════════════════════════════════════════════════
Клиентский поток: handshake + разбор WS-фреймов
═══════════════════════════════════════════════════════════════════════════ }
procedure TTCIServer.HandleClient(Client: TTCIClient);
var
Ws: TWsClient;
R, HeaderEnd: Integer;
Header, Key, AcceptKey, Response, Text: string;
Raw: array[0..4095] of Byte;
RawLen, Rest: Integer;
B0, B1: Byte;
Masked, Fin, Pending: Boolean;
PayLen, Need, i, j, Consumed: Integer;
Hi32: LongWord;
Mask: array[0..3] of Byte;
Payload: array of Byte;
Opcode, MsgOp: Byte;
Frag: string;
Bin: array of Byte; // сборка бинарного сообщения (блок потока)
BinLen: Integer;
Cmds: TStringList;
begin
Ws := Client.Ws;
RawLen := 0;
// ── HTTP-запрос: ждём конца заголовков ───────────────────────────────────
Header := '';
HeaderEnd := 0;
repeat
R := SockRecv(Ws.Socket, @Raw[RawLen], SizeOf(Raw) - RawLen, 0);
// R <= 0 здесь — это и разрыв, и истёкший TCI_HANDSHAKE_MS: молчащее
// соединение уходит само, не занимая место.
if R <= 0 then begin Ws.State := wsClosed; Break; end;
Inc(RawLen, R);
SetLength(Header, RawLen);
Move(Raw[0], Header[1], RawLen);
HeaderEnd := System.Pos(#13#10#13#10, Header);
until (HeaderEnd > 0) or (RawLen >= SizeOf(Raw));
if (Ws.State = wsClosed) or (HeaderEnd = 0) then Exit;
Consumed := HeaderEnd + 3; // длина заголовков вместе с CRLFCRLF
Header := Copy(Header, 1, Consumed);
// Путь не проверяем: клиенты ходят на '/', но протокол его не оговаривает.
if (Pos('websocket', LowerCase(TCIHttpHeader(Header, 'upgrade'))) = 0) or
(Pos('upgrade', LowerCase(TCIHttpHeader(Header, 'connection'))) = 0) then
begin
Response := 'HTTP/1.1 426 Upgrade Required'#13#10 +
'Content-Length: 0'#13#10'Connection: close'#13#10#13#10;
Ws.SendRaw(Response[1], Length(Response));
Exit;
end;
// Браузерный клиент. Origin шлют только браузеры, и он — единственный
// признак, отличающий страницу от нативной программы. Авторизации в TCI
// нет: без этой проверки открытая вкладка с чужого сайта дотянулась бы по
// ws://127.0.0.1:40001 до TRX/TUNE/VFO. Своим web-страницам нужен явный
// прокси, а не дыра по умолчанию.
if TCIHttpHeader(Header, 'origin') <> '' then
begin
Response := 'HTTP/1.1 403 Forbidden'#13#10 +
'Content-Length: 0'#13#10'Connection: close'#13#10#13#10;
Ws.SendRaw(Response[1], Length(Response));
Exit;
end;
Key := TCIHttpHeader(Header, 'sec-websocket-key');
// Пустой ключ = не WebSocket-клиент (или сломанный): Accept без ключа
// формально считается валидным, и такое «соединение» потом молча висит.
if Key = '' then
begin
Response := 'HTTP/1.1 400 Bad Request'#13#10 +
'Content-Length: 0'#13#10'Connection: close'#13#10#13#10;
Ws.SendRaw(Response[1], Length(Response));
Exit;
end;
// Слот протокола — до ответа 101: отказать после «Switching Protocols» уже
// некрасиво, клиент считал бы себя подключённым.
if not Promote(Client) then
begin
Response := 'HTTP/1.1 503 Service Unavailable'#13#10 +
'Content-Length: 0'#13#10'Connection: close'#13#10#13#10;
Ws.SendRaw(Response[1], Length(Response));
Exit;
end;
AcceptKey := Base64EncodeBytes(SHA1(Key + TCI_WS_GUID), 20);
Response := 'HTTP/1.1 101 Switching Protocols'#13#10 +
'Upgrade: websocket'#13#10 +
'Connection: Upgrade'#13#10 +
'Sec-WebSocket-Accept: ' + AcceptKey + #13#10#13#10;
if not Ws.SendRaw(Response[1], Length(Response)) then Exit;
Ws.State := wsOpen;
// Дальше клиент вправе молчать сколько угодно, но просыпаться нам надо:
// на этом же потоке уходит его очередь отправки (Flush).
SockSetRcvTimeout(Ws.Socket, TCI_POLL_MS);
// Хвост первого пакета: клиент вправе прислать первый WS-фрейм в том же
// сегменте, что и заголовки. Выбросить его — потерять первую команду.
Rest := RawLen - Consumed;
if Rest > 0 then Move(Raw[Consumed], Ws.BufData[0], Rest);
Ws.BufLen := Rest;
// Пачка инициализации + текущее состояние (§3.1) — дело адаптера. Под
// FClientLock: пока она набирается, рассылка обязана ждать. Иначе изменение,
// случившееся после строки снимка, но до Ready=True, пропадало навсегда —
// Broadcast пропускает не-Ready клиента, и тот оставался со старым значением,
// считая инициализацию завершённой. Лок держится только на укладку строк в
// очередь (сеть тут не пишется), но обработчик OnConnect по этой же причине
// НЕ имеет права звать Invoke в поток контроллера: тот может ждать этот лок.
if Assigned(FOnConnect) then
begin
FClientLock.Enter;
try
FOnConnect(Client);
finally
FClientLock.Leave;
end;
end;
// ── Цикл WS-сообщений ────────────────────────────────────────────────────
Cmds := TStringList.Create;
Frag := '';
MsgOp := 0;
BinLen := 0;
// Буфер под сборку блока потока заводим сразу: расти по ходу приёма он всё
// равно не имеет права (потолок задан протоколом), а перевыделение на
// каждый блок TX-аудио — это мусор в куче двадцать раз в секунду.
SetLength(Bin, TCI_STREAM_MAX);
Pending := Rest > 0; // хвост handshake разбираем до первого recv
try
while FRunning and (Ws.State = wsOpen) and not Client.Dead do
begin
if not Pending then
begin
R := Ws.Recv;
// R = 0 — клиент закрыл свою сторону (EOF). Это НЕ ошибка, errno при
// этом не трогается и вполне может нести EAGAIN от прошлого истёкшего
// TCI_POLL_MS — спрашивать его тут нельзя, иначе обычный TCP-разрыв
// без close-кадра выглядит как таймаут и слот клиента не освобождается
// до остановки сервера. Место в буфере есть всегда (кадр крупнее
// буфера рвётся выше), так что нулю иного смысла нет.
if R = 0 then Break;
// R < 0 — ошибка: либо разрыв, либо просто истёк TCI_POLL_MS. Второе
// штатно: просыпаемся, чтобы отдать накопившуюся очередь.
if (R < 0) and not SockRecvTimedOut then Break;
end;
Pending := False;
while Ws.BufLen >= 2 do
begin
B0 := Ws.BufData[0];
B1 := Ws.BufData[1];
Fin := (B0 and $80) <> 0;
Opcode := B0 and $0F;
Masked := (B1 and $80) <> 0;
PayLen := B1 and $7F;
// RSV1..3 без согласованных расширений обязаны быть нулями.
if (B0 and $70) <> 0 then begin Ws.State := wsClosed; Break; end;
Need := 2;
if PayLen = 126 then Inc(Need, 2)
else if PayLen = 127 then Inc(Need, 8);
if Masked then Inc(Need, 4);
if Ws.BufLen < Need then Break;
i := 2;
if PayLen = 126 then
begin
PayLen := (Ws.BufData[2] shl 8) or Ws.BufData[3];
Inc(i, 2);
end
else if PayLen = 127 then
begin
// 64-битная длина: старшие четыре байта обязаны быть нулём, иначе
// значение не помещается в Integer и превращается в отрицательное.
Hi32 := (LongWord(Ws.BufData[2]) shl 24) or (LongWord(Ws.BufData[3]) shl 16) or
(LongWord(Ws.BufData[4]) shl 8) or LongWord(Ws.BufData[5]);
if Hi32 <> 0 then begin Ws.State := wsClosed; Break; end;
Hi32 := (LongWord(Ws.BufData[6]) shl 24) or (LongWord(Ws.BufData[7]) shl 16) or
(LongWord(Ws.BufData[8]) shl 8) or LongWord(Ws.BufData[9]);
if Hi32 > LongWord(SizeOf(Raw)) then begin Ws.State := wsClosed; Break; end;
PayLen := Integer(Hi32);
Inc(i, 8);
end;
// Клиент ОБЯЗАН маскировать (RFC 6455 §5.1). Незамаскированный кадр —
// либо не клиент, либо попытка прогнать через нас чужой трафик.
if not Masked then begin Ws.State := wsClosed; Break; end;
// Управляющие кадры: только короткие и только целиком (§5.5).
if (Opcode >= $08) and ((PayLen > 125) or (not Fin)) then
begin
Ws.State := wsClosed;
Break;
end;
// Фрейм крупнее приёмного буфера TWsClient никогда не соберётся —
// BufLen упрётся в потолок и цикл встанет намертво. Рвём соединение:
// команд такой длины у TCI нет, а самый крупный законный кадр —
// блок TX-аудио (заголовок + data[16384]) — в буфер помещается.
if Need + PayLen > Ws.BufCapacity then
begin
Ws.State := wsClosed;
Break;
end;
if Ws.BufLen < Need + PayLen then Break;
Mask[0] := Ws.BufData[i]; Mask[1] := Ws.BufData[i+1];
Mask[2] := Ws.BufData[i+2]; Mask[3] := Ws.BufData[i+3];
Inc(i, 4);
SetLength(Payload, PayLen);
if PayLen > 0 then
begin
Move(Ws.BufData[i], Payload[0], PayLen);
for j := 0 to PayLen - 1 do
Payload[j] := Payload[j] xor Mask[j and 3];
end;
Consumed := i + PayLen;
if Ws.BufLen > Consumed then
Move(Ws.BufData[Consumed], Ws.BufData[0], Ws.BufLen - Consumed);
Ws.BufLen := Ws.BufLen - Consumed;
case Opcode of
$00, $01, $02: // данные: продолжение / текст / binary
begin
if Opcode = $00 then
begin
// Продолжение без начала — рассинхрон, дальше читать нечего.
if MsgOp = 0 then begin Ws.State := wsClosed; Break; end;
end
else
begin
// Новое сообщение поверх недособранного — тоже рассинхрон.
if MsgOp <> 0 then begin Ws.State := wsClosed; Break; end;
MsgOp := Opcode;
Frag := '';
BinLen := 0;
end;
if MsgOp = $01 then
begin
if Length(Frag) + PayLen > TCI_MSG_MAX then
begin
Ws.State := wsClosed;
Break;
end;
if PayLen > 0 then
begin
SetLength(Text, PayLen);
Move(Payload[0], Text[1], PayLen);
Frag := Frag + Text;
end;
end
else
begin
// Бинарное сообщение — блок потока от клиента (TX-аудио).
// Клиент вправе резать его на фрагменты, поэтому копим так же,
// как текст, но с потолком в один блок: длиннее протокол не
// определяет, и растить буфер на чужой каприз мы не обязаны.
if BinLen + PayLen > TCI_STREAM_MAX then
begin
Ws.State := wsClosed;
Break;
end;
if PayLen > 0 then
begin
Move(Payload[0], Bin[BinLen], PayLen);
Inc(BinLen, PayLen);
end;
end;
if Fin then
begin
if (MsgOp = $02) and Assigned(FOnBinary) and (BinLen > 0) then
FOnBinary(Client, @Bin[0], BinLen);
BinLen := 0;
// Текстовое сообщение обязано быть валидным UTF-8 (§5.6);
// битую последовательность RFC велит закрывать, а не молча
// скармливать разбору команд.
if (MsgOp = $01) and not TCIValidUTF8(Frag) then
begin
Ws.State := wsClosed;
Break;
end;
if (MsgOp = $01) and Assigned(FOnCommand) and (Frag <> '') then
begin
TCISplit(Frag, Cmds);
for j := 0 to Cmds.Count - 1 do
FOnCommand(Client, Cmds[j]);
end;
MsgOp := 0;
Frag := '';
end;
end;
$08: // close: RFC 6455 §5.5.1 требует ответить своим close-кадром
begin
// Полезная нагрузка close — либо пустая, либо код (2 байта) плюс
// причина. Ровно один байт невалиден: отвечать на такое нечем.
if PayLen = 1 then begin Ws.State := wsClosed; Break; end;
// В ответе — только код: причину повторять не обязаны (§5.5.1),
// а чужой текст мы наружу не пересылаем.
if PayLen >= 2 then Ws.SendWsFrame($08, Payload[0], 2)
else Ws.SendWsFrame($08, PayLen, 0);
Ws.State := wsClosed;
Break;
end;
$09: // ping → pong
if PayLen > 0 then Ws.SendWsFrame($0A, Payload[0], PayLen)
else Ws.SendWsFrame($0A, PayLen, 0);
$0A: ; // pong — ничего не ждём
else
// Незнакомый opcode: по RFC соединение обязано закрыться.
Ws.State := wsClosed;
Break;
end;
end;
// Очередь — в сокет здесь же, на потоке этого клиента: ответы на только
// что разобранные команды уходят сразу, а медленный клиент задерживает
// только себя (тик-поток очереди лишь наполняет).
if not Client.Flush then Break;
end;
finally
Cmds.Free;
end;
end;
{ ═══════════════════════════════════════════════════════════════════════════
Рассылка и обход
═══════════════════════════════════════════════════════════════════════════ }
procedure TTCIServer.Broadcast(const S: string; Skip: TTCIClient);
var i: Integer;
begin
if S = '' then Exit;
FClientLock.Enter;
try
// Send только кладёт строку в очередь клиента — лок держится микросекунды,
// сколько бы клиент ни тормозил. В сокеты пишет тик-поток.
for i := 0 to FClientCount - 1 do
if (FClients[i] <> nil) and (FClients[i] <> Skip) and FClients[i].Ready then
FClients[i].Send(S);
finally
FClientLock.Leave;
end;
end;
procedure TTCIServer.EnumClients(Proc: TTCIClientEvent);
var i: Integer;
begin
if not Assigned(Proc) then Exit;
FClientLock.Enter;
try
for i := 0 to FClientCount - 1 do
if (FClients[i] <> nil) and FClients[i].FUp and not FClients[i].FClosed then
Proc(FClients[i]);
finally
FClientLock.Leave;
end;
end;
procedure TTCIServer.ThreadDone;
begin
InterLockedDecrement(FThreadCount);
end;
function TTCIServer.ClientCount: Integer;
begin
Result := FUpCount;
end;
procedure TTCIServer.TxTickLoop;
// Планировщик TX_CHRONO. Отдельный поток и АБСОЛЮТНЫЕ дедлайны: период здесь —
// длительность одного кванта запроса (около 10.7 мс при 512 отсчётах на 48 кГц),
// и он не обязан быть кратен чему-либо ещё в сервере. Дробный остаток копится в
// микросекундах, поэтому сетка не уезжает от округления периода до миллисекунд.
var
NextDueUs: Int64;
NowUs: Int64;
PeriodUs: Int64;
SleepMs: Integer;
begin
NextDueUs := MonotonicUs;
while FRunning do
begin
// Период ставит адаптер по кванту запроса. Пока передачи нет, он нулевой —
// но тикать всё равно надо, иначе адаптеру негде будет его выставить, когда
// передача начнётся (тик и период определяют друг друга).
PeriodUs := Int64(FTxTickPeriodMs) * 1000;
if PeriodUs < TCI_TX_TICK_MIN_MS * 1000 then
PeriodUs := TCI_TX_TICK_IDLE_MS * 1000;
NowUs := MonotonicUs;
// Фронт передачи: тикаем немедленно, не дожидаясь дремотного шага.
if FTxKick then
begin
FTxKick := False;
NextDueUs := NowUs;
end;
if NextDueUs <= NowUs then
begin
if Assigned(FOnTxTick) then FOnTxTick;
Inc(NextDueUs, PeriodUs);
// Проспали больше периода (планировщик ОС, пауза процесса) — не
// отыгрываем пропущенные тики пачкой: долг всё равно считается по часам
// внутри адаптера, а пачка маркеров только раздует окно в полёте.
NowUs := MonotonicUs;
if NextDueUs < NowUs then NextDueUs := NowUs;
Continue;
end;
SleepMs := Integer((NextDueUs - NowUs) div 1000);
if SleepMs <= 0 then SleepMs := TCI_TX_TICK_MIN_MS;
// Просыпаемся либо по расписанию, либо по звонку с фронта PTT.
RTLEventWaitFor(FTxWake, SleepMs);
end;
end;
procedure TTCIServer.TickLoop;
begin
while FRunning do
begin
Sleep(TCI_TICK_MS);
if not FRunning then Break;
ReapClients; // отключившиеся — освобождаем только здесь
if (FUpCount > 0) and Assigned(FOnTick) then FOnTick;
// В сокеты пишет каждый клиент сам, на своём потоке (см. HandleClient):
// общий поток отправки означал бы, что один медленный клиент задерживает
// очередь всех остальных на свой таймаут записи.
end;
end;
end.