Class: Tina4::QueueBackends::LiteBackend
- Inherits:
-
Object
- Object
- Tina4::QueueBackends::LiteBackend
- Defined in:
- lib/tina4/queue_backends/lite_backend.rb
Overview
File-based queue backend — JSON files on disk. Zero dependencies.
Each job is stored as a separate .queue-data file under
Dequeue policy: highest priority first, ties broken oldest-first by the stored created_at (ISO-8601, so lexicographic == chronological). The file name is NOT the ordering key — the stored priority/created_at are.
Constant Summary collapse
- JOB_EXTENSION =
The store's layout — extension, topic segment, dead-letter directory — is defined once on Tina4::Queue and shared with the dev-admin panel that reads the same files.
Tina4::Queue::JOB_EXTENSION
- LEGACY_DIRNAME =
The pre-alignment Ruby store: Dir.pwd/.queue/
/ .json. Read once per backend so an upgrading app's jobs are rescued rather than stranded (see migrate_legacy_store). ".queue"- LEGACY_JOB_EXTENSION =
".json"
Instance Attribute Summary collapse
-
#max_retries ⇒ Object
Retry policy — settable so a Queue can propagate its own max_retries / retry_backoff onto a backend instance passed directly (legacy path).
-
#retry_backoff ⇒ Object
Retry policy — settable so a Queue can propagate its own max_retries / retry_backoff onto a backend instance passed directly (legacy path).
-
#visibility_timeout ⇒ Object
Retry/visibility policy is settable so a Queue can propagate its own configuration onto a backend instance passed directly (legacy path).
Instance Method Summary collapse
- #acknowledge(message) ⇒ Object
-
#clear(topic) ⇒ Object
Remove all pending jobs from a topic (and any held reservations).
-
#close ⇒ Object
No-op: the file backend holds no connection to release.
- #complete(message) ⇒ Object
- #dead_letter(message) ⇒ Object
-
#dead_letter_count(topic) ⇒ Object
Count dead-letter / failed messages for a topic.
-
#dead_letters(topic, max_retries: 3) ⇒ Object
Get dead letter jobs for a topic — messages that exceeded max retries.
- #dequeue(topic) ⇒ Object
- #dequeue_batch(topic, count) ⇒ Object
- #enqueue(message) ⇒ Object
-
#fail(job, error = "") ⇒ Object
Record a failed attempt.
-
#failed(topic, max_retries: 3) ⇒ Object
Jobs that have failed at least once but are still being retried.
-
#find_by_id(topic, id) ⇒ Object
Find a specific pending job by id, claim it (delete the file) and return it.
-
#initialize(options = {}) ⇒ LiteBackend
constructor
A new instance of LiteBackend.
-
#purge(topic, status) ⇒ Object
Delete messages by status.
- #requeue(message) ⇒ Object
-
#reserved_count(topic) ⇒ Object
Count currently-reserved (in-flight) jobs for a topic.
-
#retry(job, delay_seconds: 0) ⇒ Object
Explicit re-queue requested by the caller (job.retry()).
-
#retry_failed(topic, max_retries: 3) ⇒ Object
Revive dead-letter jobs (under max_retries) back to pending.
-
#retry_job(topic, job_id: nil, delay_seconds: 0) ⇒ Object
Revive a specific dead-letter job by id back to the pending queue.
- #size(topic) ⇒ Object
- #topics ⇒ Object
Constructor Details
#initialize(options = {}) ⇒ LiteBackend
Returns a new instance of LiteBackend.
35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 |
# File 'lib/tina4/queue_backends/lite_backend.rb', line 35 def initialize( = {}) # An explicit dir: is the caller's own store — take it as given. Only a # store we RESOLVED is one the legacy rescue may touch. resolved_store = [:dir].nil? @dir = [:dir] || Tina4::Queue.base_path @dead_letter_dir = File.join(@dir, Tina4::Queue::DEAD_LETTER_DIRNAME) # Retry policy. Mirrors the Python lite backend: a failed job is # re-enqueued while attempts < max_retries, then dead-lettered. @max_retries = [:max_retries] || 3 # Seconds to delay a job's next attempt when fail() re-enqueues it. # 0 (default) = retry on the very next pop/consume iteration. @retry_backoff = [:retry_backoff] || 0 # Reservation/visibility timeout (seconds). When a job is popped it is # held in <topic>/reserved/ with available_at = now + visibility_timeout. # If the consumer dies before complete()/fail() (crash, OOM, k8s # eviction) the next pop() reclaims it once the window expires — # incrementing attempts and re-enqueuing, or dead-lettering past # max_retries. <= 0 disables the reclaim (a reservation then lasts until # the consumer acks — the old at-most-once behaviour). @visibility_timeout = [:visibility_timeout] || 300.0 FileUtils.mkdir_p(@dir) FileUtils.mkdir_p(@dead_letter_dir) @mutex = Mutex.new migrate_legacy_store if resolved_store end |
Instance Attribute Details
#max_retries ⇒ Object
Retry policy — settable so a Queue can propagate its own max_retries / retry_backoff onto a backend instance passed directly (legacy path).
33 34 35 |
# File 'lib/tina4/queue_backends/lite_backend.rb', line 33 def max_retries @max_retries end |
#retry_backoff ⇒ Object
Retry policy — settable so a Queue can propagate its own max_retries / retry_backoff onto a backend instance passed directly (legacy path).
33 34 35 |
# File 'lib/tina4/queue_backends/lite_backend.rb', line 33 def retry_backoff @retry_backoff end |
#visibility_timeout ⇒ Object
Retry/visibility policy is settable so a Queue can propagate its own configuration onto a backend instance passed directly (legacy path).
63 64 65 |
# File 'lib/tina4/queue_backends/lite_backend.rb', line 63 def visibility_timeout @visibility_timeout end |
Instance Method Details
#acknowledge(message) ⇒ Object
134 135 136 |
# File 'lib/tina4/queue_backends/lite_backend.rb', line 134 def acknowledge() # File already deleted on dequeue — nothing to do. end |
#clear(topic) ⇒ Object
Remove all pending jobs from a topic (and any held reservations). Returns count removed.
315 316 317 318 319 320 321 322 323 324 325 326 327 |
# File 'lib/tina4/queue_backends/lite_backend.rb', line 315 def clear(topic) count = 0 [topic_path(topic), reserved_path(topic)].each do |dir| next unless Dir.exist?(dir) job_files(dir).each do |file| File.delete(file) count += 1 rescue Errno::ENOENT next end end count end |
#close ⇒ Object
No-op: the file backend holds no connection to release.
It exists so Queue#close can call ONE method on every backend instead of
testing for it, and so switching TINA4_QUEUE_BACKEND to "lite" never
turns a working close into a NoMethodError. It was the ONLY backend
without one, so every close if respond_to?(:close) guard in the tree
silently did nothing here. Idempotent by construction - nothing to drop.
194 |
# File 'lib/tina4/queue_backends/lite_backend.rb', line 194 def close; end |
#complete(message) ⇒ Object
138 139 140 141 142 143 |
# File 'lib/tina4/queue_backends/lite_backend.rb', line 138 def complete() # The pending file was claimed on dequeue and a reservation record # written; complete() is terminal, so drop the reservation. The job is # done and gone. clear_reservation(.topic, .id) end |
#dead_letter(message) ⇒ Object
174 175 176 177 178 179 |
# File 'lib/tina4/queue_backends/lite_backend.rb', line 174 def dead_letter() path = job_file(@dead_letter_dir, .id) data = .to_hash data[:status] = "dead" File.write(path, JSON.generate(data)) end |
#dead_letter_count(topic) ⇒ Object
Count dead-letter / failed messages for a topic.
204 205 206 207 208 209 210 211 212 213 214 215 |
# File 'lib/tina4/queue_backends/lite_backend.rb', line 204 def dead_letter_count(topic) return 0 unless Dir.exist?(@dead_letter_dir) count = 0 job_files(@dead_letter_dir).each do |file| data = JSON.parse(File.read(file)) count += 1 if data["topic"] == topic.to_s rescue JSON::ParserError next end count end |
#dead_letters(topic, max_retries: 3) ⇒ Object
Get dead letter jobs for a topic — messages that exceeded max retries. Returns Hashes (raw job data with status "dead").
226 227 228 229 230 231 232 233 234 235 236 237 238 239 240 241 242 243 |
# File 'lib/tina4/queue_backends/lite_backend.rb', line 226 def dead_letters(topic, max_retries: 3) return [] unless Dir.exist?(@dead_letter_dir) files = job_files(@dead_letter_dir).sort_by { |f| File.mtime(f) } jobs = [] files.each do |file| data = JSON.parse(File.read(file)) next unless data["topic"] == topic.to_s next if (data["attempts"] || 0) < max_retries data["status"] = "dead" jobs << data rescue JSON::ParserError next end jobs end |
#dequeue(topic) ⇒ Object
74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 |
# File 'lib/tina4/queue_backends/lite_backend.rb', line 74 def dequeue(topic) @mutex.synchronize do # First return any reservations whose consumer died mid-flight. reclaim_expired(topic) candidate = available_candidates(topic).first return nil unless candidate # Write the reservation BEFORE claiming the pending file, so a crash # between claim and reserve can never strand the job. Only the worker # that wins the delete owns — and returns — it. write_reserved(candidate[:data], topic) begin File.delete(candidate[:file]) rescue Errno::ENOENT return nil # already consumed by another worker end job_from_data(candidate[:data], topic) end end |
#dequeue_batch(topic, count) ⇒ Object
94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 |
# File 'lib/tina4/queue_backends/lite_backend.rb', line 94 def dequeue_batch(topic, count) @mutex.synchronize do reclaim_expired(topic) chosen = available_candidates(topic).first(count) chosen.filter_map do |c| write_reserved(c[:data], topic) begin File.delete(c[:file]) rescue Errno::ENOENT next end job_from_data(c[:data], topic) end end end |
#enqueue(message) ⇒ Object
65 66 67 68 69 70 71 72 |
# File 'lib/tina4/queue_backends/lite_backend.rb', line 65 def enqueue() @mutex.synchronize do topic_dir = topic_path(.topic) FileUtils.mkdir_p(topic_dir) path = job_file(topic_dir, .id) File.write(path, .to_json) end end |
#fail(job, error = "") ⇒ Object
Record a failed attempt. Increments attempts and stores the error. While attempts < max_retries the job is re-enqueued to pending (after the configured retry_backoff). Once attempts >= max_retries it is moved to the dead-letter store.
153 154 155 156 157 158 159 160 161 162 163 |
# File 'lib/tina4/queue_backends/lite_backend.rb', line 153 def fail(job, error = "") # Clear the reservation — the consumer acknowledged (with a failure). clear_reservation(job.topic, job.id) job.attempts += 1 job.error = error if job.attempts < @max_retries requeue_job(job, delay_seconds: @retry_backoff, error: error) else move_to_dead_letter(job, error) end end |
#failed(topic, max_retries: 3) ⇒ Object
Jobs that have failed at least once but are still being retried.
Under the auto-retry lifecycle a failed-but-retryable job lives in the pending queue (not the dead-letter dir), so this scans the topic dir for jobs with 0 < attempts < max_retries. Dead-lettered jobs are returned by dead_letters(). Returns Hashes.
335 336 337 338 339 340 341 342 343 344 345 346 347 348 |
# File 'lib/tina4/queue_backends/lite_backend.rb', line 335 def failed(topic, max_retries: 3) dir = topic_path(topic) return [] unless Dir.exist?(dir) jobs = [] job_files(dir).sort_by { |f| File.mtime(f) }.each do |file| data = JSON.parse(File.read(file)) attempts = data["attempts"] || 0 next unless attempts > 0 && attempts < max_retries jobs << data rescue JSON::ParserError next end jobs end |
#find_by_id(topic, id) ⇒ Object
Find a specific pending job by id, claim it (delete the file) and return it. Returns nil when no pending job with that id exists.
112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 |
# File 'lib/tina4/queue_backends/lite_backend.rb', line 112 def find_by_id(topic, id) @mutex.synchronize do dir = topic_path(topic) return nil unless Dir.exist?(dir) target = id.to_s job_files(dir).each do |f| data = JSON.parse(File.read(f)) next unless data["id"].to_s == target # Reserve (so a dead consumer's job is reclaimable) then claim the # pending file — mirrors dequeue. write_reserved(data, topic) File.delete(f) return job_from_data(data, topic) rescue JSON::ParserError next end nil end end |
#purge(topic, status) ⇒ Object
Delete messages by status. For "failed"/"dead"/"dead_letter", removes from the dead_letter directory. For "completed"/"pending", removes matching jobs from the topic (pending) directory. Returns count purged.
248 249 250 251 252 253 254 255 256 257 258 259 260 261 262 263 264 265 266 267 268 269 270 271 272 273 274 275 276 277 278 279 |
# File 'lib/tina4/queue_backends/lite_backend.rb', line 248 def purge(topic, status) count = 0 if dead_status?(status) return 0 unless Dir.exist?(@dead_letter_dir) job_files(@dead_letter_dir).each do |file| data = JSON.parse(File.read(file)) if data["topic"] == topic.to_s File.delete(file) count += 1 end rescue JSON::ParserError next end else dir = topic_path(topic) return 0 unless Dir.exist?(dir) job_files(dir).each do |file| data = JSON.parse(File.read(file)) if data["status"].to_s == status.to_s File.delete(file) count += 1 end rescue JSON::ParserError next end end count end |
#requeue(message) ⇒ Object
145 146 147 |
# File 'lib/tina4/queue_backends/lite_backend.rb', line 145 def requeue() enqueue() end |
#reserved_count(topic) ⇒ Object
Count currently-reserved (in-flight) jobs for a topic.
197 198 199 200 201 |
# File 'lib/tina4/queue_backends/lite_backend.rb', line 197 def reserved_count(topic) dir = reserved_path(topic) return 0 unless Dir.exist?(dir) job_files(dir).length end |
#retry(job, delay_seconds: 0) ⇒ Object
Explicit re-queue requested by the caller (job.retry()). Always re-enqueues regardless of the retry limit — a manual override, distinct from the automatic fail() path. Increments attempts, clears the error.
168 169 170 171 172 |
# File 'lib/tina4/queue_backends/lite_backend.rb', line 168 def retry(job, delay_seconds: 0) clear_reservation(job.topic, job.id) job.attempts += 1 requeue_job(job, delay_seconds: delay_seconds, error: nil) end |
#retry_failed(topic, max_retries: 3) ⇒ Object
Revive dead-letter jobs (under max_retries) back to pending. Returns the number of jobs re-queued.
283 284 285 286 287 288 289 290 291 292 293 294 295 296 297 298 299 300 301 302 303 304 305 306 307 308 309 310 311 |
# File 'lib/tina4/queue_backends/lite_backend.rb', line 283 def retry_failed(topic, max_retries: 3) return 0 unless Dir.exist?(@dead_letter_dir) dir = topic_path(topic) FileUtils.mkdir_p(dir) count = 0 job_files(@dead_letter_dir).each do |file| data = JSON.parse(File.read(file)) next unless data["topic"] == topic.to_s next if (data["attempts"] || 0) >= max_retries msg = Tina4::Job.new( topic: data["topic"], payload: data["payload"], id: data["id"], priority: data["priority"] || 0, attempts: data["attempts"] || 0, error: data["error"] ) enqueue(msg) File.delete(file) count += 1 rescue JSON::ParserError next end count end |
#retry_job(topic, job_id: nil, delay_seconds: 0) ⇒ Object
Revive a specific dead-letter job by id back to the pending queue. When job_id is nil, revives every dead-letter for the topic.
This is a manual override (Queue#retry(job_id)) — it always revives a dead-letter regardless of attempt count, mirroring job.retry. Returns true if any job was re-queued.
356 357 358 359 360 361 362 363 364 365 366 367 368 369 370 371 372 373 374 375 376 377 378 379 380 381 382 383 384 385 |
# File 'lib/tina4/queue_backends/lite_backend.rb', line 356 def retry_job(topic, job_id: nil, delay_seconds: 0) return false unless Dir.exist?(@dead_letter_dir) available_at = delay_seconds > 0 ? Time.now + delay_seconds : nil count = 0 job_files(@dead_letter_dir).each do |file| data = JSON.parse(File.read(file)) next unless data["topic"] == topic.to_s next if job_id && data["id"] != job_id.to_s msg = Tina4::Job.new( topic: data["topic"], payload: data["payload"], id: data["id"], priority: data["priority"] || 0, attempts: (data["attempts"] || 0) + 1, available_at: available_at, error: nil ) enqueue(msg) File.delete(file) count += 1 break if job_id # found the specific job, stop scanning rescue JSON::ParserError next end count > 0 end |
#size(topic) ⇒ Object
181 182 183 184 185 |
# File 'lib/tina4/queue_backends/lite_backend.rb', line 181 def size(topic) dir = topic_path(topic) return 0 unless Dir.exist?(dir) job_files(dir).length end |
#topics ⇒ Object
217 218 219 220 221 222 |
# File 'lib/tina4/queue_backends/lite_backend.rb', line 217 def topics return [] unless Dir.exist?(@dir) Dir.children(@dir) .reject { |d| d == Tina4::Queue::DEAD_LETTER_DIRNAME } .select { |d| File.directory?(File.join(@dir, d)) } end |