Class: Omnizip::Parallel::JobQueue

Inherits:
Object
  • Object
show all
Defined in:
lib/omnizip/parallel/job_queue.rb

Overview

Thread-safe job queue for parallel compression/extraction

Manages a queue of compression or extraction jobs with priority support. Jobs are ordered by priority (large files first for better load balancing).

Examples:

Create and use job queue

queue = Omnizip::Parallel::JobQueue.new(max_size: 100)
queue.push(file: 'large.dat', size: 1_000_000, priority: :high)
job = queue.pop

Size-based priority

queue.push_with_size(file: 'file.txt', size: 1024)

Defined Under Namespace

Classes: Job

Instance Attribute Summary collapse

Instance Method Summary collapse

Constructor Details

#initialize(max_size: 1000) ⇒ JobQueue

Initialize job queue

Parameters:

  • max_size (Integer) (defaults to: 1000)

    maximum number of jobs in queue



41
42
43
44
45
46
47
48
# File 'lib/omnizip/parallel/job_queue.rb', line 41

def initialize(max_size: 1000)
  @max_size = max_size
  @queue = []
  @mutex = Mutex.new
  @cond = ConditionVariable.new
  @closed = false
  @size = 0
end

Instance Attribute Details

#max_sizeInteger (readonly)

Returns maximum queue size.

Returns:

  • (Integer)

    maximum queue size



33
34
35
# File 'lib/omnizip/parallel/job_queue.rb', line 33

def max_size
  @max_size
end

#sizeInteger (readonly)

Returns current queue size.

Returns:

  • (Integer)

    current queue size



36
37
38
# File 'lib/omnizip/parallel/job_queue.rb', line 36

def size
  @size
end

Instance Method Details

#clearInteger

Clear all jobs from queue

Returns:

  • (Integer)

    number of jobs cleared



180
181
182
183
184
185
186
187
188
# File 'lib/omnizip/parallel/job_queue.rb', line 180

def clear
  @mutex.synchronize do
    count = @queue.size
    @queue.clear
    @size = 0
    @cond.broadcast
    count
  end
end

#closeObject

Close the queue

No more jobs can be pushed after closing. Pending pops will return nil.



170
171
172
173
174
175
# File 'lib/omnizip/parallel/job_queue.rb', line 170

def close
  @mutex.synchronize do
    @closed = true
    @cond.broadcast # Wake up all waiting threads
  end
end

#closed?Boolean

Check if queue is closed

Returns:

  • (Boolean)

    true if closed



162
163
164
# File 'lib/omnizip/parallel/job_queue.rb', line 162

def closed?
  @mutex.synchronize { @closed }
end

#empty?Boolean

Check if queue is empty

Returns:

  • (Boolean)

    true if empty



155
156
157
# File 'lib/omnizip/parallel/job_queue.rb', line 155

def empty?
  @mutex.synchronize { @queue.empty? }
end

#pop(timeout: nil) ⇒ Job?

Pop a job from the queue

Parameters:

  • timeout (Numeric, nil) (defaults to: nil)

    timeout in seconds, nil for no timeout

Returns:

  • (Job, nil)

    job or nil if timeout or closed



113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
# File 'lib/omnizip/parallel/job_queue.rb', line 113

def pop(timeout: nil)
  @mutex.synchronize do
    if timeout
      deadline = Time.now + timeout
      while @queue.empty? && !@closed
        remaining = deadline - Time.now
        return nil if remaining <= 0

        @cond.wait(@mutex, remaining)
      end
    else
      @cond.wait(@mutex) while @queue.empty? && !@closed
    end

    return nil if @queue.empty?

    job = @queue.shift
    @size -= 1
    @cond.signal # Signal waiting pushers
    job
  end
end

#pop_batch(count, timeout: nil) ⇒ Array<Job>

Pop multiple jobs in batch

Parameters:

  • count (Integer)

    maximum number of jobs to pop

  • timeout (Numeric, nil) (defaults to: nil)

    timeout in seconds

Returns:

  • (Array<Job>)

    array of jobs (may be empty)



141
142
143
144
145
146
147
148
149
150
# File 'lib/omnizip/parallel/job_queue.rb', line 141

def pop_batch(count, timeout: nil)
  jobs = []
  count.times do
    job = pop(timeout: timeout)
    break unless job

    jobs << job
  end
  jobs
end

#push(file:, data: nil, size: 0, priority: :normal, metadata: {}) ⇒ Job

Push a job onto the queue

Parameters:

  • file (String)

    file path

  • data (Object) (defaults to: nil)

    job data

  • size (Integer) (defaults to: 0)

    file size in bytes

  • priority (Symbol) (defaults to: :normal)

    job priority (:high, :normal, :low)

  • metadata (Hash) (defaults to: {})

    additional metadata

Returns:

  • (Job)

    the created job

Raises:



59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
# File 'lib/omnizip/parallel/job_queue.rb', line 59

def push(file:, data: nil, size: 0, priority: :normal, metadata: {})
  @mutex.synchronize do
    raise ClosedQueueError, "Queue is closed" if @closed

    # Wait if queue is full
    @cond.wait(@mutex) while @size >= @max_size && !@closed

    raise ClosedQueueError, "Queue is closed" if @closed

    job = Job.new(
      file: file,
      data: data,
      size: size,
      priority: priority,
      metadata: ,
    )

    @queue << job
    @size += 1

    # Keep queue sorted by priority
    @queue.sort!

    @cond.signal
    job
  end
end

#push_with_size(file:, size:, data: nil, metadata: {}) ⇒ Job

Push a job with automatic priority based on file size

Parameters:

  • file (String)

    file path

  • size (Integer)

    file size in bytes

  • data (Object) (defaults to: nil)

    job data

  • metadata (Hash) (defaults to: {})

    additional metadata

Returns:

  • (Job)

    the created job



94
95
96
97
98
99
100
101
102
103
104
105
106
107
# File 'lib/omnizip/parallel/job_queue.rb', line 94

def push_with_size(file:, size:, data: nil, metadata: {})
  # Determine priority based on size
  # Large files (>10MB) get high priority for better load balancing
  priority = if size > 10 * 1024 * 1024
               :high
             elsif size > 1024 * 1024
               :normal
             else
               :low
             end

  push(file: file, data: data, size: size, priority: priority,
       metadata: )
end

#statsHash

Get queue statistics

Returns:

  • (Hash)

    statistics hash



193
194
195
196
197
198
199
200
201
202
203
# File 'lib/omnizip/parallel/job_queue.rb', line 193

def stats
  @mutex.synchronize do
    {
      size: @size,
      max_size: @max_size,
      closed: @closed,
      utilization: @max_size.zero? ? 0.0 : @size.to_f / @max_size,
      priority_counts: @queue.group_by(&:priority).transform_values(&:count),
    }
  end
end