Class: Mfp::Server::Server
- Inherits:
-
Object
- Object
- Mfp::Server::Server
- 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
-
#port ⇒ Object
readonly
Returns the value of attribute port.
Instance Method Summary collapse
- #accept ⇒ Object
- #close ⇒ Object
- #connect(client) ⇒ Object
- #disconnect(client) ⇒ Object
-
#initialize(host, port, call_sign, **opts) ⇒ Server
constructor
A new instance of Server.
- #start(task) ⇒ Object
Constructor Details
#initialize(host, port, call_sign, **opts) ⇒ Server
Returns a new instance of Server.
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
#port ⇒ Object (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
#accept ⇒ Object
115 |
# File 'lib/mfp/server/server.rb', line 115 def accept = @client_queue.dequeue |
#close ⇒ Object
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 |