From 24647dfef437356df15870421e7f90d802a4bef9 Mon Sep 17 00:00:00 2001 From: Dieter Lang Date: Mon, 10 Aug 2026 10:38:14 +0200 Subject: [PATCH] Implement server-side RS232 data transport --- .gitignore | 1 + CHANGELOG.md | 64 +++ README.md | 101 +++- configs/client.json | 16 + configs/server.json | 26 + go.mod | 6 +- go.sum | 4 + internal/config/config.go | 245 +++++++++ internal/config/config_test.go | 459 +++++++++++++++++ internal/config/device.go | 83 ++++ internal/serial/serial.go | 159 ++++++ internal/serial/serial_test.go | 391 +++++++++++++++ internal/server/control.go | 508 +++++++++++++++++++ internal/server/control_test.go | 635 ++++++++++++++++++++++++ internal/server/data_connection.go | 169 +++++++ internal/server/data_connection_test.go | 314 ++++++++++++ internal/server/data_handler.go | 143 ++++++ internal/server/data_handler_test.go | 350 +++++++++++++ internal/server/data_listener.go | 146 ++++++ internal/server/data_listener_test.go | 150 ++++++ internal/server/session.go | 462 +++++++++++++++++ internal/server/session_test.go | 440 ++++++++++++++++ internal/transport/codec.go | 101 ++++ internal/transport/codec_test.go | 164 ++++++ internal/transport/control.go | 202 ++++++++ internal/transport/control_test.go | 125 +++++ 26 files changed, 5460 insertions(+), 4 deletions(-) create mode 100644 configs/client.json create mode 100644 configs/server.json create mode 100644 go.sum create mode 100644 internal/config/config.go create mode 100644 internal/config/config_test.go create mode 100644 internal/config/device.go create mode 100644 internal/serial/serial.go create mode 100644 internal/serial/serial_test.go create mode 100644 internal/server/control.go create mode 100644 internal/server/control_test.go create mode 100644 internal/server/data_connection.go create mode 100644 internal/server/data_connection_test.go create mode 100644 internal/server/data_handler.go create mode 100644 internal/server/data_handler_test.go create mode 100644 internal/server/data_listener.go create mode 100644 internal/server/data_listener_test.go create mode 100644 internal/server/session.go create mode 100644 internal/server/session_test.go create mode 100644 internal/transport/codec.go create mode 100644 internal/transport/codec_test.go create mode 100644 internal/transport/control.go create mode 100644 internal/transport/control_test.go diff --git a/.gitignore b/.gitignore index 9f9b114..08bfe85 100644 --- a/.gitignore +++ b/.gitignore @@ -1,3 +1,4 @@ +Notizen.txt # Build-Verzeichnis /build/ diff --git a/CHANGELOG.md b/CHANGELOG.md index fb8d92b..1ef8afb 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,70 @@ Alle wesentlichen Änderungen am Projekt werden in dieser Datei dokumentiert. +## [Unreleased] + +### Added + +- Implementierung der seriellen Schnittstellenschicht unter `internal/serial` +- Verwendung von `go.bug.st/serial v1.7.1` für die Ansteuerung der seriellen Schnittstellen +- Konfigurationsunterstützung für Baudrate, Datenbits, Parität und Stopbits +- Öffnen und Schließen serieller Schnittstellen über die interne Serial-Abstraktion +- Lesen und Schreiben von RS232-Daten über die Serial-Abstraktion +- Dynamische TCP-Data-Listener für konfigurierte Geräte +- Session-Verwaltung für dynamische Data-Listener +- Session-Verwaltung für aktive DataConnections +- Bidirektionale Datenübertragung zwischen TCP und RS232 +- Öffnen der seriellen Schnittstelle erst beim Aufbau einer tatsächlichen Data-Verbindung +- Begrenzung auf eine aktive Data-Verbindung pro konfiguriertem Gerät +- Integrationstest mit virtuellen seriellen Schnittstellen über `socat` +- Tests für TCP → RS232 und RS232 → TCP +- Tests für Session-Reconnect und Ressourcenverwaltung + +### Changed + +- Go-Version des Projektes auf Go `1.25.0` aktualisiert +- `go.bug.st/serial` wird in Version `v1.7.1` verwendet +- `golang.org/x/sys v0.43.0` wird als indirekte Abhängigkeit verwendet +- Die Serial-Bibliothek `go.bug.st/serial` wurde zusätzlich auf dem privaten Git-Server des Projektes als Ausfallsicherung gespiegelt: + `git.lang-dieter.de/third-party/go-serial` +- Der `DataListener` wurde nebenläufigkeitssicher implementiert +- Das Schließen eines `DataListener` kann gleichzeitig mit einem laufenden `Accept()` erfolgen + +### Tests + +- `go test ./...` erfolgreich +- `go test -race ./...` erfolgreich +- Race Condition im `DataListener` erkannt und behoben + +### Status + +Die grundlegende Server-seitige TCP-/RS232-Datenübertragung ist implementiert und durch automatisierte Tests abgesichert. + +Der Datenpfad ist derzeit: + +```text +TCP-Control-Verbindung + │ + ▼ + Session + │ + ▼ + dynamischer Data-Port + │ + ▼ + TCP-Data-Verbindung + │ + ▼ + internal/serial + │ + ▼ + RS232-Gerät +``` + +Die Client-seitige Bereitstellung einer virtuellen seriellen Schnittstelle ist noch nicht implementiert. + +--- + ## [0.0.1] - 2026-08-09 ### Added diff --git a/README.md b/README.md index 8a3fcd3..bf60c9b 100644 --- a/README.md +++ b/README.md @@ -40,6 +40,41 @@ rs2322tcp-client │ Funkgerät Rotor ``` +## Aktueller Server-Datenpfad + +Die serverseitige TCP-/RS232-Verbindung ist inzwischen implementiert: + +```text +TCP-Control + │ + ▼ + Session + │ + ├── Gerät 1 ──► dynamischer Data-Port + │ │ + │ ▼ + │ TCP Data + │ │ + │ ▼ + │ Serial Layer + │ │ + │ ▼ + │ RS232 + │ + └── Gerät 2 ──► dynamischer Data-Port +``` + +Die serielle Schnittstelle wird erst geöffnet, wenn ein Client tatsächlich eine Data-Verbindung zum entsprechenden dynamischen TCP-Port aufbaut. + +Pro konfiguriertem Gerät ist nur eine aktive Data-Verbindung vorgesehen. + +Die Datenübertragung erfolgt bidirektional: + +```text +TCP ───────────────► RS232 +TCP ◄────────────── RS232 +``` + ## Client Der Client soll gleichberechtigt unter folgenden Betriebssystemen eingesetzt werden können: @@ -49,12 +84,32 @@ Der Client soll gleichberechtigt unter folgenden Betriebssystemen eingesetzt wer Die plattformspezifische Bereitstellung der virtuellen seriellen Schnittstelle wird vom gemeinsamen Client-Kern getrennt. +Die Client-seitige virtuelle serielle Schnittstelle ist derzeit noch nicht implementiert. + ## Server Der Server ist zunächst für den Betrieb auf einem Raspberry Pi 5 vorgesehen. Mehrere USB-to-RS232-Adapter können angeschlossen werden. Anzahl, Bezeichnung und serielle Parameter der Anschlüsse werden über eine Konfigurationsdatei festgelegt. +Die serverseitige Serial-Kommunikation verwendet `go.bug.st/serial`. + +## Serielle Schnittstelle + +Für die Ansteuerung der seriellen Schnittstellen wird derzeit verwendet: + +```text +go.bug.st/serial v1.7.1 +``` + +Die Bibliothek wurde zusätzlich auf dem privaten Git-Server des Projektes gespiegelt, um bei einem Ausfall des ursprünglichen Anbieters weiterhin auf den verwendeten Quellstand zugreifen zu können: + +```text +git.lang-dieter.de/third-party/go-serial +``` + +Die verwendete Version `v1.7.1` ist auf dem Backup-Repository einschließlich Git-Tag vorhanden. + ## Netzwerk Für die Übertragung der RS232-Daten wird TCP verwendet. @@ -90,6 +145,33 @@ Insbesondere soll eine Darstellung der übertragenen Bytes möglich sein, um Pro Während der Entwicklung kann `socat` als zusätzliches Werkzeug für Tests und Diagnose eingesetzt werden. +## Tests + +Die wichtigsten Komponenten verfügen über automatisierte Tests. + +Der aktuelle Stand wird unter anderem mit folgenden Befehlen geprüft: + +```bash +go test ./... +``` + +und: + +```bash +go test -race ./... +``` + +Der Race Detector wird eingesetzt, um Probleme bei der nebenläufigen Verarbeitung von Sessions, Data-Listenern und DataConnections zu erkennen. + +Für die serielle Datenübertragung werden virtuelle serielle Schnittstellen über `socat` verwendet. Dadurch kann der Datenpfad ohne angeschlossene RS232-Hardware getestet werden. + +Dabei werden insbesondere beide Übertragungsrichtungen geprüft: + +```text +TCP → RS232 +RS232 → TCP +``` + ## Projektstruktur ```text @@ -148,11 +230,24 @@ Release-Versionen werden über Git-Tags gekennzeichnet. ## Entwicklungsstand -Das Projekt befindet sich derzeit in der frühen Entwicklungsphase. +Das Projekt befindet sich weiterhin in der frühen Entwicklungsphase. -Version `0.0.1` enthält zunächst die Projektgrundlage, die beiden Programmgerüste sowie das Buildsystem. +Die ursprüngliche Projektgrundlage aus Version `0.0.1` wurde inzwischen um eine funktionierende serverseitige TCP-/RS232-Datenübertragung erweitert. -Eine funktionierende RS232- oder TCP-Übertragung ist in dieser Version noch nicht implementiert. +Implementiert und getestet sind derzeit: + +- Server-Control-Verbindung +- Session-Verwaltung +- dynamische Data-Ports +- DataConnection +- bidirektionale TCP-/RS232-Datenübertragung +- Serial-Abstraktion +- Konfiguration der seriellen Parameter +- virtuelle serielle Integrationstests +- nebenläufigkeitssichere Data-Listener +- Race-Detection + +Noch nicht implementiert ist insbesondere die clientseitige virtuelle serielle Schnittstelle für Windows und Linux. ## Lizenz diff --git a/configs/client.json b/configs/client.json new file mode 100644 index 0000000..0be587a --- /dev/null +++ b/configs/client.json @@ -0,0 +1,16 @@ +{ + "server": { + "address": "100.64.0.10", + "port": 5000 + }, + "virtual_ports": [ + { + "port": "COM7", + "remote_device": "radio" + }, + { + "port": "COM8", + "remote_device": "rotor" + } + ] +} diff --git a/configs/server.json b/configs/server.json new file mode 100644 index 0000000..1db84ad --- /dev/null +++ b/configs/server.json @@ -0,0 +1,26 @@ +{ + "listen": { + "address": "0.0.0.0", + "port": 5000 + }, + "devices": [ + { + "id": "radio", + "name": "Funkgerät", + "serial_port": "/dev/ttyUSB0", + "baud_rate": 9600, + "data_bits": 8, + "parity": "none", + "stop_bits": 1 + }, + { + "id": "rotor", + "name": "Antennenrotor", + "serial_port": "/dev/ttyUSB1", + "baud_rate": 4800, + "data_bits": 8, + "parity": "none", + "stop_bits": 1 + } + ] +} diff --git a/go.mod b/go.mod index 10e3e44..d033741 100644 --- a/go.mod +++ b/go.mod @@ -1,3 +1,7 @@ module git.lang-dieter.de/rs2322tcp -go 1.22.2 +go 1.25.0 + +require go.bug.st/serial v1.7.1 + +require golang.org/x/sys v0.43.0 // indirect diff --git a/go.sum b/go.sum new file mode 100644 index 0000000..dc4a229 --- /dev/null +++ b/go.sum @@ -0,0 +1,4 @@ +go.bug.st/serial v1.7.1 h1:5aP8wYL0UjEYOVs3oPAGscjaSfRQLHtCvBFXNN/rwtc= +go.bug.st/serial v1.7.1/go.mod h1:d0MmS16Qt9b1m06yoYRNUXhRRTJV5Qg2S5EKqQtnayQ= +golang.org/x/sys v0.43.0 h1:Rlag2XtaFTxp19wS8MXlJwTvoh8ArU6ezoyFsMyCTNI= +golang.org/x/sys v0.43.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw= diff --git a/internal/config/config.go b/internal/config/config.go new file mode 100644 index 0000000..1d6765e --- /dev/null +++ b/internal/config/config.go @@ -0,0 +1,245 @@ +/* +Package config provides configuration types and JSON handling for rs2322tcp. + +The server configuration describes the physical serial devices available +on the remote system. The client configuration describes the connection +to the server and the user's assignment of virtual serial ports to remote +devices. + +Project: rs2322tcp +Module: git.lang-dieter.de/rs2322tcp +*/ +package config + +import ( + "encoding/json" + "fmt" + "os" +) + +/////////////////////////////////////////////////////////////////////////////// +// Server configuration +/////////////////////////////////////////////////////////////////////////////// + +// ServerConfig contains the complete server configuration. +type ServerConfig struct { + Listen ListenConfig `json:"listen"` + Devices []DeviceConfig `json:"devices"` +} + +// ListenConfig contains the network listener configuration. +type ListenConfig struct { + Address string `json:"address"` + Port int `json:"port"` +} + +// DeviceConfig describes one physical RS232 device connected to the server. +type DeviceConfig struct { + ID string `json:"id"` + Name string `json:"name"` + SerialPort string `json:"serial_port"` + BaudRate int `json:"baud_rate"` + DataBits int `json:"data_bits"` + Parity string `json:"parity"` + StopBits int `json:"stop_bits"` +} + +/////////////////////////////////////////////////////////////////////////////// +// Client configuration +/////////////////////////////////////////////////////////////////////////////// + +// ClientConfig contains the complete client configuration. +type ClientConfig struct { + Server ServerConnectionConfig `json:"server"` + VirtualPorts []VirtualPortConfig `json:"virtual_ports"` +} + +// ServerConnectionConfig contains the connection information for the +// rs2322tcp server. +type ServerConnectionConfig struct { + Address string `json:"address"` + Port int `json:"port"` +} + +// VirtualPortConfig describes one local virtual serial port and the +// remote device assigned to it. +type VirtualPortConfig struct { + Port string `json:"port"` + RemoteDevice string `json:"remote_device"` +} + +/////////////////////////////////////////////////////////////////////////////// +// JSON loading +/////////////////////////////////////////////////////////////////////////////// + +// LoadServer loads a server configuration from a JSON file. +func LoadServer(filename string) (*ServerConfig, error) { + var cfg ServerConfig + + if err := loadJSON(filename, &cfg); err != nil { + return nil, err + } + + if err := cfg.Validate(); err != nil { + return nil, err + } + + return &cfg, nil +} + +// LoadClient loads a client configuration from a JSON file. +func LoadClient(filename string) (*ClientConfig, error) { + var cfg ClientConfig + + if err := loadJSON(filename, &cfg); err != nil { + return nil, err + } + + if err := cfg.Validate(); err != nil { + return nil, err + } + + return &cfg, nil +} + +/////////////////////////////////////////////////////////////////////////////// +// JSON saving +/////////////////////////////////////////////////////////////////////////////// + +// SaveServer saves a server configuration as formatted JSON. +func SaveServer(filename string, cfg *ServerConfig) error { + if cfg == nil { + return fmt.Errorf("server configuration is nil") + } + + if err := cfg.Validate(); err != nil { + return err + } + + return saveJSON(filename, cfg) +} + +// SaveClient saves a client configuration as formatted JSON. +func SaveClient(filename string, cfg *ClientConfig) error { + if cfg == nil { + return fmt.Errorf("client configuration is nil") + } + + if err := cfg.Validate(); err != nil { + return err + } + + return saveJSON(filename, cfg) +} + +/////////////////////////////////////////////////////////////////////////////// +// Validation +/////////////////////////////////////////////////////////////////////////////// + +// Validate checks the server configuration for basic errors. +func (cfg *ServerConfig) Validate() error { + if cfg == nil { + return fmt.Errorf("server configuration is nil") + } + + if cfg.Listen.Port < 1 || cfg.Listen.Port > 65535 { + return fmt.Errorf("invalid listen port: %d", cfg.Listen.Port) + } + + ids := make(map[string]bool) + + for i, device := range cfg.Devices { + if device.ID == "" { + return fmt.Errorf("device %d: ID is empty", i) + } + + if device.SerialPort == "" { + return fmt.Errorf("device %q: serial port is empty", device.ID) + } + + if ids[device.ID] { + return fmt.Errorf("duplicate device ID: %q", device.ID) + } + + ids[device.ID] = true + + if device.BaudRate <= 0 { + return fmt.Errorf("device %q: invalid baud rate", device.ID) + } + + if device.DataBits < 5 || device.DataBits > 8 { + return fmt.Errorf("device %q: invalid data bits", device.ID) + } + + if device.StopBits < 1 || device.StopBits > 2 { + return fmt.Errorf("device %q: invalid stop bits", device.ID) + } + } + + return nil +} + +// Validate checks the client configuration for basic errors. +func (cfg *ClientConfig) Validate() error { + if cfg == nil { + return fmt.Errorf("client configuration is nil") + } + + if cfg.Server.Port < 1 || cfg.Server.Port > 65535 { + return fmt.Errorf("invalid server port: %d", cfg.Server.Port) + } + + ports := make(map[string]bool) + + for i, virtualPort := range cfg.VirtualPorts { + if virtualPort.Port == "" { + return fmt.Errorf("virtual port %d: port is empty", i) + } + + if virtualPort.RemoteDevice == "" { + return fmt.Errorf("virtual port %q: remote device is empty", + virtualPort.Port) + } + + if ports[virtualPort.Port] { + return fmt.Errorf("duplicate virtual port: %q", + virtualPort.Port) + } + + ports[virtualPort.Port] = true + } + + return nil +} + +/////////////////////////////////////////////////////////////////////////////// +// Internal JSON helpers +/////////////////////////////////////////////////////////////////////////////// + +func loadJSON(filename string, target interface{}) error { + data, err := os.ReadFile(filename) + if err != nil { + return fmt.Errorf("read configuration %q: %w", filename, err) + } + + if err := json.Unmarshal(data, target); err != nil { + return fmt.Errorf("parse configuration %q: %w", filename, err) + } + + return nil +} + +func saveJSON(filename string, value interface{}) error { + data, err := json.MarshalIndent(value, "", " ") + if err != nil { + return fmt.Errorf("encode configuration: %w", err) + } + + data = append(data, '\n') + + if err := os.WriteFile(filename, data, 0644); err != nil { + return fmt.Errorf("write configuration %q: %w", filename, err) + } + + return nil +} diff --git a/internal/config/config_test.go b/internal/config/config_test.go new file mode 100644 index 0000000..5e59cf2 --- /dev/null +++ b/internal/config/config_test.go @@ -0,0 +1,459 @@ +/* +Package config_test contains tests for the rs2322tcp configuration package. + +Project: rs2322tcp +Module: git.lang-dieter.de/rs2322tcp +*/ +package config_test + +import ( + "encoding/json" + "git.lang-dieter.de/rs2322tcp/internal/config" + "os" + "path/filepath" + "strings" + "testing" +) + +/////////////////////////////////////////////////////////////////////////////// +// Server configuration +/////////////////////////////////////////////////////////////////////////////// + +func TestLoadServer(t *testing.T) { + dir := t.TempDir() + filename := filepath.Join(dir, "server.json") + + data := `{ + "listen": { + "address": "0.0.0.0", + "port": 5000 + }, + "devices": [ + { + "id": "radio", + "name": "Funkgerät", + "serial_port": "/dev/ttyUSB0", + "baud_rate": 9600, + "data_bits": 8, + "parity": "none", + "stop_bits": 1 + } + ] + }` + + if err := os.WriteFile(filename, []byte(data), 0644); err != nil { + t.Fatalf("write test configuration: %v", err) + } + + cfg, err := config.LoadServer(filename) + if err != nil { + t.Fatalf("LoadServer() failed: %v", err) + } + + if cfg.Listen.Address != "0.0.0.0" { + t.Errorf("Listen.Address = %q, want %q", + cfg.Listen.Address, "0.0.0.0") + } + + if cfg.Listen.Port != 5000 { + t.Errorf("Listen.Port = %d, want %d", + cfg.Listen.Port, 5000) + } + + if len(cfg.Devices) != 1 { + t.Fatalf("len(Devices) = %d, want 1", len(cfg.Devices)) + } + + device := cfg.Devices[0] + + if device.ID != "radio" { + t.Errorf("Device.ID = %q, want %q", + device.ID, "radio") + } + + if device.SerialPort != "/dev/ttyUSB0" { + t.Errorf("Device.SerialPort = %q, want %q", + device.SerialPort, "/dev/ttyUSB0") + } + + if device.BaudRate != 9600 { + t.Errorf("Device.BaudRate = %d, want %d", + device.BaudRate, 9600) + } +} + +/////////////////////////////////////////////////////////////////////////////// +// Client configuration +/////////////////////////////////////////////////////////////////////////////// + +func TestLoadClient(t *testing.T) { + dir := t.TempDir() + filename := filepath.Join(dir, "client.json") + + data := `{ + "server": { + "address": "100.64.0.10", + "port": 5000 + }, + "virtual_ports": [ + { + "port": "COM7", + "remote_device": "radio" + } + ] + }` + + if err := os.WriteFile(filename, []byte(data), 0644); err != nil { + t.Fatalf("write test configuration: %v", err) + } + + cfg, err := config.LoadClient(filename) + if err != nil { + t.Fatalf("LoadClient() failed: %v", err) + } + + if cfg.Server.Address != "100.64.0.10" { + t.Errorf("Server.Address = %q, want %q", + cfg.Server.Address, "100.64.0.10") + } + + if cfg.Server.Port != 5000 { + t.Errorf("Server.Port = %d, want %d", + cfg.Server.Port, 5000) + } + + if len(cfg.VirtualPorts) != 1 { + t.Fatalf("len(VirtualPorts) = %d, want 1", + len(cfg.VirtualPorts)) + } + + virtualPort := cfg.VirtualPorts[0] + + if virtualPort.Port != "COM7" { + t.Errorf("VirtualPort.Port = %q, want %q", + virtualPort.Port, "COM7") + } + + if virtualPort.RemoteDevice != "radio" { + t.Errorf("VirtualPort.RemoteDevice = %q, want %q", + virtualPort.RemoteDevice, "radio") + } +} + +/////////////////////////////////////////////////////////////////////////////// +// Save and reload +/////////////////////////////////////////////////////////////////////////////// + +func TestSaveAndLoadClient(t *testing.T) { + dir := t.TempDir() + filename := filepath.Join(dir, "client.json") + + original := &config.ClientConfig{ + Server: config.ServerConnectionConfig{ + Address: "100.64.0.10", + Port: 5000, + }, + VirtualPorts: []config.VirtualPortConfig{ + { + Port: "COM7", + RemoteDevice: "radio", + }, + { + Port: "COM8", + RemoteDevice: "rotor", + }, + }, + } + + if err := config.SaveClient(filename, original); err != nil { + t.Fatalf("SaveClient() failed: %v", err) + } + + loaded, err := config.LoadClient(filename) + if err != nil { + t.Fatalf("LoadClient() failed: %v", err) + } + + if loaded.Server != original.Server { + t.Errorf("loaded Server differs from original") + } + + if len(loaded.VirtualPorts) != len(original.VirtualPorts) { + t.Fatalf("len(VirtualPorts) = %d, want %d", + len(loaded.VirtualPorts), + len(original.VirtualPorts)) + } + + for i := range original.VirtualPorts { + if loaded.VirtualPorts[i] != original.VirtualPorts[i] { + t.Errorf("VirtualPorts[%d] differs from original", i) + } + } +} + +/////////////////////////////////////////////////////////////////////////////// +// Validation +/////////////////////////////////////////////////////////////////////////////// + +func TestServerValidationDuplicateDeviceID(t *testing.T) { + cfg := &config.ServerConfig{ + Listen: config.ListenConfig{ + Port: 5000, + }, + Devices: []config.DeviceConfig{ + { + ID: "radio", + SerialPort: "/dev/ttyUSB0", + BaudRate: 9600, + DataBits: 8, + StopBits: 1, + }, + { + ID: "radio", + SerialPort: "/dev/ttyUSB1", + BaudRate: 4800, + DataBits: 8, + StopBits: 1, + }, + }, + } + + if err := cfg.Validate(); err == nil { + t.Fatal("Validate() succeeded, want duplicate ID error") + } +} + +func TestClientValidationDuplicateVirtualPort(t *testing.T) { + cfg := &config.ClientConfig{ + Server: config.ServerConnectionConfig{ + Port: 5000, + }, + VirtualPorts: []config.VirtualPortConfig{ + { + Port: "COM7", + RemoteDevice: "radio", + }, + { + Port: "COM7", + RemoteDevice: "rotor", + }, + }, + } + + if err := cfg.Validate(); err == nil { + t.Fatal("Validate() succeeded, want duplicate virtual port error") + } +} + +func TestServerValidationInvalidPort(t *testing.T) { + cfg := &config.ServerConfig{ + Listen: config.ListenConfig{ + Port: 70000, + }, + } + + if err := cfg.Validate(); err == nil { + t.Fatal("Validate() succeeded, want invalid port error") + } +} + +func TestClientValidationInvalidPort(t *testing.T) { + cfg := &config.ClientConfig{ + Server: config.ServerConnectionConfig{ + Port: 0, + }, + } + + if err := cfg.Validate(); err == nil { + t.Fatal("Validate() succeeded, want invalid port error") + } +} + +/////////////////////////////////////////////////////////////////////////////// +// Example configurations +/////////////////////////////////////////////////////////////////////////////// + +func TestExampleServerConfig(t *testing.T) { + cfg, err := config.LoadServer("../../configs/server.json") + if err != nil { + t.Fatalf("LoadServer() failed: %v", err) + } + + if len(cfg.Devices) != 2 { + t.Fatalf("len(Devices) = %d, want 2", len(cfg.Devices)) + } +} + +func TestExampleClientConfig(t *testing.T) { + cfg, err := config.LoadClient("../../configs/client.json") + if err != nil { + t.Fatalf("LoadClient() failed: %v", err) + } + + if len(cfg.VirtualPorts) != 2 { + t.Fatalf("len(VirtualPorts) = %d, want 2", len(cfg.VirtualPorts)) + } +} + +/////////////////////////////////////////////////////////////////////////////// +// Remote device +/////////////////////////////////////////////////////////////////////////////// + +func TestDeviceRemoteDevice(t *testing.T) { + device := config.DeviceConfig{ + ID: "radio", + Name: "Funkgerät", + SerialPort: "/dev/ttyUSB0", + BaudRate: 9600, + DataBits: 8, + Parity: "none", + StopBits: 1, + } + + remote := device.RemoteDevice() + + if remote.ID != "radio" { + t.Errorf("RemoteDevice.ID = %q, want %q", + remote.ID, "radio") + } + + if remote.Name != "Funkgerät" { + t.Errorf("RemoteDevice.Name = %q, want %q", + remote.Name, "Funkgerät") + } + + if remote.BaudRate != 9600 { + t.Errorf("RemoteDevice.BaudRate = %d, want %d", + remote.BaudRate, 9600) + } + + if remote.DataBits != 8 { + t.Errorf("RemoteDevice.DataBits = %d, want %d", + remote.DataBits, 8) + } + + if remote.Parity != "none" { + t.Errorf("RemoteDevice.Parity = %q, want %q", + remote.Parity, "none") + } + + if remote.StopBits != 1 { + t.Errorf("RemoteDevice.StopBits = %d, want %d", + remote.StopBits, 1) + } +} + +/////////////////////////////////////////////////////////////////////////////// +// Remote device list +/////////////////////////////////////////////////////////////////////////////// + +func TestServerRemoteDevices(t *testing.T) { + cfg := &config.ServerConfig{ + Devices: []config.DeviceConfig{ + { + ID: "radio", + Name: "Funkgerät", + SerialPort: "/dev/ttyUSB0", + BaudRate: 9600, + DataBits: 8, + Parity: "none", + StopBits: 1, + }, + { + ID: "rotor", + Name: "Antennenrotor", + SerialPort: "/dev/ttyUSB1", + BaudRate: 4800, + DataBits: 8, + Parity: "none", + StopBits: 1, + }, + }, + } + + list := cfg.RemoteDevices() + + if len(list.Devices) != 2 { + t.Fatalf("len(Devices) = %d, want 2", len(list.Devices)) + } + + if list.Devices[0].ID != "radio" { + t.Errorf("Devices[0].ID = %q, want %q", + list.Devices[0].ID, "radio") + } + + if list.Devices[1].ID != "rotor" { + t.Errorf("Devices[1].ID = %q, want %q", + list.Devices[1].ID, "rotor") + } + + if list.Devices[0].Name != "Funkgerät" { + t.Errorf("Devices[0].Name = %q, want %q", + list.Devices[0].Name, "Funkgerät") + } + + if list.Devices[0].BaudRate != 9600 { + t.Errorf("Devices[0].BaudRate = %d, want %d", + list.Devices[0].BaudRate, 9600) + } + + if list.Devices[1].BaudRate != 4800 { + t.Errorf("Devices[1].BaudRate = %d, want %d", + list.Devices[1].BaudRate, 4800) + } +} + +func TestServerRemoteDevicesNil(t *testing.T) { + var cfg *config.ServerConfig + + list := cfg.RemoteDevices() + + if len(list.Devices) != 0 { + t.Fatalf("len(Devices) = %d, want 0", len(list.Devices)) + } +} + +/////////////////////////////////////////////////////////////////////////////// +// Remote device JSON +/////////////////////////////////////////////////////////////////////////////// + +func TestRemoteDeviceJSON(t *testing.T) { + list := config.DeviceList{ + Devices: []config.RemoteDevice{ + { + ID: "radio", + Name: "Funkgerät", + BaudRate: 9600, + DataBits: 8, + Parity: "none", + StopBits: 1, + }, + }, + } + + data, err := json.Marshal(list) + if err != nil { + t.Fatalf("json.Marshal() failed: %v", err) + } + + jsonText := string(data) + + if strings.Contains(jsonText, "serial_port") { + t.Fatalf("JSON contains server-internal serial_port") + } + + var decoded config.DeviceList + + if err := json.Unmarshal(data, &decoded); err != nil { + t.Fatalf("json.Unmarshal() failed: %v", err) + } + + if len(decoded.Devices) != 1 { + t.Fatalf("len(Devices) = %d, want 1", len(decoded.Devices)) + } + + if decoded.Devices[0].ID != "radio" { + t.Errorf("Devices[0].ID = %q, want %q", + decoded.Devices[0].ID, "radio") + } +} diff --git a/internal/config/device.go b/internal/config/device.go new file mode 100644 index 0000000..19e113e --- /dev/null +++ b/internal/config/device.go @@ -0,0 +1,83 @@ +/* +Package config provides configuration types and JSON handling for rs2322tcp. + +This file contains the distinction between the server-internal physical +device configuration and the public device information made available +to clients. + +Project: rs2322tcp +Module: git.lang-dieter.de/rs2322tcp +*/ +package config + +/////////////////////////////////////////////////////////////////////////////// +// Remote device +/////////////////////////////////////////////////////////////////////////////// + +// RemoteDevice describes a physical RS232 device that is made available +// to a client by the rs2322tcp server. +// +// Server-internal information such as the physical serial device path +// is deliberately not included. +type RemoteDevice struct { + ID string `json:"id"` + Name string `json:"name"` + BaudRate int `json:"baud_rate"` + DataBits int `json:"data_bits"` + Parity string `json:"parity"` + StopBits int `json:"stop_bits"` +} + +/////////////////////////////////////////////////////////////////////////////// +// Remote device list +/////////////////////////////////////////////////////////////////////////////// + +// DeviceList contains the remote devices made available by the server. +type DeviceList struct { + Devices []RemoteDevice `json:"devices"` +} + +/////////////////////////////////////////////////////////////////////////////// +// Public device information +/////////////////////////////////////////////////////////////////////////////// + +// RemoteDevice returns the public description of a configured physical +// device. +// +// The physical serial port remains server-internal and is not exposed +// to the client. +func (cfg DeviceConfig) RemoteDevice() RemoteDevice { + return RemoteDevice{ + ID: cfg.ID, + Name: cfg.Name, + BaudRate: cfg.BaudRate, + DataBits: cfg.DataBits, + Parity: cfg.Parity, + StopBits: cfg.StopBits, + } +} + +/////////////////////////////////////////////////////////////////////////////// +// Remote device list +/////////////////////////////////////////////////////////////////////////////// + +// RemoteDevices returns the public device list for all configured +// physical devices. +// +// Server-internal information such as the physical serial port is +// deliberately omitted. +func (cfg *ServerConfig) RemoteDevices() DeviceList { + if cfg == nil { + return DeviceList{} + } + + devices := make([]RemoteDevice, 0, len(cfg.Devices)) + + for _, device := range cfg.Devices { + devices = append(devices, device.RemoteDevice()) + } + + return DeviceList{ + Devices: devices, + } +} diff --git a/internal/serial/serial.go b/internal/serial/serial.go new file mode 100644 index 0000000..82bbd37 --- /dev/null +++ b/internal/serial/serial.go @@ -0,0 +1,159 @@ +/* + * ============================================================================ + * Projekt.....: rs2322tcp + * Datei.......: serial.go + * Copyright (C) 2026 Dieter Lang + * + * SPDX-License-Identifier: GPL-3.0-or-later + * + * Beschreibung: + * Öffnet und verwaltet eine serielle Schnittstelle auf Basis der + * konfigurierten Geräteparameter. + * ============================================================================ + */ +package serial + +import ( + "fmt" + "io" + + bugserial "go.bug.st/serial" + + "git.lang-dieter.de/rs2322tcp/internal/config" +) + +/////////////////////////////////////////////////////////////////////////////// +// Connection +/////////////////////////////////////////////////////////////////////////////// + +// Connection represents an opened serial connection. +type Connection struct { + port bugserial.Port +} + +/////////////////////////////////////////////////////////////////////////////// +// Open +/////////////////////////////////////////////////////////////////////////////// + +// Open opens the serial device described by the configuration. +func Open(device config.DeviceConfig) (*Connection, error) { + mode, err := createMode(device) + if err != nil { + return nil, err + } + + port, err := bugserial.Open(device.SerialPort, mode) + if err != nil { + return nil, fmt.Errorf( + "open serial port %q: %w", + device.SerialPort, + err, + ) + } + + return &Connection{ + port: port, + }, nil +} + +/////////////////////////////////////////////////////////////////////////////// +// Mode +/////////////////////////////////////////////////////////////////////////////// + +func createMode(device config.DeviceConfig) (*bugserial.Mode, error) { + mode := &bugserial.Mode{ + BaudRate: device.BaudRate, + DataBits: device.DataBits, + Parity: bugserial.NoParity, + StopBits: bugserial.OneStopBit, + } + + if device.BaudRate <= 0 { + return nil, fmt.Errorf( + "invalid baud rate: %d", + device.BaudRate, + ) + } + + switch device.DataBits { + case 5, 6, 7, 8: + // Supported by go.bug.st/serial. + + default: + return nil, fmt.Errorf( + "unsupported data bits: %d", + device.DataBits, + ) + } + + switch device.Parity { + case "none": + mode.Parity = bugserial.NoParity + + case "odd": + mode.Parity = bugserial.OddParity + + case "even": + mode.Parity = bugserial.EvenParity + + default: + return nil, fmt.Errorf( + "unsupported parity: %q", + device.Parity, + ) + } + + switch device.StopBits { + case 1: + mode.StopBits = bugserial.OneStopBit + + case 2: + mode.StopBits = bugserial.TwoStopBits + + default: + return nil, fmt.Errorf( + "unsupported stop bits: %d", + device.StopBits, + ) + } + + return mode, nil +} + +/////////////////////////////////////////////////////////////////////////////// +// Read / Write +/////////////////////////////////////////////////////////////////////////////// + +// Read reads data from the serial connection. +func (c *Connection) Read(p []byte) (int, error) { + if c == nil || c.port == nil { + return 0, io.ErrClosedPipe + } + + return c.port.Read(p) +} + +// Write writes data to the serial connection. +func (c *Connection) Write(p []byte) (int, error) { + if c == nil || c.port == nil { + return 0, io.ErrClosedPipe + } + + return c.port.Write(p) +} + +/////////////////////////////////////////////////////////////////////////////// +// Close +/////////////////////////////////////////////////////////////////////////////// + +// Close closes the serial connection. +func (c *Connection) Close() error { + if c == nil || c.port == nil { + return nil + } + + err := c.port.Close() + c.port = nil + + return err +} diff --git a/internal/serial/serial_test.go b/internal/serial/serial_test.go new file mode 100644 index 0000000..182db64 --- /dev/null +++ b/internal/serial/serial_test.go @@ -0,0 +1,391 @@ +/* + * ============================================================================ + * Projekt.....: rs2322tcp + * Datei.......: serial_test.go + * Copyright (C) 2026 Dieter Lang + * + * SPDX-License-Identifier: GPL-3.0-or-later + * + * Beschreibung: + * Tests für die serielle Schnittstellenabstraktion. + * Zusätzlich wird unter Linux mit socat ein virtuelles serielles Portpaar + * erzeugt, um Read, Write und Close ohne reale RS232-Hardware zu testen. + * ============================================================================ + */ +package serial + +import ( + "fmt" + "os" + "os/exec" + "path/filepath" + "testing" + "time" + + "git.lang-dieter.de/rs2322tcp/internal/config" +) + +/////////////////////////////////////////////////////////////////////////////// +// Test helpers +/////////////////////////////////////////////////////////////////////////////// + +func testDevice() config.DeviceConfig { + return config.DeviceConfig{ + Name: "Testgerät", + SerialPort: "/dev/ttyTEST", + BaudRate: 9600, + DataBits: 8, + Parity: "none", + StopBits: 1, + } +} + +/////////////////////////////////////////////////////////////////////////////// +// Mode tests +/////////////////////////////////////////////////////////////////////////////// + +func TestCreateMode(t *testing.T) { + device := testDevice() + + mode, err := createMode(device) + if err != nil { + t.Fatalf("createMode() failed: %v", err) + } + + if mode.BaudRate != 9600 { + t.Fatalf("BaudRate = %d, want 9600", mode.BaudRate) + } + + if mode.DataBits != 8 { + t.Fatalf("DataBits = %d, want 8", mode.DataBits) + } +} + +func TestCreateModeParity(t *testing.T) { + tests := []struct { + name string + parity string + }{ + {"none", "none"}, + {"odd", "odd"}, + {"even", "even"}, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + device := testDevice() + device.Parity = tt.parity + + if _, err := createMode(device); err != nil { + t.Fatalf( + "createMode() with parity %q failed: %v", + tt.parity, + err, + ) + } + }) + } +} + +func TestCreateModeStopBits(t *testing.T) { + for _, stopBits := range []int{1, 2} { + t.Run(fmt.Sprintf("stopbits_%d", stopBits), func(t *testing.T) { + device := testDevice() + device.StopBits = stopBits + + if _, err := createMode(device); err != nil { + t.Fatalf( + "createMode() with stop bits %d failed: %v", + stopBits, + err, + ) + } + }) + } +} + +/////////////////////////////////////////////////////////////////////////////// +// Invalid configuration +/////////////////////////////////////////////////////////////////////////////// + +func TestCreateModeRejectsInvalidParity(t *testing.T) { + device := testDevice() + device.Parity = "invalid" + + if _, err := createMode(device); err == nil { + t.Fatal("createMode() succeeded with invalid parity") + } +} + +func TestCreateModeRejectsInvalidStopBits(t *testing.T) { + device := testDevice() + device.StopBits = 3 + + if _, err := createMode(device); err == nil { + t.Fatal("createMode() succeeded with invalid stop bits") + } +} + +func TestCreateModeRejectsInvalidDataBits(t *testing.T) { + device := testDevice() + device.DataBits = 9 + + if _, err := createMode(device); err == nil { + t.Fatal("createMode() succeeded with invalid data bits") + } +} + +func TestCreateModeRejectsInvalidBaudRate(t *testing.T) { + device := testDevice() + device.BaudRate = 0 + + if _, err := createMode(device); err == nil { + t.Fatal("createMode() succeeded with invalid baud rate") + } +} + +/////////////////////////////////////////////////////////////////////////////// +// Open +/////////////////////////////////////////////////////////////////////////////// + +func TestOpenInvalidPort(t *testing.T) { + device := testDevice() + + _, err := Open(device) + if err == nil { + t.Fatal("Open() succeeded with nonexistent serial port") + } +} + +/////////////////////////////////////////////////////////////////////////////// +// Nil connection +/////////////////////////////////////////////////////////////////////////////// + +func TestNilConnectionRead(t *testing.T) { + var connection *Connection + + buffer := make([]byte, 1) + + _, err := connection.Read(buffer) + if err == nil { + t.Fatal("Read() succeeded on nil connection") + } +} + +func TestNilConnectionWrite(t *testing.T) { + var connection *Connection + + _, err := connection.Write([]byte("test")) + if err == nil { + t.Fatal("Write() succeeded on nil connection") + } +} + +func TestNilConnectionClose(t *testing.T) { + var connection *Connection + + if err := connection.Close(); err != nil { + t.Fatalf("Close() failed: %v", err) + } +} + +/////////////////////////////////////////////////////////////////////////////// +// Virtual serial port +/////////////////////////////////////////////////////////////////////////////// + +func startVirtualSerialPair(t *testing.T) (string, string, func()) { + t.Helper() + + if _, err := exec.LookPath("socat"); err != nil { + t.Skip("socat not installed") + } + + dir := t.TempDir() + + portA := filepath.Join(dir, "ttyA") + portB := filepath.Join(dir, "ttyB") + + cmd := exec.Command( + "socat", + "-d", + "-d", + fmt.Sprintf("pty,raw,echo=0,link=%s", portA), + fmt.Sprintf("pty,raw,echo=0,link=%s", portB), + ) + + if err := cmd.Start(); err != nil { + t.Fatalf("failed to start socat: %v", err) + } + + cleanup := func() { + if cmd.Process != nil { + _ = cmd.Process.Kill() + _ = cmd.Wait() + } + } + + deadline := time.Now().Add(2 * time.Second) + + for { + if _, errA := os.Stat(portA); errA == nil { + if _, errB := os.Stat(portB); errB == nil { + break + } + } + + if time.Now().After(deadline) { + cleanup() + t.Fatal("timeout waiting for virtual serial ports") + } + + time.Sleep(10 * time.Millisecond) + } + + return portA, portB, cleanup +} + +/////////////////////////////////////////////////////////////////////////////// +// Read / Write integration test +/////////////////////////////////////////////////////////////////////////////// + +func TestVirtualSerialReadWrite(t *testing.T) { + portA, portB, cleanup := startVirtualSerialPair(t) + defer cleanup() + + device := testDevice() + device.SerialPort = portA + + connection, err := Open(device) + if err != nil { + t.Fatalf("Open() failed: %v", err) + } + defer connection.Close() + + peer, err := os.OpenFile( + portB, + os.O_RDWR, + 0, + ) + if err != nil { + t.Fatalf("opening peer port failed: %v", err) + } + defer peer.Close() + + //////////////////////////////////////////////////////////////////////////// + // Serial -> peer + //////////////////////////////////////////////////////////////////////////// + + outgoing := []byte("hello from serial") + + n, err := connection.Write(outgoing) + if err != nil { + t.Fatalf("serial Write() failed: %v", err) + } + + if n != len(outgoing) { + t.Fatalf( + "serial Write() wrote %d bytes, want %d", + n, + len(outgoing), + ) + } + + received := make([]byte, len(outgoing)) + + if err := readWithTimeout( + peer, + received, + 2*time.Second, + ); err != nil { + t.Fatalf("peer Read() failed: %v", err) + } + + if string(received) != string(outgoing) { + t.Fatalf( + "peer received %q, want %q", + string(received), + string(outgoing), + ) + } + + //////////////////////////////////////////////////////////////////////////// + // peer -> Serial + //////////////////////////////////////////////////////////////////////////// + + incoming := []byte("hello from peer") + + n, err = peer.Write(incoming) + if err != nil { + t.Fatalf("peer Write() failed: %v", err) + } + + if n != len(incoming) { + t.Fatalf( + "peer Write() wrote %d bytes, want %d", + n, + len(incoming), + ) + } + + received = make([]byte, len(incoming)) + + if err := readWithTimeout( + connection, + received, + 2*time.Second, + ); err != nil { + t.Fatalf("serial Read() failed: %v", err) + } + + if string(received) != string(incoming) { + t.Fatalf( + "serial received %q, want %q", + string(received), + string(incoming), + ) + } +} + +/////////////////////////////////////////////////////////////////////////////// +// Timeout helper +/////////////////////////////////////////////////////////////////////////////// + +type reader interface { + Read([]byte) (int, error) +} + +func readWithTimeout( + r reader, + buffer []byte, + timeout time.Duration, +) error { + result := make(chan error, 1) + + go func() { + n, err := r.Read(buffer) + + if err != nil { + result <- err + return + } + + if n != len(buffer) { + result <- fmt.Errorf( + "read %d bytes, want %d", + n, + len(buffer), + ) + return + } + + result <- nil + }() + + select { + case err := <-result: + return err + + case <-time.After(timeout): + return fmt.Errorf("read timeout after %s", timeout) + } +} diff --git a/internal/server/control.go b/internal/server/control.go new file mode 100644 index 0000000..c3e82f2 --- /dev/null +++ b/internal/server/control.go @@ -0,0 +1,508 @@ +/* + * ============================================================================ + * Projekt.....: rs2322tcp + * Datei.......: control.go + * Copyright (C) 2026 Dieter Lang + * + * SPDX-License-Identifier: GPL-3.0-or-later + * + * Beschreibung: + * TCP-Control-Server für rs2322tcp. + * Verwaltung der Control-Verbindungen, Client-Sessions und Bereitstellung + * der aktuellen Geräteinformationen einschließlich dynamischer Data-Ports. + * Verwaltung der dynamischen Datenverbindungen zwischen TCP und RS232. + * ============================================================================ + */ +package server + +import ( + "bufio" + "errors" + "fmt" + "io" + "log" + "net" + "sync/atomic" + + "git.lang-dieter.de/rs2322tcp/internal/config" + "git.lang-dieter.de/rs2322tcp/internal/serial" + "git.lang-dieter.de/rs2322tcp/internal/transport" +) + +/////////////////////////////////////////////////////////////////////////////// +// ControlServer +/////////////////////////////////////////////////////////////////////////////// + +// ControlServer implements the rs2322tcp control listener. +type ControlServer struct { + config *config.ServerConfig + listener net.Listener + nextSessionID uint64 +} + +/////////////////////////////////////////////////////////////////////////////// +// Constructor +/////////////////////////////////////////////////////////////////////////////// + +// NewControlServer creates a new control server using the supplied +// server configuration. +func NewControlServer(cfg *config.ServerConfig) (*ControlServer, error) { + if cfg == nil { + return nil, fmt.Errorf("server configuration is nil") + } + + if err := cfg.Validate(); err != nil { + return nil, err + } + + return &ControlServer{ + config: cfg, + }, nil +} + +/////////////////////////////////////////////////////////////////////////////// +// Listener +/////////////////////////////////////////////////////////////////////////////// + +// Listen starts the TCP control listener. +func (s *ControlServer) Listen() error { + if s == nil { + return fmt.Errorf("control server is nil") + } + + if s.listener != nil { + return fmt.Errorf("control server is already listening") + } + + address := net.JoinHostPort( + s.config.Listen.Address, + fmt.Sprintf("%d", s.config.Listen.Port), + ) + + listener, err := net.Listen("tcp", address) + if err != nil { + return fmt.Errorf("listen on %s: %w", address, err) + } + + s.listener = listener + + return nil +} + +// Addr returns the actual address of the control listener. +func (s *ControlServer) Addr() net.Addr { + if s == nil || s.listener == nil { + return nil + } + + return s.listener.Addr() +} + +// Close stops the control listener. +func (s *ControlServer) Close() error { + if s == nil || s.listener == nil { + return nil + } + + err := s.listener.Close() + s.listener = nil + + return err +} + +/////////////////////////////////////////////////////////////////////////////// +// Serve +/////////////////////////////////////////////////////////////////////////////// + +// Serve accepts incoming control connections. +// +// Each connection is handled independently. +func (s *ControlServer) Serve() error { + if s == nil { + return fmt.Errorf("control server is nil") + } + + if s.listener == nil { + return fmt.Errorf("control server is not listening") + } + + for { + conn, err := s.listener.Accept() + if err != nil { + if errors.Is(err, net.ErrClosed) { + return nil + } + + return fmt.Errorf("accept control connection: %w", err) + } + + go s.handleConnection(conn) + } +} + +/////////////////////////////////////////////////////////////////////////////// +// Connection handling +/////////////////////////////////////////////////////////////////////////////// + +// handleConnection creates a session for one control connection and +// handles all control messages belonging to that session. +func (s *ControlServer) handleConnection(conn net.Conn) *Session { + session, err := s.newSession(conn) + if err != nil { + log.Printf("create session: %v", err) + _ = conn.Close() + return nil + } + + log.Printf("session %d connected", session.ID()) + + defer func() { + if err := session.Close(); err != nil { + log.Printf( + "session %d close error: %v", + session.ID(), + err, + ) + } + + log.Printf("session %d closed", session.ID()) + }() + + reader := bufio.NewReader(conn) + + for { + var message transport.Message + + if err := transport.ReadMessage(reader, &message); err != nil { + if !errors.Is(err, io.EOF) { + log.Printf( + "session %d control connection error: %v", + session.ID(), + err, + ) + } + + return session + } + + if err := s.handleMessage(session, message); err != nil { + log.Printf( + "session %d control message error: %v", + session.ID(), + err, + ) + + _ = transport.WriteMessage( + conn, + transport.NewError( + transport.ErrorInternal, + err.Error(), + ), + ) + + return session + } + } +} + +// HandleConnectionForTest handles one control connection and returns +// the session after the connection has ended. +// +// This method exists only to allow package-level tests without opening +// a real TCP control listener. +func (s *ControlServer) HandleConnectionForTest(conn net.Conn) *Session { + return s.handleConnection(conn) +} + +/////////////////////////////////////////////////////////////////////////////// +// Session creation +/////////////////////////////////////////////////////////////////////////////// + +// newSession creates a new client session, assigns a unique ID and +// creates one dynamic data listener for every configured device. +// +// A data worker is started for every listener. The worker waits for an +// incoming TCP data connection and opens the corresponding serial device +// only when a client actually connects to the dynamic data port. +func (s *ControlServer) newSession(conn net.Conn) (*Session, error) { + if conn == nil { + return nil, fmt.Errorf("connection is nil") + } + + id := atomic.AddUint64(&s.nextSessionID, 1) + + session, err := NewSession(id, conn) + if err != nil { + return nil, err + } + + for _, device := range s.config.Devices { + listener, err := NewDataListener() + if err != nil { + _ = session.Close() + + return nil, fmt.Errorf( + "create data listener for device %q: %w", + device.ID, + err, + ) + } + + if err := session.AddDataListener( + device.ID, + listener, + ); err != nil { + _ = session.Close() + + return nil, fmt.Errorf( + "add data listener for device %q: %w", + device.ID, + err, + ) + } + + go s.runDataListener(session, device, listener) + } + + return session, nil +} + +/////////////////////////////////////////////////////////////////////////////// +// Data connection handling +/////////////////////////////////////////////////////////////////////////////// + +// runDataListener waits for incoming TCP data connections for one +// configured device. +// +// The listener remains active for the lifetime of the session. A serial +// connection is opened only after a TCP client actually connects. +// +// Only one active data connection is permitted for each device. +func (s *ControlServer) runDataListener( + session *Session, + device config.DeviceConfig, + listener *DataListener, +) { + if session == nil || listener == nil { + return + } + + for { + tcpConn, err := listener.Accept() + if err != nil { + if session.IsClosed() { + return + } + + log.Printf( + "session %d device %q data accept error: %v", + session.ID(), + device.ID, + err, + ) + + return + } + + if session.IsClosed() { + _ = tcpConn.Close() + return + } + + if _, exists := session.DataConnection(device.ID); exists { + log.Printf( + "session %d device %q already has an active data connection", + session.ID(), + device.ID, + ) + + _ = tcpConn.Close() + continue + } + + serialConn, err := serial.Open(device) + if err != nil { + log.Printf( + "session %d device %q serial open error: %v", + session.ID(), + device.ID, + err, + ) + + _ = tcpConn.Close() + continue + } + + dataConnection, err := NewDataConnection( + tcpConn, + serialConn, + ) + if err != nil { + log.Printf( + "session %d device %q create data connection error: %v", + session.ID(), + device.ID, + err, + ) + + _ = serialConn.Close() + _ = tcpConn.Close() + + continue + } + + if err := session.AddDataConnection( + device.ID, + dataConnection, + ); err != nil { + log.Printf( + "session %d device %q register data connection error: %v", + session.ID(), + device.ID, + err, + ) + + continue + } + + log.Printf( + "session %d device %q data connection established", + session.ID(), + device.ID, + ) + + err = dataConnection.Run() + + if err != nil { + log.Printf( + "session %d device %q data connection error: %v", + session.ID(), + device.ID, + err, + ) + } + + if err := session.RemoveDataConnection(device.ID); err != nil { + log.Printf( + "session %d device %q remove data connection error: %v", + session.ID(), + device.ID, + err, + ) + } + + log.Printf( + "session %d device %q data connection closed", + session.ID(), + device.ID, + ) + + if session.IsClosed() { + return + } + } +} + +/////////////////////////////////////////////////////////////////////////////// +// Message handling +/////////////////////////////////////////////////////////////////////////////// + +// handleMessage processes one control message for the specified session. +func (s *ControlServer) handleMessage( + session *Session, + message transport.Message, +) error { + if session == nil { + return fmt.Errorf("session is nil") + } + + conn := session.Conn() + if conn == nil { + return fmt.Errorf("session control connection is closed") + } + + if message.Version != transport.ProtocolVersion { + return transport.WriteMessage( + conn, + transport.NewError( + transport.ErrorUnsupportedVersion, + fmt.Sprintf( + "unsupported protocol version: %d", + message.Version, + ), + ), + ) + } + + switch message.Type { + case transport.MessageHello: + return transport.WriteMessage( + conn, + transport.NewHelloResponse(), + ) + + case transport.MessageGetDevices: + return s.writeDeviceList(session) + + default: + return transport.WriteMessage( + conn, + transport.NewError( + transport.ErrorUnknownMessage, + fmt.Sprintf( + "unknown message type: %q", + message.Type, + ), + ), + ) + } +} + +/////////////////////////////////////////////////////////////////////////////// +// Device list +/////////////////////////////////////////////////////////////////////////////// + +// writeDeviceList sends the configured remote devices together with +// their session-specific dynamic data ports. +func (s *ControlServer) writeDeviceList(session *Session) error { + if session == nil { + return fmt.Errorf("session is nil") + } + + conn := session.Conn() + if conn == nil { + return fmt.Errorf("session control connection is closed") + } + + remoteDevices := s.config.RemoteDevices() + + devices := make( + []transport.RemoteDeviceInfo, + 0, + len(remoteDevices.Devices), + ) + + for _, device := range remoteDevices.Devices { + dataPort := session.DataPort(device.ID) + + if dataPort == 0 { + return fmt.Errorf( + "no data listener for device %q", + device.ID, + ) + } + + devices = append( + devices, + transport.NewRemoteDeviceInfo( + device, + dataPort, + ), + ) + } + + return transport.WriteMessage( + conn, + transport.NewDeviceList(devices), + ) +} diff --git a/internal/server/control_test.go b/internal/server/control_test.go new file mode 100644 index 0000000..d8567c8 --- /dev/null +++ b/internal/server/control_test.go @@ -0,0 +1,635 @@ +/* + * ============================================================================ + * Projekt.....: rs2322tcp + * Datei.......: control_test.go + * Copyright (C) 2026 Dieter Lang + * + * SPDX-License-Identifier: GPL-3.0-or-later + * + * Beschreibung: + * Tests für den TCP-Control-Server einschließlich Session-Verwaltung, + * dynamischer Data-Ports und der bidirektionalen Verbindung zur + * seriellen Schnittstelle. + * ============================================================================ + */ +package server_test + +import ( + "bufio" + "fmt" + "io" + "net" + "os" + "os/exec" + "path/filepath" + "testing" + "time" + + "git.lang-dieter.de/rs2322tcp/internal/config" + "git.lang-dieter.de/rs2322tcp/internal/server" + "git.lang-dieter.de/rs2322tcp/internal/transport" +) + +/////////////////////////////////////////////////////////////////////////////// +// Test configuration +/////////////////////////////////////////////////////////////////////////////// + +func testServerConfig() *config.ServerConfig { + return &config.ServerConfig{ + Listen: config.ListenConfig{ + Address: "127.0.0.1", + Port: 5000, + }, + Devices: []config.DeviceConfig{ + { + ID: "radio", + Name: "Funkgerät", + SerialPort: "/dev/ttyUSB0", + BaudRate: 9600, + DataBits: 8, + Parity: "none", + StopBits: 1, + }, + { + ID: "rotor", + Name: "Antennenrotor", + SerialPort: "/dev/ttyUSB1", + BaudRate: 4800, + DataBits: 8, + Parity: "none", + StopBits: 1, + }, + }, + } +} + +/////////////////////////////////////////////////////////////////////////////// +// Basic server tests +/////////////////////////////////////////////////////////////////////////////// + +func TestNewControlServer(t *testing.T) { + cfg := testServerConfig() + + srv, err := server.NewControlServer(cfg) + if err != nil { + t.Fatalf("NewControlServer() failed: %v", err) + } + + if srv == nil { + t.Fatal("NewControlServer() returned nil server") + } +} + +func TestControlServerHello(t *testing.T) { + cfg := testServerConfig() + + srv, err := server.NewControlServer(cfg) + if err != nil { + t.Fatalf("NewControlServer() failed: %v", err) + } + + serverConn, clientConn := net.Pipe() + defer serverConn.Close() + defer clientConn.Close() + + done := make(chan struct{}) + + var session *server.Session + + go func() { + session = srv.HandleConnectionForTest(serverConn) + close(done) + }() + + if err := transport.WriteMessage( + clientConn, + transport.NewHello(), + ); err != nil { + t.Fatalf("WriteMessage() failed: %v", err) + } + + reader := bufio.NewReader(clientConn) + + var response transport.HelloResponseMessage + + if err := transport.ReadMessage(reader, &response); err != nil { + t.Fatalf("ReadMessage() failed: %v", err) + } + + if response.Type != transport.MessageHelloResponse { + t.Errorf( + "Type = %q, want %q", + response.Type, + transport.MessageHelloResponse, + ) + } + + clientConn.Close() + + select { + case <-done: + case <-time.After(time.Second): + t.Fatal("server connection handler did not terminate") + } + + if session == nil { + t.Fatal("server returned nil session") + } + + if session.ID() != 1 { + t.Errorf("Session ID = %d, want 1", session.ID()) + } + + if !session.IsClosed() { + t.Fatal("session is not closed") + } +} + +func TestControlServerGetDevices(t *testing.T) { + cfg := testServerConfig() + + srv, err := server.NewControlServer(cfg) + if err != nil { + t.Fatalf("NewControlServer() failed: %v", err) + } + + serverConn, clientConn := net.Pipe() + defer serverConn.Close() + defer clientConn.Close() + + done := make(chan struct{}) + + go func() { + srv.HandleConnectionForTest(serverConn) + close(done) + }() + + if err := transport.WriteMessage( + clientConn, + transport.NewGetDevices(), + ); err != nil { + t.Fatalf("WriteMessage() failed: %v", err) + } + + reader := bufio.NewReader(clientConn) + + var response transport.DeviceListMessage + + if err := transport.ReadMessage(reader, &response); err != nil { + t.Fatalf("ReadMessage() failed: %v", err) + } + + if response.Type != transport.MessageDeviceList { + t.Errorf( + "Type = %q, want %q", + response.Type, + transport.MessageDeviceList, + ) + } + + if len(response.Devices) != 2 { + t.Fatalf( + "len(Devices) = %d, want 2", + len(response.Devices), + ) + } + + if response.Devices[0].ID != "radio" { + t.Errorf( + "Devices[0].ID = %q, want %q", + response.Devices[0].ID, + "radio", + ) + } + + if response.Devices[1].ID != "rotor" { + t.Errorf( + "Devices[1].ID = %q, want %q", + response.Devices[1].ID, + "rotor", + ) + } + + if response.Devices[0].DataPort == 0 { + t.Fatal("radio data port is zero") + } + + if response.Devices[1].DataPort == 0 { + t.Fatal("rotor data port is zero") + } + + clientConn.Close() + + select { + case <-done: + case <-time.After(time.Second): + t.Fatal("server connection handler did not terminate") + } +} + +func TestControlServerReconnectCreatesNewSession(t *testing.T) { + cfg := testServerConfig() + + srv, err := server.NewControlServer(cfg) + if err != nil { + t.Fatalf("NewControlServer() failed: %v", err) + } + + // First connection. + serverConn1, clientConn1 := net.Pipe() + + done1 := make(chan *server.Session) + + go func() { + done1 <- srv.HandleConnectionForTest(serverConn1) + }() + + if err := transport.WriteMessage( + clientConn1, + transport.NewHello(), + ); err != nil { + t.Fatalf("first WriteMessage() failed: %v", err) + } + + reader1 := bufio.NewReader(clientConn1) + + var response1 transport.HelloResponseMessage + + if err := transport.ReadMessage(reader1, &response1); err != nil { + t.Fatalf("first ReadMessage() failed: %v", err) + } + + if response1.Type != transport.MessageHelloResponse { + t.Errorf( + "first response Type = %q, want %q", + response1.Type, + transport.MessageHelloResponse, + ) + } + + if err := clientConn1.Close(); err != nil { + t.Fatalf("close first client connection: %v", err) + } + + session1 := <-done1 + + if session1 == nil { + t.Fatal("first session is nil") + } + + if session1.ID() != 1 { + t.Errorf( + "first session ID = %d, want 1", + session1.ID(), + ) + } + + if !session1.IsClosed() { + t.Fatal("first session is not closed") + } + + // Second connection - simulates reconnect. + serverConn2, clientConn2 := net.Pipe() + + done2 := make(chan *server.Session) + + go func() { + done2 <- srv.HandleConnectionForTest(serverConn2) + }() + + if err := transport.WriteMessage( + clientConn2, + transport.NewHello(), + ); err != nil { + t.Fatalf("second WriteMessage() failed: %v", err) + } + + reader2 := bufio.NewReader(clientConn2) + + var response2 transport.HelloResponseMessage + + if err := transport.ReadMessage(reader2, &response2); err != nil { + t.Fatalf("second ReadMessage() failed: %v", err) + } + + if response2.Type != transport.MessageHelloResponse { + t.Errorf( + "second response Type = %q, want %q", + response2.Type, + transport.MessageHelloResponse, + ) + } + + if err := clientConn2.Close(); err != nil { + t.Fatalf("close second client connection: %v", err) + } + + session2 := <-done2 + + if session2 == nil { + t.Fatal("second session is nil") + } + + if session2.ID() != 2 { + t.Errorf( + "second session ID = %d, want 2", + session2.ID(), + ) + } + + if session2.ID() == session1.ID() { + t.Fatal("reconnect reused the old session ID") + } + + if !session2.IsClosed() { + t.Fatal("second session is not closed") + } +} + +/////////////////////////////////////////////////////////////////////////////// +// Virtual serial pair +/////////////////////////////////////////////////////////////////////////////// + +func startVirtualSerialPair(t *testing.T) (string, string, func()) { + t.Helper() + + if _, err := exec.LookPath("socat"); err != nil { + t.Skip("socat not installed") + } + + dir := t.TempDir() + + portA := filepath.Join(dir, "ttyA") + portB := filepath.Join(dir, "ttyB") + + cmd := exec.Command( + "socat", + "-d", + "-d", + fmt.Sprintf( + "pty,raw,echo=0,link=%s", + portA, + ), + fmt.Sprintf( + "pty,raw,echo=0,link=%s", + portB, + ), + ) + + if err := cmd.Start(); err != nil { + t.Fatalf("failed to start socat: %v", err) + } + + cleanup := func() { + if cmd.Process != nil { + _ = cmd.Process.Kill() + _ = cmd.Wait() + } + } + + deadline := time.Now().Add(2 * time.Second) + + for { + _, errA := os.Stat(portA) + _, errB := os.Stat(portB) + + if errA == nil && errB == nil { + break + } + + if time.Now().After(deadline) { + cleanup() + t.Fatal("timeout waiting for virtual serial ports") + } + + time.Sleep(10 * time.Millisecond) + } + + return portA, portB, cleanup +} + +/////////////////////////////////////////////////////////////////////////////// +// Read helper +/////////////////////////////////////////////////////////////////////////////// + +func readExactWithTimeout( + t *testing.T, + reader io.Reader, + buffer []byte, + timeout time.Duration, +) { + t.Helper() + + done := make(chan error, 1) + + go func() { + _, err := io.ReadFull(reader, buffer) + done <- err + }() + + select { + case err := <-done: + if err != nil { + t.Fatalf("read failed: %v", err) + } + + case <-time.After(timeout): + t.Fatalf("read timeout after %s", timeout) + } +} + +/////////////////////////////////////////////////////////////////////////////// +// Data connection integration test +/////////////////////////////////////////////////////////////////////////////// + +func TestControlServerDataConnection(t *testing.T) { + portA, portB, cleanup := startVirtualSerialPair(t) + defer cleanup() + + cfg := &config.ServerConfig{ + Listen: config.ListenConfig{ + Address: "127.0.0.1", + Port: 5000, + }, + Devices: []config.DeviceConfig{ + { + ID: "radio", + Name: "Funkgerät", + SerialPort: portA, + BaudRate: 9600, + DataBits: 8, + Parity: "none", + StopBits: 1, + }, + }, + } + + srv, err := server.NewControlServer(cfg) + if err != nil { + t.Fatalf("NewControlServer() failed: %v", err) + } + + //////////////////////////////////////////////////////////////////////////// + // Control connection + //////////////////////////////////////////////////////////////////////////// + + serverConn, clientConn := net.Pipe() + defer serverConn.Close() + defer clientConn.Close() + + done := make(chan *server.Session) + + go func() { + done <- srv.HandleConnectionForTest(serverConn) + }() + + if err := transport.WriteMessage( + clientConn, + transport.NewGetDevices(), + ); err != nil { + t.Fatalf("WriteMessage() failed: %v", err) + } + + controlReader := bufio.NewReader(clientConn) + + var deviceList transport.DeviceListMessage + + if err := transport.ReadMessage( + controlReader, + &deviceList, + ); err != nil { + t.Fatalf("ReadMessage() failed: %v", err) + } + + if len(deviceList.Devices) != 1 { + t.Fatalf( + "len(Devices) = %d, want 1", + len(deviceList.Devices), + ) + } + + if deviceList.Devices[0].ID != "radio" { + t.Fatalf( + "device ID = %q, want %q", + deviceList.Devices[0].ID, + "radio", + ) + } + + dataPort := deviceList.Devices[0].DataPort + if dataPort == 0 { + t.Fatal("data port is zero") + } + + //////////////////////////////////////////////////////////////////////////// + // Data connection + //////////////////////////////////////////////////////////////////////////// + + dataConn, err := net.DialTimeout( + "tcp", + net.JoinHostPort("127.0.0.1", fmt.Sprintf("%d", dataPort)), + time.Second, + ) + if err != nil { + t.Fatalf("connect data port failed: %v", err) + } + defer dataConn.Close() + + //////////////////////////////////////////////////////////////////////////// + // Serial peer + //////////////////////////////////////////////////////////////////////////// + + serialPeer, err := os.OpenFile( + portB, + os.O_RDWR, + 0, + ) + if err != nil { + t.Fatalf("open serial peer failed: %v", err) + } + defer serialPeer.Close() + + //////////////////////////////////////////////////////////////////////////// + // TCP -> Serial + //////////////////////////////////////////////////////////////////////////// + + tcpToSerial := []byte("hello from tcp") + + if _, err := dataConn.Write(tcpToSerial); err != nil { + t.Fatalf("TCP Write() failed: %v", err) + } + + serialReceived := make([]byte, len(tcpToSerial)) + + readExactWithTimeout( + t, + serialPeer, + serialReceived, + 2*time.Second, + ) + + if string(serialReceived) != string(tcpToSerial) { + t.Fatalf( + "serial received %q, want %q", + string(serialReceived), + string(tcpToSerial), + ) + } + + //////////////////////////////////////////////////////////////////////////// + // Serial -> TCP + //////////////////////////////////////////////////////////////////////////// + + serialToTCP := []byte("hello from serial") + + if _, err := serialPeer.Write(serialToTCP); err != nil { + t.Fatalf("serial Write() failed: %v", err) + } + + tcpReceived := make([]byte, len(serialToTCP)) + + readExactWithTimeout( + t, + dataConn, + tcpReceived, + 2*time.Second, + ) + + if string(tcpReceived) != string(serialToTCP) { + t.Fatalf( + "TCP received %q, want %q", + string(tcpReceived), + string(serialToTCP), + ) + } + + //////////////////////////////////////////////////////////////////////////// + // Close data connection + //////////////////////////////////////////////////////////////////////////// + + if err := dataConn.Close(); err != nil { + t.Fatalf("close data connection failed: %v", err) + } + + //////////////////////////////////////////////////////////////////////////// + // Close control connection + //////////////////////////////////////////////////////////////////////////// + + if err := clientConn.Close(); err != nil { + t.Fatalf("close control connection failed: %v", err) + } + + select { + case session := <-done: + if session == nil { + t.Fatal("server returned nil session") + } + + if !session.IsClosed() { + t.Fatal("session is not closed") + } + + case <-time.After(2 * time.Second): + t.Fatal("server connection handler did not terminate") + } +} diff --git a/internal/server/data_connection.go b/internal/server/data_connection.go new file mode 100644 index 0000000..199e1d6 --- /dev/null +++ b/internal/server/data_connection.go @@ -0,0 +1,169 @@ +/* + * ============================================================================ + * Projekt.....: rs2322tcp + * Datei.......: data_connection.go + * Copyright (C) 2026 Dieter Lang + * + * SPDX-License-Identifier: GPL-3.0-or-later + * + * Beschreibung: + * Bidirektionale Datenverbindung zwischen einer TCP-Verbindung und + * einem seriellen Gerätekanal. + * ============================================================================ + */ +package server + +import ( + "fmt" + "io" + "net" + "sync" +) + +/////////////////////////////////////////////////////////////////////////////// +// DataConnection +/////////////////////////////////////////////////////////////////////////////// + +// DataConnection connects one TCP connection with one serial device. +// +// Data is transferred bidirectionally: +// +// TCP -> Serial +// TCP <- Serial +// +// The serial side is represented by an io.ReadWriteCloser so that this +// server component does not depend directly on the concrete serial +// implementation. +type DataConnection struct { + tcp net.Conn + serial io.ReadWriteCloser + + closeOnce sync.Once + closeErr error +} + +/////////////////////////////////////////////////////////////////////////////// +// Constructor +/////////////////////////////////////////////////////////////////////////////// + +// NewDataConnection creates a new bidirectional data connection. +func NewDataConnection( + tcp net.Conn, + serial io.ReadWriteCloser, +) (*DataConnection, error) { + if tcp == nil { + return nil, fmt.Errorf("TCP connection is nil") + } + + if serial == nil { + return nil, fmt.Errorf("serial connection is nil") + } + + return &DataConnection{ + tcp: tcp, + serial: serial, + }, nil +} + +/////////////////////////////////////////////////////////////////////////////// +// Properties +/////////////////////////////////////////////////////////////////////////////// + +// TCPConn returns the TCP connection. +func (c *DataConnection) TCPConn() net.Conn { + if c == nil { + return nil + } + + return c.tcp +} + +// SerialConn returns the serial connection. +func (c *DataConnection) SerialConn() io.ReadWriteCloser { + if c == nil { + return nil + } + + return c.serial +} + +/////////////////////////////////////////////////////////////////////////////// +// Run +/////////////////////////////////////////////////////////////////////////////// + +// Run transfers data in both directions. +// +// Run blocks until one of the two transfer directions terminates. +// The other direction is then stopped and both connections are closed. +// +// The first non-EOF transfer error is returned. +func (c *DataConnection) Run() error { + if c == nil { + return fmt.Errorf("data connection is nil") + } + + var wg sync.WaitGroup + + errCh := make(chan error, 2) + + wg.Add(2) + + go func() { + defer wg.Done() + + _, err := io.Copy(c.serial, c.tcp) + errCh <- err + }() + + go func() { + defer wg.Done() + + _, err := io.Copy(c.tcp, c.serial) + errCh <- err + }() + + err := <-errCh + + _ = c.Close() + + wg.Wait() + + if err == nil || err == io.EOF { + return nil + } + + return err +} + +/////////////////////////////////////////////////////////////////////////////// +// Lifecycle +/////////////////////////////////////////////////////////////////////////////// + +// Close closes both sides of the data connection. +// +// Close is safe to call multiple times. +func (c *DataConnection) Close() error { + if c == nil { + return nil + } + + c.closeOnce.Do(func() { + var firstErr error + + if c.tcp != nil { + if err := c.tcp.Close(); err != nil { + firstErr = err + } + } + + if c.serial != nil { + if err := c.serial.Close(); err != nil && firstErr == nil { + firstErr = err + } + } + + c.closeErr = firstErr + }) + + return c.closeErr +} diff --git a/internal/server/data_connection_test.go b/internal/server/data_connection_test.go new file mode 100644 index 0000000..314755d --- /dev/null +++ b/internal/server/data_connection_test.go @@ -0,0 +1,314 @@ +/* + * ============================================================================ + * Projekt.....: rs2322tcp + * Datei.......: data_connection_test.go + * Copyright (C) 2026 Dieter Lang + * + * SPDX-License-Identifier: GPL-3.0-or-later + * + * Beschreibung: + * Tests für die bidirektionale TCP-/RS232-Datenverbindung. + * ============================================================================ + */ +package server_test + +import ( + "bytes" + "io" + "net" + "sync" + "testing" + "time" + + "git.lang-dieter.de/rs2322tcp/internal/server" +) + +/////////////////////////////////////////////////////////////////////////////// +// Test serial connection +/////////////////////////////////////////////////////////////////////////////// + +type testSerialConnection struct { + reader *bytes.Reader + writer bytes.Buffer + + mu sync.Mutex + closed bool + + closeCh chan struct{} + closeOnce sync.Once +} + +func newTestSerialConnection(data []byte) *testSerialConnection { + return &testSerialConnection{ + reader: bytes.NewReader(data), + closeCh: make(chan struct{}), + } +} + +func (s *testSerialConnection) Read(p []byte) (int, error) { + s.mu.Lock() + + if s.closed { + s.mu.Unlock() + return 0, io.ErrClosedPipe + } + + n, err := s.reader.Read(p) + + s.mu.Unlock() + + if err != io.EOF { + return n, err + } + + // A real serial connection normally waits for further data instead + // of immediately returning EOF. Block until the connection is closed. + <-s.closeCh + + return 0, io.ErrClosedPipe +} + +func (s *testSerialConnection) Write(p []byte) (int, error) { + s.mu.Lock() + defer s.mu.Unlock() + + if s.closed { + return 0, io.ErrClosedPipe + } + + return s.writer.Write(p) +} + +func (s *testSerialConnection) Close() error { + s.closeOnce.Do(func() { + s.mu.Lock() + s.closed = true + s.mu.Unlock() + + close(s.closeCh) + }) + + return nil +} + +func (s *testSerialConnection) WrittenData() []byte { + s.mu.Lock() + defer s.mu.Unlock() + + data := make([]byte, s.writer.Len()) + + copy(data, s.writer.Bytes()) + + return data +} + +/////////////////////////////////////////////////////////////////////////////// +// Constructor +/////////////////////////////////////////////////////////////////////////////// + +func TestNewDataConnection(t *testing.T) { + tcpServer, tcpClient := net.Pipe() + defer tcpServer.Close() + defer tcpClient.Close() + + serial := newTestSerialConnection(nil) + + connection, err := server.NewDataConnection( + tcpServer, + serial, + ) + if err != nil { + t.Fatalf("NewDataConnection() failed: %v", err) + } + + if connection.TCPConn() == nil { + t.Fatal("TCPConn() returned nil") + } + + if connection.SerialConn() == nil { + t.Fatal("SerialConn() returned nil") + } +} + +func TestNewDataConnectionRejectsNilTCP(t *testing.T) { + serial := newTestSerialConnection(nil) + + connection, err := server.NewDataConnection( + nil, + serial, + ) + + if err == nil { + t.Fatal("NewDataConnection() succeeded, want error") + } + + if connection != nil { + t.Fatal("NewDataConnection() returned connection despite error") + } +} + +func TestNewDataConnectionRejectsNilSerial(t *testing.T) { + tcpServer, tcpClient := net.Pipe() + defer tcpServer.Close() + defer tcpClient.Close() + + connection, err := server.NewDataConnection( + tcpServer, + nil, + ) + + if err == nil { + t.Fatal("NewDataConnection() succeeded, want error") + } + + if connection != nil { + t.Fatal("NewDataConnection() returned connection despite error") + } +} + +/////////////////////////////////////////////////////////////////////////////// +// TCP -> Serial +/////////////////////////////////////////////////////////////////////////////// + +func TestDataConnectionTCPToSerial(t *testing.T) { + tcpServer, tcpClient := net.Pipe() + defer tcpClient.Close() + + serial := newTestSerialConnection(nil) + + connection, err := server.NewDataConnection( + tcpServer, + serial, + ) + if err != nil { + t.Fatalf("NewDataConnection() failed: %v", err) + } + + done := make(chan error, 1) + + go func() { + done <- connection.Run() + }() + + testData := []byte("hello serial") + + if _, err := tcpClient.Write(testData); err != nil { + t.Fatalf("TCP Write() failed: %v", err) + } + + deadline := time.Now().Add(time.Second) + + for { + if bytes.Equal(serial.WrittenData(), testData) { + break + } + + if time.Now().After(deadline) { + t.Fatalf( + "serial data = %q, want %q", + serial.WrittenData(), + testData, + ) + } + + time.Sleep(time.Millisecond) + } + + _ = tcpClient.Close() + + select { + case <-done: + case <-time.After(time.Second): + t.Fatal("DataConnection.Run() did not terminate") + } +} + +/////////////////////////////////////////////////////////////////////////////// +// Serial -> TCP +/////////////////////////////////////////////////////////////////////////////// + +func TestDataConnectionSerialToTCP(t *testing.T) { + tcpServer, tcpClient := net.Pipe() + defer tcpClient.Close() + + testData := []byte("hello tcp") + + serial := newTestSerialConnection(testData) + + connection, err := server.NewDataConnection( + tcpServer, + serial, + ) + if err != nil { + t.Fatalf("NewDataConnection() failed: %v", err) + } + + done := make(chan error, 1) + + go func() { + done <- connection.Run() + }() + + buffer := make([]byte, len(testData)) + + if _, err := io.ReadFull(tcpClient, buffer); err != nil { + t.Fatalf("TCP Read() failed: %v", err) + } + + if !bytes.Equal(buffer, testData) { + t.Fatalf( + "TCP data = %q, want %q", + buffer, + testData, + ) + } + + _ = tcpClient.Close() + + select { + case <-done: + case <-time.After(time.Second): + t.Fatal("DataConnection.Run() did not terminate") + } +} + +/////////////////////////////////////////////////////////////////////////////// +// Close +/////////////////////////////////////////////////////////////////////////////// + +func TestDataConnectionClose(t *testing.T) { + tcpServer, tcpClient := net.Pipe() + defer tcpClient.Close() + + serial := newTestSerialConnection(nil) + + connection, err := server.NewDataConnection( + tcpServer, + serial, + ) + if err != nil { + t.Fatalf("NewDataConnection() failed: %v", err) + } + + if err := connection.Close(); err != nil { + t.Fatalf("Close() failed: %v", err) + } + + if err := connection.Close(); err != nil { + t.Fatalf("second Close() failed: %v", err) + } + + buffer := make([]byte, 1) + + if _, err := tcpClient.Read(buffer); err == nil { + t.Fatal("TCP connection is still open") + } + + serial.mu.Lock() + closed := serial.closed + serial.mu.Unlock() + + if !closed { + t.Fatal("serial connection is still open") + } +} diff --git a/internal/server/data_handler.go b/internal/server/data_handler.go new file mode 100644 index 0000000..274ad5c --- /dev/null +++ b/internal/server/data_handler.go @@ -0,0 +1,143 @@ +/* + * ============================================================================ + * Projekt.....: rs2322tcp + * Datei.......: data_handler.go + * Copyright (C) 2026 Dieter Lang + * + * SPDX-License-Identifier: GPL-3.0-or-later + * + * Beschreibung: + * Verbindet eingehende TCP-Datenverbindungen eines Data-Listeners mit + * einem vom Server bereitgestellten seriellen Gerätekanal. + * ============================================================================ + */ +package server + +import ( + "fmt" + "io" + "log" + "net" +) + +/////////////////////////////////////////////////////////////////////////////// +// SerialFactory +/////////////////////////////////////////////////////////////////////////////// + +// SerialFactory creates a serial connection for one device. +// +// The concrete implementation is provided by the serial package later. +// Keeping the factory as a function type prevents the server package from +// depending directly on the concrete serial implementation. +type SerialFactory func() (io.ReadWriteCloser, error) + +/////////////////////////////////////////////////////////////////////////////// +// DataHandler +/////////////////////////////////////////////////////////////////////////////// + +// DataHandler accepts TCP data connections and connects them to a serial +// device. +type DataHandler struct { + listener *DataListener + serialFactory SerialFactory +} + +/////////////////////////////////////////////////////////////////////////////// +// Constructor +/////////////////////////////////////////////////////////////////////////////// + +// NewDataHandler creates a new data handler. +// +// The listener provides the TCP endpoint. The serial factory is called for +// every accepted TCP connection. +func NewDataHandler( + listener *DataListener, + serialFactory SerialFactory, +) (*DataHandler, error) { + if listener == nil { + return nil, fmt.Errorf("data listener is nil") + } + + if serialFactory == nil { + return nil, fmt.Errorf("serial factory is nil") + } + + return &DataHandler{ + listener: listener, + serialFactory: serialFactory, + }, nil +} + +/////////////////////////////////////////////////////////////////////////////// +// Serve +/////////////////////////////////////////////////////////////////////////////// + +// Serve waits for incoming TCP data connections. +// +// Every accepted connection gets its own serial connection and +// DataConnection. The handler runs independently for each client. +// +// Serve terminates when the listener is closed. +func (h *DataHandler) Serve() error { + if h == nil { + return fmt.Errorf("data handler is nil") + } + + for { + tcpConn, err := h.listener.Accept() + if err != nil { + return err + } + + go h.handleConnection(tcpConn) + } +} + +/////////////////////////////////////////////////////////////////////////////// +// Connection handling +/////////////////////////////////////////////////////////////////////////////// + +// handleConnection connects one accepted TCP connection to one serial +// connection. +func (h *DataHandler) handleConnection(tcpConn net.Conn) { + if tcpConn == nil { + return + } + + serialConn, err := h.serialFactory() + if err != nil { + log.Printf("create serial connection: %v", err) + _ = tcpConn.Close() + return + } + + dataConnection, err := NewDataConnection( + tcpConn, + serialConn, + ) + if err != nil { + log.Printf("create data connection: %v", err) + + _ = serialConn.Close() + _ = tcpConn.Close() + + return + } + + if err := dataConnection.Run(); err != nil { + log.Printf("data connection terminated: %v", err) + } +} + +/////////////////////////////////////////////////////////////////////////////// +// Test support +/////////////////////////////////////////////////////////////////////////////// + +// HandleConnectionForTest handles one TCP connection using the configured +// serial factory. +// +// This method is intentionally provided for package-level tests without +// requiring a real TCP listener. +func (h *DataHandler) HandleConnectionForTest(tcpConn net.Conn) { + h.handleConnection(tcpConn) +} diff --git a/internal/server/data_handler_test.go b/internal/server/data_handler_test.go new file mode 100644 index 0000000..ca40d56 --- /dev/null +++ b/internal/server/data_handler_test.go @@ -0,0 +1,350 @@ +/* + * ============================================================================ + * Projekt.....: rs2322tcp + * Datei.......: data_handler_test.go + * Copyright (C) 2026 Dieter Lang + * + * SPDX-License-Identifier: GPL-3.0-or-later + * + * Beschreibung: + * Tests für die Verbindung von DataListener, TCP-Datenverbindung und + * serieller Geräteverbindung. + * ============================================================================ + */ +package server_test + +import ( + "bytes" + "io" + "net" + "sync" + "testing" + "time" + + "git.lang-dieter.de/rs2322tcp/internal/server" +) + +/////////////////////////////////////////////////////////////////////////////// +// Test serial connection +/////////////////////////////////////////////////////////////////////////////// + +type handlerTestSerial struct { + reader *bytes.Reader + + mu sync.Mutex + writer bytes.Buffer + closed bool + + closeCh chan struct{} + closeOnce sync.Once +} + +func newHandlerTestSerial(data []byte) *handlerTestSerial { + return &handlerTestSerial{ + reader: bytes.NewReader(data), + closeCh: make(chan struct{}), + } +} + +func (s *handlerTestSerial) Read(p []byte) (int, error) { + s.mu.Lock() + + if s.closed { + s.mu.Unlock() + return 0, io.ErrClosedPipe + } + + n, err := s.reader.Read(p) + + s.mu.Unlock() + + if err != io.EOF { + return n, err + } + + <-s.closeCh + + return 0, io.ErrClosedPipe +} + +func (s *handlerTestSerial) Write(p []byte) (int, error) { + s.mu.Lock() + defer s.mu.Unlock() + + if s.closed { + return 0, io.ErrClosedPipe + } + + return s.writer.Write(p) +} + +func (s *handlerTestSerial) Close() error { + s.closeOnce.Do(func() { + s.mu.Lock() + s.closed = true + s.mu.Unlock() + + close(s.closeCh) + }) + + return nil +} + +func (s *handlerTestSerial) WrittenData() []byte { + s.mu.Lock() + defer s.mu.Unlock() + + data := make([]byte, s.writer.Len()) + + copy(data, s.writer.Bytes()) + + return data +} + +/////////////////////////////////////////////////////////////////////////////// +// Constructor +/////////////////////////////////////////////////////////////////////////////// + +func TestNewDataHandler(t *testing.T) { + listener, err := server.NewDataListener() + if err != nil { + t.Fatalf("NewDataListener() failed: %v", err) + } + defer listener.Close() + + serial := newHandlerTestSerial(nil) + + handler, err := server.NewDataHandler( + listener, + func() (io.ReadWriteCloser, error) { + return serial, nil + }, + ) + if err != nil { + t.Fatalf("NewDataHandler() failed: %v", err) + } + + if handler == nil { + t.Fatal("NewDataHandler() returned nil") + } +} + +func TestNewDataHandlerRejectsNilListener(t *testing.T) { + handler, err := server.NewDataHandler( + nil, + func() (io.ReadWriteCloser, error) { + return newHandlerTestSerial(nil), nil + }, + ) + + if err == nil { + t.Fatal("NewDataHandler() succeeded, want error") + } + + if handler != nil { + t.Fatal("NewDataHandler() returned handler despite error") + } +} + +func TestNewDataHandlerRejectsNilFactory(t *testing.T) { + listener, err := server.NewDataListener() + if err != nil { + t.Fatalf("NewDataListener() failed: %v", err) + } + defer listener.Close() + + handler, err := server.NewDataHandler( + listener, + nil, + ) + + if err == nil { + t.Fatal("NewDataHandler() succeeded, want error") + } + + if handler != nil { + t.Fatal("NewDataHandler() returned handler despite error") + } +} + +/////////////////////////////////////////////////////////////////////////////// +// TCP -> Serial +/////////////////////////////////////////////////////////////////////////////// + +func TestDataHandlerTCPToSerial(t *testing.T) { + listener, err := server.NewDataListener() + if err != nil { + t.Fatalf("NewDataListener() failed: %v", err) + } + defer listener.Close() + + serial := newHandlerTestSerial(nil) + + handler, err := server.NewDataHandler( + listener, + func() (io.ReadWriteCloser, error) { + return serial, nil + }, + ) + if err != nil { + t.Fatalf("NewDataHandler() failed: %v", err) + } + + serveDone := make(chan error, 1) + + go func() { + serveDone <- handler.Serve() + }() + + tcpConn, err := net.Dial("tcp", listener.Addr().String()) + if err != nil { + t.Fatalf("net.Dial() failed: %v", err) + } + + testData := []byte("handler tcp to serial") + + if _, err := tcpConn.Write(testData); err != nil { + t.Fatalf("TCP Write() failed: %v", err) + } + + deadline := time.Now().Add(time.Second) + + for { + if bytes.Equal(serial.WrittenData(), testData) { + break + } + + if time.Now().After(deadline) { + t.Fatalf( + "serial data = %q, want %q", + serial.WrittenData(), + testData, + ) + } + + time.Sleep(time.Millisecond) + } + + _ = tcpConn.Close() + _ = listener.Close() + + select { + case <-serveDone: + case <-time.After(time.Second): + t.Fatal("DataHandler.Serve() did not terminate") + } +} + +/////////////////////////////////////////////////////////////////////////////// +// Serial -> TCP +/////////////////////////////////////////////////////////////////////////////// + +func TestDataHandlerSerialToTCP(t *testing.T) { + listener, err := server.NewDataListener() + if err != nil { + t.Fatalf("NewDataListener() failed: %v", err) + } + defer listener.Close() + + testData := []byte("handler serial to tcp") + + serial := newHandlerTestSerial(testData) + + handler, err := server.NewDataHandler( + listener, + func() (io.ReadWriteCloser, error) { + return serial, nil + }, + ) + if err != nil { + t.Fatalf("NewDataHandler() failed: %v", err) + } + + serveDone := make(chan error, 1) + + go func() { + serveDone <- handler.Serve() + }() + + tcpConn, err := net.Dial("tcp", listener.Addr().String()) + if err != nil { + t.Fatalf("net.Dial() failed: %v", err) + } + + buffer := make([]byte, len(testData)) + + if _, err := io.ReadFull(tcpConn, buffer); err != nil { + t.Fatalf("TCP Read() failed: %v", err) + } + + if !bytes.Equal(buffer, testData) { + t.Fatalf( + "TCP data = %q, want %q", + buffer, + testData, + ) + } + + _ = tcpConn.Close() + _ = listener.Close() + + select { + case <-serveDone: + case <-time.After(time.Second): + t.Fatal("DataHandler.Serve() did not terminate") + } +} + +/////////////////////////////////////////////////////////////////////////////// +// Serial factory error +/////////////////////////////////////////////////////////////////////////////// + +func TestDataHandlerSerialFactoryError(t *testing.T) { + listener, err := server.NewDataListener() + if err != nil { + t.Fatalf("NewDataListener() failed: %v", err) + } + defer listener.Close() + + handler, err := server.NewDataHandler( + listener, + func() (io.ReadWriteCloser, error) { + return nil, io.ErrClosedPipe + }, + ) + if err != nil { + t.Fatalf("NewDataHandler() failed: %v", err) + } + + serveDone := make(chan error, 1) + + go func() { + serveDone <- handler.Serve() + }() + + tcpConn, err := net.Dial("tcp", listener.Addr().String()) + if err != nil { + t.Fatalf("net.Dial() failed: %v", err) + } + + // The handler must close the TCP connection when the serial factory + // fails. + buffer := make([]byte, 1) + + _ = tcpConn.SetReadDeadline(time.Now().Add(time.Second)) + + _, readErr := tcpConn.Read(buffer) + + if readErr == nil { + t.Fatal("TCP connection remained open after serial factory error") + } + + _ = tcpConn.Close() + _ = listener.Close() + + select { + case <-serveDone: + case <-time.After(time.Second): + t.Fatal("DataHandler.Serve() did not terminate") + } +} diff --git a/internal/server/data_listener.go b/internal/server/data_listener.go new file mode 100644 index 0000000..189914b --- /dev/null +++ b/internal/server/data_listener.go @@ -0,0 +1,146 @@ +/* + * ============================================================================ + * Projekt.....: rs2322tcp + * Datei.......: data_listener.go + * Copyright (C) 2026 Dieter Lang + * + * SPDX-License-Identifier: GPL-3.0-or-later + * + * Beschreibung: + * Dynamischer TCP-Listener für die RS232-Datenverbindungen. + * Der Listener verwendet einen vom Betriebssystem zugewiesenen freien Port. + * Der Zugriff auf den zugrunde liegenden Netzwerk-Listener ist + * nebenläufigkeitssicher. + * ============================================================================ + */ +package server + +import ( + "fmt" + "net" + "sync" +) + +/////////////////////////////////////////////////////////////////////////////// +// DataListener +/////////////////////////////////////////////////////////////////////////////// + +// DataListener represents one TCP listener for an RS232 data connection. +// +// The listener uses a dynamically assigned TCP port. The actual port is +// available through Port(). +type DataListener struct { + mu sync.Mutex + listener net.Listener +} + +/////////////////////////////////////////////////////////////////////////////// +// Constructor +/////////////////////////////////////////////////////////////////////////////// + +// NewDataListener creates a TCP listener on a dynamically assigned port. +// +// The listener is bound to all local interfaces. +func NewDataListener() (*DataListener, error) { + listener, err := net.Listen("tcp", ":0") + if err != nil { + return nil, fmt.Errorf("listen on dynamic data port: %w", err) + } + + return &DataListener{ + listener: listener, + }, nil +} + +/////////////////////////////////////////////////////////////////////////////// +// Properties +/////////////////////////////////////////////////////////////////////////////// + +// Addr returns the actual network address of the listener. +func (l *DataListener) Addr() net.Addr { + if l == nil { + return nil + } + + l.mu.Lock() + defer l.mu.Unlock() + + if l.listener == nil { + return nil + } + + return l.listener.Addr() +} + +// Port returns the dynamically assigned TCP port. +// +// It returns 0 if the listener is nil or has already been closed. +func (l *DataListener) Port() int { + if l == nil { + return 0 + } + + l.mu.Lock() + defer l.mu.Unlock() + + if l.listener == nil { + return 0 + } + + tcpAddr, ok := l.listener.Addr().(*net.TCPAddr) + if !ok { + return 0 + } + + return tcpAddr.Port +} + +/////////////////////////////////////////////////////////////////////////////// +// Accept +/////////////////////////////////////////////////////////////////////////////// + +// Accept waits for an incoming data connection. +func (l *DataListener) Accept() (net.Conn, error) { + if l == nil { + return nil, fmt.Errorf("data listener is closed") + } + + l.mu.Lock() + + listener := l.listener + + l.mu.Unlock() + + if listener == nil { + return nil, fmt.Errorf("data listener is closed") + } + + return listener.Accept() +} + +/////////////////////////////////////////////////////////////////////////////// +// Lifecycle +/////////////////////////////////////////////////////////////////////////////// + +// Close closes the data listener. +// +// Close is safe to call more than once and may safely run concurrently +// with Accept(). +func (l *DataListener) Close() error { + if l == nil { + return nil + } + + l.mu.Lock() + + listener := l.listener + l.listener = nil + + l.mu.Unlock() + + if listener == nil { + return nil + } + + return listener.Close() +} diff --git a/internal/server/data_listener_test.go b/internal/server/data_listener_test.go new file mode 100644 index 0000000..005a60d --- /dev/null +++ b/internal/server/data_listener_test.go @@ -0,0 +1,150 @@ +/* + * ============================================================================ + * Projekt.....: rs2322tcp + * Datei.......: data_listener_test.go + * Copyright (C) 2026 Dieter Lang + * + * SPDX-License-Identifier: GPL-3.0-or-later + * + * Beschreibung: + * Tests für den dynamischen TCP-Data-Listener. + * ============================================================================ + */ +package server_test + +import ( + "fmt" + "net" + "testing" + "time" + + "git.lang-dieter.de/rs2322tcp/internal/server" +) + +func TestNewDataListener(t *testing.T) { + listener, err := server.NewDataListener() + if err != nil { + t.Fatalf("NewDataListener() failed: %v", err) + } + defer listener.Close() + + if listener.Addr() == nil { + t.Fatal("Addr() returned nil") + } + + if listener.Port() == 0 { + t.Fatal("Port() returned 0") + } +} + +func TestDataListenerAccept(t *testing.T) { + listener, err := server.NewDataListener() + if err != nil { + t.Fatalf("NewDataListener() failed: %v", err) + } + defer listener.Close() + + done := make(chan error, 1) + + go func() { + conn, err := listener.Accept() + if conn != nil { + _ = conn.Close() + } + + done <- err + }() + + conn, err := net.Dial("tcp", listener.Addr().String()) + if err != nil { + t.Fatalf("net.Dial() failed: %v", err) + } + defer conn.Close() + + select { + case err := <-done: + if err != nil { + t.Fatalf("Accept() failed: %v", err) + } + + case <-time.After(time.Second): + t.Fatal("Accept() did not return") + } +} + +func TestDataListenerClose(t *testing.T) { + listener, err := server.NewDataListener() + if err != nil { + t.Fatalf("NewDataListener() failed: %v", err) + } + + if listener.Port() == 0 { + t.Fatal("Port() returned 0") + } + + if err := listener.Close(); err != nil { + t.Fatalf("Close() failed: %v", err) + } + + if listener.Port() != 0 { + t.Errorf("Port() after Close() = %d, want 0", + listener.Port()) + } + + if err := listener.Close(); err != nil { + t.Fatalf("second Close() failed: %v", err) + } +} + +func TestDataListenerGetsDifferentPorts(t *testing.T) { + listener1, err := server.NewDataListener() + if err != nil { + t.Fatalf("NewDataListener() #1 failed: %v", err) + } + defer listener1.Close() + + listener2, err := server.NewDataListener() + if err != nil { + t.Fatalf("NewDataListener() #2 failed: %v", err) + } + defer listener2.Close() + + port1 := listener1.Port() + port2 := listener2.Port() + + if port1 == 0 { + t.Fatal("listener1 Port() returned 0") + } + + if port2 == 0 { + t.Fatal("listener2 Port() returned 0") + } + + if port1 == port2 { + t.Fatalf("both listeners use port %d", port1) + } +} + +func TestDataListenerCloseReleasesPort(t *testing.T) { + listener, err := server.NewDataListener() + if err != nil { + t.Fatalf("NewDataListener() failed: %v", err) + } + + port := listener.Port() + + if err := listener.Close(); err != nil { + t.Fatalf("Close() failed: %v", err) + } + + address := fmt.Sprintf("127.0.0.1:%d", port) + + testListener, err := net.Listen("tcp", address) + if err != nil { + t.Fatalf("port %d was not released: %v", port, err) + } + + if err := testListener.Close(); err != nil { + t.Fatalf("closing test listener failed: %v", err) + } +} diff --git a/internal/server/session.go b/internal/server/session.go new file mode 100644 index 0000000..fdb0ae9 --- /dev/null +++ b/internal/server/session.go @@ -0,0 +1,462 @@ +/* + * ============================================================================ + * Projekt.....: rs2322tcp + * Datei.......: session.go + * Copyright (C) 2026 Dieter Lang + * + * SPDX-License-Identifier: GPL-3.0-or-later + * + * Beschreibung: + * Verwaltung des Lebenszyklus einer Client-Session einschließlich + * der zugehörigen Netzwerkressourcen, Data-Listener und aktiven + * Datenverbindungen. + * ============================================================================ + */ +package server + +import ( + "fmt" + "io" + "net" + "sync" +) + +/////////////////////////////////////////////////////////////////////////////// +// Session +/////////////////////////////////////////////////////////////////////////////// + +// Session represents one active client connection. +// +// All resources belonging to one client connection are associated with +// the session. This includes the control connection, dynamic data +// listeners and active data connections belonging to the configured +// remote devices. +type Session struct { + id uint64 + conn net.Conn + + mu sync.Mutex + closed bool + resources []io.Closer + dataListeners map[string]*DataListener + dataConnections map[string]*DataConnection +} + +/////////////////////////////////////////////////////////////////////////////// +// Constructor +/////////////////////////////////////////////////////////////////////////////// + +// NewSession creates a new client session. +func NewSession(id uint64, conn net.Conn) (*Session, error) { + if conn == nil { + return nil, fmt.Errorf("connection is nil") + } + + if id == 0 { + return nil, fmt.Errorf("session ID must not be zero") + } + + return &Session{ + id: id, + conn: conn, + resources: make([]io.Closer, 0), + dataListeners: make(map[string]*DataListener), + dataConnections: make(map[string]*DataConnection), + }, nil +} + +/////////////////////////////////////////////////////////////////////////////// +// Properties +/////////////////////////////////////////////////////////////////////////////// + +// ID returns the unique session ID. +func (s *Session) ID() uint64 { + if s == nil { + return 0 + } + + return s.id +} + +// Conn returns the control connection belonging to the session. +func (s *Session) Conn() net.Conn { + if s == nil { + return nil + } + + s.mu.Lock() + defer s.mu.Unlock() + + if s.closed { + return nil + } + + return s.conn +} + +// IsClosed reports whether the session has already been closed. +func (s *Session) IsClosed() bool { + if s == nil { + return true + } + + s.mu.Lock() + defer s.mu.Unlock() + + return s.closed +} + +/////////////////////////////////////////////////////////////////////////////// +// Resources +/////////////////////////////////////////////////////////////////////////////// + +// AddResource adds a resource to the session. +// +// The resource must implement io.Closer. It will automatically be closed +// when the session is closed. +// +// If the session is already closed, the resource is closed immediately +// and an error is returned. +func (s *Session) AddResource(resource io.Closer) error { + if s == nil { + if resource != nil { + _ = resource.Close() + } + + return fmt.Errorf("session is nil") + } + + if resource == nil { + return fmt.Errorf("resource is nil") + } + + s.mu.Lock() + + if s.closed { + s.mu.Unlock() + + _ = resource.Close() + + return fmt.Errorf("session is already closed") + } + + s.resources = append(s.resources, resource) + + s.mu.Unlock() + + return nil +} + +// ResourceCount returns the number of resources currently owned by the +// session. +// +// This method is primarily useful for diagnostics and tests. +func (s *Session) ResourceCount() int { + if s == nil { + return 0 + } + + s.mu.Lock() + defer s.mu.Unlock() + + return len(s.resources) +} + +/////////////////////////////////////////////////////////////////////////////// +// Data listeners +/////////////////////////////////////////////////////////////////////////////// + +// AddDataListener associates a dynamic data listener with a remote device. +// +// The listener becomes a session resource and is therefore automatically +// closed when the session is closed. +func (s *Session) AddDataListener( + deviceID string, + listener *DataListener, +) error { + if s == nil { + if listener != nil { + _ = listener.Close() + } + + return fmt.Errorf("session is nil") + } + + if deviceID == "" { + if listener != nil { + _ = listener.Close() + } + + return fmt.Errorf("device ID is empty") + } + + if listener == nil { + return fmt.Errorf("data listener is nil") + } + + s.mu.Lock() + + if s.closed { + s.mu.Unlock() + + _ = listener.Close() + + return fmt.Errorf("session is already closed") + } + + if _, exists := s.dataListeners[deviceID]; exists { + s.mu.Unlock() + + _ = listener.Close() + + return fmt.Errorf( + "data listener already exists for device %q", + deviceID, + ) + } + + s.resources = append(s.resources, listener) + s.dataListeners[deviceID] = listener + + s.mu.Unlock() + + return nil +} + +// DataListener returns the data listener associated with a device. +// +// The second return value reports whether a listener exists for the +// specified device. +func (s *Session) DataListener( + deviceID string, +) (*DataListener, bool) { + if s == nil { + return nil, false + } + + s.mu.Lock() + defer s.mu.Unlock() + + if s.closed { + return nil, false + } + + listener, ok := s.dataListeners[deviceID] + + return listener, ok +} + +// DataPort returns the dynamic TCP data port associated with a device. +// +// It returns 0 if the device has no listener or the session is closed. +func (s *Session) DataPort(deviceID string) int { + listener, ok := s.DataListener(deviceID) + if !ok { + return 0 + } + + return listener.Port() +} + +// DataListenerCount returns the number of data listeners currently +// associated with the session. +func (s *Session) DataListenerCount() int { + if s == nil { + return 0 + } + + s.mu.Lock() + defer s.mu.Unlock() + + return len(s.dataListeners) +} + +/////////////////////////////////////////////////////////////////////////////// +// Active data connections +/////////////////////////////////////////////////////////////////////////////// + +// AddDataConnection associates an active data connection with a remote +// device. +// +// Only one active data connection is allowed per device. The data +// connection becomes a session resource and is therefore automatically +// closed when the session is closed. +func (s *Session) AddDataConnection( + deviceID string, + connection *DataConnection, +) error { + if s == nil { + if connection != nil { + _ = connection.Close() + } + + return fmt.Errorf("session is nil") + } + + if deviceID == "" { + if connection != nil { + _ = connection.Close() + } + + return fmt.Errorf("device ID is empty") + } + + if connection == nil { + return fmt.Errorf("data connection is nil") + } + + s.mu.Lock() + + if s.closed { + s.mu.Unlock() + + _ = connection.Close() + + return fmt.Errorf("session is already closed") + } + + if _, exists := s.dataConnections[deviceID]; exists { + s.mu.Unlock() + + _ = connection.Close() + + return fmt.Errorf( + "data connection already exists for device %q", + deviceID, + ) + } + + s.resources = append(s.resources, connection) + s.dataConnections[deviceID] = connection + + s.mu.Unlock() + + return nil +} + +// DataConnection returns the active data connection associated with a +// device. +// +// The second return value reports whether an active connection exists. +func (s *Session) DataConnection( + deviceID string, +) (*DataConnection, bool) { + if s == nil { + return nil, false + } + + s.mu.Lock() + defer s.mu.Unlock() + + if s.closed { + return nil, false + } + + connection, ok := s.dataConnections[deviceID] + + return connection, ok +} + +// RemoveDataConnection removes an active data connection from the +// session. +// +// The connection is closed before it is removed from the session. +// Removing a connection that is not registered is harmless. +func (s *Session) RemoveDataConnection(deviceID string) error { + if s == nil { + return nil + } + + s.mu.Lock() + + connection, exists := s.dataConnections[deviceID] + if !exists { + s.mu.Unlock() + return nil + } + + delete(s.dataConnections, deviceID) + + for i, resource := range s.resources { + if resource == connection { + s.resources = append( + s.resources[:i], + s.resources[i+1:]..., + ) + break + } + } + + s.mu.Unlock() + + return connection.Close() +} + +// DataConnectionCount returns the number of currently active data +// connections. +func (s *Session) DataConnectionCount() int { + if s == nil { + return 0 + } + + s.mu.Lock() + defer s.mu.Unlock() + + return len(s.dataConnections) +} + +/////////////////////////////////////////////////////////////////////////////// +// Lifecycle +/////////////////////////////////////////////////////////////////////////////// + +// Close terminates the session and releases all resources belonging +// to the session. +// +// Close is safe to call multiple times. +// +// Data listeners and other session resources are closed before the +// control connection is closed. +func (s *Session) Close() error { + if s == nil { + return nil + } + + s.mu.Lock() + + if s.closed { + s.mu.Unlock() + return nil + } + + s.closed = true + + conn := s.conn + s.conn = nil + + resources := s.resources + s.resources = nil + s.dataListeners = nil + s.dataConnections = nil + + s.mu.Unlock() + + var firstErr error + + for _, resource := range resources { + if resource == nil { + continue + } + + if err := resource.Close(); err != nil && firstErr == nil { + firstErr = err + } + } + + if conn != nil { + if err := conn.Close(); err != nil && firstErr == nil { + firstErr = err + } + } + + return firstErr +} diff --git a/internal/server/session_test.go b/internal/server/session_test.go new file mode 100644 index 0000000..cab5b09 --- /dev/null +++ b/internal/server/session_test.go @@ -0,0 +1,440 @@ +/* + * ============================================================================ + * Projekt.....: rs2322tcp + * Datei.......: session_test.go + * Copyright (C) 2026 Dieter Lang + * + * SPDX-License-Identifier: GPL-3.0-or-later + * + * Beschreibung: + * Tests für den Lebenszyklus einer Client-Session, die Verwaltung + * der zugehörigen Ressourcen und der dynamischen Data-Listener. + * ============================================================================ + */ +package server_test + +import ( + "errors" + "net" + "testing" + + "git.lang-dieter.de/rs2322tcp/internal/server" +) + +/////////////////////////////////////////////////////////////////////////////// +// Test resource +/////////////////////////////////////////////////////////////////////////////// + +type testResource struct { + closed bool + err error +} + +func (r *testResource) Close() error { + r.closed = true + return r.err +} + +/////////////////////////////////////////////////////////////////////////////// +// Constructor +/////////////////////////////////////////////////////////////////////////////// + +func TestNewSession(t *testing.T) { + serverConn, clientConn := net.Pipe() + defer serverConn.Close() + defer clientConn.Close() + + session, err := server.NewSession(1, serverConn) + if err != nil { + t.Fatalf("NewSession() failed: %v", err) + } + + if session.ID() != 1 { + t.Errorf("ID() = %d, want 1", session.ID()) + } + + if session.Conn() == nil { + t.Fatal("Conn() returned nil") + } + + if session.IsClosed() { + t.Fatal("new session is already closed") + } + + if session.ResourceCount() != 0 { + t.Errorf("ResourceCount() = %d, want 0", + session.ResourceCount()) + } + + if session.DataListenerCount() != 0 { + t.Errorf("DataListenerCount() = %d, want 0", + session.DataListenerCount()) + } +} + +func TestNewSessionRejectsNilConnection(t *testing.T) { + session, err := server.NewSession(1, nil) + + if err == nil { + t.Fatal("NewSession() succeeded, want error") + } + + if session != nil { + t.Fatal("NewSession() returned a session despite error") + } +} + +func TestNewSessionRejectsZeroID(t *testing.T) { + serverConn, clientConn := net.Pipe() + defer serverConn.Close() + defer clientConn.Close() + + session, err := server.NewSession(0, serverConn) + + if err == nil { + t.Fatal("NewSession() succeeded, want error") + } + + if session != nil { + t.Fatal("NewSession() returned a session despite error") + } +} + +/////////////////////////////////////////////////////////////////////////////// +// Resources +/////////////////////////////////////////////////////////////////////////////// + +func TestSessionAddResource(t *testing.T) { + serverConn, clientConn := net.Pipe() + defer serverConn.Close() + defer clientConn.Close() + + session, err := server.NewSession(1, serverConn) + if err != nil { + t.Fatalf("NewSession() failed: %v", err) + } + + resource1 := &testResource{} + resource2 := &testResource{} + + if err := session.AddResource(resource1); err != nil { + t.Fatalf("AddResource(resource1) failed: %v", err) + } + + if err := session.AddResource(resource2); err != nil { + t.Fatalf("AddResource(resource2) failed: %v", err) + } + + if session.ResourceCount() != 2 { + t.Errorf("ResourceCount() = %d, want 2", + session.ResourceCount()) + } + + if resource1.closed { + t.Fatal("resource1 was closed before session.Close()") + } + + if resource2.closed { + t.Fatal("resource2 was closed before session.Close()") + } +} + +func TestAddNilResource(t *testing.T) { + serverConn, clientConn := net.Pipe() + defer serverConn.Close() + defer clientConn.Close() + + session, err := server.NewSession(1, serverConn) + if err != nil { + t.Fatalf("NewSession() failed: %v", err) + } + + if err := session.AddResource(nil); err == nil { + t.Fatal("AddResource(nil) succeeded, want error") + } +} + +func TestAddResourceAfterClose(t *testing.T) { + serverConn, clientConn := net.Pipe() + defer clientConn.Close() + + session, err := server.NewSession(1, serverConn) + if err != nil { + t.Fatalf("NewSession() failed: %v", err) + } + + if err := session.Close(); err != nil { + t.Fatalf("Close() failed: %v", err) + } + + resource := &testResource{} + + if err := session.AddResource(resource); err == nil { + t.Fatal("AddResource() succeeded, want error") + } + + if !resource.closed { + t.Fatal("resource was not closed after rejected AddResource()") + } + + if session.ResourceCount() != 0 { + t.Errorf("ResourceCount() = %d, want 0", + session.ResourceCount()) + } +} + +/////////////////////////////////////////////////////////////////////////////// +// Data listeners +/////////////////////////////////////////////////////////////////////////////// + +func TestSessionAddDataListener(t *testing.T) { + serverConn, clientConn := net.Pipe() + defer clientConn.Close() + + session, err := server.NewSession(1, serverConn) + if err != nil { + t.Fatalf("NewSession() failed: %v", err) + } + + listener, err := server.NewDataListener() + if err != nil { + t.Fatalf("NewDataListener() failed: %v", err) + } + + port := listener.Port() + + if port == 0 { + t.Fatal("DataListener.Port() returned 0") + } + + if err := session.AddDataListener("radio", listener); err != nil { + t.Fatalf("AddDataListener() failed: %v", err) + } + + if session.DataListenerCount() != 1 { + t.Errorf("DataListenerCount() = %d, want 1", + session.DataListenerCount()) + } + + if session.ResourceCount() != 1 { + t.Errorf("ResourceCount() = %d, want 1", + session.ResourceCount()) + } + + storedListener, ok := session.DataListener("radio") + if !ok { + t.Fatal("DataListener(radio) not found") + } + + if storedListener != listener { + t.Fatal("stored listener is not the same listener") + } + + if session.DataPort("radio") != port { + t.Errorf("DataPort(radio) = %d, want %d", + session.DataPort("radio"), port) + } + + if session.DataPort("unknown") != 0 { + t.Errorf("DataPort(unknown) = %d, want 0", + session.DataPort("unknown")) + } +} + +func TestSessionRejectsDuplicateDataListener(t *testing.T) { + serverConn, clientConn := net.Pipe() + defer clientConn.Close() + + session, err := server.NewSession(1, serverConn) + if err != nil { + t.Fatalf("NewSession() failed: %v", err) + } + + listener1, err := server.NewDataListener() + if err != nil { + t.Fatalf("NewDataListener() #1 failed: %v", err) + } + + listener2, err := server.NewDataListener() + if err != nil { + t.Fatalf("NewDataListener() #2 failed: %v", err) + } + + if err := session.AddDataListener("radio", listener1); err != nil { + t.Fatalf("AddDataListener() #1 failed: %v", err) + } + + if err := session.AddDataListener("radio", listener2); err == nil { + t.Fatal("AddDataListener() #2 succeeded, want error") + } + + if listener2.Port() != 0 { + t.Errorf("duplicate listener port = %d, want 0", + listener2.Port()) + } + + if session.DataListenerCount() != 1 { + t.Errorf("DataListenerCount() = %d, want 1", + session.DataListenerCount()) + } +} + +func TestSessionRejectsEmptyDeviceID(t *testing.T) { + serverConn, clientConn := net.Pipe() + defer clientConn.Close() + + session, err := server.NewSession(1, serverConn) + if err != nil { + t.Fatalf("NewSession() failed: %v", err) + } + + listener, err := server.NewDataListener() + if err != nil { + t.Fatalf("NewDataListener() failed: %v", err) + } + + if err := session.AddDataListener("", listener); err == nil { + t.Fatal("AddDataListener() succeeded, want error") + } + + if listener.Port() != 0 { + t.Errorf("listener port = %d, want 0", + listener.Port()) + } +} + +func TestSessionClosesDataListener(t *testing.T) { + serverConn, clientConn := net.Pipe() + defer clientConn.Close() + + session, err := server.NewSession(1, serverConn) + if err != nil { + t.Fatalf("NewSession() failed: %v", err) + } + + dataListener, err := server.NewDataListener() + if err != nil { + t.Fatalf("NewDataListener() failed: %v", err) + } + + if err := session.AddDataListener("radio", dataListener); err != nil { + t.Fatalf("AddDataListener() failed: %v", err) + } + + if session.DataListenerCount() != 1 { + t.Fatalf("DataListenerCount() = %d, want 1", + session.DataListenerCount()) + } + + if err := session.Close(); err != nil { + t.Fatalf("Session.Close() failed: %v", err) + } + + if dataListener.Port() != 0 { + t.Errorf("DataListener.Port() after Session.Close() = %d, want 0", + dataListener.Port()) + } + + if session.DataListenerCount() != 0 { + t.Errorf("DataListenerCount() after Session.Close() = %d, want 0", + session.DataListenerCount()) + } + + if session.DataPort("radio") != 0 { + t.Errorf("DataPort(radio) after Close() = %d, want 0", + session.DataPort("radio")) + } +} + +/////////////////////////////////////////////////////////////////////////////// +// Close +/////////////////////////////////////////////////////////////////////////////// + +func TestSessionClose(t *testing.T) { + serverConn, clientConn := net.Pipe() + defer clientConn.Close() + + session, err := server.NewSession(1, serverConn) + if err != nil { + t.Fatalf("NewSession() failed: %v", err) + } + + resource1 := &testResource{} + resource2 := &testResource{} + + if err := session.AddResource(resource1); err != nil { + t.Fatalf("AddResource(resource1) failed: %v", err) + } + + if err := session.AddResource(resource2); err != nil { + t.Fatalf("AddResource(resource2) failed: %v", err) + } + + if err := session.Close(); err != nil { + t.Fatalf("Close() failed: %v", err) + } + + if !session.IsClosed() { + t.Fatal("session is not closed") + } + + if session.Conn() != nil { + t.Fatal("Conn() returned a connection after Close()") + } + + if session.ResourceCount() != 0 { + t.Errorf("ResourceCount() after Close() = %d, want 0", + session.ResourceCount()) + } + + if !resource1.closed { + t.Fatal("resource1 was not closed") + } + + if !resource2.closed { + t.Fatal("resource2 was not closed") + } + + if err := session.Close(); err != nil { + t.Fatalf("second Close() failed: %v", err) + } +} + +func TestSessionCloseResourceError(t *testing.T) { + serverConn, clientConn := net.Pipe() + defer clientConn.Close() + + session, err := server.NewSession(1, serverConn) + if err != nil { + t.Fatalf("NewSession() failed: %v", err) + } + + expectedErr := errors.New("test resource error") + + resource := &testResource{ + err: expectedErr, + } + + if err := session.AddResource(resource); err != nil { + t.Fatalf("AddResource() failed: %v", err) + } + + err = session.Close() + + if !errors.Is(err, expectedErr) { + t.Fatalf("Close() error = %v, want %v", + err, expectedErr) + } + + if !resource.closed { + t.Fatal("resource was not closed") + } +} + +func TestSessionCloseNil(t *testing.T) { + var session *server.Session + + if err := session.Close(); err != nil { + t.Fatalf("Close() on nil session failed: %v", err) + } +} diff --git a/internal/transport/codec.go b/internal/transport/codec.go new file mode 100644 index 0000000..2550ccc --- /dev/null +++ b/internal/transport/codec.go @@ -0,0 +1,101 @@ +/* +Package transport contains the network transport definitions for rs2322tcp. + +This file implements the framing of control protocol messages. Control +messages are encoded as one JSON object per line (JSON Lines). + +Project: rs2322tcp +Module: git.lang-dieter.de/rs2322tcp +*/ +package transport + +import ( + "bufio" + "encoding/json" + "fmt" + "io" +) + +/////////////////////////////////////////////////////////////////////////////// +// Constants +/////////////////////////////////////////////////////////////////////////////// + +const ( + // MaximumControlMessageSize limits the size of one control message. + MaximumControlMessageSize = 64 * 1024 +) + +/////////////////////////////////////////////////////////////////////////////// +// Writer +/////////////////////////////////////////////////////////////////////////////// + +// WriteMessage writes one control message as a JSON object followed by +// a newline. +func WriteMessage(w io.Writer, message interface{}) error { + if w == nil { + return fmt.Errorf("writer is nil") + } + + if message == nil { + return fmt.Errorf("message is nil") + } + + data, err := json.Marshal(message) + if err != nil { + return fmt.Errorf("encode control message: %w", err) + } + + if len(data) > MaximumControlMessageSize { + return fmt.Errorf("control message too large: %d bytes", + len(data)) + } + + data = append(data, '\n') + + if _, err := w.Write(data); err != nil { + return fmt.Errorf("write control message: %w", err) + } + + return nil +} + +/////////////////////////////////////////////////////////////////////////////// +// Reader +/////////////////////////////////////////////////////////////////////////////// + +// ReadMessage reads exactly one JSON Lines control message. +// +// The target must be a pointer to the expected message structure. +func ReadMessage(r *bufio.Reader, target interface{}) error { + if r == nil { + return fmt.Errorf("reader is nil") + } + + if target == nil { + return fmt.Errorf("target is nil") + } + + data, err := r.ReadBytes('\n') + if err != nil { + if err == io.EOF && len(data) == 0 { + return io.EOF + } + + if err == io.EOF { + return fmt.Errorf("incomplete control message: %w", err) + } + + return fmt.Errorf("read control message: %w", err) + } + + if len(data) > MaximumControlMessageSize { + return fmt.Errorf("control message too large: %d bytes", + len(data)) + } + + if err := json.Unmarshal(data, target); err != nil { + return fmt.Errorf("decode control message: %w", err) + } + + return nil +} diff --git a/internal/transport/codec_test.go b/internal/transport/codec_test.go new file mode 100644 index 0000000..3c804fc --- /dev/null +++ b/internal/transport/codec_test.go @@ -0,0 +1,164 @@ +package transport_test + +import ( + "bufio" + "bytes" + "io" + "strings" + "testing" + + "git.lang-dieter.de/rs2322tcp/internal/transport" +) + +func TestWriteAndReadMessage(t *testing.T) { + original := transport.NewHello() + + var buffer bytes.Buffer + + if err := transport.WriteMessage(&buffer, original); err != nil { + t.Fatalf("WriteMessage() failed: %v", err) + } + + reader := bufio.NewReader(&buffer) + + var decoded transport.HelloMessage + + if err := transport.ReadMessage(reader, &decoded); err != nil { + t.Fatalf("ReadMessage() failed: %v", err) + } + + if decoded.Version != original.Version { + t.Errorf("Version = %d, want %d", + decoded.Version, original.Version) + } + + if decoded.Type != original.Type { + t.Errorf("Type = %q, want %q", + decoded.Type, original.Type) + } +} + +func TestMultipleMessages(t *testing.T) { + var buffer bytes.Buffer + + if err := transport.WriteMessage(&buffer, transport.NewHello()); err != nil { + t.Fatalf("WriteMessage(hello) failed: %v", err) + } + + if err := transport.WriteMessage(&buffer, transport.NewGetDevices()); err != nil { + t.Fatalf("WriteMessage(get_devices) failed: %v", err) + } + + reader := bufio.NewReader(&buffer) + + var hello transport.HelloMessage + + if err := transport.ReadMessage(reader, &hello); err != nil { + t.Fatalf("ReadMessage(hello) failed: %v", err) + } + + if hello.Type != transport.MessageHello { + t.Errorf("first Type = %q, want %q", + hello.Type, transport.MessageHello) + } + + var getDevices transport.GetDevicesMessage + + if err := transport.ReadMessage(reader, &getDevices); err != nil { + t.Fatalf("ReadMessage(get_devices) failed: %v", err) + } + + if getDevices.Type != transport.MessageGetDevices { + t.Errorf("second Type = %q, want %q", + getDevices.Type, transport.MessageGetDevices) + } +} + +func TestMessageFraming(t *testing.T) { + var buffer bytes.Buffer + + if err := transport.WriteMessage(&buffer, transport.NewHello()); err != nil { + t.Fatalf("WriteMessage() failed: %v", err) + } + + text := buffer.String() + + if !strings.HasSuffix(text, "\n") { + t.Fatalf("message does not end with newline: %q", text) + } + + if strings.Count(text, "\n") != 1 { + t.Fatalf("newline count = %d, want 1", + strings.Count(text, "\n")) + } +} + +func TestReadMessageEOF(t *testing.T) { + reader := bufio.NewReader(strings.NewReader("")) + + var message transport.HelloMessage + + err := transport.ReadMessage(reader, &message) + + if err != io.EOF { + t.Fatalf("ReadMessage() error = %v, want io.EOF", err) + } +} + +func TestReadMessageInvalidJSON(t *testing.T) { + reader := bufio.NewReader( + strings.NewReader(`{"version":1,"type":invalid}` + "\n"), + ) + + var message transport.HelloMessage + + if err := transport.ReadMessage(reader, &message); err == nil { + t.Fatal("ReadMessage() succeeded, want JSON error") + } +} + +func TestReadMessageIncomplete(t *testing.T) { + reader := bufio.NewReader( + strings.NewReader(`{"version":1,"type":"hello"}`), + ) + + var message transport.HelloMessage + + err := transport.ReadMessage(reader, &message) + + if err == nil { + t.Fatal("ReadMessage() succeeded, want incomplete message error") + } +} + +func TestNilWriter(t *testing.T) { + if err := transport.WriteMessage(nil, transport.NewHello()); err == nil { + t.Fatal("WriteMessage() succeeded, want error") + } +} + +func TestNilMessage(t *testing.T) { + var buffer bytes.Buffer + + if err := transport.WriteMessage(&buffer, nil); err == nil { + t.Fatal("WriteMessage() succeeded, want error") + } +} + +func TestNilReader(t *testing.T) { + var message transport.HelloMessage + + if err := transport.ReadMessage(nil, &message); err == nil { + t.Fatal("ReadMessage() succeeded, want error") + } +} + +func TestNilTarget(t *testing.T) { + reader := bufio.NewReader(strings.NewReader( + `{"version":1,"type":"hello"}` + "\n", + )) + + if err := transport.ReadMessage(reader, nil); err == nil { + t.Fatal("ReadMessage() succeeded, want error") + } +} diff --git a/internal/transport/control.go b/internal/transport/control.go new file mode 100644 index 0000000..e8a74b9 --- /dev/null +++ b/internal/transport/control.go @@ -0,0 +1,202 @@ +/* +Package transport contains the network transport definitions for rs2322tcp. + +This file defines the control protocol messages. It deliberately does not +contain any network I/O. The actual TCP implementation will be added later. + +Project: rs2322tcp +Module: git.lang-dieter.de/rs2322tcp +*/ +package transport + +import ( + "encoding/json" + "fmt" + + "git.lang-dieter.de/rs2322tcp/internal/config" +) + +/////////////////////////////////////////////////////////////////////////////// +// Protocol +/////////////////////////////////////////////////////////////////////////////// + +// ProtocolVersion is the current control protocol version. +const ProtocolVersion = 1 + +// MessageType identifies a control protocol message. +type MessageType string + +const ( + MessageHello MessageType = "hello" + MessageHelloResponse MessageType = "hello_response" + MessageGetDevices MessageType = "get_devices" + MessageDeviceList MessageType = "device_list" + MessageError MessageType = "error" +) + +/////////////////////////////////////////////////////////////////////////////// +// Generic message +/////////////////////////////////////////////////////////////////////////////// + +// Message is the common envelope for all control protocol messages. +type Message struct { + Version int `json:"version"` + Type MessageType `json:"type"` +} + +/////////////////////////////////////////////////////////////////////////////// +// Hello +/////////////////////////////////////////////////////////////////////////////// + +// HelloMessage starts a control session. +type HelloMessage struct { + Message +} + +// HelloResponseMessage confirms that the server accepts the protocol +// version used by the client. +type HelloResponseMessage struct { + Message +} + +/////////////////////////////////////////////////////////////////////////////// +// Get devices +/////////////////////////////////////////////////////////////////////////////// + +// GetDevicesMessage requests the devices currently available from the +// server. +type GetDevicesMessage struct { + Message +} + +/////////////////////////////////////////////////////////////////////////////// +// Device list +/////////////////////////////////////////////////////////////////////////////// + +// RemoteDeviceInfo describes a device returned by the server. +// +// DataPort is runtime information belonging to the current server session. +// It is therefore deliberately not part of config.RemoteDevice. +type RemoteDeviceInfo struct { + ID string `json:"id"` + Name string `json:"name"` + BaudRate int `json:"baud_rate"` + DataBits int `json:"data_bits"` + Parity string `json:"parity"` + StopBits int `json:"stop_bits"` + DataPort int `json:"data_port"` +} + +// DeviceListMessage contains the devices currently available from the +// server and their session-specific data ports. +type DeviceListMessage struct { + Message + Devices []RemoteDeviceInfo `json:"devices"` +} + +/////////////////////////////////////////////////////////////////////////////// +// Error +/////////////////////////////////////////////////////////////////////////////// + +// ErrorCode identifies a protocol error. +type ErrorCode string + +const ( + ErrorUnknownMessage ErrorCode = "unknown_message" + ErrorUnsupportedVersion ErrorCode = "unsupported_version" + ErrorDeviceNotFound ErrorCode = "device_not_found" + ErrorDeviceUnavailable ErrorCode = "device_unavailable" + ErrorInternal ErrorCode = "internal_error" +) + +// ErrorMessage reports a control protocol error. +type ErrorMessage struct { + Message + Code ErrorCode `json:"code"` + MessageText string `json:"message"` +} + +/////////////////////////////////////////////////////////////////////////////// +// Constructors +/////////////////////////////////////////////////////////////////////////////// + +// NewHello creates a hello message. +func NewHello() HelloMessage { + return HelloMessage{ + Message: Message{ + Version: ProtocolVersion, + Type: MessageHello, + }, + } +} + +// NewHelloResponse creates a hello response message. +func NewHelloResponse() HelloResponseMessage { + return HelloResponseMessage{ + Message: Message{ + Version: ProtocolVersion, + Type: MessageHelloResponse, + }, + } +} + +// NewGetDevices creates a get-devices message. +func NewGetDevices() GetDevicesMessage { + return GetDevicesMessage{ + Message: Message{ + Version: ProtocolVersion, + Type: MessageGetDevices, + }, + } +} + +// NewDeviceList creates a device-list message from a list of remote +// devices and their session-specific data ports. +func NewDeviceList(devices []RemoteDeviceInfo) DeviceListMessage { + return DeviceListMessage{ + Message: Message{ + Version: ProtocolVersion, + Type: MessageDeviceList, + }, + Devices: devices, + } +} + +// NewError creates an error message. +func NewError(code ErrorCode, message string) ErrorMessage { + return ErrorMessage{ + Message: Message{ + Version: ProtocolVersion, + Type: MessageError, + }, + Code: code, + MessageText: message, + } +} + +// NewRemoteDeviceInfo converts a public configuration device description +// into the transport representation and adds the session-specific data port. +func NewRemoteDeviceInfo(device config.RemoteDevice, dataPort int) RemoteDeviceInfo { + return RemoteDeviceInfo{ + ID: device.ID, + Name: device.Name, + BaudRate: device.BaudRate, + DataBits: device.DataBits, + Parity: device.Parity, + StopBits: device.StopBits, + DataPort: dataPort, + } +} + +/////////////////////////////////////////////////////////////////////////////// +// JSON +/////////////////////////////////////////////////////////////////////////////// + +// EncodeMessage encodes a control message as JSON. +func EncodeMessage(message interface{}) ([]byte, error) { + if message == nil { + return nil, fmt.Errorf("message is nil") + } + + return json.Marshal(message) +} diff --git a/internal/transport/control_test.go b/internal/transport/control_test.go new file mode 100644 index 0000000..a967b7f --- /dev/null +++ b/internal/transport/control_test.go @@ -0,0 +1,125 @@ +package transport_test + +import ( + "encoding/json" + "strings" + "testing" + + "git.lang-dieter.de/rs2322tcp/internal/config" + "git.lang-dieter.de/rs2322tcp/internal/transport" +) + +func TestNewHello(t *testing.T) { + message := transport.NewHello() + + if message.Version != transport.ProtocolVersion { + t.Errorf("Version = %d, want %d", + message.Version, transport.ProtocolVersion) + } + + if message.Type != transport.MessageHello { + t.Errorf("Type = %q, want %q", + message.Type, transport.MessageHello) + } +} + +func TestNewGetDevices(t *testing.T) { + message := transport.NewGetDevices() + + if message.Version != transport.ProtocolVersion { + t.Errorf("Version = %d, want %d", + message.Version, transport.ProtocolVersion) + } + + if message.Type != transport.MessageGetDevices { + t.Errorf("Type = %q, want %q", + message.Type, transport.MessageGetDevices) + } +} + +func TestDeviceListJSON(t *testing.T) { + message := transport.NewDeviceList([]transport.RemoteDeviceInfo{ + { + ID: "radio", + Name: "Funkgerät", + BaudRate: 9600, + DataBits: 8, + Parity: "none", + StopBits: 1, + DataPort: 43721, + }, + }) + + data, err := transport.EncodeMessage(message) + if err != nil { + t.Fatalf("EncodeMessage() failed: %v", err) + } + + jsonText := string(data) + + if !strings.Contains(jsonText, `"type":"device_list"`) { + t.Fatalf("JSON does not contain device_list: %s", jsonText) + } + + if !strings.Contains(jsonText, `"data_port":43721`) { + t.Fatalf("JSON does not contain data port: %s", jsonText) + } + + if strings.Contains(jsonText, "serial_port") { + t.Fatalf("JSON contains server-internal serial_port") + } +} + +func TestNewRemoteDeviceInfo(t *testing.T) { + device := config.RemoteDevice{ + ID: "radio", + Name: "Funkgerät", + BaudRate: 9600, + DataBits: 8, + Parity: "none", + StopBits: 1, + } + + info := transport.NewRemoteDeviceInfo(device, 43721) + + if info.ID != "radio" { + t.Errorf("ID = %q, want %q", info.ID, "radio") + } + + if info.DataPort != 43721 { + t.Errorf("DataPort = %d, want %d", info.DataPort, 43721) + } +} + +func TestErrorMessageJSON(t *testing.T) { + message := transport.NewError( + transport.ErrorUnknownMessage, + "unknown message type", + ) + + data, err := transport.EncodeMessage(message) + if err != nil { + t.Fatalf("EncodeMessage() failed: %v", err) + } + + var decoded map[string]interface{} + + if err := json.Unmarshal(data, &decoded); err != nil { + t.Fatalf("json.Unmarshal() failed: %v", err) + } + + if decoded["type"] != string(transport.MessageError) { + t.Errorf("type = %v, want %q", + decoded["type"], transport.MessageError) + } + + if decoded["code"] != string(transport.ErrorUnknownMessage) { + t.Errorf("code = %v, want %q", + decoded["code"], transport.ErrorUnknownMessage) + } + + if decoded["message"] != "unknown message type" { + t.Errorf("message = %v, want %q", + decoded["message"], "unknown message type") + } +}