Class: Skiplock::Job
- Inherits:
-
ActiveRecord::Base
- Object
- ActiveRecord::Base
- Skiplock::Job
- Defined in:
- lib/skiplock/job.rb
Class Method Summary collapse
- .dispatch(purge_completion: true, max_retries: 20) ⇒ Object
- .enqueue(activejob) ⇒ Object
- .enqueue_at(activejob, timestamp) ⇒ Object
-
.flush ⇒ Object
resynchronize jobs that could not commit to database and reset any abandoned jobs for retry.
- .reset_retry_schedules ⇒ Object
Instance Method Summary collapse
Class Method Details
.dispatch(purge_completion: true, max_retries: 20) ⇒ Object
10 11 12 13 14 15 16 17 18 |
# File 'lib/skiplock/job.rb', line 10 def self.dispatch(purge_completion: true, max_retries: 20) job = nil self.connection.transaction do job = self.find_by_sql("SELECT id, scheduled_at FROM skiplock.jobs WHERE running = FALSE AND expired_at IS NULL AND finished_at IS NULL ORDER BY scheduled_at ASC NULLS FIRST, priority ASC NULLS LAST, created_at ASC FOR UPDATE SKIP LOCKED LIMIT 1").first return if job.nil? || job.scheduled_at.to_f > Time.now.to_f job = self.find_by_sql("UPDATE skiplock.jobs SET running = TRUE, updated_at = NOW() WHERE id = '#{job.id}' RETURNING *").first end self.dispatch(purge_completion: purge_completion, max_retries: max_retries) if job.execute(purge_completion: purge_completion, max_retries: max_retries) end |
.enqueue(activejob) ⇒ Object
20 21 22 |
# File 'lib/skiplock/job.rb', line 20 def self.enqueue(activejob) self.enqueue_at(activejob, nil) end |
.enqueue_at(activejob, timestamp) ⇒ Object
24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 |
# File 'lib/skiplock/job.rb', line 24 def self.enqueue_at(activejob, ) = activejob.instance_variable_get('@skiplock_options') || {} = Time.at() if if Thread.current[:skiplock_job].try(:id) == activejob.job_id Thread.current[:skiplock_job].activejob_error = [:error] Thread.current[:skiplock_job].data['activejob_error'] = true Thread.current[:skiplock_job].executions = activejob.executions Thread.current[:skiplock_job].exception_executions = activejob.exception_executions Thread.current[:skiplock_job].scheduled_at = Thread.current[:skiplock_job] else serialize = activejob.serialize self.create!(serialize.slice(*self.column_names).merge('id' => serialize['job_id'], 'data' => { 'arguments' => serialize['arguments'], 'options' => }, 'scheduled_at' => )) end end |
.flush ⇒ Object
resynchronize jobs that could not commit to database and reset any abandoned jobs for retry
41 42 43 44 45 46 47 48 49 50 51 52 53 |
# File 'lib/skiplock/job.rb', line 41 def self.flush Dir.mkdir('tmp/skiplock') unless Dir.exist?('tmp/skiplock') Dir.glob('tmp/skiplock/*').each do |f| disposed = true if self.exists?(id: File.basename(f), running: true) job = YAML.load_file(f) rescue nil disposed = job.dispose if job.is_a?(Skiplock::Job) end (File.delete(f) rescue nil) if disposed end self.where(running: true).where.not(worker_id: Worker.ids).update_all(running: false, worker_id: nil) true end |
.reset_retry_schedules ⇒ Object
55 56 57 |
# File 'lib/skiplock/job.rb', line 55 def self.reset_retry_schedules self.where('scheduled_at > NOW() AND executions > 0 AND expired_at IS NULL AND finished_at IS NULL').update_all(scheduled_at: nil, updated_at: Time.now) end |
Instance Method Details
#dispose ⇒ Object
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 |
# File 'lib/skiplock/job.rb', line 59 def dispose return unless self.max_retries yaml = self.to_yaml purging = false self.running = false self.worker_id = nil self.updated_at = Time.now > self.updated_at ? Time.now : self.updated_at + 1 # in case of clock drifting if self.exception self.exception_executions["[#{self.exception.class.name}]"] = self.exception_executions["[#{self.exception.class.name}]"].to_i + 1 unless self.data.key?('activejob_error') if (self.executions.to_i >= self.max_retries + 1) || self.data.key?('activejob_error') || self.exception.is_a?(Skiplock::Extension::ProxyError) self.expired_at = Time.now else self.scheduled_at = Time.now + (5 * 2**self.executions.to_i) end elsif self.finished_at if self.cron self.data['cron'] ||= {} self.data['cron']['executions'] = self.data['cron']['executions'].to_i + 1 self.data['cron']['last_finished_at'] = self.finished_at.utc.to_s self.data['cron']['last_result'] = self.data['result'] next_cron_at = Cron.next_schedule_at(self.cron) if next_cron_at # update job to record completions counter before resetting finished_at to nil self.update_columns(self.attributes.slice(*(self.changes.keys & self.class.column_names))) self.finished_at = nil self.executions = nil self.exception_executions = nil self.data.delete('result') self.scheduled_at = Time.at(next_cron_at) else Skiplock.logger.error("[Skiplock] ERROR: Invalid CRON '#{self.cron}' for Job #{self.job_class}") if Skiplock.logger purging = true end elsif self.purge == true purging = true end end purging ? self.delete : self.update_columns(self.attributes.slice(*(self.changes.keys & self.class.column_names))) rescue Exception => e File.write("tmp/skiplock/#{self.id}", yaml) rescue nil if Skiplock.logger Skiplock.logger.error(e.to_s) Skiplock.logger.error(e.backtrace.join("\n")) end Skiplock.on_errors.each { |p| p.call(e) } nil end |
#execute(purge_completion: true, max_retries: 20) ⇒ Object
107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 |
# File 'lib/skiplock/job.rb', line 107 def execute(purge_completion: true, max_retries: 20) raise 'Job has already been completed' if self.finished_at self.update_columns(running: true, updated_at: Time.now) unless self.running Skiplock.logger.info("[Skiplock] Performing #{self.job_class} (#{self.id}) from queue '#{self.queue_name || 'default'}'...") if Skiplock.logger self.data ||= {} self.data.delete('result') self.exception_executions ||= {} self.activejob_error = nil self.max_retries = (self.data['options'].key?('max_retries') ? self.data['options']['max_retries'].to_i : max_retries) rescue max_retries self.max_retries = 20 if self.max_retries < 0 || self.max_retries > 20 self.purge = (self.data['options'].key?('purge') ? self.data['options']['purge'] : purge_completion) rescue purge_completion job_data = self.attributes.slice('job_class', 'queue_name', 'locale', 'timezone', 'priority', 'executions', 'exception_executions').merge('job_id' => self.id, 'enqueued_at' => self.updated_at, 'arguments' => (self.data['arguments'] || [])) self.executions = self.executions.to_i + 1 Thread.current[:skiplock_job] = self start_time = Process.clock_gettime(Process::CLOCK_MONOTONIC) begin self.data['result'] = ActiveJob::Base.execute(job_data) rescue Exception => ex self.exception = ex Skiplock.on_errors.each { |p| p.call(ex) } end if Skiplock.logger if self.exception || self.activejob_error Skiplock.logger.error("[Skiplock] Job #{self.job_class} (#{self.id}) was interrupted by an exception#{ ' (rescued and retried by ActiveJob)' if self.activejob_error }") if self.exception Skiplock.logger.error(self.exception.to_s) Skiplock.logger.error(self.exception.backtrace.join("\n")) end else end_time = Process.clock_gettime(Process::CLOCK_MONOTONIC) job_name = self.job_class if self.job_class == 'Skiplock::Extension::ProxyJob' target, method_name = ::YAML.load(self.data['arguments'].first) job_name = "'#{target.name}.#{method_name}'" end Skiplock.logger.info "[Skiplock] Performed #{job_name} (#{self.id}) from queue '#{self.queue_name || 'default'}' in #{end_time - start_time} seconds" end end self.exception || self.activejob_error || self.data['result'] ensure self.finished_at ||= Time.now if self.data.key?('result') && !self.activejob_error self.dispose end |