🚀 Cosmonats
Background jobs + real-time event streaming for Ruby — unified, in one gem, backed by NATS.
No Redis. No DB polling. Disk-backed, horizontally scalable — no message is ever silently dropped.
⚡ Taste it
# Define a job with a familiar look
class SendEmailJob
include Cosmo::Job
stream: :default, retry: 3, dead: true
def perform(user_id, template)
EmailService.send(user_id, template)
end
end
# Enqueue it
SendEmailJob.perform_async(123, "welcome")
SendEmailJob.perform_in(1.day, 123, "followup")
# Process a continuous real-time event stream
class ClicksProcessor
include Cosmo::Stream
stream: :clickstream, batch_size: 100,
consumer: { subjects: ["events.clicks.>"] }
def process_one
Analytics.track(.data)
.ack
end
end
ClicksProcessor.publish({ user_id: 123, page: "/home" }, subject: "events.clicks.homepage")
bundle exec cosmo -C config/cosmo.yml -c 20 # Run jobs + streams with 20 threads
bundle exec cosmo -C config/cosmo.yml -c 20 jobs # Jobs only
bundle exec cosmo -C config/cosmo.yml -c 20 streams # Streams only

📖 Index
- Why?
- Features
- Installation
- Quick Start
- Core Concepts
- Advanced Usage
- CLI Reference
- Deployment
- Monitoring
- Examples
🎯 Why?
Most background job libraries use Redis or Postgres — tools that were never designed for this. Think of NATS as Redis — but Redis is KV first then messaging; NATS is messaging first, then KV. What NATS is:
- ~20 MB binary, ~10 MB RAM at idle Trivial to run anywhere.
- Disk-backed persistent streams Messages survive restarts, don't require RAM to fit.
- True horizontal clustering Lose a node — other nodes take over, zero message loss.
- Multilingual Official clients for Ruby, Go, Python, Rust, Java, .NET, and more. Any service can publish or consume.
One NATS server replaces your message broker, job queue, and KV store — with lower operational overhead.
| Redis/DB-backed | NATS/Cosmonats | |
|---|---|---|
| Persistence | In-memory / DB bloat | Disk-backed, TB-scale |
| Scaling | Sentinel only / Vertical only | True horizontal clustering |
| Background jobs | Yes | Yes |
| Real-time stream | No | Yes |
| Zero message loss | No | Yes |
| Message replay | No | Yes |
| Backpressure | No, grow unbounded | Yes |
| Multi-DC | Complex setup | Native geo-distribution |
Killer Features:
— Jobs + Streams, unified in one gem.
Most Ruby gems handle exactly that — background jobs. If you also need to consume a continuous event feed, that's a second system, second config, second set of
worker processes, second Dockerfile entry. Cosmonats is the only Ruby gem with a first-class Job primitive and a first-class Stream primitive, sharing
one server, one config, one CLI, one monitoring endpoint.
— Message replay and time-travel debugging.
NATS persists messages to disk and lets any consumer rewind to any point — beginning of time, a specific timestamp, or only new messages.
- Incident recovery — your pipeline crashed for 3 hours. Replay from the crash timestamp.
- New consumer bootstrap — a new service needs historical events. Start it from the beginning.
- Bug reproduction — replay the exact sequence of messages that caused a production issue.
— Multi-datacenter queues, natively.
NATS has a first-class cluster + leaf-node architecture for geo-distribution. Spanning multiple regions or datacenters is a config block — not a separate product or a third-party replication tool. NATS was built for edge computing, IoT, and satellite communication — multi-DC is a first-class concern, not an afterthought.
— Transport-level deduplication + built-in KV. No extra infrastructure.
NATS deduplicates messages at the broker — same-ID messages within the configured window are dropped before they ever reach a worker. No uniqueness gems, no advisory locks, no extra round-trips. It also ships a built-in Key/Value store usable for distributed locks and rate limiting — no Redis, no Memcached, nothing else to run.
✨ Features
🎪 Job Processing
- Familiar API —
perform_async,perform_in,perform_at - Priority queues — critical, high, default, low with weighted round-robin
- Scheduled jobs — execute at a specific time or after a delay
- Automatic retries — exponential backoff, configurable attempts
- Dead letter queue — capture permanently failed jobs
- Job uniqueness — prevent duplicate execution
- Concurrency limits — cap simultaneous executions per class or per key
- Cron scheduling — recurring jobs manageable live from the web UI
🌊 Stream Processing
- Real-time event streams — process continuous data feeds
- Batch processing — handle multiple messages in one go
- Message replay — reprocess from any point in time
- Consumer groups — load-balanced across workers
- Custom serialization — JSON, MessagePack, Protobuf
- Pause / resume — stop and restart a stream's processing without losing its position
📦 Installation
# Gemfile
gem "cosmonats"
Requirements: Ruby ≥ 3.1, NATS Server (install guide)
Spin up NATS instantly with Docker — one command, that's it:
docker run -p 4222:4222 -p 8222:8222 nats:alpine -js
Or add it to your existing docker-compose.yml:
services:
nats:
image: nats:alpine
command: -js
ports:
- "4222:4222"
- "8222:8222"
Mount the monitoring UI in your Rack app:
require "cosmo/web"
# Rails
mount Cosmo::Web => "/cosmo"
# Any Rack app (config.ru)
map "/cosmo" { run Cosmo::Web }
🚀 Quick Start
1. Create config/cosmo.yml
concurrency: 5 # Number of worker threads
consumers: # Declare consumer groups for streams, things that pull messages and process them
jobs: # Consumer configs for jobs (or streams)
default: # Stream name
ack_policy: explicit # Acknowledgment required for each message, can be explicit, none, or all
max_deliver: 10 # Max retry attempts before sending to a dead stream
max_ack_pending: 10 # Max messages waiting for ack, if exceeded, the server will stop delivering new messages until some are acked
ack_wait: 15 # Seconds to wait for ack before redelivering
subject: jobs.%{name}.> # Subject pattern for this consumer, %{name} replaced with stream name, becomes `jobs.default.>`
setup: # Initial stream creation only `cosmo -S`
jobs: # Stream configs for jobs (or streams)
default: # Stream name
storage: file # Storage type (file or memory)
retention: workqueue # Retention policy (limits, interest, workqueue). workqueue - deletes acked/nacked, limits - append only
subjects: ["jobs.%{name}.>"] # Subject pattern for this stream, %{name} replaced with stream name
allow_direct: true # Allow direct messages to stream (required for web UI)
2. Create streams in NATS (one-time), grabs config from setup section of config/cosmo.yml
bundle exec cosmo -S
3. Define a job in app/jobs/
class SendEmailJob
include Cosmo::Job
stream: :default, retry: 3, dead: true
def perform(user_id, email_type)
UserMailer.send(email_type, user_id).deliver_now
end
end
4. Enqueue & run
SendEmailJob.perform_async(42, "welcome")
bundle exec cosmo -C config/cosmo.yml -c 10 -r ./app/jobs jobs
💡 Core Concepts
Jobs
class ReportJob
include Cosmo::Job
(
stream: :critical, # Stream name
retry: 5, # Retry attempts
dead: true # Send to dead letter queue on final failure
)
def perform(report_id)
logger.info "Processing report #{report_id}"
Report.find(report_id).generate!
rescue StandardError => e
logger.error "Failed: #{e.}"
raise # Triggers retry with exponential backoff
end
end
ReportJob.perform_async(42) # Enqueue now
ReportJob.perform_in(30.minutes, 42) # Delayed
ReportJob.perform_at(Time.parse("2026-01-25 10:00"), 42) # Scheduled
ReportJob.perform_sync(42) # Inline, no NATS (great for tests)
Streams
class ClicksProcessor
include Cosmo::Stream
(
stream: :clickstream,
batch_size: 100,
start_position: :last, # :first, :last, :new, or timestamp
consumer: {
ack_policy: "explicit",
max_deliver: 3,
max_ack_pending: 100,
subjects: ["events.clicks.>"]
}
)
# Process one message at a time
def process_one
Analytics.track_click(.data)
.ack
end
# OR process a batch
def process()
Analytics.bulk_track(.map(&:data))
.each(&:ack)
end
end
# Publishing
ClicksProcessor.publish({ user_id: 123, page: "/home" }, subject: "events.clicks.homepage")
# Acknowledgment strategies
.ack # Success
.nack(delay: 5_000_000_000) # Retry in 5 seconds (nanoseconds)
.term # Permanent failure, no retry
Configuration
NATS subjects follow a dot-separated hierarchy (events.clicks.homepage).
The > wildcard matches everything after that prefix. Think of subjects as topic names — flexible routing with no extra configuration.
Full config/cosmo.yml example:
timeout: 25 # Shutdown timeout in seconds
concurrency: &concurrency 1 # Number of worker threads
max_retries: &max_retries 3 # Default max retries
stream_config: &stream_config
storage: file # storage type (file or memory)
retention: workqueue # retention policy (limits, interest, workqueue)
duplicate_window: 120 # time window for duplicate message detection in seconds
discard: old # discard new messages when stream is full (discard new or old)
allow_direct: true # allow direct messages to stream, required for web UI
subjects:
- jobs.%{name}.> # subject pattern for stream, %{name} will be replaced with stream name
consumer_config: &consumer_config
ack_policy: explicit # ack policy (explicit, none, all), each individual message must be acknowledged
max_deliver: 10 # maximum number of times a message will be delivered before it's considered failed
max_ack_pending: 20 # maximum number of messages with pending ack for this consumer
ack_wait: 60 # time in seconds to wait for an ack before redelivering the message
subject: jobs.%{name}.> # subject pattern for consumer, %{name} will be replaced with stream name
consumers:
jobs:
critical:
<<: *consumer_config
priority: 50
high:
<<: *consumer_config
priority: 30
default:
<<: *consumer_config
priority: 15
low:
<<: *consumer_config
priority: 5
scheduled:
<<: *consumer_config
max_deliver: 1
max_ack_pending: 100
ack_wait: 10
setup:
jobs:
critical:
<<: *stream_config
description: Very critical priority jobs
high:
<<: *stream_config
description: Higher priority jobs
default:
<<: *stream_config
description: Default priority jobs
low:
<<: *stream_config
description: Lower priority jobs
scheduled:
<<: *stream_config
description: Scheduled jobs
dead:
<<: *stream_config
retention: limits
max_msgs: 10000
max_age: 604800 # 7d
description: Broken jobs (DLQ)
development:
verbose: false
concurrency: *concurrency
staging:
verbose: true
concurrency: 3
production:
concurrency: 3
Programmatic:
Cosmo::Config.set(:concurrency, 20)
Cosmo::Config.set(:setup, :streams, :custom, { storage: "file", subjects: ["custom.>"] })
Environment variables:
export NATS_URL=nats://localhost:4222
export COSMO_JOBS_FETCH_TIMEOUT=0.1
export COSMO_STREAMS_FETCH_TIMEOUT=0.1
🔧 Advanced Usage
Cron
Recurring jobs, without a separate scheduler process. A schedule is just a message parked in the job's own NATS stream (requires NATS Server 2.14+) — NATS fires it on the cron expression, and it lands back in the stream as a regular job. Deploy it once; whatever's in NATS is exactly what runs and exactly what shows up in the web UI's Crons tab, where each entry can be inspected, run immediately, or deleted.
Declare schedules right in config/cosmo.yml:
setup:
cron:
daily_report:
class: ReportJob
schedule: "@daily" # @-shortcuts are passed straight through to NATS
stream: default
weekday_digest:
class: ReportJob
schedule: "0 9 * * 1-5" # 6-field NATS cron (seconds first); 5-field UNIX cron is auto-normalized
stream: default
args: ["daily"]
timezone: America/New_York # optional, cron expressions only
cosmo -C config/cosmo.yml -S syncs it — whatever's in the file is exactly what ends up scheduled in NATS, same as streams.
Prefer to manage schedules at runtime instead? The same operations are available from Ruby:
Cosmo::API::Cron.instance.upsert!(
class_name: "ReportJob", stream: "default", schedule: "0 9 * * 1-5",
args: ["daily"], timezone: "America/New_York", name: "weekday_report"
)
Cosmo::API::Cron.instance.all # every schedule currently deployed
Cosmo::API::Cron.instance.run_now!("cosmo.cron.default.report_job.weekday_report") # bypass the timer
Cosmo::API::Cron.instance.delete!("cosmo.cron.default.report_job.weekday_report") # stop future firings
Priority Queues
class UrgentJob
include Cosmo::Job
stream: :critical # priority: 50 in config — polled most frequently
end
Concurrency Limiting
class ThirdPartyApiJob
include Cosmo::Job
# At most 3 instances of this job run at once, cluster-wide.
# Jobs that lose the race are NAK'd with a delay equal to `duration`
# so they aren't redelivered until a slot is guaranteed free.
limit: { duration: 30, concurrency: 3 }
end
class PerAccountSyncJob
include Cosmo::Job
# Scope the cap per key instead of class-wide — e.g. one concurrent sync per account.
limit: { duration: 30, concurrency: { to: 1, key: ->(account_id) { account_id } } }
def perform(account_id)
Account.find(account_id).sync!
end
end
Custom Serializers
module MessagePackSerializer
def self.serialize(data) = MessagePack.pack(data)
def self.deserialize(payload) = MessagePack.unpack(payload)
end
class FastStream
include Cosmo::Stream
publisher: { serializer: MessagePackSerializer }
end
Error Handling
class ResilientJob
include Cosmo::Job
retry: 5, dead: true
def perform(data)
process_data(data)
rescue RetryableError => e
logger.warn "Retryable: #{e.}"
raise # Will retry with exponential backoff
rescue FatalError => e
logger.error "Fatal: #{e.}"
# Don't raise — won't retry, won't go to DLQ
end
end
Testing
# Synchronous — no NATS needed
SendEmailJob.perform_sync(123, "test")
# Async — returns a job ID
jid = SendEmailJob.perform_async(123, "welcome")
assert_kind_of String, jid
Integrations
ActiveJob:
# config/application.rb
config.active_job.queue_adapter = :cosmonats
The ActiveJob queue name maps directly to a Cosmo stream. Use cosmo_options for anything
Cosmo-specific — retries, DLQ behavior, or overriding the target stream:
class ReportJob < ApplicationJob
retry: 5, dead: false, stream: :critical
def perform(report_id)
Report.find(report_id).generate!
end
end
Inside a Rails app this is wired up automatically by the bundled Railtie — it registers the
adapter and loads config/cosmo.yml if present. Outside Rails:
require "cosmo/active_job"
ActiveJob::Base.queue_adapter = Cosmo::ActiveJobAdapter::Adapter.new
Sentry:
require "cosmo/sentry/auto"
Wraps every job execution in a Sentry transaction (queue.cosmonats) and captures unhandled
exceptions with the job's id, stream, subject, and retry count attached as context — no other
setup beyond having sentry-ruby initialized.
🖥️ CLI Reference
cosmo -C config/cosmo.yml --setup # Create streams in NATS (idempotent)
cosmo -C config/cosmo.yml -c 20 -r ./app/jobs jobs # Jobs only
cosmo -C config/cosmo.yml -c 20 streams # Streams only
cosmo -C config/cosmo.yml -c 20 # Both
| Flag | Description | Example |
|---|---|---|
-C, --config PATH |
Config file path | -C config/cosmo.yml |
-c, --concurrency INT |
Worker threads | -c 20 |
-r, --require PATH |
Auto-require directory | -r ./app/jobs |
-t, --timeout NUM |
Shutdown timeout (sec) | -t 60 |
-S, --setup |
Setup streams & exit | --setup |
🚢 Deployment
NATS Cluster config:
# nats-server.conf
port: 4222
jetstream {
store_dir: /var/lib/nats
max_file: 10G
}
cluster {
name: cosmo-cluster
listen: 0.0.0.0:6222
routes: [nats://nats-2:6222, nats://nats-3:6222]
}
Docker Compose:
services:
nats:
image: nats:latest
command: -js -c /etc/nats/nats-server.conf
volumes:
- ./nats.conf:/etc/nats/nats-server.conf
- nats-data:/var/lib/nats
worker:
build: .
environment:
NATS_URL: nats://nats:4222
command: bundle exec cosmo -C config/cosmo.yml -c 20 jobs
deploy:
replicas: 3
Systemd Service:
# /etc/systemd/system/cosmo.service
[Unit]
Description=Cosmo Background Processor
After=network.target
[Service]
Type=simple
User=deploy
WorkingDirectory=/var/www/myapp
Environment=RAILS_ENV=production
Environment=NATS_URL=nats://localhost:4222
ExecStart=/usr/local/bin/bundle exec cosmo -C config/cosmo.yml -c 20 jobs
Restart=always
RestartSec=10
StandardOutput=syslog
StandardError=syslog
SyslogIdentifier=cosmo
[Install]
WantedBy=multi-user.target
sudo systemctl enable cosmo && sudo systemctl start cosmo
📊 Monitoring
Web UI — mount Cosmo::Web (see Installation) for a live, htmx-powered dashboard:
- Jobs — enqueued, scheduled, busy, and dead views, with per-job retry and delete
- Streams — per-stream state (messages, bytes, consumers) with pause/resume
- Crons — every schedule deployed in NATS, with run-now and delete
- Summary counters (processed / failed / busy / enqueued / retries / scheduled / dead) backed by a NATS KV counter, no separate metrics store needed
Structured logs:
2026-01-23T10:15:30.123Z INFO pid=12345 tid=abc jid=def: start
2026-01-23T10:15:32.456Z INFO pid=12345 tid=abc jid=def elapsed=2.333: done
Stream Metrics:
client = Cosmo::Client.instance
info = client.stream_info("default")
info.state. # Total messages
info.state.bytes # Total bytes
info.state.consumer_count # Number of consumers
Prometheus — NATS exposes metrics at :8222/metrics:
jetstream_server_store_msgs— Messages in streamjetstream_consumer_delivered_msgs— Delivered messagesjetstream_consumer_ack_pending— Pending acknowledgments
💼 Examples
Email queue with scheduling:
class EmailJob
include Cosmo::Job
stream: :default, retry: 3
def perform(user_id, template)
user = User.find(user_id)
EmailService.send(user.email, template)
end
end
EmailJob.perform_async(123, "welcome")
EmailJob.perform_in(1.day, 123, "followup")
Image Processing Pipeline:
class ImageProcessor
include Cosmo::Stream
(
stream: :images,
consumer: { subjects: ["images.uploaded.>"] }
)
def process_one
processed = ImageService.process(.data["url"])
publish(processed, subject: "images.processed.optimized")
.ack
rescue => e
logger.error "Processing failed: #{e.}"
.nack(delay: 30_000_000_000) # retry in 30s
end
end
ImageProcessor.publish({ url: "https://example.com/image.jpg" }, subject: "images.uploaded.user")
Real-Time Analytics:
class AnalyticsAggregator
include Cosmo::Stream
batch_size: 1000, consumer: { subjects: ["events.*.>"] }
def process()
aggregates = .map(&:data).group_by { |e| e["type"] }.transform_values(&:count)
Analytics.bulk_insert(aggregates)
.each(&:ack)
end
end
Made with ❤️ for Ruby
Blast off Cosmonats! 🚀