diff options
| author | alex <[email protected]> | 2026-07-10 20:28:47 +0200 |
|---|---|---|
| committer | alex <[email protected]> | 2026-07-10 20:28:47 +0200 |
| commit | f18210c30e67de4bede7c6080d85629f7086676c (patch) | |
| tree | c9aedf21ed9f905424f8ab6fdb633254e127fbc6 /pkg | |
| parent | 36d70e64340033f1e5b928524436c0b2f44c831e (diff) | |
| download | redis-clone-f18210c30e67de4bede7c6080d85629f7086676c.tar.xz redis-clone-f18210c30e67de4bede7c6080d85629f7086676c.zip | |
separated commands into their own files TODO: refactor all commands to work like this
Diffstat (limited to 'pkg')
| -rw-r--r-- | pkg/commands/command-registry.go | 58 | ||||
| -rw-r--r-- | pkg/commands/del.go | 15 | ||||
| -rw-r--r-- | pkg/commands/get.go | 28 | ||||
| -rw-r--r-- | pkg/commands/set.go | 10 | ||||
| -rw-r--r-- | pkg/core/redisDB.go | 61 | ||||
| -rw-r--r-- | pkg/core/shards.go | 134 |
6 files changed, 306 insertions, 0 deletions
diff --git a/pkg/commands/command-registry.go b/pkg/commands/command-registry.go new file mode 100644 index 0000000..0193da1 --- /dev/null +++ b/pkg/commands/command-registry.go @@ -0,0 +1,58 @@ +package commands + +import ( + "errors" + "redisClone/pkg/core" +) +func SingleKey(args []string) []string { + return []string{args[1]} +} + +func AllSubsequentKeys(args []string) []string { + return args[1:] +} + +func NoKeys(args []string) []string { + return nil +} + +type CommandDef struct{ + MinArgs int + ExtractKeys func(args []string)[]string + Execute func(args []string, getShard func(k string) *core.Shard) interface{} +} + +func DispatchBaseCommand(db core.RedisDB ,commandName string, args []string) interface{} { + cmd, exists := Registry[commandName] + if !exists { + return nil + } + + if len(args) < cmd.MinArgs { + return errors.New("wrong number of arguments") + } + + keys := cmd.ExtractKeys(args) + + db.Lock(keys) + defer db.Unlock(keys) + return cmd.Execute(args, db.GetShard) +} + +var Registry = map[string]CommandDef{ + "GET": { + MinArgs: 2, + ExtractKeys: SingleKey, + Execute: Get, + }, + "SET": { + MinArgs: 3, + ExtractKeys: SingleKey, + Execute: Set, + }, + "DEL": { + MinArgs: 2, + ExtractKeys: AllSubsequentKeys, + Execute: Del, + }, +}
\ No newline at end of file diff --git a/pkg/commands/del.go b/pkg/commands/del.go new file mode 100644 index 0000000..bc7d3f9 --- /dev/null +++ b/pkg/commands/del.go @@ -0,0 +1,15 @@ +package commands + +import "redisClone/pkg/core" + +func Del(args []string, getShard func(k string) *core.Shard) interface{} { + count := 0 + for _, key := range args[1:] { + shard := getShard(key) + if _, exists := shard.Data[key]; exists { + delete(shard.Data, key) + count++ + } + } + return count + }
\ No newline at end of file diff --git a/pkg/commands/get.go b/pkg/commands/get.go new file mode 100644 index 0000000..9dc2bcd --- /dev/null +++ b/pkg/commands/get.go @@ -0,0 +1,28 @@ +package commands + +import ( + "errors" + "redisClone/pkg/core" + "time" +) + +func Get(args []string, getShard func(k string) *core.Shard) interface{} { + key := args[1] + s := getShard(key) + item, exists := s.Data[key] + if !exists { + return nil + } + + if item.ExpiresAt != nil && time.Now().After(*item.ExpiresAt) { + delete(s.Data, key) + return nil + } + + strValue, ok := item.Value.(string) + if !ok { + err := errors.New("WRONGTYPE Operation against a key holding the wrong kind of value") + return err + } + return strValue + }
\ No newline at end of file diff --git a/pkg/commands/set.go b/pkg/commands/set.go new file mode 100644 index 0000000..19175f1 --- /dev/null +++ b/pkg/commands/set.go @@ -0,0 +1,10 @@ +package commands + +import "redisClone/pkg/core" + +func Set(args []string, getShard func(k string) *core.Shard) interface{} { + key := args[1] + item := core.Item{Value: args[2]} + getShard(key).Data[key] = item + return core.SimpleString("OK") + }
\ No newline at end of file diff --git a/pkg/core/redisDB.go b/pkg/core/redisDB.go new file mode 100644 index 0000000..cbd737c --- /dev/null +++ b/pkg/core/redisDB.go @@ -0,0 +1,61 @@ +package core + +import ( + "fmt" + "time" +) + + +type Item struct { + Value any + ExpiresAt *time.Time +} + +type RedisDB struct { + Shards []*Shard +} + +type SimpleString string + +func SerializeRESP(v any) []byte { + switch val := v.(type) { + case string: + bulkString := fmt.Sprintf("$%d\r\n%s\r\n", len(val), val) + return []byte(bulkString) + case int: + integerString := fmt.Sprintf(":%d\r\n",val) + return []byte(integerString) + case nil: + return []byte("$-1\r\n") + case error: + errorString := fmt.Sprintf("-ERR %s\r\n",val.Error()) + return []byte(errorString) + case []string: + size := len(val) + result := []byte(fmt.Sprintf("*%d\r\n", size)) + for i := 0; i < size; i++ { + result = append(result, SerializeRESP(val[i])...) + } + return result + case []any: + size := len(val) + result := []byte(fmt.Sprintf("*%d\r\n", size)) + for i := 0; i < size; i++ { + result = append(result, SerializeRESP(val[i])...) + } + return result + case map[string]string: + size := len(val) + result := []byte(fmt.Sprintf("*%d\r\n", size*2)) + for key, value := range val { + result = append(result, SerializeRESP(key)...) + result = append(result, SerializeRESP(value)...) + } + return result + case SimpleString: + simpleStr := fmt.Sprintf("+%s\r\n", val) + return []byte(simpleStr) + default: + return []byte("-ERR internal server error: unknown type\r\n") + } +} diff --git a/pkg/core/shards.go b/pkg/core/shards.go new file mode 100644 index 0000000..5ede61c --- /dev/null +++ b/pkg/core/shards.go @@ -0,0 +1,134 @@ +package core + +import ( + "hash/fnv" + "sort" + "sync" +) +const NumShards = 16 + +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) +} + +func ( db *RedisDB) Lock(keys []string){ + if len(keys) == 1{ + shard := db.GetShard(keys[0]) + shard.Mu.Lock() + }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() + } + } +} +func ( db *RedisDB) Unlock(keys []string){ + 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 i := len(sortedIDs)-1; i >= 0; i --{ + shardMap[sortedIDs[i]].Mu.Unlock() + } +}
\ No newline at end of file |
