Class: Takagi::Server::Tcp

Inherits:
Object
  • Object
show all
Defined in:
lib/takagi/server/tcp.rb,
sig/takagi/server/tcp.rbs

Overview

TCP server implementation for CoAP over TCP

Instance Method Summary collapse

Constructor Details

#initialize(port: 5683, worker_threads: 2, middleware_stack: nil, router: nil, logger: nil, watcher: nil, sender: nil) ⇒ Tcp

Returns a new instance of Tcp.

Parameters:

  • port: (::Integer) (defaults to: 5683)
  • worker_threads: (::Integer) (defaults to: 2)
  • middleware_stack: (Object, nil) (defaults to: nil)
  • router: (Object, nil) (defaults to: nil)
  • logger: (Object, nil) (defaults to: nil)
  • watcher: (Object, nil) (defaults to: nil)
  • sender: (Object, nil) (defaults to: nil)


10
11
12
13
14
15
16
17
18
19
20
21
22
23
# File 'lib/takagi/server/tcp.rb', line 10

def initialize(port: 5683, worker_threads: 2,
               middleware_stack: nil, router: nil, logger: nil, watcher: nil, sender: nil)
  @port = port
  @worker_threads = worker_threads
  @middleware_stack = middleware_stack || Takagi::MiddlewareStack.instance
  @router = router || Takagi::Router.instance
  @logger = logger || Takagi.logger
  @watcher = watcher || Takagi::Observer::Watcher.new(interval: 1)

  Initializer.run!

  @server = TCPServer.new('0.0.0.0', @port)
  @sender = sender || Takagi::Network::TcpSender.instance
end

Instance Method Details

#build_response(inbound_request) ⇒ Object

Parameters:

  • inbound_request (Object)

Returns:

  • (Object)


134
135
136
137
138
139
140
# File 'lib/takagi/server/tcp.rb', line 134

def build_response(inbound_request)
  # TCP already runs each connection in its own thread, so we don't need
  # to delegate to controller pools - just process synchronously.
  # Controller pools are primarily for UDP where we have fixed worker threads.
  result = @middleware_stack.call(inbound_request)
  ResponseBuilder.build(inbound_request, result, logger: @logger)
end

#handle_connection(sock) ⇒ Object

Parameters:

  • sock (Object)

Returns:

  • (Object)


81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
# File 'lib/takagi/server/tcp.rb', line 81

def handle_connection(sock)
  # RFC 8323 §5.3: Read client CSM first, then send server CSM
  csm_received = false

  loop do
    inbound_request = read_request(sock)
    break unless inbound_request

    @logger.debug "Received request from client: #{inbound_request.inspect}"

    case inbound_request.code
    when CoAP::Registries::Signaling::CSM
      @logger.debug "Received CSM from client"
      unless csm_received
        # Send our CSM in response to client's CSM
        send_csm(sock)
        csm_received = true
      end
      next
    when CoAP::Registries::Signaling::PING
      @logger.debug "Received PING from client"
      send_pong(sock, inbound_request)
      next
    when CoAP::Registries::Signaling::RELEASE, CoAP::Registries::Signaling::ABORT
      @logger.debug "Received #{Takagi::CoAP::Registries::Signaling.name_for(inbound_request.code)} from client, closing connection"
      break
    end

    # Process regular CoAP requests
    response = build_response(inbound_request)
    transmit_response(sock, response)
  end
  @logger.debug "Client connection closed gracefully"
rescue StandardError => e
  @logger.error "TCP handle_connection failed: #{e.message}"
  @logger.debug e.backtrace.join("\n")
ensure
  sock.close unless sock.closed?
end

#read_request(sock) ⇒ nil, untyped

Read request using RFC 8323 §3.3 variable-length framing Uses the new Network::Framing::Tcp module

Parameters:

  • sock (Object)

Returns:

  • (nil, untyped)


123
124
125
126
127
128
129
130
131
132
# File 'lib/takagi/server/tcp.rb', line 123

def read_request(sock)
  # NEW: Use transport framing module
  data = Takagi::Network::Framing::Tcp.read_from_socket(sock, logger: @logger)
  return nil unless data

  Takagi::Message::Inbound.new(data, transport: :tcp)
rescue IOError, Errno::ECONNRESET => e
  @logger.debug "read_request: Socket error (#{e.class}: #{e.message})"
  nil
end

#run!Object

Returns:

  • (Object)


25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
# File 'lib/takagi/server/tcp.rb', line 25

def run!
  Takagi::Hooks.emit(:server_starting, protocol: :tcp, port: @port)
  @logger.info "Starting Takagi TCP server on port #{@port}"
  @workers = []
  @watcher.start

  # Set flag instead of calling shutdown! directly from trap context
  # This avoids "can't be called from trap context" errors with logger
  trap('INT') { @shutdown_requested = true }

  loop do
    break if @shutdown_called || @shutdown_requested

    begin
      @logger.debug "Waiting for client connection..."
      client = @server.accept
      @logger.debug "Client connected from #{client.peeraddr.inspect}"
    rescue IOError, SystemCallError => e
      @logger.error "TCP server accept failed: #{e.class}: #{e.message}"
      @logger.debug "TCP server accept loop exiting: #{e.message}" if @shutdown_called
      break
    end

    @logger.debug "Spawning handler thread for client"
    Thread.new(client) do |sock|
      begin
        handle_connection(sock)
      rescue => e
        @logger.error "Handler thread crashed: #{e.class}: #{e.message}"
        @logger.debug e.backtrace.join("\n")
      end
    end
  end

  # Call shutdown if it was requested by signal
  shutdown! if @shutdown_requested

  @logger.info "TCP server stopped"
  Takagi::Hooks.emit(:server_stopped, protocol: :tcp, port: @port)
end

#shutdown!nil, untyped

Returns:

  • (nil, untyped)


66
67
68
69
70
71
72
73
74
75
76
77
# File 'lib/takagi/server/tcp.rb', line 66

def shutdown!
  return if @shutdown_called

  @shutdown_called = true
  @watcher.stop
  @server.close if @server && !@server.closed?

  # Join the server thread if it was spawned
  if defined?(@server_thread) && @server_thread&.alive?
    @server_thread.join(5) # Wait up to 5 seconds
  end
end

#transmit_response(sock, response) ⇒ Object

Parameters:

  • sock (Object)
  • response (Object)

Returns:

  • (Object)


142
143
144
145
146
147
148
# File 'lib/takagi/server/tcp.rb', line 142

def transmit_response(sock, response)
  # NEW: to_bytes now returns fully framed data from transport registry
  framed = response.to_bytes(transport: :tcp)
  written = sock.write(framed)
  sock.flush
  @logger.debug "Sent #{framed.bytesize} bytes to client (wrote #{written} bytes)"
end