Skip to content

Commit 3f6b0e5

Browse files
authored
Merge pull request #15 from bluesky-social/jc/zset
Subset of Commands for Ordered Sets
2 parents 5793ea6 + 81b4f46 commit 3f6b0e5

10 files changed

Lines changed: 1144 additions & 188 deletions

File tree

.gitignore

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -30,3 +30,6 @@ go.work.sum
3030
# Editor/IDE
3131
# .idea/
3232
# .vscode/
33+
34+
# Delve build artifacts
35+
**__debug_bin**

README.md

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -68,3 +68,4 @@ FoundationDB provides strictly serializable transactions atop an ordered key-val
6868
|`redis_v0/<user_id>/obj/<obj_id>`|Storage of data objects|
6969
|`redis_v0/<user_id>/list/<list_id>`|Metadata protobuf for information on each list|
7070
|`redis_v0/<user_id>/list/<list_id>/<item_id>`|Metadata protobuf for information on each item in a list|
71+
|`redis_v0/<user_id>/zset/<set_id>`|Stores an ordered list of item scores for ordered sets|

internal/server/redis/auth.go

Lines changed: 7 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -117,13 +117,19 @@ func (s *session) setSessionUser(user *types.User) error {
117117
return fmt.Errorf("failed to initialize user list directory: %w", err)
118118
}
119119

120+
zsetDir, err := userDir.CreateOrOpen(s.fdb, []string{"zset"}, nil)
121+
if err != nil {
122+
return fmt.Errorf("failed to initialize user zset directory: %w", err)
123+
}
124+
120125
s.userMu.Lock()
121-
s.user = &sessionUser{
126+
s.user = &userSession{
122127
objDir: objDir,
123128
metaDir: metaDir,
124129
uidDir: uidDir,
125130
reverseUIDDir: reverseUIDDir,
126131
listDir: listDir,
132+
zsetDir: zsetDir,
127133
user: user,
128134
}
129135
s.userMu.Unlock()

internal/server/redis/list.go

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -97,7 +97,7 @@ func (s *session) handlePush(ctx context.Context, args []resp.Value, left bool)
9797
}
9898

