Class: Cosmo::Job::Processor
- Includes:
- Sentry::JobProcessorMiddleware
- Defined in:
- lib/cosmo/job/processor.rb,
sig/cosmo/job/processor.rbs
Constant Summary
Constants included from Sentry::JobProcessorMiddleware
Sentry::JobProcessorMiddleware::NAME_PREFIX, Sentry::JobProcessorMiddleware::OP_NAME, Sentry::JobProcessorMiddleware::SPAN_ORIGIN, Sentry::JobProcessorMiddleware::STATUS_FAIL, Sentry::JobProcessorMiddleware::STATUS_OK
Constants inherited from Processor
Processor::STREAMS_PAUSED_IDLE_SLEEP, Processor::STREAM_EMPTY_BACKOFF_MAX, Processor::STREAM_PAUSED_RECHECK_TTL
Instance Method Summary collapse
-
#acquire_concurrency_slot(worker_class, message, data) ⇒ ::String, false
Tries to acquire a concurrency slot for the job.
- #consumers ⇒ Array[untyped]
- #drop_message(message, data) ⇒ void
- #fetch_subjects(config) ⇒ Object
- #fetch_timeout(_config) ⇒ Float
-
#handle_failure(message, data) ⇒ Boolean
rubocop:disable Naming/PredicateMethod.
- #move_message(message, data = nil) ⇒ void
-
#process(messages, _) ⇒ void
rubocop:disable Metrics/AbcSize, Metrics/MethodLength, Metrics/CyclomaticComplexity, Metrics/PerceivedComplexity.
-
#schedule_loop ⇒ void
rubocop:disable Metrics/CyclomaticComplexity, Metrics/PerceivedComplexity, Metrics/MethodLength, Metrics/AbcSize.
- #scheduler? ⇒ Boolean
- #setup ⇒ void
- #subscribe(stream_name, config) ⇒ [untyped, Hash[Symbol, untyped], nil]
- #with_stats(message) { ... } ⇒ void
Methods inherited from Processor
#client, #consumer_state, #fetch, #initialize, #lock, run, #run, #run_loop, #running?, #stop, #stopwatch, #work_loop
Constructor Details
This class inherits a constructor from Cosmo::Processor
Instance Method Details
#acquire_concurrency_slot(worker_class, message, data) ⇒ ::String, false
Tries to acquire a concurrency slot for the job.
Returns the slot key (String) on success, or false if all slots are
taken (message is NAK'd with a delay equal to duration before returning).
111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 |
# File 'lib/cosmo/job/processor.rb', line 111 def acquire_concurrency_slot(worker_class, , data) = worker_class. key = worker_class.concurrency_key(data[:args]) slot = Limit.instance.acquire(key, jid: data[:jid], limit: [:limit], duration: [:duration]) return slot if slot .nak(delay: [:duration] * Config::NANO) Logger.debug "concurrency limit reached for #{data[:class]}, re-queueing back #{data[:jid]}" false rescue NATS::Error => e # Unexpected KV failure (e.g. transient NATS error). NAK immediately so # the message is retried rather than stuck in-flight until ack_wait expires. Logger.error e .nak false end |
#consumers ⇒ Array[untyped]
172 173 174 175 |
# File 'lib/cosmo/job/processor.rb', line 172 def consumers @weights ||= @consumers.filter_map { |(_, c, _)| [c[:stream]] * [c[:priority].to_i, 1].max }.flatten @weights.shuffle.map { |s| @consumers.find { |(_, c, _)| c[:stream] == s } } end |
#drop_message(message, data) ⇒ void
This method returns an undefined value.
155 156 157 158 |
# File 'lib/cosmo/job/processor.rb', line 155 def (, data) .term Logger.debug "job dropped #{data[:jid]}" end |
#fetch_subjects(config) ⇒ Object
177 178 179 |
# File 'lib/cosmo/job/processor.rb', line 177 def fetch_subjects(config) config[:subject] end |
#fetch_timeout(_config) ⇒ Float
181 182 183 |
# File 'lib/cosmo/job/processor.rb', line 181 def fetch_timeout(_config) ENV.fetch("COSMO_JOBS_FETCH_TIMEOUT", 0.1).to_f end |
#handle_failure(message, data) ⇒ Boolean
rubocop:disable Naming/PredicateMethod
129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 |
# File 'lib/cosmo/job/processor.rb', line 129 def handle_failure(, data) # rubocop:disable Naming/PredicateMethod current_attempt = ..num_delivered max_retries = data[:retry].to_i + 1 if current_attempt < max_retries # NATS will auto-retry with delay (exponential backoff based on current attempt). # When max_deliver is reached, NATS stops redelivering the message and marks it as "max deliveries exceeded". # The message is effectively abandoned by NATS — it stays in the stream (consuming a slot) but will never be delivered again to that consumer. delay_ns = ((current_attempt**4) + 15) * Config::NANO .nak(delay: delay_ns) return false end data[:dead] ? (, data) : (, data) true end |
#move_message(message, data = nil) ⇒ void
This method returns an undefined value.
160 161 162 163 164 165 166 |
# File 'lib/cosmo/job/processor.rb', line 160 def (, data = nil) klass = data ? Utils::String.underscore(data[:class]) : "default" headers = { "X-Stream" => ..stream, "X-Subject" => .subject } Client.instance.publish("jobs.dead.#{klass}", .data, header: headers) .ack Logger.debug "job moved #{data&.dig(:jid)} to DLQ" end |
#process(messages, _) ⇒ void
This method returns an undefined value.
rubocop:disable Metrics/AbcSize, Metrics/MethodLength, Metrics/CyclomaticComplexity, Metrics/PerceivedComplexity
54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 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 |
# File 'lib/cosmo/job/processor.rb', line 54 def process(, _) # rubocop:disable Metrics/AbcSize, Metrics/MethodLength, Metrics/CyclomaticComplexity, Metrics/PerceivedComplexity = .first Logger.debug "received messages #{.inspect}" data = Utils::Json.parse(.data) unless data Logger.error ArgumentError.new("malformed payload") () return end worker_class = Utils::String.safe_constantize(data[:class]) unless worker_class Logger.error ArgumentError.new("#{data[:class]} class not found") (, data) return end if worker_class.limits_concurrency? slot = acquire_concurrency_slot(worker_class, , data) return if slot == false end duration = worker_class.[:limit]&.dig(:duration)&.to_i with_stats() do sw = stopwatch Logger.with(jid: data[:jid]) Logger.info "start" instance = worker_class.new.tap { |w| w.jid = data[:jid] } perform_job(instance, data: data, message: , duration: duration) .ack Logger.with(elapsed: sw.elapsed_seconds) { Logger.info "done" } true rescue Timeout::Error Logger.with(elapsed: sw.elapsed_seconds) { Logger.info "fail[timeout]" } dropped = handle_failure(, data) false if dropped rescue StandardError => e Logger.debug e Logger.with(elapsed: sw.elapsed_seconds) { Logger.info "fail[error]" } dropped = handle_failure(, data) false if dropped rescue Exception # rubocop:disable Lint/RescueException Logger.with(elapsed: sw.elapsed_seconds) { Logger.info "fail[exception]" } raise end ensure Limit.instance.release(slot) if slot Logger.without(:jid) Logger.debug "processed message #{.inspect}" end |
#schedule_loop ⇒ void
This method returns an undefined value.
rubocop:disable Metrics/CyclomaticComplexity, Metrics/PerceivedComplexity, Metrics/MethodLength, Metrics/AbcSize
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 51 52 |
# File 'lib/cosmo/job/processor.rb', line 24 def schedule_loop # rubocop:disable Metrics/CyclomaticComplexity, Metrics/PerceivedComplexity, Metrics/MethodLength, Metrics/AbcSize config = Config.dig(:consumers, :jobs, :scheduled) return unless config subscription, = subscribe(:scheduled, config) while running? break unless running? now = Time.now.to_i timeout = ENV.fetch("COSMO_JOBS_SCHEDULER_FETCH_TIMEOUT", 5).to_f = fetch(subscription, batch_size: 100, timeout:) &.each do || headers = .header.except("X-Stream", "X-Subject", "X-Execute-At", "Nats-Expected-Stream") stream, subject, execute_at = .header.values_at("X-Stream", "X-Subject", "X-Execute-At") headers["Nats-Expected-Stream"] = stream execute_at = execute_at.to_i if now >= execute_at client.publish(subject, .data, headers: headers) .ack else delay_ns = (execute_at - now) * 1_000_000_000 .nak(delay: delay_ns) end end break unless running? end end |
#scheduler? ⇒ Boolean
168 169 170 |
# File 'lib/cosmo/job/processor.rb', line 168 def scheduler? true end |
#setup ⇒ void
This method returns an undefined value.
10 11 12 13 14 15 16 17 18 19 20 21 22 |
# File 'lib/cosmo/job/processor.rb', line 10 def setup # Initialize singletons before starting to process messages API::Busy.instance API::Counter.instance Limit.instance jobs_config = Config.dig(:consumers, :jobs) jobs_config&.each do |stream_name, config| next if stream_name == :scheduled # scheduled jobs are handled in schedule_loop @consumers << subscribe(stream_name, config) end end |
#subscribe(stream_name, config) ⇒ [untyped, Hash[Symbol, untyped], nil]
146 147 148 149 150 151 152 153 |
# File 'lib/cosmo/job/processor.rb', line 146 def subscribe(stream_name, config) config = config.dup config[:batch_size] = 1 config[:stream] = stream_name consumer_name = "consumer-#{stream_name}" subscription = client.subscribe(config[:subject], consumer_name, config.except(:subject, :priority, :stream, :batch_size)) [subscription, config, nil] end |