From f18210c30e67de4bede7c6080d85629f7086676c Mon Sep 17 00:00:00 2001 From: alex Date: Fri, 10 Jul 2026 20:28:47 +0200 Subject: separated commands into their own files TODO: refactor all commands to work like this --- pkg/core/redisDB.go | 61 ++++++++++++++++++++++++ pkg/core/shards.go | 134 ++++++++++++++++++++++++++++++++++++++++++++++++++++ 2 files changed, 195 insertions(+) create mode 100644 pkg/core/redisDB.go create mode 100644 pkg/core/shards.go (limited to 'pkg/core') 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 -- cgit v1.2.3