_fixSyncSeqRollback will correct a rolled back _sync:seq document in the bucket
(ctx context.Context, prevAllocTo, expectedValue uint64)
| 451 | |
| 452 | // _fixSyncSeqRollback will correct a rolled back _sync:seq document in the bucket |
| 453 | func (s *sequenceAllocator) _fixSyncSeqRollback(ctx context.Context, prevAllocTo, expectedValue uint64) (allocatedToSeq uint64, err error) { |
| 454 | base.WarnfCtx(ctx, "rollback of _sync:seq document detected. Allocated to %d but expected value of at least %d", prevAllocTo, expectedValue) |
| 455 | // find diff between current _sync:seq value and what we expected it to be + correction value |
| 456 | correctionIncrValue := (expectedValue - prevAllocTo) + syncSeqCorrectionValue |
| 457 | |
| 458 | worker := func() (bool, error, interface{}) { |
| 459 | // grab _sync:seq value + its current cas value |
| 460 | var result uint64 |
| 461 | cas, err := s.datastore.Get(s.metaKeys.SyncSeqKey(), &result) |
| 462 | if err != nil { |
| 463 | return false, err, 0 |
| 464 | } |
| 465 | // set the value to _sync:seq current value + incr value if result above is below that value |
| 466 | setVal := result + correctionIncrValue |
| 467 | if result < setVal { |
| 468 | _, err = s.datastore.WriteCas(s.metaKeys.SyncSeqKey(), 0, cas, setVal, 0) |
| 469 | if base.IsCasMismatch(err) { |
| 470 | // retry on cas error |
| 471 | return true, err, nil |
| 472 | } |
| 473 | if err == nil { |
| 474 | // successfully corrected _sync:seq value above |
| 475 | base.DebugfCtx(ctx, base.KeyCRUD, "corrected _sync:seq document from value %d by the value of %d", prevAllocTo, setVal) |
| 476 | return false, nil, setVal |
| 477 | } |
| 478 | } |
| 479 | // if we get here we either had error above that is not a cas mismatch thus we need to exit with failure or the result |
| 480 | // from the fetch of _sync:seq was larger than expected (_sync:seq may have been fixed by another node) |
| 481 | return false, err, 0 |
| 482 | } |
| 483 | |
| 484 | retryErr, _ := base.RetryLoop(ctx, "Fix _sync:seq Value", worker, base.CreateDoublingSleeperFunc(1000, 5)) |
| 485 | if retryErr != nil { |
| 486 | base.WarnfCtx(ctx, "error: %v in retry loop to correct _sync:seq value", retryErr) |
| 487 | return 0, retryErr |
| 488 | } |
| 489 | |
| 490 | base.DebugfCtx(ctx, base.KeyCRUD, "_sync:seq value successfully corrected, entering normal sequence batch processing for this node") |
| 491 | // if _sync:seq has been fixed successfully just increment by batch size to get new unique batch for this node |
| 492 | allocatedToSeq, err = s._incrementSequence(s.sequenceBatchSize) |
| 493 | if err != nil { |
| 494 | return 0, err |
| 495 | } |
| 496 | |
| 497 | return allocatedToSeq, err |
| 498 | } |
no test coverage detected