508 lines
11 KiB
Go
508 lines
11 KiB
Go
/*
|
|
* ============================================================================
|
|
* 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),
|
|
)
|
|
}
|