Class: OMQ::Rust::Engine::RoutingStub
- Inherits:
-
Object
- Object
- OMQ::Rust::Engine::RoutingStub
- Defined in:
- lib/omq/rust/engine.rb
Instance Method Summary collapse
-
#initialize(engine) ⇒ RoutingStub
constructor
A new instance of RoutingStub.
- #join(group) ⇒ Object
- #leave(group) ⇒ Object
- #replay_pending(socket) ⇒ Object
- #subscribe(prefix) ⇒ Object
- #subscriber_joined ⇒ Object
- #unsubscribe(prefix) ⇒ Object
Constructor Details
#initialize(engine) ⇒ RoutingStub
Returns a new instance of RoutingStub.
404 405 406 407 408 |
# File 'lib/omq/rust/engine.rb', line 404 def initialize(engine) @engine = engine @pending_subscribe = [] @pending_join = [] end |
Instance Method Details
#join(group) ⇒ Object
431 432 433 434 435 436 437 438 |
# File 'lib/omq/rust/engine.rb', line 431 def join(group) socket = @engine.instance_variable_get(:@socket) if @engine.instance_variable_get(:@materialized) socket.join(group) else @pending_join << group end end |
#leave(group) ⇒ Object
441 442 443 |
# File 'lib/omq/rust/engine.rb', line 441 def leave(group) @engine.instance_variable_get(:@socket).leave(group) end |
#replay_pending(socket) ⇒ Object
446 447 448 449 450 451 |
# File 'lib/omq/rust/engine.rb', line 446 def replay_pending(socket) @pending_subscribe.each { |p| socket.subscribe(p) } @pending_subscribe.clear @pending_join.each { |g| socket.join(g) } @pending_join.clear end |
#subscribe(prefix) ⇒ Object
416 417 418 419 420 421 422 423 |
# File 'lib/omq/rust/engine.rb', line 416 def subscribe(prefix) socket = @engine.instance_variable_get(:@socket) if @engine.instance_variable_get(:@materialized) socket.subscribe(prefix.b) else @pending_subscribe << prefix.b end end |
#subscriber_joined ⇒ Object
411 412 413 |
# File 'lib/omq/rust/engine.rb', line 411 def subscriber_joined @engine.subscriber_joined end |
#unsubscribe(prefix) ⇒ Object
426 427 428 |
# File 'lib/omq/rust/engine.rb', line 426 def unsubscribe(prefix) @engine.instance_variable_get(:@socket).unsubscribe(prefix.b) end |