File: worker.rb

package info (click to toggle)
ruby-sidekiq 6.4.1%2Bdfsg-1
  • links: PTS, VCS
  • area: main
  • in suites: bookworm
  • size: 792 kB
  • sloc: ruby: 4,582; makefile: 20; sh: 6
file content (362 lines) | stat: -rw-r--r-- 11,851 bytes parent folder | download
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
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
53
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
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
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
# frozen_string_literal: true

require "sidekiq/client"

module Sidekiq
  ##
  # Include this module in your worker class and you can easily create
  # asynchronous jobs:
  #
  #   class HardWorker
  #     include Sidekiq::Worker
  #     sidekiq_options queue: 'critical', retry: 5
  #
  #     def perform(*args)
  #       # do some work
  #     end
  #   end
  #
  # Then in your Rails app, you can do this:
  #
  #   HardWorker.perform_async(1, 2, 3)
  #
  # Note that perform_async is a class method, perform is an instance method.
  #
  # Sidekiq::Worker also includes several APIs to provide compatibility with
  # ActiveJob.
  #
  #   class SomeWorker
  #     include Sidekiq::Worker
  #     queue_as :critical
  #
  #     def perform(...)
  #     end
  #   end
  #
  #   SomeWorker.set(wait_until: 1.hour).perform_async(123)
  #
  # Note that arguments passed to the job must still obey Sidekiq's
  # best practice for simple, JSON-native data types. Sidekiq will not
  # implement ActiveJob's more complex argument serialization. For
  # this reason, we don't implement `perform_later` as our call semantics
  # are very different.
  #
  module Worker
    ##
    # The Options module is extracted so we can include it in ActiveJob::Base
    # and allow native AJs to configure Sidekiq features/internals.
    module Options
      def self.included(base)
        base.extend(ClassMethods)
        base.sidekiq_class_attribute :sidekiq_options_hash
        base.sidekiq_class_attribute :sidekiq_retry_in_block
        base.sidekiq_class_attribute :sidekiq_retries_exhausted_block
      end

      module ClassMethods
        ACCESSOR_MUTEX = Mutex.new

        ##
        # Allows customization for this type of Worker.
        # Legal options:
        #
        #   queue - name of queue to use for this job type, default *default*
        #   retry - enable retries for this Worker in case of error during execution,
        #      *true* to use the default or *Integer* count
        #   backtrace - whether to save any error backtrace in the retry payload to display in web UI,
        #      can be true, false or an integer number of lines to save, default *false*
        #
        # In practice, any option is allowed.  This is the main mechanism to configure the
        # options for a specific job.
        def sidekiq_options(opts = {})
          opts = opts.transform_keys(&:to_s) # stringify
          self.sidekiq_options_hash = get_sidekiq_options.merge(opts)
        end

        def sidekiq_retry_in(&block)
          self.sidekiq_retry_in_block = block
        end

        def sidekiq_retries_exhausted(&block)
          self.sidekiq_retries_exhausted_block = block
        end

        def get_sidekiq_options # :nodoc:
          self.sidekiq_options_hash ||= Sidekiq.default_worker_options
        end

        def sidekiq_class_attribute(*attrs)
          instance_reader = true
          instance_writer = true

          attrs.each do |name|
            synchronized_getter = "__synchronized_#{name}"

            singleton_class.instance_eval do
              undef_method(name) if method_defined?(name) || private_method_defined?(name)
            end

            define_singleton_method(synchronized_getter) { nil }
            singleton_class.class_eval do
              private(synchronized_getter)
            end

            define_singleton_method(name) { ACCESSOR_MUTEX.synchronize { send synchronized_getter } }

            ivar = "@#{name}"

            singleton_class.instance_eval do
              m = "#{name}="
              undef_method(m) if method_defined?(m) || private_method_defined?(m)
            end
            define_singleton_method("#{name}=") do |val|
              singleton_class.class_eval do
                ACCESSOR_MUTEX.synchronize do
                  undef_method(synchronized_getter) if method_defined?(synchronized_getter) || private_method_defined?(synchronized_getter)
                  define_method(synchronized_getter) { val }
                end
              end

              if singleton_class?
                class_eval do
                  undef_method(name) if method_defined?(name) || private_method_defined?(name)
                  define_method(name) do
                    if instance_variable_defined? ivar
                      instance_variable_get ivar
                    else
                      singleton_class.send name
                    end
                  end
                end
              end
              val
            end

            if instance_reader
              undef_method(name) if method_defined?(name) || private_method_defined?(name)
              define_method(name) do
                if instance_variable_defined?(ivar)
                  instance_variable_get ivar
                else
                  self.class.public_send name
                end
              end
            end

            if instance_writer
              m = "#{name}="
              undef_method(m) if method_defined?(m) || private_method_defined?(m)
              attr_writer name
            end
          end
        end
      end
    end

    attr_accessor :jid

    def self.included(base)
      raise ArgumentError, "Sidekiq::Worker cannot be included in an ActiveJob: #{base.name}" if base.ancestors.any? { |c| c.name == "ActiveJob::Base" }

      base.include(Options)
      base.extend(ClassMethods)
    end

    def logger
      Sidekiq.logger
    end

    # This helper class encapsulates the set options for `set`, e.g.
    #
    #     SomeWorker.set(queue: 'foo').perform_async(....)
    #
    class Setter
      include Sidekiq::JobUtil

      def initialize(klass, opts)
        @klass = klass
        @opts = opts

        # ActiveJob compatibility
        interval = @opts.delete(:wait_until) || @opts.delete(:wait)
        at(interval) if interval
      end

      def set(options)
        interval = options.delete(:wait_until) || options.delete(:wait)
        @opts.merge!(options)
        at(interval) if interval
        self
      end

      def perform_async(*args)
        if @opts["sync"] == true
          perform_inline(*args)
        else
          @klass.client_push(@opts.merge("args" => args, "class" => @klass))
        end
      end

      # Explicit inline execution of a job. Returns nil if the job did not
      # execute, true otherwise.
      def perform_inline(*args)
        raw = @opts.merge("args" => args, "class" => @klass).transform_keys(&:to_s)

        # validate and normalize payload
        item = normalize_item(raw)
        queue = item["queue"]

        # run client-side middleware
        result = Sidekiq.client_middleware.invoke(item["class"], item, queue, Sidekiq.redis_pool) do
          item
        end
        return nil unless result

        # round-trip the payload via JSON
        msg = Sidekiq.load_json(Sidekiq.dump_json(item))

        # prepare the job instance
        klass = msg["class"].constantize
        job = klass.new
        job.jid = msg["jid"]
        job.bid = msg["bid"] if job.respond_to?(:bid)

        # run the job through server-side middleware
        result = Sidekiq.server_middleware.invoke(job, msg, msg["queue"]) do
          # perform it
          job.perform(*msg["args"])
          true
        end
        return nil unless result
        # jobs do not return a result. they should store any
        # modified state.
        true
      end
      alias_method :perform_sync, :perform_inline

      def perform_bulk(args, batch_size: 1_000)
        hash = @opts.transform_keys(&:to_s)
        pool = Thread.current[:sidekiq_via_pool] || @klass.get_sidekiq_options["pool"] || Sidekiq.redis_pool
        client = Sidekiq::Client.new(pool)
        result = args.each_slice(batch_size).flat_map do |slice|
          client.push_bulk(hash.merge("class" => @klass, "args" => slice))
        end

        result.is_a?(Enumerator::Lazy) ? result.force : result
      end

      # +interval+ must be a timestamp, numeric or something that acts
      #   numeric (like an activesupport time interval).
      def perform_in(interval, *args)
        at(interval).perform_async(*args)
      end
      alias_method :perform_at, :perform_in

      private

      def at(interval)
        int = interval.to_f
        now = Time.now.to_f
        ts = (int < 1_000_000_000 ? now + int : int)
        # Optimization to enqueue something now that is scheduled to go out now or in the past
        @opts["at"] = ts if ts > now
        self
      end
    end

    module ClassMethods
      def delay(*args)
        raise ArgumentError, "Do not call .delay on a Sidekiq::Worker class, call .perform_async"
      end

      def delay_for(*args)
        raise ArgumentError, "Do not call .delay_for on a Sidekiq::Worker class, call .perform_in"
      end

      def delay_until(*args)
        raise ArgumentError, "Do not call .delay_until on a Sidekiq::Worker class, call .perform_at"
      end

      def queue_as(q)
        sidekiq_options("queue" => q.to_s)
      end

      def set(options)
        Setter.new(self, options)
      end

      def perform_async(*args)
        Setter.new(self, {}).perform_async(*args)
      end

      # Inline execution of job's perform method after passing through Sidekiq.client_middleware and Sidekiq.server_middleware
      def perform_inline(*args)
        Setter.new(self, {}).perform_inline(*args)
      end

      ##
      # Push a large number of jobs to Redis, while limiting the batch of
      # each job payload to 1,000. This method helps cut down on the number
      # of round trips to Redis, which can increase the performance of enqueueing
      # large numbers of jobs.
      #
      # +items+ must be an Array of Arrays.
      #
      # For finer-grained control, use `Sidekiq::Client.push_bulk` directly.
      #
      # Example (3 Redis round trips):
      #
      #     SomeWorker.perform_async(1)
      #     SomeWorker.perform_async(2)
      #     SomeWorker.perform_async(3)
      #
      # Would instead become (1 Redis round trip):
      #
      #     SomeWorker.perform_bulk([[1], [2], [3]])
      #
      def perform_bulk(*args, **kwargs)
        Setter.new(self, {}).perform_bulk(*args, **kwargs)
      end

      # +interval+ must be a timestamp, numeric or something that acts
      #   numeric (like an activesupport time interval).
      def perform_in(interval, *args)
        int = interval.to_f
        now = Time.now.to_f
        ts = (int < 1_000_000_000 ? now + int : int)

        item = {"class" => self, "args" => args}

        # Optimization to enqueue something now that is scheduled to go out now or in the past
        item["at"] = ts if ts > now

        client_push(item)
      end
      alias_method :perform_at, :perform_in

      ##
      # Allows customization for this type of Worker.
      # Legal options:
      #
      #   queue - use a named queue for this Worker, default 'default'
      #   retry - enable the RetryJobs middleware for this Worker, *true* to use the default
      #      or *Integer* count
      #   backtrace - whether to save any error backtrace in the retry payload to display in web UI,
      #      can be true, false or an integer number of lines to save, default *false*
      #   pool - use the given Redis connection pool to push this type of job to a given shard.
      #
      # In practice, any option is allowed.  This is the main mechanism to configure the
      # options for a specific job.
      def sidekiq_options(opts = {})
        super
      end

      def client_push(item) # :nodoc:
        pool = Thread.current[:sidekiq_via_pool] || get_sidekiq_options["pool"] || Sidekiq.redis_pool
        stringified_item = item.transform_keys(&:to_s)

        Sidekiq::Client.new(pool).push(stringified_item)
      end
    end
  end
end