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