Class: Cosmo::API::Cron
- Inherits:
-
Object
- Object
- Cosmo::API::Cron
- Defined in:
- lib/cosmo/api/cron.rb,
lib/cosmo/api/cron/entry.rb,
sig/cosmo/api/cron.rbs,
sig/cosmo/api/cron/entry.rbs
Overview
Web-facing API for cron schedules. Single interface for all cron NATS operations.
Derives the schedule list entirely from NATS. Whatever is deployed in NATS is exactly what appears in the UI.
Schedule templates live in the same job stream they target (e.g. default),
stored at subjects matching cosmo.cron.<stream>.>. NATS 2.14 fires each
template by publishing the body to Nats-Schedule-Target as a regular
JetStream message that accumulates alongside pending jobs.
Defined Under Namespace
Classes: Entry
Class Method Summary collapse
Instance Method Summary collapse
-
#all ⇒ Array<Hash>
Every cron schedule currently deployed in NATS.
- #build_from_nats(stream_name, subject) ⇒ Hash[Symbol, untyped]?
-
#delete!(subject) ⇒ void
Purge the schedule message from NATS (stops future firings).
-
#name_from_subject(subject) ⇒ ::String?
"cosmo.cron.default.report_job" → nil "cosmo.cron.default.report_job.monthly" → "monthly".
-
#run_now!(schedule_subject) ⇒ void
Dispatch the job immediately to the target stream, bypassing the timer.
- #schedules_from_stream(stream_name) ⇒ ::Array[Hash[Symbol, untyped]]
-
#upsert!(class_name: nil, stream: nil, schedule: nil, args: [], timezone: nil, name: nil) ⇒ Hash?
Publish (or replace) a schedule message in NATS.
Class Method Details
.instance ⇒ Object
18 19 20 |
# File 'lib/cosmo/api/cron.rb', line 18 def self.instance @instance ||= new end |
Instance Method Details
#all ⇒ Array<Hash>
Returns every cron schedule currently deployed in NATS.
23 24 25 26 27 |
# File 'lib/cosmo/api/cron.rb', line 23 def all Stream.jobs.flat_map { |s| schedules_from_stream(s.name) } rescue StandardError [] end |
#build_from_nats(stream_name, subject) ⇒ Hash[Symbol, untyped]?
88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 |
# File 'lib/cosmo/api/cron.rb', line 88 def build_from_nats(stream_name, subject) msg = client.(stream_name, subject: subject) return unless msg headers = msg.headers || {} body = Utils::Json.parse(msg.data) || {} { class: body[:class], stream: stream_name, schedule: headers["Nats-Schedule"], timezone: headers["Nats-Schedule-Time-Zone"], args: body[:args] || [], name: name_from_subject(subject), schedule_subject: subject, target_subject: headers["Nats-Schedule-Target"], registry_key: subject.split(".").drop(2).join("/") } rescue StandardError nil end |
#delete!(subject) ⇒ void
This method returns an undefined value.
Purge the schedule message from NATS (stops future firings).
45 46 47 48 49 50 |
# File 'lib/cosmo/api/cron.rb', line 45 def delete!(subject) stream_name = subject.to_s.split(".")[2] client.purge(stream_name, subject) rescue NATS::JetStream::Error::NotFound, NATS::IO::Timeout nil end |
#name_from_subject(subject) ⇒ ::String?
"cosmo.cron.default.report_job" → nil "cosmo.cron.default.report_job.monthly" → "monthly"
112 113 114 115 |
# File 'lib/cosmo/api/cron.rb', line 112 def name_from_subject(subject) parts = subject.to_s.split(".") parts.length > 4 ? parts.drop(4).join(".") : nil end |
#run_now!(schedule_subject) ⇒ void
This method returns an undefined value.
Dispatch the job immediately to the target stream, bypassing the timer.
54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 |
# File 'lib/cosmo/api/cron.rb', line 54 def run_now!(schedule_subject) # rubocop:disable Metrics/AbcSize, Metrics/CyclomaticComplexity, Metrics/PerceivedComplexity stream_name = schedule_subject.to_s.split(".")[2] msg = client.(stream_name, subject: schedule_subject) return unless msg headers = msg.headers || {} body = Utils::Json.parse(msg.data) || {} target = headers["Nats-Schedule-Target"] return unless target && body[:class] payload = Utils::Json.dump({ jid: SecureRandom.hex(12), class: body[:class], args: body[:args] || [], retry: body[:retry] || Job::Data::DEFAULTS[:retry], dead: body[:dead].nil? ? Job::Data::DEFAULTS[:dead] : body[:dead] }) client.publish(target, payload, stream: stream_name) rescue NATS::JetStream::Error::NotFound nil end |
#schedules_from_stream(stream_name) ⇒ ::Array[Hash[Symbol, untyped]]
82 83 84 85 86 |
# File 'lib/cosmo/api/cron.rb', line 82 def schedules_from_stream(stream_name) filter = "#{Entry::SUBJECT_PREFIX}.#{stream_name}.>" subjects = client.cron_subjects_in_stream(stream_name, filter) subjects.filter_map { |subj| build_from_nats(stream_name, subj) } end |
#upsert!(class_name: nil, stream: nil, schedule: nil, args: [], timezone: nil, name: nil) ⇒ Hash?
Publish (or replace) a schedule message in NATS.
31 32 33 34 35 36 37 38 39 40 41 |
# File 'lib/cosmo/api/cron.rb', line 31 def upsert!(class_name: nil, stream: nil, schedule: nil, args: [], timezone: nil, name: nil) e = Entry.new(class_name: class_name, stream: stream, expression: schedule, args: args, timezone: timezone, name: name) headers = { "Nats-Schedule" => e.expression, "Nats-Schedule-Target" => e.target_subject } headers["Nats-Schedule-Time-Zone"] = e.timezone if e.timezone client.publish(e.schedule_subject, e.job_payload, stream: e.stream, header: headers) build_from_nats(e.stream, e.schedule_subject) end |