File size: 3,334 Bytes
b110593
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
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
//                           _       _
// __      _____  __ ___   ___  __ _| |_ ___
// \ \ /\ / / _ \/ _` \ \ / / |/ _` | __/ _ \
//  \ V  V /  __/ (_| |\ V /| | (_| | ||  __/
//   \_/\_/ \___|\__,_| \_/ |_|\__,_|\__\___|
//
//  Copyright © 2016 - 2024 Weaviate B.V. All rights reserved.
//
//  CONTACT: [email protected]
//

package common

import "sync"

const DefaultShardedLocksCount = 512

type ShardedLocks struct {
	// number of locks
	count int
	// ensures single LockAll and multiple RLockAll, Lock and RLock
	// LockAll is exclusive to RLockAll, Lock, RLock
	writeAll *sync.RWMutex
	// indicates whether write of any single shard is ongoing, exclusive with readAll
	writeAny *sync.RWMutex
	// indicates whether read of all shards is ongoing, exclusive with writeAny
	readAll *sync.RWMutex
	// allows safe transition between writeAny and readAll
	change *sync.RWMutex
	// sharded locks
	shards []*sync.RWMutex
}

func NewDefaultShardedLocks() *ShardedLocks {
	return NewShardedLocks(DefaultShardedLocksCount)
}

func NewShardedLocks(count int) *ShardedLocks {
	if count < 2 {
		count = 2
	}

	writeAll := new(sync.RWMutex)
	writeAny := new(sync.RWMutex)
	readAll := new(sync.RWMutex)
	change := new(sync.RWMutex)
	shards := make([]*sync.RWMutex, count)
	for i := 0; i < count; i++ {
		shards[i] = new(sync.RWMutex)
	}

	return &ShardedLocks{
		count:    count,
		writeAll: writeAll,
		readAll:  readAll,
		writeAny: writeAny,
		change:   change,
		shards:   shards,
	}
}

func (sl *ShardedLocks) LockAll() {
	sl.writeAll.Lock()
}

func (sl *ShardedLocks) UnlockAll() {
	sl.writeAll.Unlock()
}

func (sl *ShardedLocks) LockedAll(callback func()) {
	sl.LockAll()
	defer sl.UnlockAll()

	callback()
}

func (sl *ShardedLocks) Lock(id uint64) {
	sl.writeAll.RLock()
	sl.markOngoingWriteAny()
	sl.shards[sl.mid(id)].Lock()
}

func (sl *ShardedLocks) Unlock(id uint64) {
	sl.shards[sl.mid(id)].Unlock()
	sl.writeAny.RUnlock()
	sl.writeAll.RUnlock()
}

func (sl *ShardedLocks) Locked(id uint64, callback func()) {
	sl.Lock(id)
	defer sl.Unlock(id)

	callback()
}

func (sl *ShardedLocks) RLockAll() {
	sl.writeAll.RLock()
	sl.markOngoingReadAll()
}

func (sl *ShardedLocks) RUnlockAll() {
	sl.readAll.RUnlock()
	sl.writeAll.RUnlock()
}

func (sl *ShardedLocks) RLockedAll(callback func()) {
	sl.RLockAll()
	defer sl.RUnlockAll()

	callback()
}

func (sl *ShardedLocks) RLock(id uint64) {
	sl.writeAll.RLock()
	sl.shards[sl.mid(id)].RLock()
}

func (sl *ShardedLocks) RUnlock(id uint64) {
	sl.shards[sl.mid(id)].RUnlock()
	sl.writeAll.RUnlock()
}

func (sl *ShardedLocks) RLocked(id uint64, callback func()) {
	sl.RLock(id)
	defer sl.RUnlock(id)

	callback()
}

func (sl *ShardedLocks) mid(id uint64) uint64 {
	return id % uint64(sl.count)
}

func (sl *ShardedLocks) markOngoingWriteAny() {
	sl.change.RLock()
	defer sl.change.RUnlock()

	// wait until no ongoing readAll
	sl.readAll.Lock()
	// mark ongoing writeAny
	sl.writeAny.RLock()
	sl.readAll.Unlock()
}

func (sl *ShardedLocks) markOngoingReadAll() {
	sl.change.Lock()
	defer sl.change.Unlock()

	// wait until no ongoing writeAny
	sl.writeAny.Lock()
	// mark ongoing readAll
	sl.readAll.RLock()
	sl.writeAny.Unlock()
}