aboutsummaryrefslogtreecommitdiff
path: root/pkg/core
diff options
context:
space:
mode:
authoralex <[email protected]>2026-07-10 20:28:47 +0200
committeralex <[email protected]>2026-07-10 20:28:47 +0200
commitf18210c30e67de4bede7c6080d85629f7086676c (patch)
treec9aedf21ed9f905424f8ab6fdb633254e127fbc6 /pkg/core
parent36d70e64340033f1e5b928524436c0b2f44c831e (diff)
downloadredis-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/core')
-rw-r--r--pkg/core/redisDB.go61
-rw-r--r--pkg/core/shards.go134
2 files changed, 195 insertions, 0 deletions
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