From 7d08b5c5528b4a93b5be2074195d222d00dfbe33 Mon Sep 17 00:00:00 2001 From: Dan Buch Date: Tue, 14 Oct 2025 12:53:16 -0400 Subject: [PATCH 1/2] Add deadline API for tracked messages to get a slice of tracked message values that have exceeded their deadline within a given time range into the past. --- queue/client.go | 41 +++++- queue/client_test.go | 287 ++++++++++++++++++++++++++++++++++++++++ queue/queue.go | 7 +- queue/writetracking.lua | 7 + 4 files changed, 336 insertions(+), 6 deletions(-) diff --git a/queue/client.go b/queue/client.go index 8858daf..5a0f0bf 100644 --- a/queue/client.go +++ b/queue/client.go @@ -171,11 +171,23 @@ func (c *Client) gcProcessBatch(ctx context.Context, f OnGCFunc, idsToDelete, ke } } - return c.rdb.HDel( - ctx, - MetaCancelationHash, - keysToDelete..., - ).Result() + nDeleted, err := c.rdb.HDel(ctx, MetaCancelationHash, keysToDelete...).Result() + if err != nil { + return nDeleted, err + } + + // NOTE: ZRem requires an explicit []any which cannot be automatically + // converted from a []string. + zremArgs := make([]any, len(idsToDelete)) + for i, id := range idsToDelete { + zremArgs[i] = id + } + + if err := c.rdb.ZRem(ctx, MetaDeadlinesZSet, zremArgs...).Err(); err != nil { + return nDeleted, err + } + + return nDeleted, nil } func (c *Client) callOnGC(ctx context.Context, f OnGCFunc, idsToDelete []string) error { @@ -213,6 +225,25 @@ func (c *Client) callOnGC(ctx context.Context, f OnGCFunc, idsToDelete []string) return f(ctx, trackValues) } +// DeadlineExceeded returns a slice of "track values" that have exceeded their deadline +// within a given duration into the past. The times are truncated to the second because +// the deadlines scored set uses unix timestamps as scores. +func (c *Client) DeadlineExceeded(ctx context.Context, within time.Duration) ([]string, error) { + start, err := c.rdb.Time(ctx).Result() + if err != nil { + return []string{}, err + } + + return c.rdb.ZRangeByScore( + ctx, + MetaDeadlinesZSet, + &redis.ZRangeBy{ + Min: strconv.Itoa(int(start.Add(-within).Unix())), + Max: strconv.Itoa(int(start.Add(1 * time.Second).Unix())), + }, + ).Result() +} + // Len calculates the aggregate length (XLEN) of the queue. It adds up the // lengths of all the streams in the queue. func (c *Client) Len(ctx context.Context, name string) (int64, error) { diff --git a/queue/client_test.go b/queue/client_test.go index aeb46cf..c12ec74 100644 --- a/queue/client_test.go +++ b/queue/client_test.go @@ -501,6 +501,293 @@ func TestClientGCIntegration(t *testing.T) { }) } +func TestClientDeadlineExceededIntegration(t *testing.T) { + ctx := test.Context(t) + rdb := test.Redis(ctx, t) + + ttl := 24 * time.Hour + client := queue.NewTrackingClient(rdb, ttl, "tracking_id") + require.NoError(t, client.Prepare(ctx)) + + t.Run("empty queue", func(t *testing.T) { + require.NoError(t, rdb.FlushAll(ctx).Err()) + + exceeded, err := client.DeadlineExceeded(ctx, 1*time.Hour) + require.NoError(t, err) + assert.Empty(t, exceeded) + }) + + t.Run("no deadlines exceeded", func(t *testing.T) { + require.NoError(t, rdb.FlushAll(ctx).Err()) + + now, err := rdb.Time(ctx).Result() + require.NoError(t, err) + + t.Logf("Writing messages with deadlines in the future") + for i := range 5 { + trackID, err := uuid.NewV7() + require.NoError(t, err) + + _, err = client.Write(ctx, &queue.WriteArgs{ + Name: "testqueue", + Streams: 2, + StreamsPerShard: 1, + ShardKey: []byte(fmt.Sprintf("key-%d", i)), + Values: map[string]any{ + "tracking_id": trackID.String(), + "data": fmt.Sprintf("message-%d", i), + }, + Deadline: now.Add(2 * time.Hour), + }) + require.NoError(t, err) + } + + t.Logf("Checking for exceeded deadlines within past hour - should be empty") + exceeded, err := client.DeadlineExceeded(ctx, 1*time.Hour) + require.NoError(t, err) + assert.Empty(t, exceeded) + }) + + t.Run("recently exceeded deadlines", func(t *testing.T) { + require.NoError(t, rdb.FlushAll(ctx).Err()) + + now, err := rdb.Time(ctx).Result() + require.NoError(t, err) + + trackIDs := []string{} + + t.Logf("Writing messages with deadlines recently exceeded (30 minutes in the past)") + for i := range 5 { + trackID, err := uuid.NewV7() + require.NoError(t, err) + trackIDs = append(trackIDs, trackID.String()) + + _, err = client.Write(ctx, &queue.WriteArgs{ + Name: "testqueue", + Streams: 2, + StreamsPerShard: 1, + ShardKey: []byte(fmt.Sprintf("key-%d", i)), + Values: map[string]any{ + "tracking_id": trackID.String(), + "data": fmt.Sprintf("message-%d", i), + }, + Deadline: now.Add(-30 * time.Minute), + }) + require.NoError(t, err) + } + + t.Logf("Checking for deadlines exceeded within past hour - should get all 5") + exceeded, err := client.DeadlineExceeded(ctx, 1*time.Hour) + require.NoError(t, err) + assert.Len(t, exceeded, 5) + + t.Logf("Verifying all track IDs are present") + exceededSet := make(map[string]bool) + for _, id := range exceeded { + exceededSet[id] = true + } + for _, trackID := range trackIDs { + assert.True(t, exceededSet[trackID], "expected track ID %s to be in exceeded list", trackID) + } + }) + + t.Run("mixed deadlines", func(t *testing.T) { + require.NoError(t, rdb.FlushAll(ctx).Err()) + + now, err := rdb.Time(ctx).Result() + require.NoError(t, err) + + recentTrackIDs := []string{} + oldTrackIDs := []string{} + + t.Logf("Writing 3 messages with recently exceeded deadlines (30 minutes in the past)") + for i := range 3 { + trackID, err := uuid.NewV7() + require.NoError(t, err) + recentTrackIDs = append(recentTrackIDs, trackID.String()) + + _, err = client.Write(ctx, &queue.WriteArgs{ + Name: "testqueue", + Streams: 2, + StreamsPerShard: 1, + ShardKey: []byte(fmt.Sprintf("recent-key-%d", i)), + Values: map[string]any{ + "tracking_id": trackID.String(), + "data": fmt.Sprintf("recent-message-%d", i), + }, + Deadline: now.Add(-30 * time.Minute), + }) + require.NoError(t, err) + } + + t.Logf("Writing 2 messages with old exceeded deadlines (2 hours in the past)") + for i := range 2 { + trackID, err := uuid.NewV7() + require.NoError(t, err) + oldTrackIDs = append(oldTrackIDs, trackID.String()) + + _, err = client.Write(ctx, &queue.WriteArgs{ + Name: "testqueue", + Streams: 2, + StreamsPerShard: 1, + ShardKey: []byte(fmt.Sprintf("old-key-%d", i)), + Values: map[string]any{ + "tracking_id": trackID.String(), + "data": fmt.Sprintf("old-message-%d", i), + }, + Deadline: now.Add(-2 * time.Hour), + }) + require.NoError(t, err) + } + + t.Logf("Checking within past hour - should only get the recent ones") + exceeded, err := client.DeadlineExceeded(ctx, 1*time.Hour) + require.NoError(t, err) + assert.Len(t, exceeded, 3) + + t.Logf("Verifying only recent track IDs are present") + exceededSet := make(map[string]bool) + for _, id := range exceeded { + exceededSet[id] = true + } + for _, trackID := range recentTrackIDs { + assert.True(t, exceededSet[trackID], "expected recent track ID %s to be in exceeded list", trackID) + } + for _, trackID := range oldTrackIDs { + assert.False(t, exceededSet[trackID], "expected old track ID %s to not be in exceeded list", trackID) + } + }) + + t.Run("within window filters correctly", func(t *testing.T) { + require.NoError(t, rdb.FlushAll(ctx).Err()) + + now, err := rdb.Time(ctx).Result() + require.NoError(t, err) + + recentTrackIDs := []string{} + oldTrackIDs := []string{} + + t.Logf("Writing messages with deadlines 30 minutes in the past") + for i := range 3 { + trackID, err := uuid.NewV7() + require.NoError(t, err) + recentTrackIDs = append(recentTrackIDs, trackID.String()) + + _, err = client.Write(ctx, &queue.WriteArgs{ + Name: "testqueue", + Streams: 2, + StreamsPerShard: 1, + ShardKey: []byte(fmt.Sprintf("recent-key-%d", i)), + Values: map[string]any{ + "tracking_id": trackID.String(), + "data": fmt.Sprintf("recent-message-%d", i), + }, + Deadline: now.Add(-30 * time.Minute), + }) + require.NoError(t, err) + } + + t.Logf("Writing messages with deadlines 2 hours in the past") + for i := range 2 { + trackID, err := uuid.NewV7() + require.NoError(t, err) + oldTrackIDs = append(oldTrackIDs, trackID.String()) + + _, err = client.Write(ctx, &queue.WriteArgs{ + Name: "testqueue", + Streams: 2, + StreamsPerShard: 1, + ShardKey: []byte(fmt.Sprintf("old-key-%d", i)), + Values: map[string]any{ + "tracking_id": trackID.String(), + "data": fmt.Sprintf("old-message-%d", i), + }, + Deadline: now.Add(-2 * time.Hour), + }) + require.NoError(t, err) + } + + t.Logf("Checking within past 1 hour - should only get the recent ones") + exceeded, err := client.DeadlineExceeded(ctx, 1*time.Hour) + require.NoError(t, err) + assert.Len(t, exceeded, 3) + + exceededSet := make(map[string]bool) + for _, id := range exceeded { + exceededSet[id] = true + } + for _, trackID := range recentTrackIDs { + assert.True(t, exceededSet[trackID], "expected recent track ID %s to be in exceeded list", trackID) + } + for _, trackID := range oldTrackIDs { + assert.False(t, exceededSet[trackID], "expected old track ID %s to not be in exceeded list", trackID) + } + + t.Logf("Checking within past 3 hours - should get all messages") + exceeded, err = client.DeadlineExceeded(ctx, 3*time.Hour) + require.NoError(t, err) + assert.Len(t, exceeded, 5) + }) + + t.Run("zero duration", func(t *testing.T) { + require.NoError(t, rdb.FlushAll(ctx).Err()) + + now, err := rdb.Time(ctx).Result() + require.NoError(t, err) + + t.Logf("Writing a message with a deadline 1 second in the past") + trackID, err := uuid.NewV7() + require.NoError(t, err) + + _, err = client.Write(ctx, &queue.WriteArgs{ + Name: "testqueue", + Streams: 2, + StreamsPerShard: 1, + ShardKey: []byte("test-key"), + Values: map[string]any{ + "tracking_id": trackID.String(), + "data": "test-message", + }, + Deadline: now.Add(-1 * time.Second), + }) + require.NoError(t, err) + + t.Logf("Checking with zero duration - should not return anything") + exceeded, err := client.DeadlineExceeded(ctx, 0) + require.NoError(t, err) + assert.Empty(t, exceeded) + }) + + t.Run("boundary at now", func(t *testing.T) { + require.NoError(t, rdb.FlushAll(ctx).Err()) + + now, err := rdb.Time(ctx).Result() + require.NoError(t, err) + + justExpiredID, err := uuid.NewV7() + require.NoError(t, err) + + t.Logf("Writing a message with deadline right around now") + _, err = client.Write(ctx, &queue.WriteArgs{ + Name: "testqueue", + Streams: 2, + StreamsPerShard: 1, + ShardKey: []byte("test-key"), + Values: map[string]any{ + "tracking_id": justExpiredID.String(), + "data": "test-message", + }, + Deadline: now, + }) + require.NoError(t, err) + + t.Logf("Checking with 5 second window - should include items at or just before now") + exceeded, err := client.DeadlineExceeded(ctx, 5*time.Second) + require.NoError(t, err) + assert.Contains(t, exceeded, justExpiredID.String()) + }) +} + // TestPickupLatencyIntegration runs a test with a mostly-empty queue -- by // running artificially slow producers and full-speed consumers -- to ensure // that the blocking read operation has low latency. diff --git a/queue/queue.go b/queue/queue.go index c69aafe..0a78736 100644 --- a/queue/queue.go +++ b/queue/queue.go @@ -65,7 +65,11 @@ var ( writeTrackingCmd string writeTrackingScript = redis.NewScript( strings.ReplaceAll( - writeTrackingCmd, + strings.ReplaceAll( + writeTrackingCmd, + "__META_DEADLINES_ZSET__", + MetaDeadlinesZSet, + ), "__META_CANCELATION_HASH__", MetaCancelationHash, ), @@ -74,6 +78,7 @@ var ( const ( MetaCancelationHash = "meta:cancelation" + MetaDeadlinesZSet = "meta:deadlines" ) func prepare(ctx context.Context, rdb redis.Cmdable) error { diff --git a/queue/writetracking.lua b/queue/writetracking.lua index 340a03f..e2b7255 100644 --- a/queue/writetracking.lua +++ b/queue/writetracking.lua @@ -138,6 +138,13 @@ if track_value ~= '' then cancelation_expiry_key, '1' ) + + redis.call( + 'ZADD', + '__META_DEADLINES_ZSET__', + deadline, + track_value + ) end -- Set expiry on selected stream + meta/notifications keys From 1699e64f720d28b79b2362800426cb22a245a053 Mon Sep 17 00:00:00 2001 From: Dan Buch Date: Tue, 14 Oct 2025 16:52:12 -0400 Subject: [PATCH 2/2] Switch name and impl to "cordon" track values that have exceeded deadline so that the `GC` can delete them without double checking. --- queue/client.go | 85 ++++++++++++++++++++++++++++++++++++++------ queue/client_test.go | 31 ++++++++++------ queue/queue.go | 5 +-- 3 files changed, 98 insertions(+), 23 deletions(-) diff --git a/queue/client.go b/queue/client.go index 5a0f0bf..04a3a45 100644 --- a/queue/client.go +++ b/queue/client.go @@ -69,6 +69,36 @@ type OnGCFunc func(ctx context.Context, trackValues []string) error // of the current server time as a way of limiting the keyspace scanned. As a special // case, any value <= -1 will result in all keys being scanned. func (c *Client) GC(ctx context.Context, nTimeDigits int, f OnGCFunc) (uint64, uint64, error) { + pipe := c.rdb.Pipeline() + presortedGarbageCmd := pipe.SMembers(ctx, MetaPresortedGarbageSet) + pipe.Del(ctx, MetaPresortedGarbageSet) + + if _, err := pipe.Exec(ctx); err != nil { + return 0, 0, err + } + + presortedGarbage := presortedGarbageCmd.Val() + + if len(presortedGarbage) > 0 { + toDelete := []string{} + for _, key := range presortedGarbage { + trackValue, deadline, ok := strings.Cut(key, ":") + if !ok { + continue + } + + toDelete = append( + toDelete, + trackValue, + fmt.Sprintf("%s:expiry:%s", trackValue, deadline), + ) + } + + if err := c.rdb.HDel(ctx, MetaCancelationHash, toDelete...).Err(); err != nil { + return 0, 0, err + } + } + now, err := c.rdb.Time(ctx).Result() if err != nil { return 0, 0, err @@ -225,23 +255,58 @@ func (c *Client) callOnGC(ctx context.Context, f OnGCFunc, idsToDelete []string) return f(ctx, trackValues) } -// DeadlineExceeded returns a slice of "track values" that have exceeded their deadline -// within a given duration into the past. The times are truncated to the second because -// the deadlines scored set uses unix timestamps as scores. -func (c *Client) DeadlineExceeded(ctx context.Context, within time.Duration) ([]string, error) { +// CordonDeadlineExceeded selects a chunk of "track values" that have exceeded their +// deadline within a given duration into the past, moves them into the "presorted garbage" +// set for use with `GC`, and returns them as a slice. The times are truncated to the +// second because the deadlines scored set uses unix timestamps as scores. +func (c *Client) CordonDeadlineExceeded(ctx context.Context, within time.Duration) ([]string, error) { start, err := c.rdb.Time(ctx).Result() if err != nil { return []string{}, err } - return c.rdb.ZRangeByScore( + pipe := c.rdb.Pipeline() + zRangeOpts := &redis.ZRangeBy{ + Min: strconv.Itoa(int(start.Add(-within).Unix())), + Max: strconv.Itoa(int(start.Add(1 * time.Second).Unix())), + } + + zRangeCmd := c.rdb.ZRangeByScoreWithScores( ctx, MetaDeadlinesZSet, - &redis.ZRangeBy{ - Min: strconv.Itoa(int(start.Add(-within).Unix())), - Max: strconv.Itoa(int(start.Add(1 * time.Second).Unix())), - }, - ).Result() + zRangeOpts, + ) + + c.rdb.ZRemRangeByScore( + ctx, + MetaDeadlinesZSet, + zRangeOpts.Min, + zRangeOpts.Max, + ) + + if _, err := pipe.Exec(ctx); err != nil { + return []string{}, err + } + + trackValues := zRangeCmd.Val() + + if len(trackValues) == 0 { + return []string{}, nil + } + + ret := make([]string, len(trackValues)) + sMembers := make([]any, len(trackValues)) + + for i, trackValue := range trackValues { + sMembers[i] = fmt.Sprintf("%s:%d", trackValue.Member, int(trackValue.Score)) + ret[i] = fmt.Sprintf("%v", trackValue.Member) + } + + if err := c.rdb.SAdd(ctx, MetaPresortedGarbageSet, sMembers...).Err(); err != nil { + return ret, err + } + + return ret, nil } // Len calculates the aggregate length (XLEN) of the queue. It adds up the diff --git a/queue/client_test.go b/queue/client_test.go index c12ec74..5caed5a 100644 --- a/queue/client_test.go +++ b/queue/client_test.go @@ -501,7 +501,7 @@ func TestClientGCIntegration(t *testing.T) { }) } -func TestClientDeadlineExceededIntegration(t *testing.T) { +func TestClientCordonDeadlineExceededIntegration(t *testing.T) { ctx := test.Context(t) rdb := test.Redis(ctx, t) @@ -512,7 +512,7 @@ func TestClientDeadlineExceededIntegration(t *testing.T) { t.Run("empty queue", func(t *testing.T) { require.NoError(t, rdb.FlushAll(ctx).Err()) - exceeded, err := client.DeadlineExceeded(ctx, 1*time.Hour) + exceeded, err := client.CordonDeadlineExceeded(ctx, 1*time.Hour) require.NoError(t, err) assert.Empty(t, exceeded) }) @@ -543,7 +543,7 @@ func TestClientDeadlineExceededIntegration(t *testing.T) { } t.Logf("Checking for exceeded deadlines within past hour - should be empty") - exceeded, err := client.DeadlineExceeded(ctx, 1*time.Hour) + exceeded, err := client.CordonDeadlineExceeded(ctx, 1*time.Hour) require.NoError(t, err) assert.Empty(t, exceeded) }) @@ -577,7 +577,7 @@ func TestClientDeadlineExceededIntegration(t *testing.T) { } t.Logf("Checking for deadlines exceeded within past hour - should get all 5") - exceeded, err := client.DeadlineExceeded(ctx, 1*time.Hour) + exceeded, err := client.CordonDeadlineExceeded(ctx, 1*time.Hour) require.NoError(t, err) assert.Len(t, exceeded, 5) @@ -641,7 +641,7 @@ func TestClientDeadlineExceededIntegration(t *testing.T) { } t.Logf("Checking within past hour - should only get the recent ones") - exceeded, err := client.DeadlineExceeded(ctx, 1*time.Hour) + exceeded, err := client.CordonDeadlineExceeded(ctx, 1*time.Hour) require.NoError(t, err) assert.Len(t, exceeded, 3) @@ -708,7 +708,7 @@ func TestClientDeadlineExceededIntegration(t *testing.T) { } t.Logf("Checking within past 1 hour - should only get the recent ones") - exceeded, err := client.DeadlineExceeded(ctx, 1*time.Hour) + exceeded, err := client.CordonDeadlineExceeded(ctx, 1*time.Hour) require.NoError(t, err) assert.Len(t, exceeded, 3) @@ -723,10 +723,19 @@ func TestClientDeadlineExceededIntegration(t *testing.T) { assert.False(t, exceededSet[trackID], "expected old track ID %s to not be in exceeded list", trackID) } - t.Logf("Checking within past 3 hours - should get all messages") - exceeded, err = client.DeadlineExceeded(ctx, 3*time.Hour) + t.Logf("Checking within past 3 hours - should get remaining 2 old messages") + exceeded, err = client.CordonDeadlineExceeded(ctx, 3*time.Hour) require.NoError(t, err) - assert.Len(t, exceeded, 5) + assert.Len(t, exceeded, 2) + + t.Logf("Verifying only old track IDs are present") + exceededSet = make(map[string]bool) + for _, id := range exceeded { + exceededSet[id] = true + } + for _, trackID := range oldTrackIDs { + assert.True(t, exceededSet[trackID], "expected old track ID %s to be in exceeded list", trackID) + } }) t.Run("zero duration", func(t *testing.T) { @@ -753,7 +762,7 @@ func TestClientDeadlineExceededIntegration(t *testing.T) { require.NoError(t, err) t.Logf("Checking with zero duration - should not return anything") - exceeded, err := client.DeadlineExceeded(ctx, 0) + exceeded, err := client.CordonDeadlineExceeded(ctx, 0) require.NoError(t, err) assert.Empty(t, exceeded) }) @@ -782,7 +791,7 @@ func TestClientDeadlineExceededIntegration(t *testing.T) { require.NoError(t, err) t.Logf("Checking with 5 second window - should include items at or just before now") - exceeded, err := client.DeadlineExceeded(ctx, 5*time.Second) + exceeded, err := client.CordonDeadlineExceeded(ctx, 5*time.Second) require.NoError(t, err) assert.Contains(t, exceeded, justExpiredID.String()) }) diff --git a/queue/queue.go b/queue/queue.go index 0a78736..e531989 100644 --- a/queue/queue.go +++ b/queue/queue.go @@ -77,8 +77,9 @@ var ( ) const ( - MetaCancelationHash = "meta:cancelation" - MetaDeadlinesZSet = "meta:deadlines" + MetaCancelationHash = "meta:cancelation" + MetaDeadlinesZSet = "meta:deadlines" + MetaPresortedGarbageSet = "meta:presortedgarbage" ) func prepare(ctx context.Context, rdb redis.Cmdable) error {