unit TCIServer; { TCIServer.pas — транспорт TCI: WebSocket-сервер (роль сервера играем мы, как ExpertSDR3; клиенты — логгеры, скиммеры, программы цифровых видов). Что делает: • слушает TCP-порт (умолчание 40001), принимает HTTP-Upgrade на WebSocket по любому пути (клиенты ходят на ws://host:40001/); • режет входящие текстовые фреймы на команды и отдаёт их наверх (OnCommand) прямо в потоке клиента — маршалинг в поток контроллера делает TCIAdapter, как это устроено у CAT; • рассылает строки всем клиентам (Broadcast) — сервер TCI обязан синхронизировать всех подключённых (§3.5); • тикает OnTick (умолчание 20 мс) — по нему адаптер шлёт показания измерителей с индивидуальным для каждого клиента периодом. Бинарные фреймы (потоки IQ/аудио, §3.4) пока не обрабатываются: этап 2, см. doc/TCI.md. Приходящие от клиента binary-фреймы молча отбрасываются. Сокеты и WS-фреймы переиспользованы из веб-подсистемы (WebUtils/WsClient): тот же код handshake и та же схема «поток на клиента + accept-поток», что в WebServer. Авторизации у TCI нет by design. Порт слушается там, где сказано в настройках; умолчание — 127.0.0.1, чтобы наружу он не торчал без спроса. } {$IFDEF FPC} {$MODE Delphi} {$LONGSTRINGS ON} {$ENDIF} interface uses Classes, SysUtils, WebUtils, WsClient, TCIProtocol {$IFDEF WINDOWS}, Windows, WinSock2{$ELSE}, Sockets{$ENDIF}, SyncObjs; // ← после платформенных юнитов (конфликт идентификатора Create) const TCI_MAX_CLIENTS = 8; TCI_TICK_MS = 20; // период OnTick (сенсоры троттлятся адаптером) TCI_WS_GUID = '258EAFA5-E914-47DA-95CA-C5AB0DC85B11'; TCI_SEND_TIMEOUT = 1000; // мс на SockSend, иначе клиент считается мёртвым type TTCIServer = class; { Один подключённый клиент: WS-сокет + его личные подписки. Подписки на сенсоры в TCI индивидуальны (RX_SENSORS_ENABLE «отправляется только клиентом»), поэтому живут здесь, а не в адаптере. } TTCIClient = class private FWs: TWsClient; FReady: Boolean; // пачка инициализации отправлена FRxSensors: Boolean; FRxSensorsMs: Integer; FRxSensorsAt: QWord; // тик последней отправки FTxSensors: Boolean; FTxSensorsMs: Integer; FTxSensorsAt: QWord; public constructor Create(AWs: TWsClient); { Текстовый фрейм клиенту. False — соединение уже мертво. } function Send(const S: string): Boolean; property Ws: TWsClient read FWs; property Ready: Boolean read FReady write FReady; property RxSensors: Boolean read FRxSensors write FRxSensors; property RxSensorsMs: Integer read FRxSensorsMs write FRxSensorsMs; property RxSensorsAt: QWord read FRxSensorsAt write FRxSensorsAt; property TxSensors: Boolean read FTxSensors write FTxSensors; property TxSensorsMs: Integer read FTxSensorsMs write FTxSensorsMs; property TxSensorsAt: QWord read FTxSensorsAt write FTxSensorsAt; end; TTCIClientEvent = procedure(Client: TTCIClient) of object; TTCICommandEvent = procedure(Client: TTCIClient; const Cmd: string) of object; TTCIServer = class private FListenSock: TSocket; FClients: array[0..TCI_MAX_CLIENTS-1] of TTCIClient; FClientCount: Integer; FClientLock: TCriticalSection; FAcceptThread: TThread; FTickThread: TThread; FThreadCount: LongInt; // живых клиентских потоков (Interlocked*) FRunning: Boolean; FPort: Word; FBindIP: string; FOnCommand: TTCICommandEvent; FOnConnect: TTCIClientEvent; FOnDisconnect: TTCIClientEvent; FOnTick: TThreadMethod; function InitListen: Boolean; procedure RemoveClient(Client: TTCIClient); public constructor Create; destructor Destroy; override; { Настройка слушателя. Применяется при следующем Start. } procedure Configure(APort: Word; const ABindIP: string); function Start: Boolean; procedure Stop; function Running: Boolean; { Всем клиентам, прошедшим инициализацию. Skip — кого пропустить (обычно автора изменения не пропускаем: сервер отвечает и ему тоже, это и есть подтверждение установки). } procedure Broadcast(const S: string; Skip: TTCIClient = nil); { Обход клиентов под локом — для рассылки с индивидуальным периодом. } procedure EnumClients(Proc: TTCIClientEvent); function ClientCount: Integer; { Внутреннее (зовётся потоками сервера). } procedure AcceptLoop; procedure TickLoop; procedure HandleClient(Client: TTCIClient); procedure ThreadDone; // клиентский поток отработал property Port: Word read FPort; property BindIP: string read FBindIP; property OnCommand: TTCICommandEvent read FOnCommand write FOnCommand; property OnConnect: TTCIClientEvent read FOnConnect write FOnConnect; property OnDisconnect: TTCIClientEvent read FOnDisconnect write FOnDisconnect; property OnTick: TThreadMethod read FOnTick write FOnTick; end; implementation type TTCIAcceptThread = 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 TTCIClientThread.Create(AServer: TTCIServer; AClient: TTCIClient); begin inherited Create(True); FServer := AServer; FClient := AClient; FreeOnTerminate := True; end; procedure TTCIClientThread.Execute; begin try FServer.HandleClient(FClient); finally FServer.ThreadDone; // Stop ждёт обнуления счётчика перед зачисткой end; end; { IPv4 из строки в сетевом порядке. Свой, потому что WebServer держит такой же в implementation и наружу не отдаёт. } function TCIParseIPv4(const S: string): LongWord; var Oct: array[0..3] of LongWord; N, i, Start: Integer; Part: string; begin Result := 0; // INADDR_ANY if (S = '') or (S = '0.0.0.0') then Exit; 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; Part := Copy(S, Start, i - Start); Oct[N] := LongWord(StrToIntDef(Part, 0)) and $FF; Inc(N); Start := i + 1; end; if N <> 4 then Exit; // Сетевой порядок байт: первый октет — младший байт in_addr. Result := Oct[0] or (Oct[1] shl 8) or (Oct[2] shl 16) or (Oct[3] shl 24); end; { ═══════════════════════════════════════════════════════════════════════════ TTCIClient ═══════════════════════════════════════════════════════════════════════════ } constructor TTCIClient.Create(AWs: TWsClient); begin inherited Create; FWs := AWs; FReady := False; FRxSensors := False; FRxSensorsMs := 200; FTxSensors := False; FTxSensorsMs := 200; end; function TTCIClient.Send(const S: string): Boolean; begin Result := (FWs <> nil) and (FWs.State = wsOpen) and FWs.SendText(S); 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; FClientLock := TCriticalSection.Create; FPort := TCI_DEFAULT_PORT; FBindIP := '127.0.0.1'; end; destructor TTCIServer.Destroy; begin Stop; FClientLock.Free; {$IFDEF WINDOWS} WSACleanup; {$ENDIF} inherited; end; procedure TTCIServer.Configure(APort: Word; const ABindIP: string); begin FPort := APort; FBindIP := ABindIP; end; function TTCIServer.Running: Boolean; begin Result := FRunning; end; function TTCIServer.InitListen: Boolean; var Addr: {$IFDEF WINDOWS}TSockAddrIn{$ELSE}TInetSockAddr{$ENDIF}; One: Integer; begin Result := False; {$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 := TCIParseIPv4(FBindIP); 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 := TCIParseIPv4(FBindIP); 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; FRunning := True; FAcceptThread := TTCIAcceptThread.Create(Self); TTCIAcceptThread(FAcceptThread).Start; FTickThread := TTCITickThread.Create(Self); TTCITickThread(FTickThread).Start; Result := True; end; procedure TTCIServer.Stop; var i, Waited: Integer; begin if not FRunning then Exit; 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: будим клиентские потоки, висящие в recv. FClientLock.Enter; try for i := 0 to FClientCount - 1 do if FClients[i] <> nil then begin FClients[i].Ws.State := wsClosed; SockShutdown(FClients[i].Ws.Socket); end; finally FClientLock.Leave; end; // Шаг 3: свои потоки (клиентские — FreeOnTerminate, ждём их отдельно). if FAcceptThread <> nil then begin FAcceptThread.WaitFor; FreeAndNil(FAcceptThread); end; if FTickThread <> nil then begin FTickThread.WaitFor; FreeAndNil(FTickThread); end; // Шаг 4: даём клиентским потокам выйти самим. Освобождает клиента ТОТ, // кто вынул его из массива (RemoveClient в клиентском потоке), — иначе // Stop освободил бы объект из-под работающего потока. Waited := 0; while (FThreadCount > 0) and (Waited < 2000) do begin Sleep(10); Inc(Waited, 10); end; // Вырожденный случай: поток завис (не должно случаться — сокеты закрыты). // Чистим остатки, чтобы не течь; объекты уже никем не используются. FClientLock.Enter; try for i := 0 to FClientCount - 1 do if FClients[i] <> nil then begin FClients[i].Ws.Free; FreeAndNil(FClients[i]); end; FClientCount := 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; 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 FClientCount >= TCI_MAX_CLIENTS then begin SockClose(CSock); Continue; end; SockSetSndTimeout(CSock, TCI_SEND_TIMEOUT); Client := TTCIClient.Create(TWsClient.Create(CSock)); FClientLock.Enter; try FClients[FClientCount] := Client; Inc(FClientCount); finally FClientLock.Leave; end; InterLockedIncrement(FThreadCount); T := TTCIClientThread.Create(Self, Client); T.Start; end; end; procedure TTCIServer.RemoveClient(Client: TTCIClient); var i, j: Integer; begin if Client = nil then Exit; if Assigned(FOnDisconnect) then FOnDisconnect(Client); FClientLock.Enter; try for i := 0 to FClientCount - 1 do if FClients[i] = Client then begin for j := i to FClientCount - 2 do FClients[j] := FClients[j + 1]; FClients[FClientCount - 1] := nil; Dec(FClientCount); Break; end; finally FClientLock.Leave; end; Client.Ws.Free; // закрывает сокет Client.Free; end; { ═══════════════════════════════════════════════════════════════════════════ Клиентский поток: handshake + разбор WS-фреймов ═══════════════════════════════════════════════════════════════════════════ } procedure TTCIServer.HandleClient(Client: TTCIClient); var Ws: TWsClient; R, HeaderEnd: Integer; Header, HeaderLC, Key, AcceptKey, Response, Text: string; Raw: array[0..4095] of Byte; RawLen: Integer; B0, B1: Byte; Masked: Boolean; PayLen, Need, i, j, Consumed, KPos, KEnd: Integer; Mask: array[0..3] of Byte; Payload: array of Byte; Opcode: Byte; Cmds: TStringList; begin Ws := Client.Ws; RawLen := 0; // ── HTTP-запрос: ждём конца заголовков ─────────────────────────────────── Header := ''; HeaderEnd := 0; repeat R := SockRecv(Ws.Socket, @Raw[RawLen], SizeOf(Raw) - RawLen, 0); 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 begin RemoveClient(Client); Exit; end; Header := Copy(Header, 1, HeaderEnd + 3); HeaderLC := LowerCase(Header); // Путь не проверяем: клиенты ходят на '/', но протокол его не оговаривает. if System.Pos('upgrade: websocket', HeaderLC) = 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)); RemoveClient(Client); Exit; end; Key := ''; KPos := System.Pos('sec-websocket-key: ', HeaderLC); if KPos > 0 then begin Key := Copy(Header, KPos + 19, 100); KEnd := System.Pos(#13, Key); if KEnd > 0 then Key := Copy(Key, 1, KEnd - 1); Key := Trim(Key); 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 begin RemoveClient(Client); Exit; end; Ws.State := wsOpen; // Пачка инициализации + текущее состояние (§3.1) — дело адаптера. if Assigned(FOnConnect) then FOnConnect(Client); // ── Цикл WS-сообщений ──────────────────────────────────────────────────── Cmds := TStringList.Create; try Ws.BufLen := 0; while FRunning and (Ws.State = wsOpen) do begin R := Ws.Recv; if R <= 0 then Break; while Ws.BufLen >= 2 do begin B0 := Ws.BufData[0]; B1 := Ws.BufData[1]; Opcode := B0 and $0F; Masked := (B1 and $80) <> 0; PayLen := B1 and $7F; 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 PayLen := (Ws.BufData[6] shl 24) or (Ws.BufData[7] shl 16) or (Ws.BufData[8] shl 8) or Ws.BufData[9]; Inc(i, 8); end; // Фрейм крупнее приёмного буфера TWsClient никогда не соберётся — // BufLen упрётся в потолок и цикл встанет намертво. Рвём соединение: // команд такой длины у TCI нет, а бинарные потоки от клиента (TX-аудио) // мы пока не принимаем. if Need + PayLen > 4096 then begin Ws.State := wsClosed; Break; end; if Ws.BufLen < Need + PayLen then Break; if Masked then begin 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); end; SetLength(Payload, PayLen); if PayLen > 0 then begin Move(Ws.BufData[i], Payload[0], PayLen); if Masked then 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 $01: // текст — одна или несколько команд в одном фрейме begin SetLength(Text, PayLen); if PayLen > 0 then Move(Payload[0], Text[1], PayLen); if Assigned(FOnCommand) then begin TCISplit(Text, Cmds); for j := 0 to Cmds.Count - 1 do FOnCommand(Client, Cmds[j]); end; end; $02: ; // binary: TX-аудио от клиента — этап 2, пока игнорируем $08: // close begin Ws.State := wsClosed; Break; end; $09: // ping → pong if PayLen > 0 then Ws.SendWsFrame($0A, Payload[0], PayLen) else Ws.SendWsFrame($0A, PayLen, 0); end; end; end; finally Cmds.Free; end; RemoveClient(Client); end; { ═══════════════════════════════════════════════════════════════════════════ Рассылка и обход ═══════════════════════════════════════════════════════════════════════════ } procedure TTCIServer.Broadcast(const S: string; Skip: TTCIClient); var i: Integer; begin if S = '' then Exit; FClientLock.Enter; try 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 then Proc(FClients[i]); finally FClientLock.Leave; end; end; procedure TTCIServer.ThreadDone; begin InterLockedDecrement(FThreadCount); end; function TTCIServer.ClientCount: Integer; begin Result := FClientCount; end; procedure TTCIServer.TickLoop; begin while FRunning do begin Sleep(TCI_TICK_MS); if not FRunning then Break; if (FClientCount > 0) and Assigned(FOnTick) then FOnTick; end; end; end.