Module: RedisQueuedLocks::Acquirer::LockSeriesPoC Private

Extended by:
AcquireLock::YieldExpire
Defined in:
lib/redis_queued_locks/acquirer/lock_series_poc.rb

Overview

This module is part of a private API. You should avoid using this module if possible, as it may be removed or be changed in the future.

NOTE: Lock Series PoC steep:ignore rubocop:disable all

Since:

  • 1.16.0

Defined Under Namespace

Modules: InstrVisitor, LogVisitor

Constant Summary

Constants included from AcquireLock::YieldExpire

AcquireLock::YieldExpire::DECREASE_LOCK_PTTL

Class Method Summary collapse

Methods included from AcquireLock::YieldExpire

yield_expire

Class Method Details

.lock_series_poc(redis, lock_names, detailed_result:, process_id:, thread_id:, fiber_id:, ractor_id:, ttl:, queue_ttl:, timeout:, timed:, retry_count:, retry_delay:, retry_jitter:, raise_errors:, instrumenter:, identity:, fail_fast:, meta:, detailed_acq_timeout_error:, instrument:, logger:, log_lock_try:, conflict_strategy:, read_write_mode:, access_strategy:, log_sampling_enabled:, log_sampling_percent:, log_sampler:, log_sample_this:, instr_sampling_enabled:, instr_sampling_percent:, instr_sampler:, instr_sample_this:, &block) ⇒ Object

This method is part of a private API. You should avoid using this method if possible, as it may be removed or be changed in the future.

Since:

  • 1.16.0

Version:

  • 1.16.2



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
# File 'lib/redis_queued_locks/acquirer/lock_series_poc.rb', line 19

