Class: OMQ::Backend::Libzmq::Engine

Inherits:
Object
  • Object
show all
Defined in:
lib/omq/backend/libzmq/engine.rb

Overview

Engine — wraps a libzmq socket to implement the OMQ Engine contract.

A dedicated I/O thread owns the zmq_socket exclusively (libzmq sockets are not thread-safe). Send and recv flow through queues, with an IO pipe to wake the Async fiber scheduler.

Defined Under Namespace

Classes: RoutingStub

Constant Summary collapse

L =
Native

Instance Attribute Summary collapse

Class Method Summary collapse

Instance Method Summary collapse

Constructor Details

#initialize(socket_type, options) ⇒ Engine

Returns a new instance of Engine.

Parameters:

  • socket_type (Symbol)

    e.g. :REQ, :PAIR

  • options (Options)


110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
# File 'lib/omq/backend/libzmq/engine.rb', line 110

def initialize(socket_type, options)
  @socket_type    = socket_type
  @options        = options
  @peer_connected = Async::Promise.new
  @all_peers_gone = Async::Promise.new
  @connections    = []
  @closed         = false
  @parent_task    = nil
  @on_io_thread   = false
  @fatal_error    = nil

  @zmq_socket = L.zmq_socket(OMQ::Backend::Libzmq.context, L::SOCKET_TYPES.fetch(@socket_type))
  raise "zmq_socket failed: #{L.zmq_strerror(L.zmq_errno)}" if @zmq_socket.null?

  apply_options

  @routing = RoutingStub.new(self)

  # Queues for cross-thread communication
  @send_queue = Thread::Queue.new   # main → io thread
  @recv_queue = Thread::Queue.new   # io thread → main
  @cmd_queue  = Thread::Queue.new   # control commands → io thread

  # Signal pipe: io thread → Async fiber (message received)
  @recv_signal_r, @recv_signal_w = IO.pipe
  # Wake pipe: main thread → io thread (send/cmd enqueued)
  @wake_r, @wake_w = IO.pipe

  @io_thread = nil
end

Instance Attribute Details

#all_peers_goneAsync::Promise (readonly)

Returns resolved when all peers have disconnected.

Returns:

  • (Async::Promise)

    resolved when all peers have disconnected



28
29
30
# File 'lib/omq/backend/libzmq/engine.rb', line 28

def all_peers_gone
  @all_peers_gone
end

#connectionsArray (readonly)

Returns active connections.

Returns:

  • (Array)

    active connections



22
23
24
# File 'lib/omq/backend/libzmq/engine.rb', line 22

def connections
  @connections
end

#monitor_queue=(value) ⇒ Object (writeonly)

Note:

Monitor events are not yet emitted by the libzmq backend; these writers exist so Socket#monitor can attach without raising. Wiring libzmq's zmq_socket_monitor is a TODO.



41
42
43
# File 'lib/omq/backend/libzmq/engine.rb', line 41

def monitor_queue=(value)
  @monitor_queue = value
end

#on_io_threadBoolean (readonly) Also known as: on_io_thread?

Returns true when the engine's parent task lives on the shared Reactor IO thread (i.e. not created under an Async task). Writable/Readable check this to pick the fast path.

Returns:

  • (Boolean)

    true when the engine's parent task lives on the shared Reactor IO thread (i.e. not created under an Async task). Writable/Readable check this to pick the fast path.



34
35
36
# File 'lib/omq/backend/libzmq/engine.rb', line 34

def on_io_thread
  @on_io_thread
end

#optionsOptions (readonly)

Returns socket options.

Returns:

  • (Options)

    socket options



20
21
22
# File 'lib/omq/backend/libzmq/engine.rb', line 20

def options
  @options
end

#parent_taskAsync::Task? (readonly)

Returns root of the engine's task tree.

Returns:

  • (Async::Task, nil)

    root of the engine's task tree



30
31
32
# File 'lib/omq/backend/libzmq/engine.rb', line 30

def parent_task
  @parent_task
end

#peer_connectedAsync::Promise (readonly)

Returns resolved when the first peer connects.

Returns:

  • (Async::Promise)

    resolved when the first peer connects



26
27
28
# File 'lib/omq/backend/libzmq/engine.rb', line 26

def peer_connected
  @peer_connected
end

#reconnect_enabled=(value) ⇒ Object (writeonly)

