aboutsummaryrefslogtreecommitdiff
path: root/shards.go
blob: ef9f9ef72320f3a8ba69c763f3f4d1b69d95b5ba (plain)
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
package main

import (
	"hash/fnv"
	"sort"
	"sync"
)
type Shard struct{
	mu sync.RWMutex
	id int
	data map[string]Item
}

func (db *RedisDB) getShard(key string)*Shard{
	hash := fnv.New64a()
	hash.Write([]byte(key))
	val := hash.Sum64()

	return db.shards[val%NumShards]
}

func (db *RedisDB) Execute(key string, fn func(*Shard) interface{}) interface{} {
    shard := db.getShard(key)
    shard.mu.Lock()
    defer shard.mu.Unlock()
    return fn(shard)
}
func (db *RedisDB) ExecuteRead(key string, fn func(*Shard) interface{}) interface{} {
    shard := db.getShard(key)
    shard.mu.RLock()
    defer shard.mu.RUnlock()
    return fn(shard)
}
func ( db *RedisDB) ExecuteMulti(keys []string, fn func([]*Shard) interface{})interface{}{
	if len(keys) == 1{
		shard := db.getShard(keys[0])
		shard.mu.Lock()
		defer shard.mu.Unlock()
		shards := []*Shard{shard}
		return fn(shards)
	}else{
		shardMap := make(map[int]*Shard)
		for _, k := range keys {
			shard := db.getShard(k)
			shardMap[shard.id] = shard
		}
		var sortedIDs []int
		for id := range shardMap {
			sortedIDs = append(sortedIDs, id)
		}
		sort.Ints(sortedIDs)
		for _, id := range sortedIDs{
			shardMap[id].mu.Lock()
		}
		defer func(){
			for i := len(sortedIDs)-1; i >= 0; i --{
				shardMap[sortedIDs[i]].mu.Unlock()
			}
		}()
		
		shards := make([]*Shard, 0, len(shardMap))
		for _, id := range sortedIDs {
			shards = append(shards, shardMap[id])
		}
		return fn(shards)
	}
}

func ( db *RedisDB) ExecuteAll(fn func([]*Shard) interface{})interface{}{
	for _, shard := range db.shards {
        shard.mu.Lock()
    }

	defer func(){
		for i := len(db.shards)-1;i >= 0; i --{
			db.shards[i].mu.Unlock()
		}
	}()

	return fn(db.shards)
}

func ( db *RedisDB) ExecuteReadAll(fn func([]*Shard) interface{})interface{}{
	for _, shard := range db.shards {
        shard.mu.RLock()
    }

	defer func(){
		for i := len(db.shards)-1;i >= 0; i --{
			db.shards[i].mu.RUnlock()
		}
	}()

	return fn(db.shards)
}