Class: Tina4::QueueBackends::MongoBackend
- Inherits:
-
Object
- Object
- Tina4::QueueBackends::MongoBackend
- Defined in:
- lib/tina4/queue_backends/mongo_backend.rb
Instance Attribute Summary collapse
-
#max_retries ⇒ Object
Reservation/visibility + retry policy (settable so a Queue can propagate its own onto a backend instance passed directly — legacy path).
-
#retry_backoff ⇒ Object
Reservation/visibility + retry policy (settable so a Queue can propagate its own onto a backend instance passed directly — legacy path).
-
#visibility_timeout ⇒ Object
Reservation/visibility + retry policy (settable so a Queue can propagate its own onto a backend instance passed directly — legacy path).
Instance Method Summary collapse
- #acknowledge(message) ⇒ Object
-
#clear(topic) ⇒ Object
Remove every pending job for a topic.
-
#close ⇒ Object
Close the MongoClient and release its connection pool.
-
#complete(message) ⇒ Object
Terminal success — the job is done and removed (mirrors the lite backend's complete()).
- #dead_letter(message) ⇒ Object
- #dead_letters(topic, max_retries: 3) ⇒ Object
- #dequeue(topic) ⇒ Object
- #enqueue(message) ⇒ Object
-
#fail(job, error = "") ⇒ Object
Record a failed attempt (mirrors the lite backend + the Python Mongo adapter).
-
#failed(topic, max_retries: 3) ⇒ Object
Jobs that failed but are still eligible for retry (under max_retries).
-
#find_by_id(topic, id) ⇒ Object
Claim ONE specific job by id, the same way dequeue claims the head.
-
#initialize(options = {}) ⇒ MongoBackend
constructor
A new instance of MongoBackend.
- #purge(topic, status) ⇒ Object
-
#reclaim_expired(topic, max_retries) ⇒ Object
Return reservations whose visibility window expired (at-least-once).
- #requeue(message) ⇒ Object
- #resolve_visibility_timeout(option) ⇒ Object
-
#retry(job, delay_seconds: 0) ⇒ Object
Explicit re-queue requested by the caller (job.retry()).
-
#retry_failed(topic, max_retries: 3) ⇒ Object
Re-queue failed-but-retryable jobs back to pending.
-
#retry_job(topic, job_id: nil, delay_seconds: 0) ⇒ Object
Move ONE dead-lettered job back to its main topic as pending.
- #size(topic) ⇒ Object
Constructor Details
#initialize(options = {}) ⇒ MongoBackend
Returns a new instance of MongoBackend.
10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 |
# File 'lib/tina4/queue_backends/mongo_backend.rb', line 10 def initialize( = {}) require "mongo" @max_retries = [:max_retries] || 3 # Seconds to delay a requeued (failed/retried) job before it is eligible # again. Default 0 = available on the very next dequeue, so a fail()'d job # retries immediately (matching the lite backend) instead of waiting out # the visibility window. @retry_backoff = ([:retry_backoff] || 0).to_f uri = [:uri] || ENV["TINA4_MONGO_URI"] host = [:host] || ENV.fetch("TINA4_MONGO_HOST", "localhost") port = ([:port] || ENV.fetch("TINA4_MONGO_PORT", 27017)).to_i username = [:username] || ENV["TINA4_MONGO_USERNAME"] password = [:password] || ENV["TINA4_MONGO_PASSWORD"] db_name = [:db] || ENV.fetch("TINA4_MONGO_DB", "tina4") @collection_name = [:collection] || ENV.fetch("TINA4_MONGO_COLLECTION", "tina4_queue") # Reservation/visibility timeout (seconds): a dequeued message is held # reserved (status "processing") with available_at = now + timeout; # reclaim_expired returns it once that passes (consumer died mid-flight). # <= 0 disables reclaim. @visibility_timeout = resolve_visibility_timeout([:visibility_timeout]) if uri # Honour the explicit db: / TINA4_MONGO_DB even when a uri is given. # Mongo::Client.new(uri) with no database path defaults to "admin", so # without passing :database the requested db_name was silently dropped # and every job/dead-letter landed in admin (a data-isolation footgun). # The explicit db_name always wins over the URI's (often absent) default. @client = Mongo::Client.new(uri, database: db_name) else = { database: db_name } [:user] = username if username [:password] = password if password @client = Mongo::Client.new(["#{host}:#{port}"], ) end @db = @client.database create_indexes rescue LoadError raise "MongoDB backend requires the 'mongo' gem. Install with: gem install mongo" end |
Instance Attribute Details
#max_retries ⇒ Object
Reservation/visibility + retry policy (settable so a Queue can propagate its own onto a backend instance passed directly — legacy path).
8 9 10 |
# File 'lib/tina4/queue_backends/mongo_backend.rb', line 8 def max_retries @max_retries end |
#retry_backoff ⇒ Object
Reservation/visibility + retry policy (settable so a Queue can propagate its own onto a backend instance passed directly — legacy path).
8 9 10 |
# File 'lib/tina4/queue_backends/mongo_backend.rb', line 8 def retry_backoff @retry_backoff end |
#visibility_timeout ⇒ Object
Reservation/visibility + retry policy (settable so a Queue can propagate its own onto a backend instance passed directly — legacy path).
8 9 10 |
# File 'lib/tina4/queue_backends/mongo_backend.rb', line 8 def visibility_timeout @visibility_timeout end |
Instance Method Details
#acknowledge(message) ⇒ Object
166 167 168 |
# File 'lib/tina4/queue_backends/mongo_backend.rb', line 166 def acknowledge() collection.delete_one(_id: .id) end |
#clear(topic) ⇒ Object
Remove every pending job for a topic. Queue#clear used to return 0 silently here, because this backend had no clear at all and the caller guarded on respond_to?(:clear) - so clearing a mongo-backed queue was a no-op that looked like success.
243 244 245 246 |
# File 'lib/tina4/queue_backends/mongo_backend.rb', line 243 def clear(topic) result = collection.delete_many(topic: topic, status: "pending") result.deleted_count end |
#close ⇒ Object
Close the MongoClient and release its connection pool.
IDEMPOTENT by construction: the handles are dropped in an ensure, so a second close finds nothing and returns. Before 3.13.95 the ivars were left set, so a shutdown path that ran twice (an explicit close plus an at_exit / ensure) closed an already-closed client.
352 353 354 355 356 357 |
# File 'lib/tina4/queue_backends/mongo_backend.rb', line 352 def close @client&.close ensure @client = nil @db = nil end |
#complete(message) ⇒ Object
Terminal success — the job is done and removed (mirrors the lite backend's complete()). Without this, job.complete() was a no-op on MongoDB and the document stayed "processing" forever.
173 174 175 |
# File 'lib/tina4/queue_backends/mongo_backend.rb', line 173 def complete() collection.delete_one(_id: .id) end |
#dead_letter(message) ⇒ Object
231 232 233 234 235 236 237 |
# File 'lib/tina4/queue_backends/mongo_backend.rb', line 231 def dead_letter() collection.find_one_and_update( { _id: .id }, { "$set" => { status: "dead", topic: "#{.topic}.dead_letter" } }, upsert: true ) end |
#dead_letters(topic, max_retries: 3) ⇒ Object
289 290 291 292 |
# File 'lib/tina4/queue_backends/mongo_backend.rb', line 289 def dead_letters(topic, max_retries: 3) collection.find(topic: "#{topic}.dead_letter", status: "dead") .map { |doc| job_from_doc(doc) } end |
#dequeue(topic) ⇒ Object
81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 |
# File 'lib/tina4/queue_backends/mongo_backend.rb', line 81 def dequeue(topic) # Reclaim any reservations whose consumer died before acking, then take # the next available message (at-least-once delivery). reclaim_expired(topic, @max_retries) now = Time.now.utc # The claim advances available_at to now + visibility_timeout and records # reserved_at so reclaim_expired can return the job if the consumer dies # before acknowledge/complete. This is the fix for the "reserved forever" # bug — previously available_at was left unchanged. # A job is claimable only once available_at has passed. The $or arm is not # optional: documents enqueued before available_at was written have no # such field, and a bare { "$lte" => now } would strand every one of them # in the collection forever. doc = collection.find_one_and_update( { topic: topic, status: "pending", "$or" => [ { available_at: nil }, { available_at: { "$exists" => false } }, { available_at: { "$lte" => now } } ] }, { "$set" => { status: "processing", reserved_at: now, available_at: now + (@visibility_timeout || 0) } }, # Highest priority first, ties broken oldest-first — the same ordering # policy LiteBackend applies. Sorting on created_at alone made the # backend pure FIFO, so an urgent job queued behind a backlog waited # for all of it. sort: { priority: -1, created_at: 1 }, return_document: :after ) return nil unless doc Tina4::Job.new( topic: doc["topic"], payload: doc["payload"], id: doc["_id"], priority: doc["priority"] || 0, attempts: doc["attempts"] || 0 ) end |
#enqueue(message) ⇒ Object
60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 |
# File 'lib/tina4/queue_backends/mongo_backend.rb', line 60 def enqueue() # Queue#push already resolved delay_seconds into an available_at (nil when # undelayed). Persisting it is what makes a delayed job invisible until # its time: dequeue filters on it below. Before 3.13.95 this field was # never written, so a delayed job on Mongo fired immediately while the # same code delayed correctly on the file backend. collection.insert_one( _id: .id, topic: .topic, payload: .payload, created_at: .created_at.utc, available_at: .available_at&.utc, attempts: .attempts, # Stored so the dequeue sort can order on it. It was read back on the # way out (doc["priority"] || 0) but never written, so every job # scored 0 and the queue was pure FIFO. priority: .priority, status: "pending" ) end |
#fail(job, error = "") ⇒ Object
Record a failed attempt (mirrors the lite backend + the Python Mongo adapter). Increments attempts; while attempts < max_retries the job is re-queued to pending (visible again immediately, or after retry_backoff), otherwise it is dead-lettered. The Queue/Job lifecycle expects fail() to route the requeue here — previously MongoBackend had no fail(), so job.fail() degraded to in-memory bookkeeping and never touched Mongo.
183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 |
# File 'lib/tina4/queue_backends/mongo_backend.rb', line 183 def fail(job, error = "") job.attempts += 1 if job.attempts >= @max_retries collection.find_one_and_update( { _id: job.id }, # attempts MUST be persisted here. fail() increments the in-memory # counter only; without writing it the dead-letter document keeps # the value it held BEFORE the final attempt, so dead_letters() # under-reported by one on every Mongo dead letter. { "$set" => { status: "dead", topic: "#{job.topic}.dead_letter", attempts: job.attempts, error: error, reserved_at: nil } }, upsert: true ) else requeue_with_error(job, error) end end |
#failed(topic, max_retries: 3) ⇒ Object
Jobs that failed but are still eligible for retry (under max_retries).
Found by the ATTEMPTS COUNTER, not by a "failed" status. fail() under max_retries re-queues the job as "pending" (that is what makes the next dequeue redeliver it), so no document ever carries status="failed" on the normal path. This method did not exist at all, and Queue#failed silently returned [] for it — indistinguishable from "nothing has failed" (ADR-0022 decision 7).
281 282 283 284 285 286 287 |
# File 'lib/tina4/queue_backends/mongo_backend.rb', line 281 def failed(topic, max_retries: 3) collection.find( topic: topic, status: { "$in" => %w[pending failed] }, attempts: { "$gt" => 0, "$lt" => max_retries } ).map { |doc| job_from_doc(doc) } end |
#find_by_id(topic, id) ⇒ Object
Claim ONE specific job by id, the same way dequeue claims the head. Queue#pop_by_id used to return nil silently for the same reason.
250 251 252 253 254 255 256 257 258 259 260 261 262 263 264 265 266 267 |
# File 'lib/tina4/queue_backends/mongo_backend.rb', line 250 def find_by_id(topic, id) now = Time.now.utc doc = collection.find_one_and_update( { _id: id, topic: topic, status: "pending" }, { "$set" => { status: "processing", reserved_at: now, available_at: now + (@visibility_timeout || 0) } }, return_document: :after ) return nil unless doc Tina4::Job.new( topic: doc["topic"], payload: doc["payload"], id: doc["_id"], priority: doc["priority"] || 0, attempts: doc["attempts"] || 0 ) end |
#purge(topic, status) ⇒ Object
294 295 296 297 |
# File 'lib/tina4/queue_backends/mongo_backend.rb', line 294 def purge(topic, status) result = collection.delete_many(topic: topic, status: status.to_s) result.deleted_count end |
#reclaim_expired(topic, max_retries) ⇒ Object
Return reservations whose visibility window expired (at-least-once).
A message left "processing" with available_at <= now had a consumer die before acknowledging. Each is atomically flipped back to "pending" with attempts incremented (so the next dequeue re-delivers it); once attempts
max_retries it is dead-lettered instead. Returns the number reclaimed.
Disabled when visibility_timeout <= 0.
131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 |
# File 'lib/tina4/queue_backends/mongo_backend.rb', line 131 def reclaim_expired(topic, max_retries) return 0 if @visibility_timeout.nil? || @visibility_timeout <= 0 reclaimed = 0 loop do now = Time.now.utc doc = collection.find_one_and_update( { topic: topic, status: "processing", available_at: { "$lte" => now } }, { "$set" => { status: "pending", available_at: now, reserved_at: nil }, "$inc" => { attempts: 1 } }, sort: { available_at: 1 }, return_document: :after ) break unless doc reclaimed += 1 next if (doc["attempts"] || 0) < max_retries # Out of retries — move it to the dead-letter queue and remove the # original so it is not re-delivered. collection.insert_one( _id: "#{doc["_id"]}.dead_letter", topic: "#{topic}.dead_letter", payload: doc["payload"], status: "dead", priority: doc["priority"] || 0, attempts: doc["attempts"] || 0, error: "reservation timed out - consumer did not acknowledge within the visibility timeout", created_at: Time.now.utc ) collection.delete_one(_id: doc["_id"], topic: topic) end reclaimed end |
#requeue(message) ⇒ Object
216 217 218 219 220 221 222 223 224 225 226 227 228 229 |
# File 'lib/tina4/queue_backends/mongo_backend.rb', line 216 def requeue() # Reset available_at so the requeued job is visible again right away (or # after retry_backoff) and clear reserved_at. dequeue() pushed # available_at out to the reservation expiry; leaving it there stranded a # requeued job for the full visibility window instead of retrying it on # the next pop(). collection.find_one_and_update( { _id: .id }, { "$set" => { status: "pending", reserved_at: nil, available_at: requeue_available_at(@retry_backoff) }, "$inc" => { attempts: 1 } }, upsert: true ) end |
#resolve_visibility_timeout(option) ⇒ Object
52 53 54 55 56 57 58 |
# File 'lib/tina4/queue_backends/mongo_backend.rb', line 52 def resolve_visibility_timeout(option) return option.to_f unless option.nil? Float(ENV.fetch("TINA4_QUEUE_VISIBILITY_TIMEOUT", "300")) rescue ArgumentError, TypeError 300.0 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.
205 206 207 208 209 210 211 212 213 214 |
# File 'lib/tina4/queue_backends/mongo_backend.rb', line 205 def retry(job, delay_seconds: 0) job.attempts += 1 backoff = delay_seconds.to_f > 0 ? delay_seconds.to_f : @retry_backoff collection.find_one_and_update( { _id: job.id }, { "$set" => { status: "pending", error: nil, reserved_at: nil, available_at: requeue_available_at(backoff) } }, upsert: true ) end |
#retry_failed(topic, max_retries: 3) ⇒ Object
Re-queue failed-but-retryable jobs back to pending. Returns the count.
Matched on the ATTEMPTS COUNTER for the same reason failed() is: the old query was status="failed", which requeue_with_error never writes, so this matched nothing and always returned 0 — silently reporting that there was nothing to retry.
305 306 307 308 309 310 311 312 313 314 315 316 317 318 319 320 321 322 323 324 325 326 327 |
# File 'lib/tina4/queue_backends/mongo_backend.rb', line 305 def retry_failed(topic, max_retries: 3) result = collection.update_many( { topic: topic, # BOTH statuses. "pending" is where the auto-retry fail() path # leaves a still-retryable job (it re-queues rather than marking # it failed), and matching only "failed" is why this used to # return 0 on every real failure. "failed" is still matched # because reject(requeue: false) writes it, so a job explicitly # parked as failed is retryable too. status: { "$in" => %w[pending failed] }, attempts: { "$gt" => 0, "$lt" => max_retries } }, # Reset available_at so re-queued failed jobs are visible again — they # were reserved with available_at in the future at dequeue. Clear # reserved_at too. (Same Bug B reason as requeue/fail.) # status MUST be set back to "pending". The match now also picks up # docs explicitly parked as "failed" (reject(requeue: false)), and # those have to be flipped back or retry_failed reports a re-queue it # never performed. A doc already pending is unaffected. { "$set" => { status: "pending", error: nil, reserved_at: nil, available_at: requeue_available_at(@retry_backoff) } } ) result.modified_count end |
#retry_job(topic, job_id: nil, delay_seconds: 0) ⇒ Object
Move ONE dead-lettered job back to its main topic as pending. Returns true when a job was revived, false when the id was not found.
331 332 333 334 335 336 337 338 339 340 341 342 343 344 |
# File 'lib/tina4/queue_backends/mongo_backend.rb', line 331 def retry_job(topic, job_id: nil, delay_seconds: 0) filter = { topic: "#{topic}.dead_letter", status: "dead" } filter[:_id] = job_id if job_id doc = collection.find(filter).first return false unless doc available = delay_seconds.to_f > 0 ? (Time.now.utc + delay_seconds.to_f) : Time.now.utc collection.update_one( { _id: doc["_id"] }, { "$set" => { topic: topic, status: "pending", error: nil, reserved_at: nil, available_at: available } } ) true end |
#size(topic) ⇒ Object
269 270 271 |
# File 'lib/tina4/queue_backends/mongo_backend.rb', line 269 def size(topic) collection.count_documents(topic: topic, status: "pending") end |