Parameters:

  • value (Boolean)

    enables or disables automatic reconnection



37
38
39
# File 'lib/omq/backend/libzmq/engine.rb', line 37

def reconnect_enabled=(value)
  @reconnect_enabled = value
end

#routingRoutingStub (readonly)

Returns subscription/group routing interface.

Returns:

  • (RoutingStub)

    subscription/group routing interface



24
25
26
# File 'lib/omq/backend/libzmq/engine.rb', line 24

def routing
  @routing
end

#verbose_monitor=(value) ⇒ Object (writeonly)

Note:

Monitor events are not yet emitted by the libzmq backend; these writers exist so Socket#monitor can attach without raising. Wiring libzmq's zmq_socket_monitor is a TODO.



41
42
43
# File 'lib/omq/backend/libzmq/engine.rb', line 41

def verbose_monitor=(value)
  @verbose_monitor = value
end

Class Method Details

.linger_to_zmq_ms(linger) ⇒ Integer

Maps an OMQ linger value (seconds, or +nil+/+Float::INFINITY+ for "wait forever") to libzmq's ZMQ_LINGER int milliseconds (-1 = infinite, 0 = drop, N = N ms).

Parameters:

  • linger (Numeric, nil)

Returns:

  • (Integer)


101
102
103
104
# File 'lib/omq/backend/libzmq/engine.rb', line 101

def self.linger_to_zmq_ms(linger)
  return -1 if linger.nil? || linger == Float::INFINITY
  (linger * 1000).to_i
end

Instance Method Details

#bind(endpoint) ⇒ URI::Generic

Binds the socket to the given endpoint.

Parameters:

  • endpoint (String)

    ZMQ endpoint URL (e.g. "tcp://*:5555")

Returns:

  • (URI::Generic)

    resolved endpoint URI (with auto-selected port for "tcp://host:0")



148
149
150
151
152
153
154
155
156
# File 'lib/omq/backend/libzmq/engine.rb', line 148

def bind(endpoint)
  sync_identity
  send_cmd(:bind, endpoint)
  resolved = get_string_option(L::ZMQ_LAST_ENDPOINT)
  @connections << :libzmq
  @peer_connected.resolve(:libzmq) unless @peer_connected.resolved?
  resolve_subscriber_joined_without_monitor
  URI.parse(resolved)
end

#capture_parent_task(parent: nil) ⇒ void

This method returns an undefined value.

Captures the current Async task as the parent for I/O scheduling. parent: is accepted for API compatibility with the pure-Ruby engine but has no effect: the libzmq backend runs its own I/O thread and doesn't participate in the Async barrier tree.



257
258
259
260
261
262
263
264
265
266
267
268
# File 'lib/omq/backend/libzmq/engine.rb', line 257

def capture_parent_task(parent: nil)
  return if @parent_task
  if parent
    @parent_task = parent
  elsif Async::Task.current?
    @parent_task = Async::Task.current
  else
    @parent_task  = Reactor.root_task
    @on_io_thread = true
    Reactor.track_linger(@options.linger)
  end
end

#closevoid

This method returns an undefined value.

Closes the socket and shuts down the I/O thread.

Honors options.linger:

nil → wait forever for Ruby-side queue to drain into libzmq
    and for libzmq's own LINGER to flush to the network
0   → drop anything not yet in libzmq's kernel buffers, close fast
N   → up to N seconds for drain + N + 1s grace for join


225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
# File 'lib/omq/backend/libzmq/engine.rb', line 225

def close
  return if @closed
  @closed = true
  if @io_thread
    @cmd_queue.push([:stop])
    wake_io_thread
    linger = @options.linger
    if linger.nil?
      @io_thread.join
    elsif linger.zero?
      @io_thread.join(0.5) # fast path: zmq_close is non-blocking with LINGER=0
    else
      @io_thread.join(linger + 1.0)
    end
    @io_thread.kill if @io_thread.alive? # hard stop if deadline exceeded
  else
    # IO thread never started — close socket directly
    L.zmq_close(@zmq_socket)
  end
  @recv_signal_r&.close rescue nil
  @recv_signal_w&.close rescue nil
  @wake_r&.close rescue nil
  @wake_w&.close rescue nil
end

#connect(endpoint) ⇒ URI::Generic

Connects the socket to the given endpoint.