def lock_series_poc( # steep:ignore
  redis,
  lock_names,
  detailed_result:,
  process_id:,
  thread_id:,
  fiber_id:,
  ractor_id:,
  ttl:,
  queue_ttl:,
  timeout:,
  timed:,
  retry_count:,
  retry_delay:,
  retry_jitter:,
  raise_errors:,
  instrumenter:,
  identity:,
  fail_fast:,
  meta:,
  detailed_acq_timeout_error:,
  instrument:,
  logger:,
  log_lock_try:,
  conflict_strategy:,
  read_write_mode:,
  access_strategy:,
  log_sampling_enabled:,
  log_sampling_percent:,
  log_sampler:,
  log_sample_this:,
  instr_sampling_enabled:,
  instr_sampling_percent:,
  instr_sampler:,
  instr_sample_this:,
  &block
)
  case meta
  when Hash, NilClass then nil
  else
    raise(
      RedisQueuedLocks::ArgumentError,
      "`:meta` argument should be a type of NilClass or Hash, got #{meta.class}."
    )
  end

  if meta.is_a?(::Hash) && (meta.any? do |key, _value|
    key == 'acq_id' ||
    key == 'hst_id' ||
    key == 'ts' ||
    key == 'ini_ttl' ||
    key == 'lock_key' ||
    key == 'rem_ttl' ||
    key == 'spc_ext_ttl' ||
    key == 'spc_cnt' ||
    key == 'l_spc_ext_ini_ttl' ||
    key == 'l_spc_ext_ts' ||
    key == 'l_spc_ts'
  end)
    raise(
      RedisQueuedLocks::ArgumentError,
      '`:meta` keys can not overlap reserved lock data keys ' \
      '"acq_id", "hst_id", "ts", "ini_ttl", "lock_key", "rem_ttl", "spc_cnt", ' \
      '"spc_ext_ttl", "l_spc_ext_ini_ttl", "l_spc_ext_ts", "l_spc_ts"'
    )
  end

  locks_count = lock_names.size

  lock_acquirement_operations = lock_names.each_with_index.map do |lock_name, position|
    lock_ttl =
      ttl * (locks_count - position) +
      retry_delay * (locks_count - position) +
      retry_jitter * (locks_count - position)

    { lock_name: lock_name, ttl: lock_ttl }
  end

  log_sampled = RedisQueuedLocks::Logging.should_log?(
    log_sampling_enabled,
    log_sample_this,
    log_sampling_percent,
    log_sampler
  )

  instr_sampled = RedisQueuedLocks::Instrument.should_instrument?(
    instr_sampling_enabled,
    instr_sample_this,
    instr_sampling_percent,
    instr_sampler
  )

  lock_keys_for_instrumentation = lock_names.map do |lock_name|
    RedisQueuedLocks::Resource.prepare_lock_key(lock_name)
  end

  acquirer_id_for_instrumentation = RedisQueuedLocks::Resource.acquirer_identifier(
    process_id,
    thread_id,
    fiber_id,
    ractor_id,
    identity
  )

  host_id_for_instrumentation = RedisQueuedLocks::Resource.host_identifier(
    process_id,
    thread_id,
    ractor_id,
    identity
  )

  RedisQueuedLocks::Acquirer::LockSeriesPoC::LogVisitor.start_lock_series_obtaining( # steep:ignore
    logger, log_sampled, lock_keys_for_instrumentation,
    queue_ttl, acquirer_id_for_instrumentation, host_id_for_instrumentation, access_strategy
  )

  acq_start_time = RedisQueuedLocks::Utilities.clock_gettime
  successfully_acquired_locks = [] # steep:ignore
  failed_locks_and_errors = {} # steep:ignore
  failed_on_lock = nil

  lock_acquirement_operations.map do |lock_operation_options|
    result =
      begin
        RedisQueuedLocks::Acquirer::AcquireLock.acquire_lock(
          redis,
          lock_operation_options[:lock_name],
          process_id:,
          thread_id:,
          fiber_id:,
          ractor_id:,
          ttl: lock_operation_options[:ttl],
          queue_ttl:,
          timeout:,
          timed:,
          retry_count:,
          retry_delay:,
          retry_jitter:,
          raise_errors:,
          instrumenter:,
          identity:,
          fail_fast:,
          meta:,
          detailed_acq_timeout_error:,
          instrument:,
          logger:,
          log_lock_try:,
          conflict_strategy:,
          read_write_mode:,
          access_strategy:,
          log_sampling_enabled:,
          log_sampling_percent:,
          log_sampler:,
          log_sample_this:,
          instr_sampling_enabled:,
          instr_sampling_percent:,
          instr_sampler:,
          instr_sample_this:
        )
      rescue RedisQueuedLocks::LockAlreadyObtainedError,
             RedisQueuedLocks::LockAcquirementIntermediateTimeoutError,
             RedisQueuedLocks::LockAcquirementTimeoutError,
             RedisQueuedLocks::LockAcquirementRetryLimitError,
             RedisQueuedLocks::ConflictLockObtainError => error
        if successfully_acquired_locks.any?
          # NOTE: release all previously acquired locks if any next lock is already locked
          successfully_acquired_locks.each do |operation_result|
            lock_key = RedisQueuedLocks::Resource.prepare_lock_key(operation_result[:lock_name])
            redis.with do |conn|
              conn.multi(watch: [lock_key]) do |transact|
                transact.call('DEL', lock_key)
              end
            end
          end
        end

        raise(error) if raise_errors
      end

    if result[:ok]
      successfully_acquired_locks << {
        lock_name: lock_operation_options[:lock_name],
        ok: result[:ok],
        result: result[:result]
      }
    else
      failed_on_lock = lock_operation_options[:lock_name]
      failed_locks_and_errors[lock_operation_options[:lock_name]] = result
      break
    end
  end

  if (successfully_acquired_locks.size == lock_names.size && (successfully_acquired_locks.all? { |res| res[:ok] }))
    acq_end_time = RedisQueuedLocks::Utilities.clock_gettime
    acq_time = ((acq_end_time - acq_start_time) / 1_000.0).ceil(2)
    ts = Time.now.to_s

    RedisQueuedLocks::Acquirer::LockSeriesPoC::LogVisitor.lock_series_obtained( # steep:ignore
      logger, log_sampled, lock_keys_for_instrumentation,
      queue_ttl, acquirer_id_for_instrumentation, host_id_for_instrumentation,
      acq_time, access_strategy
    )

    RedisQueuedLocks::Acquirer::LockSeriesPoC::InstrVisitor.lock_series_obtained( # steep:ignore
      instrumenter, instr_sampled, lock_keys_for_instrumentation,
      ttl, acquirer_id_for_instrumentation, host_id_for_instrumentation, ts, acq_time, instrument
    )

    yield_time = RedisQueuedLocks::Utilities.clock_gettime
    ttl_shift = (
      (yield_time - acq_end_time) / 1_000.0 -
      RedisQueuedLocks::Resource::REDIS_TIMESHIFT_ERROR
    ).ceil(2)

    yield_result = nil
    is_lock_manually_released = nil
    hold_time = nil

    begin
      yield_result = yield_expire( # steep:ignore
        redis,
        logger,
        lock_keys_for_instrumentation.last,
        acquirer_id_for_instrumentation,
        host_id_for_instrumentation,
        access_strategy,
        timed,
        ttl_shift,
        ttl,
        queue_ttl,
        meta,
        log_sampled,
        instr_sampled,
        false, # should_expire (expire manually)
        false, # should_decrease (expire manually)
        &block
      )
    ensure
      is_lock_manually_released = false

      # expire locks manually
      if block_given?
        redis.with do |conn|
          # use transaction in order to exclude any cross-locking during the group expiration
          conn.multi(watch: lock_keys_for_instrumentation) do |transaction|
            lock_keys_for_instrumentation.each do |lock_key|
              transaction.call('EXPIRE', lock_key, '0')
            end
          end
        end
        is_lock_manually_released = true
      end

      rel_time = RedisQueuedLocks::Utilities.clock_gettime
      hold_time = ((rel_time - acq_end_time) / 1_000.0).ceil(2)
      ts = Time.now.to_f

      RedisQueuedLocks::Acquirer::LockSeriesPoC::LogVisitor.expire_lock_series( # steep:ignore
        logger, log_sampled, lock_keys_for_instrumentation,
        queue_ttl, acquirer_id_for_instrumentation, host_id_for_instrumentation, access_strategy
      ) if is_lock_manually_released

      RedisQueuedLocks::Acquirer::LockSeriesPoC::InstrVisitor.lock_series_hold_and_release( # steep:ignore
        instrumenter,
        instr_sampled,
        lock_keys_for_instrumentation,
        ttl,
        acquirer_id_for_instrumentation,
        host_id_for_instrumentation,
        ts,
        acq_time,
        hold_time,
        instrument
      ) if is_lock_manually_released
    end

    if detailed_result
      {
        yield_result: yield_result,
        locks_release_strategy: block_given? ? :immediate_release_after_yield : :redis_key_ttl,
        locks_released_at: block_given? ? ts : nil,
        locks_acq_time: acq_time,
        locks_hold_time: is_lock_manually_released ? hold_time : nil,
        lock_series: lock_names,
        rql_lock_series: lock_keys_for_instrumentation
      }
    else
      yield_result
    end
  else
    acquired_locks = successfully_acquired_locks.map { |state| state[:lock_name] }
    missing_locks = lock_names - acquired_locks

    {
      ok: false,
      result: {
        error: :failed_to_acquire_lock_series,
        detailed_errors: failed_locks_and_errors,
        lock_series: lock_names,
        acquired_locks:,
        missing_locks:,
        failed_on_lock:,
      }
    }
  end
end