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

Class Method Details

.connect_timeout(options) ⇒ Object



103
104
105
106
# File 'lib/omq/transport/zstd_tcp/transport.rb', line 103

def connect_timeout(options)
  ri = options.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.

Parameters:

  • endpoint (String)

    e.g. "zstd+tcp://127.0.0.1:5555"

  • engine (Engine)
  • level (Integer) (defaults to: -3,)

    Zstd compression level

  • dict (String, nil) (defaults to: nil)

    user-supplied dictionary bytes

  • auto_dict (true, Hash, nil) (defaults to: nil)

    enable automatic dictionary training. Pass true for defaults or { capacity: N }.

Returns:



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.

Parameters:

  • endpoint (String)

    e.g. "zstd+tcp://127.0.0.1:5555"

  • engine (Engine)
  • level (Integer) (defaults to: -3,)

    Zstd compression level

  • dict (String, nil) (defaults to: nil)

    user-supplied dictionary bytes

  • auto_dict (true, Hash, nil) (defaults to: nil)

    enable automatic dictionary training. Pass true for defaults or { capacity: N }.

Returns:



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