9999
_, err = s.fdb.Transact(func(tx fdb.Transaction) (any, error) {
100-
listMetaKey, objMeta, err := s.getObjectMeta(ctx, tx, key)
100+
listMetaKey, objMeta, err := s.getMeta(ctx, tx, key)
101101
if err != nil {
102102
return nil, fmt.Errorf("failed to get list meta: %w", err)
103103
}

internal/server/redis/object.go

Lines changed: 27 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,7 @@ type objectKind int64
2323
const (
2424
objectKindBasic objectKind = 1 << iota
2525
objectKindSet
26+
objectKindSortedSet
2627
objectKindList
2728
objectKindListItem
2829
)
@@ -44,6 +45,9 @@ func getNumChunks(meta *types.ObjectMeta, kind gt.Option[objectKind]) (uint32, e
4445
case *types.ObjectMeta_Set:
4546
objKind = objectKindSet
4647
numChunks = typ.Set.NumChunks
48+
case *types.ObjectMeta_SortedSet:
49+
objKind = objectKindSortedSet
50+
numChunks = typ.SortedSet.NumChunks
4751
case *types.ObjectMeta_ListItem:
4852
objKind = objectKindListItem
4953
numChunks = typ.ListItem.NumChunks
@@ -64,7 +68,7 @@ func (s *session) getObject(ctx context.Context, tx fdb.ReadTransaction, kind ob
6468
ctx, span := s.tracer.Start(ctx, "getObject")
6569
defer span.End()
6670

67-
_, meta, err := s.getObjectMeta(ctx, tx, id)
71+
_, meta, err := s.getMeta(ctx, tx, id)
6872
if err != nil {
6973
span.RecordError(err)
7074
return nil, nil, err
@@ -116,7 +120,7 @@ func (s *session) writeObject(ctx context.Context, tx fdb.Transaction, id string
116120
now := timestamppb.Now()
117121

118122
// check if the object already exists and should be overwritten
119-
metaKey, meta, err := s.getObjectMeta(ctx, tx, id)
123+
metaKey, meta, err := s.getMeta(ctx, tx, id)
120124
if err != nil {
121125
span.RecordError(err)
122126
return err
@@ -137,6 +141,8 @@ func (s *session) writeObject(ctx context.Context, tx fdb.Transaction, id string
137141
meta.Type = &types.ObjectMeta_Basic{Basic: &types.BasicObjectMeta{}}
138142
case objectKindSet:
139143
meta.Type = &types.ObjectMeta_Set{Set: &types.SetMeta{}}
144+
case objectKindSortedSet:
145+
meta.Type = &types.ObjectMeta_SortedSet{SortedSet: &types.SortedSetMeta{}}
140146
case objectKindListItem:
141147
meta.Type = &types.ObjectMeta_ListItem{ListItem: &types.ListItemMeta{}}
142148
}
@@ -157,6 +163,8 @@ func (s *session) writeObject(ctx context.Context, tx fdb.Transaction, id string
157163
typ.Basic.NumChunks = numChunks
158164
case *types.ObjectMeta_Set:
159165
typ.Set.NumChunks = numChunks
166+
case *types.ObjectMeta_SortedSet:
167+
typ.SortedSet.NumChunks = numChunks
160168
case *types.ObjectMeta_ListItem:
161169
typ.ListItem.NumChunks = numChunks
162170
}
@@ -217,10 +225,7 @@ func (s *session) deleteObject(ctx context.Context, tx fdb.Transaction, id strin
217225
span.RecordError(err)
218226
return fmt.Errorf("failed to get object start key: %w", err)
219227
}
220-
tx.ClearRange(fdb.KeyRange{
221-
Begin: begin,
222-
End: end,
223-
})
228+
tx.ClearRange(fdb.KeyRange{Begin: begin, End: end})
224229

225230
metrics.SpanOK(span)
226231
return nil
@@ -275,7 +280,7 @@ func (s *session) handleExists(ctx context.Context, args []resp.Value) (string,
275280
}
276281

277282
existsAny, err := s.fdb.ReadTransact(func(tx fdb.ReadTransaction) (any, error) {
278-
_, meta, err := s.getObjectMeta(ctx, tx, key)
283+
_, meta, err := s.getMeta(ctx, tx, key)
279284
if err != nil {
280285
return nil, err
281286
}
@@ -337,7 +342,7 @@ func (s *session) handleDelete(ctx context.Context, args []resp.Value) (string,
337342
}
338343

339344
existsAny, err := s.fdb.Transact(func(tx fdb.Transaction) (any, error) {
340-
_, meta, err := s.getObjectMeta(ctx, tx, key)
345+
_, meta, err := s.getMeta(ctx, tx, key)
341346
if err != nil {
342347
return false, err
343348
}
@@ -350,6 +355,19 @@ func (s *session) handleDelete(ctx context.Context, args []resp.Value) (string,
350355
return false, err
351356
}
352357

358+
// if we're deleting a sorted set, also clear its score directory
359+
if _, ok := meta.Type.(*types.ObjectMeta_SortedSet); ok {
360+
scoreDir, err := s.sortedSetScoreDir(key)
361+
if err != nil {
362+
return false, fmt.Errorf("failed to get sorted set score directory: %w", err)
363+
}
364+
365+
_, err = scoreDir.Remove(tx, []string{})
366+
if err != nil {
367+
return false, fmt.Errorf("failed to remove score directory: %w", err)
368+
}
369+
}
370+
353371
return true, nil
354372
})
355373
if err != nil {
@@ -578,7 +596,7 @@ func (s *session) handleExpire(ctx context.Context, args []resp.Value) (string,
578596
delta := time.Duration(secs) * time.Second
579597

580598
resAny, err := s.fdb.Transact(func(tx fdb.Transaction) (any, error) {
581-
metaKey, meta, err := s.getObjectMeta(ctx, tx, key)
599+
metaKey, meta, err := s.getMeta(ctx, tx, key)
582600
if err != nil {
583601
return 0, fmt.Errorf("failed to get object meta: %w", err)
584602
}

internal/server/redis/redis.go

Lines changed: 48 additions & 43 deletions
Original file line numberDiff line numberDiff line change
@@ -151,16 +151,17 @@ type session struct {
151151
dirs *Directories
152152

153153
// Only set once the user has authenticated
154-
user *sessionUser
154+
user *userSession
155155
userMu *sync.RWMutex
156156
}
157157

158-
type sessionUser struct {
158+
type userSession struct {
159159
objDir directory.DirectorySubspace
160160
metaDir directory.DirectorySubspace
161161
uidDir directory.DirectorySubspace
162162
reverseUIDDir directory.DirectorySubspace
163163
listDir directory.DirectorySubspace
164+
zsetDir directory.DirectorySubspace
164165

165166
user *types.User
166167
}
@@ -247,6 +248,18 @@ func (s *session) uidKey(id string) (fdb.Key, error) {
247248
return s.user.uidDir.Pack(tuple.Tuple{id}), nil
248249
}
249250

251+
// Returns the FDB directory
252+
func (s *session) sortedSetScoreDir(key string) (directory.DirectorySubspace, error) {
253+
s.userMu.RLock()
254+
defer s.userMu.RUnlock()
255+
256+
if s.user == nil {
257+
return nil, fmt.Errorf("authentication is required")
258+
}
259+
260+
return s.user.zsetDir.CreateOrOpen(s.fdb, []string{key}, nil)
261+
}
262+
250263
// Returns the FDB key of an object in the per-user uid directory
251264
func (s *session) reverseUIDKey(id string) (fdb.Key, error) {
252265
s.userMu.RLock()
@@ -382,6 +395,14 @@ func (s *session) handleCommand(ctx context.Context, cmd *resp.Command) string {
382395
res, err = s.handleSetUnion(ctx, cmd.Args)
383396
case "sdiff":
384397
res, err = s.handleSetDiff(ctx, cmd.Args)
398+
case "zadd":
399+
res, err = s.handleZAdd(ctx, cmd.Args)
400+
case "zcount":
401+
res, err = s.handleZCount(ctx, cmd.Args)
402+
case "zremrangebyscore":
403+
res, err = s.handleZRemRangeByScore(ctx, cmd.Args)
404+
case "zcard":
405+
res, err = s.handleSetCard(ctx, cmd.Args)
385406
case "llen":
386407
res, err = s.handleLLen(ctx, cmd.Args)
387408
case "lpush":
@@ -552,61 +573,43 @@ func userIsAdmin(user *types.User) bool {
552573
return false
553574
}
554575

555-
func (s *session) getMeta(ctx context.Context, tx fdb.ReadTransaction, id string, meta proto.Message) (bool, fdb.Key, error) {
576+
func (s *session) getMeta(ctx context.Context, tx fdb.ReadTransaction, id string) (fdb.Key, *types.ObjectMeta, error) {
556577
ctx, span := s.tracer.Start(ctx, "getMeta") // nolint
557578
defer span.End()
558579

559580
metaKey, err := s.metaKey(id)
560581
if err != nil {
561582
span.RecordError(err)
562-
return false, nil, fmt.Errorf("failed to get meta key: %w", err)
583+
return nil, nil, fmt.Errorf("failed to get meta key: %w", err)
563584
}
564585

565586
metaBuf, err := tx.Get(metaKey).Get()
566587
if err != nil {
567588
span.RecordError(err)
568-
return false, metaKey, fmt.Errorf("failed to get object meta: %w", err)
589+
return nil, nil, fmt.Errorf("failed to get object meta: %w", err)
569590
}
570591

571592
if len(metaBuf) == 0 {
572593
metrics.SpanOK(span)
573-
return false, metaKey, nil
574-
}
575-
576-
if err = proto.Unmarshal(metaBuf, meta); err != nil {
577-
span.RecordError(err)
578-
return false, metaKey, fmt.Errorf("failed to proto unmarshal object meta: %w", err)
594+
return metaKey, nil, nil
579595
}
580596

581-
metrics.SpanOK(span)
582-
return true, metaKey, nil
583-
}
584-
585-
// Returns the associated metadata for an object of any type. Returns nil if it not exist.
586-
func (s *session) getObjectMeta(ctx context.Context, tx fdb.ReadTransaction, id string) (fdb.Key, *types.ObjectMeta, error) {
587-
ctx, span := s.tracer.Start(ctx, "getObjectMeta")
588-
defer span.End()
589-
590597
meta := &types.ObjectMeta{}
591-
exists, key, err := s.getMeta(ctx, tx, id, meta)
592-
if err != nil {
598+
if err = proto.Unmarshal(metaBuf, meta); err != nil {
593599
span.RecordError(err)
594-
return nil, nil, err
595-
}
596-
if !exists {
597-
meta = nil
600+
return nil, nil, fmt.Errorf("failed to proto unmarshal object meta: %w", err)
598601
}
599602

600603
metrics.SpanOK(span)
601-
return key, meta, nil
604+
return metaKey, meta, nil
602605
}
603606

604607
// Returns the associated metadata for an object. Returns nil if the object does not exist.
605608
func (s *session) getListMeta(ctx context.Context, tx fdb.ReadTransaction, id string) (fdb.Key, *types.ListMeta, error) {
606609
ctx, span := s.tracer.Start(ctx, "getListMeta")
607610
defer span.End()
608611

609-
key, objMeta, err := s.getObjectMeta(ctx, tx, id)
612+
key, objMeta, err := s.getMeta(ctx, tx, id)
610613
if err != nil {
611614
span.RecordError(err)
612615
return nil, nil, fmt.Errorf("failed to get object meta: %w", err)
@@ -698,12 +701,13 @@ func (s *session) allocateNewUID(ctx context.Context, tx fdb.Transaction) (uint6
698701
return newUID, nil
699702
}
700703

701-
// Returns the UID for the given member string, creating a new one if it does not exist
702-
func (s *session) getOrAllocateUID(ctx context.Context, tx fdb.Transaction, member string) (uint64, error) {
704+
// Returns the UID for the given member string, creating a new one if it does not exist. If peek is true, it
705+
// will refuse to create a new UID if one does not already exist.
706+
func (s *session) getOrAllocateUID(ctx context.Context, tx fdb.Transaction, member *types.SetMember) (uint64, error) {
703707
ctx, span := s.tracer.Start(ctx, "getOrAllocateUID")
704708
defer span.End()
705709

706-
memberToUIDKey, err := s.reverseUIDKey(member)
710+
memberToUIDKey, err := s.reverseUIDKey(member.Member)
707711
if err != nil {
708712
span.RecordError(err)
709713
return 0, fmt.Errorf("failed to get uid key: %w", err)
@@ -718,13 +722,13 @@ func (s *session) getOrAllocateUID(ctx context.Context, tx fdb.Transaction, memb
718722

719723
if len(val) == 0 {
720724
// allocate a new UID for this member string
721-
uid, err := s.allocateNewUID(ctx, tx)
725+
member.Uid, err = s.allocateNewUID(ctx, tx)
722726
if err != nil {
723727
span.RecordError(err)
724728
return 0, fmt.Errorf("failed to allocate new uid: %w", err)
725729
}
726730

727-
uidStr := strconv.FormatUint(uid, 10)
731+
uidStr := strconv.FormatUint(member.Uid, 10)
728732
uidToMemberKey, err := s.uidKey(uidStr)
729733
if err != nil {
730734
span.RecordError(err)
@@ -733,7 +737,9 @@ func (s *session) getOrAllocateUID(ctx context.Context, tx fdb.Transaction, memb
733737

734738
// store the bi-directional mapping
735739
tx.Set(memberToUIDKey, []byte(uidStr))
736-
tx.Set(uidToMemberKey, []byte(member))
740+
if err := setProtoItem(tx, uidToMemberKey, member); err != nil {
741+
return 0, fmt.Errorf("failed to set uid to member: %w", err)
742+
}
737743

738744
val = []byte(uidStr)
739745
}
@@ -780,30 +786,29 @@ func (s *session) peekUID(ctx context.Context, tx fdb.ReadTransaction, member st
780786
return uid, nil
781787
}
782788

783-
func (s *session) memberFromUID(ctx context.Context, tx fdb.ReadTransaction, uid uint64) (string, error) {
789+
func (s *session) memberFromUID(ctx context.Context, tx fdb.ReadTransaction, uid uint64) (*types.SetMember, error) {
784790
ctx, span := s.tracer.Start(ctx, "memberFromUID") // nolint
785791
defer span.End()
786792

787793
key, err := s.uidKey(strconv.FormatUint(uid, 10))
788794
if err != nil {
789795
span.RecordError(err)
790-
return "", fmt.Errorf("failed to get uid key: %w", err)
796+
return nil, fmt.Errorf("failed to get uid key: %w", err)
791797
}
792798

793-
member, err := tx.Get(key).Get()
799+
member := &types.SetMember{}
800+
exists, err := getProtoItem(tx, key, member)
794801
if err != nil {
795-
span.RecordError(err)
796-
return "", fmt.Errorf("failed to get member for UID %d: %w", uid, err)
802+
return nil, err
797803
}
798-
799-
if len(member) == 0 {
804+
if !exists {
800805
err := fmt.Errorf("uid %d not found", uid)
801806
span.RecordError(err)
802-
return "", err
807+
return nil, err
803808
}
804809

805810
metrics.SpanOK(span)
806-
return string(member), nil
811+
return member, nil
807812
}
808813

809814
func parseVariadicArguments(args []resp.Value) ([]string, error) {

0 commit comments

Comments
 (0)