Work that outlives the request.
Define typed jobs, run workers and schedules, publish notifications, and cache strings in your application’s PostgreSQL database.
Define a job
A job is a struct that extends Caramel::ColdBrew::Job and implements perform. Each param declares a typed field and its getter. The params form the JSON payload stored with the job. A default value makes the param optional when enqueueing.
struct SendInvitation < Caramel::ColdBrew::Job
queue "mailers"
retry_on IO::TimeoutError, attempts: 5, backoff: :exponential, base: 2.seconds
param invite_id : Int64
param note : String = ""
def perform
# Load the invite by invite_id and send it.
end
endqueue takes a literal name of 1 to 63 characters from a-z, digits, _, ., : and -. A job without one uses default. Declaring the queue twice, declaring a param twice, or naming a param run_at or priority fails compilation. Param types must round-trip through JSON. Keep payloads small: store an id and load the record in perform.
Retry failures
When perform raises, its transaction rolls back and Cold Brew chooses a retry rule. retry_on takes an exception class, attempts, backoff and base. attempts counts every run, the first included. After the last run, the job is failed.
:exponentialwaitsbase * 2**(attempt - 1).:linearwaitsbase * attempt.
Rules on the job come first, then rules on its abstract parents, then global rules. The first matching rule wins. With no match, a job runs 3 times with exponential backoff from 1 second. A job whose class is not compiled into the binary fails at once.
Caramel::ColdBrew::Job.retry_on Socket::ConnectError,
attempts: 6, backoff: :linear, base: 3.secondsA failed row keeps the error’s class, message and first backtrace lines in last_error, up to 4,000 characters. Status and hooks expose only the class.
Enqueue a job
enqueue takes one keyword per param, plus run_at and priority. It returns the job’s Int64 id. An unknown keyword or a mistyped value fails compilation at the call.
SugarORM::Repo.transaction do
invite = Invite.create!(email: "elena@acme.com")
SendInvitation.enqueue(invite_id: invite.id)
end
id = SendInvitation.enqueue(invite_id: 42, run_at: 10.minutes.from_now, priority: 5)
SendInvitation.enqueue(db, invite_id: 42, note: "resend")The row is written through SugarORM::Repo’s current connection. Inside SugarORM::Repo.transaction, the job commits or rolls back with the business write. run_at defaults to now, and a job runs once it is due. Higher priority runs first; the default is 0, and equal priorities run in id order. The overload that takes a SugarORM::Handle first writes on that handle, such as a spec’s db.
Run workers
The application binary’s serve command starts Cold Brew beside the HTTP server. Its work command starts the same workers, maintenance and schedules without HTTP. Both read queue settings from the environment, and work’s flags replace them.
CARAMEL_WORKER_QUEUESComma-separated queue names, each once; defaults to defaultCARAMEL_WORKER_CONCURRENCYFibers per queue; defaults to 4--queues=NAMESwork only: replaces CARAMEL_WORKER_QUEUES--concurrency=Nwork only: replaces CARAMEL_WORKER_CONCURRENCY--no-schedulerwork only: leaves every schedule to another process./bin/bookshelf work --queues=default,mailers --concurrency=4
./bin/bookshelf work --queues=default --no-schedulerWorkers use a connection pool of their own, so jobs never wait for request connections. The pool holds one connection per fiber plus one, up to 32. Concurrency is therefore limited to 31 divided by the number of queues. Invalid settings stop the process before it connects.
On SIGTERM or SIGINT, workers stop claiming jobs, and in-flight jobs and schedule blocks finish before the process exits. work refuses pending migrations and CARAMEL_ENV=test. Under test, serve starts no workers.
Maintenance runs every 60 seconds in every process, even with --no-scheduler. It releases locks older than 15 minutes whose database session has ended, so those jobs run again. It deletes expired cache entries and creates daily job partitions from today through seven days ahead, in UTC. It drops partitions that ended more than 7 days ago and hold only finished or failed jobs; status then returns nil for those jobs.
Schedule recurring work
Caramel::ColdBrew.every(span, name) declares a block that runs once per period. Declare schedules at boot, before serve or work starts.
Caramel::ColdBrew.every(1.hour, "nightly-cleanup") { CleanupJob.enqueue }Every process ticks each schedule, but a database lease lets one process run the block per period. The block runs inside the lease’s transaction, so its writes commit with the lease. The span must be at least 1 second, and the name must be present and unique; otherwise every raises ArgumentError. A raising block is logged and tried again after its span or one minute, whichever is shorter.
Publish and subscribe
publish(channel, payload) sends a string with PostgreSQL pg_notify on the current Repo connection. Inside a transaction, the notification is delivered only when the transaction commits.
SugarORM::Repo.transaction do
Caramel::ColdBrew.publish("board_#{board_id}", board_id.to_s)
endChannel names use the queue-name rules: 1 to 63 characters from a-z, digits, _, ., : and -. Payloads hold at most 8,000 bytes, so publish an id and read the record instead. Breaking either limit raises ArgumentError.
subscribe has two forms. The block form creates a Channel(String) and unsubscribes when the block ends. subscribe(channel, subscriber) delivers into your channel until unsubscribe, a closed channel, or the end of the fiber that subscribed. A streaming action’s subscription therefore ends with its stream.
stream "text/event-stream" do |io|
Caramel::ColdBrew.subscribe("board_#{id}") do |updates|
loop { Caramel::SSE.write(io, event: "BoardUpdated", data: updates.receive) }
end
endSubscribing needs the process’s broker, which serve and work start. Elsewhere, set Caramel::ColdBrew.broker = Caramel::ColdBrew::Broker.new(database_url) first. Delivery is at most once. Each subscriber receives payloads in publish order. A payload it does not receive within 1 second is dropped, as are payloads beyond 10,000 waiting and notifications sent while the broker reconnects.
Cache strings
Caramel::Cache stores string values in the unlogged caramel_cache table, through the current Repo connection.
report = Caramel::Cache.fetch("report:2026-10", expires_in: 10.minutes) do
"computed report"
end
Caramel::Cache.write("banner", "Maintenance at noon")
Caramel::Cache.read("banner") # => "Maintenance at noon"
Caramel::Cache.delete("banner") # => truewrite(key, value, expires_in)Inserts or replaces; nil expires_in keeps the entryread(key)The value, or nil when missing or expiredfetch(key, expires_in) { … }The cached value, else the block’s value, writtendelete(key)True when an entry was removedclearDeletes every entryvacuumDeletes expired entries and returns the countA zero or negative expires_in raises ArgumentError. Maintenance vacuums expired entries every minute. Keep only values you can compute again.
Drain queues in specs
Specs run no workers. Caramel::ColdBrew.drain_queue!(db, queue) runs the queue’s due jobs synchronously on db, including jobs that those jobs enqueue. It returns the number of runs.
Caramel::ColdBrew.drain_queue!(db, "mailers")
Caramel::ColdBrew.status(db, id).not_nil!.state
.should eq(Caramel::ColdBrew::JobState::Finished)On a spec’s transaction, each job runs in a savepoint. A failed job follows its retry rule and does not run again in the same drain. drain_queue! then raises DrainFailure, listing every failed run; drain_queue does not raise. Pass include_scheduled: true to run jobs that are not yet due. A drain that has run 1,000 jobs while more are still due raises DrainLimitExceeded.