Parameters:

  • endpoint (String)

    ZMQ endpoint URL

Returns:

  • (URI::Generic)

    parsed endpoint URI



163
164
165
166
167
168
169
170
# File 'lib/omq/backend/libzmq/engine.rb', line 163

def connect(endpoint)
  sync_identity
  send_cmd(:connect, endpoint)
  @connections << :libzmq
  @peer_connected.resolve(:libzmq) unless @peer_connected.resolved?
  resolve_subscriber_joined_without_monitor
  URI.parse(endpoint)
end

#dequeue_recvArray<String>

Dequeues the next received message, blocking until one is available.

Returns:

  • (Array<String>)

    multipart message



290
291
292
293
294
295
296
# File 'lib/omq/backend/libzmq/engine.rb', line 290

def dequeue_recv
  raise_if_dead!
  ensure_io_thread
  msg = wait_for_message
  raise_if_dead! if msg.nil?
  msg
end

#dequeue_recv_sentinelvoid

This method returns an undefined value.

Pushes a nil sentinel into the recv queue to unblock a waiting consumer.



302
303
304
305
# File 'lib/omq/backend/libzmq/engine.rb', line 302

def dequeue_recv_sentinel
  @recv_queue.push(nil)
  @recv_signal_w.write_nonblock(".", exception: false) rescue nil
end

#disconnect(endpoint) ⇒ void

This method returns an undefined value.

Disconnects from the given endpoint.

Parameters:

  • endpoint (String)

    ZMQ endpoint URL



177
178
179
# File 'lib/omq/backend/libzmq/engine.rb', line 177

def disconnect(endpoint)
  send_cmd(:disconnect, endpoint)
end

#enqueue_send(parts) ⇒ void

This method returns an undefined value.

Enqueues a multipart message for sending via the I/O thread.

Parameters:

  • parts (Array<String>)

    message frames



277
278
279
280
281
282
# File 'lib/omq/backend/libzmq/engine.rb', line 277

def enqueue_send(parts)
  raise_if_dead!
  ensure_io_thread
  @send_queue.push(parts)
  wake_io_thread
end

#send_cmd(cmd, *args) ⇒ Object

This method is part of a private API. You should avoid using this method if possible, as it may be removed or be changed in the future.

Send a control command to the I/O thread.



311
312
313
314
315
316
317
318
319
320
# File 'lib/omq/backend/libzmq/engine.rb', line 311

def send_cmd(cmd, *args)
  raise_if_dead!
  ensure_io_thread
  result = Thread::Queue.new
  @cmd_queue.push([cmd, args, result])
  wake_io_thread
  r = result.pop
  raise r if r.is_a?(Exception)
  r
end

#subscribe(prefix) ⇒ void

This method returns an undefined value.

Subscribes to a topic prefix (SUB/XSUB). Delegates to the routing stub for API parity with the pure-Ruby Engine.

Parameters:

  • prefix (String)


196
197
198
# File 'lib/omq/backend/libzmq/engine.rb', line 196

def subscribe(prefix)
  @routing.subscribe(prefix)
end

#subscriber_joinedAsync::Promise

Returns resolved when a subscriber joins (PUB/XPUB).

Returns:

  • (Async::Promise)

    resolved when a subscriber joins (PUB/XPUB).



211
212
213
# File 'lib/omq/backend/libzmq/engine.rb', line 211

def subscriber_joined
  @routing.subscriber_joined
end

#unbind(endpoint) ⇒ void

This method returns an undefined value.

Unbinds from the given endpoint.

Parameters:

  • endpoint (String)

    ZMQ endpoint URL



186
187
188
# File 'lib/omq/backend/libzmq/engine.rb', line 186

def unbind(endpoint)
  send_cmd(:unbind, endpoint)
end

#unsubscribe(prefix) ⇒ void

This method returns an undefined value.

Unsubscribes from a topic prefix (SUB/XSUB).

Parameters:

  • prefix (String)


205
206
207
# File 'lib/omq/backend/libzmq/engine.rb', line 205

def unsubscribe(prefix)
  @routing.unsubscribe(prefix)
end

#wake_io_threadvoid

This method returns an undefined value.

Wakes the I/O thread via the internal pipe.



326
327
328
# File 'lib/omq/backend/libzmq/engine.rb', line 326

def wake_io_thread
  @wake_w.write_nonblock(".", exception: false)
end