aboutsummaryrefslogtreecommitdiff
diff options
context:
space:
mode:
-rw-r--r--handlers.go16
-rw-r--r--shards.go28
2 files changed, 36 insertions, 8 deletions
diff --git a/handlers.go b/handlers.go
index 23efa9a..0b5267c 100644
--- a/handlers.go
+++ b/handlers.go
@@ -183,21 +183,21 @@ func handlePing(conn net.Conn, args []string) {
conn.Write(serializeRESP(SimpleString("PONG")))
}
-/*
func handleFlushall(conn net.Conn, args []string) {
- db.mu.Lock()
- defer db.mu.Unlock()
-
if len(args) != 1 {
err := errors.New("wrong number of arguments for 'FLUSHALL'")
conn.Write(serializeRESP(err))
return
}
-
- db.data = make(map[string]Item)
- conn.Write(serializeRESP(SimpleString("OK")))
+ result := db.ExecuteAll(func(shards []*Shard)interface{}{
+ for _, shard := range shards{
+ shard.data = make(map[string]Item)
+ }
+ return SimpleString("OK")
+ })
+ conn.Write(serializeRESP(result))
}
-
+/*
func handleRpush(conn net.Conn, args []string) {
key := args[1]
newItems := args[2:]
diff --git a/shards.go b/shards.go
index 279d9b1..ef9f9ef 100644
--- a/shards.go
+++ b/shards.go
@@ -65,3 +65,31 @@ func ( db *RedisDB) ExecuteMulti(keys []string, fn func([]*Shard) interface{})in
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)
+} \ No newline at end of file