Class: Mfp::Server::Server

Inherits:
Object
  • Object
show all
Defined in:
lib/mfp/server/server.rb

Constant Summary collapse

DEFAULT_OPTIONS =
{
  max_streams: 0,
  max_payload: 1 << 24,
  idle_timeout_secs: 30,
  ping_timeout_secs: 10,
  cert: nil,
  key: nil,
  ca_certs: [],
  client_auth: nil,
  first_frame_timeout_secs: 5,
  logger: Logrb.noop,
}.freeze

Instance Attribute Summary collapse

Instance Method Summary collapse

Constructor Details

#initialize(host, port, call_sign, **opts) ⇒ Server

Returns a new instance of Server.

Raises:

  • (ArgumentError)


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
65
66
67
68
69
70
# File 'lib/mfp/server/server.rb', line 30

def initialize(host, port, call_sign, **opts)
  raise ArgumentError, "call_sign must have 4 bytes" if call_sign.bytesize != 4

  opts = DEFAULT_OPTIONS.merge(opts)
  @client_id = 0
  @clients = {}
  @client_queue = Async::Queue.new
  @call_sign = call_sign

  @logger = opts[:logger]

  @listener = if opts[:cert] && opts[:key]
    @logger.debug("Starting TLS server", host:, port:)
    ctx = Mfp.default_tls_context
    ctx.cert = OpenSSL::X509::Certificate.new(opts[:cert])
    ctx.key = OpenSSL::PKey.read(opts[:key])
    store = OpenSSL::X509::Store.new
    store.set_default_paths
    opts[:ca_certs].each do |raw|
      store.add_cert(OpenSSL::X509::Certificate.new(raw))
    end
    ctx.cert_store = store
    ctx.verify_mode = opts[:client_auth] if opts[:client_auth]
    IO::Endpoint.ssl(host, port, ssl_context: ctx)
  else
    @logger.debug("Starting plain server", host:, port:)
    IO::Endpoint.tcp(host, port)
  end

  @settings = Settings.new(
    opts[:max_streams],
    opts[:max_payload],
    opts[:idle_timeout_secs],
    opts[:ping_timeout_secs],
    opts[:first_frame_timeout_secs],
    opts[:logger],
  )

  @sockets = @listener.bind
  @port = @sockets.first.local_address.ip_port
end

Instance Attribute Details

#portObject (readonly)

Returns the value of attribute port.



28
29
30
# File 'lib/mfp/server/server.rb', line 28

def port
  @port
end

Instance Method Details

#acceptObject



115
# File 'lib/mfp/server/server.rb', line 115

def accept = @client_queue.dequeue

#closeObject



103
104
105
106
# File 'lib/mfp/server/server.rb', line 103

def close
  @accept_tasks&.each(&:stop)
  @sockets.each(&:close)
end

#connect(client) ⇒ Object



108
# File 'lib/mfp/server/server.rb', line 108

def connect(client) = @client_queue.enqueue(client)

#disconnect(client) ⇒ Object



110
111
112
113
# File 'lib/mfp/server/server.rb', line 110

def disconnect(client)
  @clients.delete(client.id)
  @logger.debug("Disconnected client", id: client.id, connected_clients: @clients.length)
end

#start(task) ⇒ Object



72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
# File 'lib/mfp/server/server.rb', line 72

def start(task)
  return if @running

  @running = true
  @accept_tasks = @sockets.map do |sock|
    task.async do |t|
      loop do
        peer = nil
        begin
          peer, = sock.accept
          peer.accept if peer.respond_to?(:accept) && peer.is_a?(OpenSSL::SSL::SSLSocket)
        rescue StandardError => e
          @logger.debug("Error accepting connection", err: e.to_s)
          peer&.close
          next
        end
        client_addr = peer.remote_address.inspect_sockaddr
        @logger.debug("Accepting connection", addr: client_addr)
        id = @client_id
        @client_id += 1
        settings = Settings.new(**@settings.to_h, logger: @logger.with_fields(id:))
        conn = Conn.new(self, id, peer, @call_sign, settings)
        @logger.debug("Created client", addr: client_addr, id:)
        @clients[id] = conn
        @logger.debug("Starting client in background", id:)
        conn.start(t)
      end
    end
  end
end