Class: Wurk::Queue

Inherits:
Object
  • Object
show all
Includes:
Enumerable, API::Fast::QueueExt
Defined in:
lib/wurk/queue.rb

Overview

Inspects / iterates one named queue against the canonical Sidekiq schema: LIST at queue:<name>, membership in the queues SET. Wire-compat is sacred — every Redis call below matches OSS exactly.

Pause/unpause is a Pro feature upstream; Wurk ships it free. Membership of the paused SET drives both paused? and the fetcher's queue filter (see Wurk::Fetcher::Reliable#queues_cmd). In-flight jobs continue to completion — pausing only stops new fetches.

Fetchers answer from a PAUSED_TTL cache rather than re-reading the SET every pass, so pause!/unpause! expire this process's copy inline and other processes pick the change up within the TTL. paused? is a reporting call and always reads Redis.

Spec: docs/target/sidekiq-free.md §19.2, sidekiq-pro.md §6.

Constant Summary collapse

PAGE_SIZE =

Page size for LRANGE traversal. Matches upstream so dashboards observing Redis traffic see the same pattern.

50

Instance Attribute Summary collapse

Class Method Summary collapse

Instance Method Summary collapse

Constructor Details

#initialize(name = 'default') ⇒ Queue

Returns a new instance of Queue.



38
39
40
41
# File 'lib/wurk/queue.rb', line 38

def initialize(name = 'default')
  @name = name.to_s
  @rname = Keys.queue(@name)
end

Instance Attribute Details

#nameObject (readonly) Also known as: id

Returns the value of attribute name.



29
30
31
# File 'lib/wurk/queue.rb', line 29

def name
  @name
end

Class Method Details

.allArray<Queue>

Returns one per known queue, sorted by name.

Returns:

  • (Array<Queue>)

    one per known queue, sorted by name.



33
34
35
36
# File 'lib/wurk/queue.rb', line 33

def self.all
  names = Wurk.redis(idempotent: true) { |conn| conn.call('SMEMBERS', Keys::QUEUES_SET) }
  names.sort.map { |n| new(n) }
end

Instance Method Details

#as_json(_options = nil) ⇒ Object



115
116
117
# File 'lib/wurk/queue.rb', line 115

def as_json(_options = nil)
  { name: @name }
end

#clearObject

UNLINK the list + drop the queue from the queues set. Pipelined so a partial failure leaves at most one of the two ops applied. Method name is Sidekiq wire-compat — clear? would break the alias.

Unlike pause!/unpause!, this one can't claim apply-safety: a replay after a lost reply would UNLINK whatever a producer enqueued in between.



105
106
107
108
109
110
111
112
113
# File 'lib/wurk/queue.rb', line 105

def clear # rubocop:disable Naming/PredicateMethod
  Wurk.redis do |conn|
    conn.pipelined do |pipe|
      pipe.call('UNLINK', @rname)
      pipe.call('SREM', Keys::QUEUES_SET, @name)
    end
  end
  true
end

#delete_by_class(klass) ⇒ Integer Originally defined in module API::Fast::QueueExt

Returns payloads removed.

Parameters:

  • klass (Class, String, Symbol)

Returns:

  • (Integer)

    payloads removed.

Raises:

  • (ArgumentError)

#delete_job(jid) ⇒ Integer Originally defined in module API::Fast::QueueExt

Returns payloads removed (0 when jid absent; >0 only in the corner case of duplicate-jid corruption).

Returns:

  • (Integer)

    payloads removed (0 when jid absent; >0 only in the corner case of duplicate-jid corruption).

Raises:

  • (ArgumentError)

#eachObject

Paged LRANGE traversal. Yields JobRecord per payload. Continues paging until Redis returns < PAGE_SIZE rows.



80
81
82
83
84
85
86
87
88
89
90
91
# File 'lib/wurk/queue.rb', line 80

def each
  page = 0
  loop do
    start  = page * PAGE_SIZE
    stop   = start + PAGE_SIZE - 1
    slice  = Wurk.redis(idempotent: true) { |conn| conn.call('LRANGE', @rname, start, stop) }
    slice.each { |value| yield JobRecord.new(value, @name) }
    break if slice.size < PAGE_SIZE

    page += 1
  end
end

#find_job(jid) ⇒ Object

O(n) scan. Returns nil if no job matches.



94
95
96
97
# File 'lib/wurk/queue.rb', line 94

def find_job(jid)
  each { |record| return record if record.jid == jid }
  nil
end

#latencyObject

Seconds since the oldest job (tail of LIST) was enqueued. 0.0 when empty.



48
49
50
51
52
53
54
55
# File 'lib/wurk/queue.rb', line 48

def latency
  payload = Wurk.redis(idempotent: true) { |conn| conn.call('LRANGE', @rname, -1, -1).first }
  return 0.0 if payload.nil?

  JobRecord.latency_from(Wurk.load_json(payload)['enqueued_at'])
rescue ::JSON::ParserError
  0.0
end

#pause!Object

Pause new fetches against this queue. Idempotent — SADD returns 0 when the name was already present. In-flight jobs are untouched.



65
66
67
68
69
# File 'lib/wurk/queue.rb', line 65

def pause! # rubocop:disable Naming/PredicateMethod
  Wurk.redis(idempotent: true) { |conn| conn.call('SADD', Keys::PAUSED_SET, @name) }
  Fetcher::Reliable.invalidate_paused_cache!
  true
end

#paused?Boolean

True iff this queue's name is a member of the paused SET. Wurk implements the Pro contract for free; fetchers consult the same set.

Returns:

  • (Boolean)


59
60
61
# File 'lib/wurk/queue.rb', line 59

def paused?
  Wurk.redis(idempotent: true) { |conn| conn.call('SISMEMBER', Keys::PAUSED_SET, @name) } == 1
end

#sizeObject



43
44
45
# File 'lib/wurk/queue.rb', line 43

def size
  Wurk.redis(idempotent: true) { |conn| conn.call('LLEN', @rname) }
end

#unpause!Object

Resume fetches. Idempotent.



72
73
74
75
76
# File 'lib/wurk/queue.rb', line 72

def unpause! # rubocop:disable Naming/PredicateMethod
  Wurk.redis(idempotent: true) { |conn| conn.call('SREM', Keys::PAUSED_SET, @name) }
  Fetcher::Reliable.invalidate_paused_cache!
  true
end