Module: RedisQueuedLocks::Acquirer::AcquireLock::TryToLock Private

Included in:
RedisQueuedLocks::Acquirer::AcquireLock
Defined in:
lib/redis_queued_locks/acquirer/acquire_lock/try_to_lock.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.

rubocop:disable Metrics/ModuleLength

Since:

  • 1.0.0

Version:

  • 1.7.0

Defined Under Namespace

Modules: LogVisitor

Constant Summary collapse

EXTEND_LOCK_PTTL =

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

Returns:

  • (String)

Since:

  • 1.3.0

"local new_lock_pttl = redis.call(\"PTTL\", KEYS[1]) + ARGV[1];\nreturn redis.call(\"PEXPIRE\", KEYS[1], new_lock_pttl);\n".strip.tr("\n", '').freeze

Instance Method Summary collapse

Instance Method Details

#try_to_lock(redis, logger, log_lock_try, lock_key, read_write_mode, lock_key_queue, read_lock_key_queue, write_lock_key_queue, acquirer_id, host_id, acquirer_position, ttl, queue_ttl, fail_fast, conflict_strategy, access_strategy, meta, log_sampled, instr_sampled) ⇒ Hash<Symbol,Any>

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.

rubocop:disable Metrics/MethodLength

Parameters:

  • redis (RedisClient)
  • logger (::Logger, #debug)
  • log_lock_try (Boolean)
  • lock_key (String)
  • read_write_mode (Symbol)
  • lock_key_queue (String)
  • read_lock_key_queue (String)
  • write_lock_key_queue (String)
  • acquirer_id (String)
  • host_id (String)
  • acquirer_position (Numeric)
  • ttl (Integer)
  • queue_ttl (Integer)
  • fail_fast (Boolean)
  • conflict_strategy (Symbol)
  • access_strategy (Symbol)
  • meta (NilClass, Hash<String|Symbol,Any>)
  • log_sampled (Boolean)
  • instr_sampled (Boolean)

Returns:

  • (Hash<Symbol,Any>) —

    Format: { ok: true/false, result: Symbol|Hash<Symbol,Any> }

Since:

  • 1.0.0

Version:

  • 1.13.0



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
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
# File 'lib/redis_queued_locks/acquirer/acquire_lock/try_to_lock.rb', line 44

def try_to_lock(
  redis,
  logger,
  log_lock_try,
  lock_key,
  read_write_mode,
  lock_key_queue,
  read_lock_key_queue,
  write_lock_key_queue,
  acquirer_id,
  host_id,
  acquirer_position,
  ttl,
  queue_ttl,
  fail_fast,
  conflict_strategy,
  access_strategy,
  meta,
  log_sampled,
  instr_sampled
)
  # Step X: intermediate invocation results
  # @type var inter_result: Symbol?
  inter_result = nil
  # @type var timestamp: Float?
  timestamp = nil
  # @type var spc_processed_timestamp: Float?
  spc_processed_timestamp = nil

  LogVisitor.start(
    logger, log_sampled, log_lock_try, lock_key,
    queue_ttl, acquirer_id, host_id, access_strategy
  )

  # Step X: start to work with lock acquiring
  result = redis.with do |rconn|
    LogVisitor.rconn_fetched(
      logger, log_sampled, log_lock_try, lock_key,
      queue_ttl, acquirer_id, host_id, access_strategy
    )

    # Step 0:
    #   watch the lock key changes (and discard acquirement if lock is already
    #   obtained by another acquirer during the current lock acquiremntt)
    rconn.multi(watch: [lock_key]) do |transact|
      # SP-Conflict status PREPARING: get the current lock obtainer
      current_lock_obtainer = rconn.call('HGET', lock_key, 'acq_id')
      # SP-Conflict status PREPARING: status flag variable
      sp_conflict_status = nil

      # SP-Conflict Step X1: calculate the current deadlock status
      if current_lock_obtainer != nil && acquirer_id == current_lock_obtainer
        LogVisitor.same_process_conflict_detected(
          logger, log_sampled, log_lock_try, lock_key,
          queue_ttl, acquirer_id, host_id, access_strategy
        )

        # SP-Conflict Step X2: self-process dead lock moment started.
        # SP-Conflict CHECK (Step CHECK): check chosen strategy and flag the current status
        case conflict_strategy
        when :work_through
          # <SP-Conflict Moment>: work through => exit
          sp_conflict_status = :conflict_work_through
        when :extendable_work_through
          # <SP-Conflict Moment>: extendable_work_through => extend the lock pttl and exit
          sp_conflict_status = :extendable_conflict_work_through
        when :wait_for_lock
          # <SP-Conflict Moment>: wait_for_lock => obtain a lock in classic way
          sp_conflict_status = :conflict_wait_for_lock
        when :dead_locking
          # <SP-Conflict Moment>: dead_locking => exit and fail
          sp_conflict_status = :conflict_dead_lock
        else
          # <SP-Conflict Moment>:
          #   - unknown status => work in classic way <wait_for_lock>
          #   - it is a case when the new status is added to the code base in the past
          #     but is forgotten to be added here;
          sp_conflict_status = :conflict_wait_for_lock
        end
        LogVisitor.same_process_conflict_analyzed(
          logger, log_sampled, log_lock_try, lock_key,
          queue_ttl, acquirer_id, host_id, access_strategy, sp_conflict_status
        )
      end

      # SP-Conflict-Step X2: switch to conflict-based logic or not
      if sp_conflict_status == :extendable_conflict_work_through
        # SP-Conflict-Step FINAL (SPCF): extend the lock and work through
        #   - extend the lock ttl;
        #   - store extensions in lock metadata;
        # SPCF Step 1: extend the lock pttl
        #   - [REDIS RESULT]: in normal cases should return the last script command value
        #     - for the current script should return:
        #       <1> => timeout was set;
        #       <0> => timeount was not set;
        transact.call('EVAL', EXTEND_LOCK_PTTL, 1, lock_key, ttl)
        # SPCF Step <Meta>: store conflict-state additionals in lock metadata:
        # SPCF Step 2: (lock meta-data)
        #   - add the added ttl to reflect the real lock TTL in info;
        #   - [REDIS RESULT]: in normal cases should return the value of <ttl> key
        #     - for non-existent key value starts from <0> (zero)
        transact.call('HINCRBY', lock_key, 'spc_ext_ttl', ttl)
        # SPCF Step 3: (lock meta-data)
        #   - increment the conflcit counter in order to remember
        #     how many times dead lock happened;
        #   - [REDIS RESULT]: in normal cases should return the value of <spc_cnt> key
        #     - for non-existent key starts from 0
        transact.call('HINCRBY', lock_key, 'spc_cnt', 1)
        # SPCF Step 4: (lock meta-data)
        #   - remember the last ext-timestamp and the last ext-initial ttl;
        #   - [REDIS RESULT]: for normal cases should return the number of fields
        #     were added/changed;
        transact.call(
          'HSET',
          lock_key,
          'l_spc_ext_ts', spc_processed_timestamp = Time.now.to_f,
          'l_spc_ext_ini_ttl', ttl
        )
        inter_result = :extendable_conflict_work_through

        # @type var sp_conflict_status: Symbol
        # @type var spc_processed_timestamp: Float
        LogVisitor.reentrant_lock__extend_and_work_through(
          logger, log_sampled, log_lock_try, lock_key,
          queue_ttl, acquirer_id, host_id, access_strategy,
          sp_conflict_status, ttl, spc_processed_timestamp
        )
      # SP-Conflict-Step X2: switch to dead lock logic or not
      elsif sp_conflict_status == :conflict_work_through
        inter_result = :conflict_work_through

        # SPCF Step X: (lock meta-data)
        #   - increment the conflcit counter in order to remember
        #     how many times dead lock happened;
        #   - [REDIS RESULT]: in normal cases should return the value of <spc_cnt> key
        #     - for non-existent key starts from 0
        transact.call('HINCRBY', lock_key, 'spc_cnt', 1)
        # SPCF Step 4: (lock meta-data)
        #   - remember the last ext-timestamp and the last ext-initial ttl;
        #   - [REDIS RESULT]: for normal cases should return the number of fields
        #     were added/changed;
        transact.call(
          'HSET',
          lock_key,
          'l_spc_ts', spc_processed_timestamp = Time.now.to_f
        )

        # @type var sp_conflict_status: Symbol
        # @type var spc_processed_timestamp: Float
        LogVisitor.reentrant_lock__work_through(
          logger, log_sampled, log_lock_try, lock_key,
          queue_ttl, acquirer_id, host_id, access_strategy,
          sp_conflict_status, spc_processed_timestamp
        )
      # SP-Conflict-Step X2: switch to dead lock logic or not
      elsif sp_conflict_status == :conflict_dead_lock
        inter_result = :conflict_dead_lock
        spc_processed_timestamp = Time.now.to_f

        # @type var sp_conflict_status: Symbol
        # @type var spc_processed_timestamp: Float
        LogVisitor.single_process_lock_conflict__dead_lock(
          logger, log_sampled, log_lock_try, lock_key,
          queue_ttl, acquirer_id, host_id, access_strategy,
          sp_conflict_status, spc_processed_timestamp
        )
      # Reached the SP-Non-Conflict Mode (NOTE):
      #   - in other sp-conflict cases we are in <wait_for_lock> (non-conflict) status and should
      #     continue to work in classic way (next lines of code):
      elsif fail_fast && current_lock_obtainer != nil # Fast-Step X0: fail-fast check
        # Fast-Step X1: lock is already obtained. fail fast leads to "no try".
        inter_result = :fail_fast_no_try
      else
        # Step 1: add an acquirer to the lock acquirement queue
        # NOTE:
        #   'NX' means "Only add new elements. Don't update already existing elements."
        #   that works as:
        #     1. (enqueue) <<if you are already in the queue - do nothing and wait your time>>
        #     2. (requeue) or <<add to the right pre-calculated position if you are
        #       not in the queue now>>;
        rconn.call('ZADD', lock_key_queue, 'NX', acquirer_position, acquirer_id)

        LogVisitor.acq_added_to_queue(
          logger, log_sampled, log_lock_try, lock_key,
          queue_ttl, acquirer_id, host_id, access_strategy
        )

        # Step 2.1: drop expired acquirers from the lock queue
        rconn.call(
          'ZREMRANGEBYSCORE',
          lock_key_queue,
          '-inf',
          RedisQueuedLocks::Resource.acquirer_dead_score(queue_ttl)
        )

        LogVisitor.remove_expired_acqs(
          logger, log_sampled, log_lock_try, lock_key,
          queue_ttl, acquirer_id, host_id, access_strategy
        )

        # Step 3: get the actual acquirer waiting in the queue
        waiting_acquirer = Array(rconn.call('ZRANGE', lock_key_queue, '0', '0')).first

        LogVisitor.get_first_from_queue(
          logger, log_sampled, log_lock_try, lock_key,
          queue_ttl, acquirer_id, host_id, access_strategy, waiting_acquirer
        )

        # Step PRE-4.x: check if the request time limit is reached
        #   (when the current try self-removes itself from queue (queue ttl has come))
        if waiting_acquirer == nil
          LogVisitor.exit__queue_ttl_reached(
            logger, log_sampled, log_lock_try, lock_key,
            queue_ttl, acquirer_id, host_id, access_strategy
          )

          inter_result = :dead_score_reached
          # Step STRATEGY: check the stragegy and corresponding preventing factor
          # Step STRATEGY (queued): check the actual acquirer: is it ours? are we aready to lock?
        elsif access_strategy == :queued && waiting_acquirer != acquirer_id
          # Step ROLLBACK 1.1: our time hasn't come yet. retry!

          LogVisitor.exit__no_first(
            logger, log_sampled, log_lock_try, lock_key,
            queue_ttl, acquirer_id, host_id, access_strategy, waiting_acquirer,
            rconn.call('HGETALL', lock_key).to_h
          )
          inter_result = :acquirer_is_not_first_in_queue
          # Step STRAGEY: successfull (:queued OR :random)
        elsif (access_strategy == :queued && waiting_acquirer == acquirer_id) ||
              (access_strategy == :random)
          # NOTE: our time has come! let's try to acquire the lock!

          # Step 5: find the lock -> check if the our lock is already acquired
          locked_by_acquirer = rconn.call('HGET', lock_key, 'acq_id')

          if locked_by_acquirer
            # Step ROLLBACK 2: required lock is stil acquired. retry!

            LogVisitor.exit__lock_still_obtained(
              logger, log_sampled, log_lock_try, lock_key,
              queue_ttl, acquirer_id, host_id, access_strategy,
              waiting_acquirer, locked_by_acquirer,
              rconn.call('HGETALL', lock_key).to_h
            )
            inter_result = :lock_is_still_acquired
          else
            # NOTE: required lock is free and ready to be acquired! acquire!

            # Step 6.1: remove our acquirer from waiting queue
            transact.call('ZREM', lock_key_queue, acquirer_id)

            # Step 6.2: acquire a lock and store an info about the acquirer and host
            transact.call(
              'HSET',
              lock_key,
              'acq_id', acquirer_id,
              'hst_id', host_id,
              'ts', timestamp = Time.now.to_f,
              'ini_ttl', ttl,
              *(meta.to_a if meta != nil) # steep:ignore
            )

            # Step 6.3: set the lock expiration time in order to prevent "infinite locks"
            transact.call('PEXPIRE', lock_key, ttl) # NOTE: in milliseconds

            LogVisitor.obtain__free_to_acquire(
              logger, log_sampled, log_lock_try, lock_key,
              queue_ttl, acquirer_id, host_id, access_strategy
            )
          end
        end
      end
    end
  end

  # Step 7: Analyze the aquirement attempt:
  # rubocop:disable Lint/DuplicateBranch
  case
  when inter_result == :extendable_conflict_work_through
    # Step 7.same_process_conflict.A:
    #   - extendable_conflict_work_through case => yield <block> without lock realesing/extending;
    #   - lock is extended in logic above;
    #   - if <result == nil> => the lock was changed during an extention:
    #     it is the fail case => go retry.
    #   - else: let's go! :))
    if result.is_a?(::Array) && result.size == 4 # NOTE: four commands should be processed
      # TODO:
      #   => (!) analyze the command result and do actions with the depending on it
      #   1. EVAL (extend lock pttl) (OK for != nil)
      #   2. HINCRBY (ttl extension) (OK for != nil)
      #   3. HINCRBY (increased spc count) (OK for != nil)
      #   4. HSET (store the last spc time and ttl data) (OK for == 2 or != nil)
      if result[0] != nil && result[1] != nil && result[2] != nil && result[3] != nil
        {
          ok: true,
          result: {
            process: :extendable_conflict_work_through,
            lock_key: lock_key,
            acq_id: acquirer_id,
            hst_id: host_id,
            ts: spc_processed_timestamp,
            ttl: ttl
          }
        }
      elsif result[0] != nil
        # NOTE: that is enough to the fact that the lock is extended but <TODO>
        # TODO: add detalized overview (log? some in-line code clarifications?) of the result
        {
          ok: true,
          result: {
            process: :extendable_conflict_work_through,
            lock_key: lock_key,
            acq_id: acquirer_id,
            hst_id: host_id,
            ts: spc_processed_timestamp,
            ttl: ttl
          }
        }
      else
        # NOTE: unknown behaviour :thinking:
        { ok: false, result: :unknown }
      end
    elsif result == nil || (result.is_a?(::Array) && result.empty?)
      # NOTE: the lock key was changed durign an SPC logic execution
      { ok: false, result: :lock_is_acquired_during_acquire_race }
    else
      # NOTE: unknown behaviour :thinking:. this part is not reachable at the moment.
      { ok: false, result: :unknown }
    end
  when inter_result == :conflict_work_through
    # Step 7.same_process_conflict.B:
    #   - conflict_work_through case => yield <block> without lock realesing/extending
    {
      ok: true,
      result: {
        process: :conflict_work_through,
        lock_key: lock_key,
        acq_id: acquirer_id,
        hst_id: host_id,
        ts: spc_processed_timestamp,
        ttl: ttl
      }
    }
  when inter_result == :conflict_dead_lock
    # Step 7.same_process_conflict.C:
    #  - deadlock. should fail in acquirement logic;
    { ok: false, result: :conflict_dead_lock }
  when fail_fast && inter_result == :fail_fast_no_try
    # Step 7.a: lock is still acquired and we should exit from the logic as soon as possible
    { ok: false, result: :fail_fast_no_try }
  when inter_result == :dead_score_reached
    { ok: false, result: :dead_score_reached }
  when inter_result == :lock_is_still_acquired
    # Step 7.b: lock is still acquired by another process => failed to acquire
    { ok: false, result: :lock_is_still_acquired }
  when inter_result == :acquirer_is_not_first_in_queue
    # Step 7.c: lock is still acquired by another process => failed to acquire
    { ok: false, result: :acquirer_is_not_first_in_queue }
  when result == nil || (result.is_a?(::Array) && result.empty?)
    # Step 7.d: lock is already acquired durign the acquire race => failed to acquire
    { ok: false, result: :lock_is_acquired_during_acquire_race }
  when result.is_a?(::Array) && result.size == 3 # NOTE: 3 is a count of redis lock commands
    # TODO:
    #   => (!) analyze the command result and do actions with the depending on it;
    #   => (*) at this moment we accept that all comamnds are completed successfully;
    #   => (!) need to analyze:
    #   1. zrem shoud return ? (?)
    #   2. hset should return 3 as minimum
    #      (lock key is added to the redis as a hashmap with 3 fields as minimum)
    #   3. pexpire should return 1 (expiration time is successfully applied)

    # Step 7.d: locked! :) let's go! => successfully acquired
    {
      ok: true,
      result: {
        process: :lock_obtaining,
        lock_key: lock_key,
        acq_id: acquirer_id,
        hst_id: host_id,
        ts: timestamp,
        ttl: ttl
      }
    }
  else
    # Ste 7.3: unknown behaviour :thinking:
    { ok: false, result: :unknown }
  end
  # rubocop:enable Lint/DuplicateBranch
end