Module: OMQ::Transport::ZstdTcp
- Defined in:
- lib/omq/transport/zstd_tcp.rb,
lib/omq/transport/zstd_tcp/codec.rb,
lib/omq/transport/zstd_tcp/transport.rb,
lib/omq/transport/zstd_tcp/connection.rb
Defined Under Namespace
Classes: Codec, Dialer, Listener, ProtocolError, ZstdConnection
Constant Summary collapse
- SCHEME =
"zstd+tcp"
Class Method Summary collapse
- .connect_timeout(options) ⇒ Object
-
.dialer(endpoint, engine, level: -3,, dict: nil, auto_dict: nil) ⇒ Dialer
Creates a zstd+tcp dialer for an endpoint.
-
.listener(endpoint, engine, level: -3,, dict: nil, auto_dict: nil) ⇒ Listener
Creates a bound zstd+tcp listener.
- .normalize_bind_host(host) ⇒ Object
- .normalize_connect_host(host) ⇒ Object
- .parse_endpoint(endpoint) ⇒ Object
- .validate_endpoint!(endpoint) ⇒ Object
Class Method Details
.connect_timeout(options) ⇒ Object
103 104 105 106 |
# File 'lib/omq/transport/zstd_tcp/transport.rb', line 103 def connect_timeout() ri = .reconnect_interval ri.is_a?(Range) ? ri.end : [ri * 10, 30].min end |
.dialer(endpoint, engine, level: -3,, dict: nil, auto_dict: nil) ⇒ Dialer
Creates a zstd+tcp dialer for an endpoint.
67 68 69 70 71 72 73 74 75 76 |
# File 'lib/omq/transport/zstd_tcp/transport.rb', line 67 def dialer(endpoint, engine, level: -3, dict: nil, auto_dict: nil, **) validate_auto_dict!(auto_dict, dict) codec = codec_for( engine, level: level, dict: dict, auto_dict: normalize_auto_dict(auto_dict), ) Dialer.new(endpoint, engine, codec) end |
.listener(endpoint, engine, level: -3,, dict: nil, auto_dict: nil) ⇒ Listener
Creates a bound zstd+tcp listener.
31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 |
# File 'lib/omq/transport/zstd_tcp/transport.rb', line 31 def listener(endpoint, engine, level: -3, dict: nil, auto_dict: nil, **) validate_auto_dict!(auto_dict, dict) codec = codec_for( engine, level: level, dict: dict, auto_dict: normalize_auto_dict(auto_dict), ) host, port = parse_endpoint(endpoint) host = normalize_bind_host(host) servers = ::Socket.tcp_server_sockets(host, port) if servers.empty? raise ::Socket::ResolutionError, "no addresses for #{host.inspect}" end actual_port = servers.first.local_address.ip_port display_host = host || "*" host_part = display_host.include?(":") ? "[#{display_host}]" : display_host resolved = "#{SCHEME}://#{host_part}:#{actual_port}" Listener.new(resolved, servers, actual_port, engine, codec) end |
.normalize_bind_host(host) ⇒ Object
92 93 94 95 |
# File 'lib/omq/transport/zstd_tcp/transport.rb', line 92 def normalize_bind_host(host) return nil if host == "*" host end |
.normalize_connect_host(host) ⇒ Object
98 99 100 |
# File 'lib/omq/transport/zstd_tcp/transport.rb', line 98 def normalize_connect_host(host) host == "*" ? "127.0.0.1" : host end |
.parse_endpoint(endpoint) ⇒ Object
86 87 88 89 |
# File 'lib/omq/transport/zstd_tcp/transport.rb', line 86 def parse_endpoint(endpoint) uri = URI.parse(endpoint) [uri.hostname, uri.port] end |
.validate_endpoint!(endpoint) ⇒ Object
79 80 81 82 83 |
# File 'lib/omq/transport/zstd_tcp/transport.rb', line 79 def validate_endpoint!(endpoint) host, _port = parse_endpoint(endpoint) host = normalize_connect_host(host) Addrinfo.getaddrinfo(host, nil, nil, :STREAM) if host end |