Module: OMQ::QoS::GatherExt
- Defined in:
- lib/omq/qos/routing_ext.rb
Overview
Same pattern as PullExt — GATHER uses fair-recv too.
Instance Method Summary collapse
Instance Method Details
#connection_added(conn) ⇒ Object
418 419 420 421 422 423 424 425 426 427 428 429 430 431 432 433 434 435 436 437 438 439 440 441 442 443 444 445 446 447 448 449 450 451 452 453 454 455 456 457 458 459 460 461 462 463 464 465 466 467 |
# File 'lib/omq/qos/routing_ext.rb', line 418 def connection_added(conn) qos = @engine..qos return super if qos.nil? return super if QoS.reliable_transport?(conn) algo = QoS.algo_for(conn) recv_queue = @recv_queue engine = @engine case qos.level when 1 engine.start_recv_pump(conn, recv_queue) do |msg| recv_queue.enqueue(msg) conn.send_command(QoS.ack_command(msg, algorithm: algo)) nil end when 2 dedup = qos.dedup_set_for(conn) QoS.install_qos2_receiver_handler(conn, dedup) engine.start_recv_pump(conn, recv_queue) do |msg| digest = QoS.digest(msg, algorithm: algo) ack = Protocol::ZMTP::Codec::Command.ack(digest, algorithm: algo) if dedup.seen?(digest) conn.send_command(ack) else dedup.add(digest) recv_queue.enqueue(msg) conn.send_command(ack) end nil end when 3 dedup = qos.dedup_set_for(conn) QoS.install_qos2_receiver_handler(conn, dedup) engine.start_recv_pump(conn, recv_queue) do |msg| digest = QoS.digest(msg, algorithm: algo) if dedup.seen?(digest) conn.send_command(Protocol::ZMTP::Codec::Command.comp(digest, algorithm: algo)) else recv_queue.enqueue(QoS::Envelope.new(parts: msg, conn: conn, digest: digest, algo: algo)) end nil end end end |