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
# 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?
  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.



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

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


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

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



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

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

#dequeue_recvArray<String>

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

Returns:

  • (Array<String>)

    multipart message



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

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.



300
301
302
303
# File 'lib/omq/backend/libzmq/engine.rb', line 300

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



175
176
177
# File 'lib/omq/backend/libzmq/engine.rb', line 175

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



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

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.



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

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)


194
195
196
# File 'lib/omq/backend/libzmq/engine.rb', line 194

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).



209
210
211
# File 'lib/omq/backend/libzmq/engine.rb', line 209

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



184
185
186
# File 'lib/omq/backend/libzmq/engine.rb', line 184

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)


203
204
205
# File 'lib/omq/backend/libzmq/engine.rb', line 203

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.



324
325
326
# File 'lib/omq/backend/libzmq/engine.rb', line 324

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