/* * ============================================================================ * 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, s.config.HardwareErrorResponse, ) if err != nil { log.Printf( "session %d device %q create data connection error: %v", session.ID(), device.ID, err, ) _ = serialConn.Close() _ = tcpConn.Close() continue } dataConnection.SetSerialMonitor( s.config.SerialMonitor, device.SerialPort, ) 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, ) _ = dataConnection.Close() 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), ) }