Class: Cosmo::ActiveJobAdapter::Adapter

Inherits:
Object
  • Object
show all
Defined in:
lib/cosmo/active_job/adapter.rb,
sig/cosmo/active_job/adapter.rbs

Overview

ActiveJob queue adapter that enqueues jobs via NATS JetStream.

Usage:

config.active_job.queue_adapter = :cosmonats
# or explicitly:
config.active_job.queue_adapter = Cosmo::ActiveJobAdapter::Adapter.new

The ActiveJob queue name maps directly to the Cosmo stream name.

Instance Method Summary collapse

Instance Method Details

#enqueue(job) ⇒ String

Enqueue a job to be run as soon as possible.

Parameters:

  • job (ActiveJob::Base)

Returns:

  • (String)


16
17
18
# File 'lib/cosmo/active_job/adapter.rb', line 16

def enqueue(job)
  publish(job, nil)
end

#enqueue_at(job, timestamp) ⇒ String

Enqueue a job to be run at (or after) a given time.

Parameters:

  • job (ActiveJob::Base)
  • timestamp (Numeric)

    Unix timestamp (seconds, float)

Returns:

  • (String)


23
24
25
# File 'lib/cosmo/active_job/adapter.rb', line 23

def enqueue_at(job, timestamp)
  publish(job, timestamp)
end

#job_cosmo_options(job) ⇒ Hash[Symbol, untyped]

Returns Cosmo-specific options declared on the job class via cosmo_options, falling back to an empty hash.

Parameters:

  • job (Object)

Returns:

  • (Hash[Symbol, untyped])


41
42
43
# File 'lib/cosmo/active_job/adapter.rb', line 41

def job_cosmo_options(job)
  job.class.respond_to?(:get_cosmo_options) ? job.class.get_cosmo_options : {}
end

#publish(job, timestamp) ⇒ String

Parameters:

  • job (Object)
  • timestamp (Numeric, nil)

Returns:

  • (String)


29
30
31
32
33
34
35
36
37
# File 'lib/cosmo/active_job/adapter.rb', line 29

def publish(job, timestamp)
  cosmo_opts = job_cosmo_options(job)
  stream     = cosmo_opts.delete(:stream) || job.queue_name.to_sym
  options    = { stream: stream }.merge(cosmo_opts)
  options[:at] = timestamp if timestamp

  data = Job::Data.new(Executor.name, [job.serialize], options)
  Publisher.publish_job(data)
end