Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
28 changes: 28 additions & 0 deletions platform/extension/messagequeue/mysql/mock_stores.go

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

20 changes: 20 additions & 0 deletions platform/extension/messagequeue/mysql/offset_store.go
Original file line number Diff line number Diff line change
Expand Up @@ -132,3 +132,23 @@ func (s *sqloffsetStore) GetMinAckedOffset(ctx context.Context, topic string, pa

return minOffset, true, nil
}

// DeleteOffset removes one consumer group's offset row for a partition.
// Idempotent — see the offsetStore interface doc.
func (s *sqloffsetStore) DeleteOffset(ctx context.Context, topic string, partitionKey string, consumerGroup string) (retErr error) {
op := metrics.Begin(s.scope, "delete_offset", metrics.StorageLatencyBuckets,
metrics.NewTag("topic", topic),
metrics.NewTag("partition_key", partitionKey),
metrics.NewTag("consumer_group", consumerGroup))
defer func() { op.Complete(retErr) }()

_, err := s.db.ExecContext(ctx, fmt.Sprintf(`
DELETE FROM %s WHERE consumer_group = ? AND topic = ? AND partition_key = ?
`, OffsetsTableName), consumerGroup, topic, partitionKey)

if err != nil {
return fmt.Errorf("delete offset topic=%s partition=%s: %w", topic, partitionKey, err)
}

return nil
}
51 changes: 51 additions & 0 deletions platform/extension/messagequeue/mysql/offset_store_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -188,3 +188,54 @@ func TestOffsetStore_GetMinAckedOffset(t *testing.T) {
})
}
}

func TestOffsetStore_DeleteOffset(t *testing.T) {
tests := []struct {
name string
setup func(mock sqlmock.Sqlmock)
wantErr bool
}{
{
name: "deletes the consumer group's offset row",
setup: func(mock sqlmock.Sqlmock) {
mock.ExpectExec("DELETE FROM queue_offsets").
WithArgs(testConsumerGroup, "test_topic", "part-1").
WillReturnResult(sqlmock.NewResult(0, 1))
},
},
{
name: "idempotent - row already gone",
setup: func(mock sqlmock.Sqlmock) {
mock.ExpectExec("DELETE FROM queue_offsets").
WithArgs(testConsumerGroup, "test_topic", "part-1").
WillReturnResult(sqlmock.NewResult(0, 0))
},
},
{
name: "database error",
setup: func(mock sqlmock.Sqlmock) {
mock.ExpectExec("DELETE FROM queue_offsets").
WithArgs(testConsumerGroup, "test_topic", "part-1").
WillReturnError(fmt.Errorf("db error"))
},
wantErr: true,
},
}

for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
db, mock, store := setupoffsetStoreTest(t)
defer db.Close()

tt.setup(mock)

err := store.DeleteOffset(context.Background(), "test_topic", "part-1", testConsumerGroup)
if tt.wantErr {
require.Error(t, err)
} else {
require.NoError(t, err)
}
require.NoError(t, mock.ExpectationsWereMet())
})
}
}
30 changes: 30 additions & 0 deletions platform/extension/messagequeue/mysql/partition_lease_store.go
Original file line number Diff line number Diff line change
Expand Up @@ -228,6 +228,36 @@ func (s *sqlpartitionLeaseStore) GetAllLeases(ctx context.Context, topic string,
return leases, nil
}

// PurgeStale deletes lease rows not renewed within olderThanMs. See the
// partitionLeaseStore interface doc.
func (s *sqlpartitionLeaseStore) PurgeStale(ctx context.Context, topic string, consumerGroup string, olderThanMs int64) (retErr error) {
op := metrics.Begin(s.scope, "purge_stale", metrics.StorageLatencyBuckets, metrics.NewTag("topic", topic))
defer func() { op.Complete(retErr) }()

threshold := currentTimeMillis() - olderThanMs

result, err := s.db.ExecContext(ctx, fmt.Sprintf(`
DELETE FROM %s
WHERE consumer_group = ? AND topic = ? AND lease_renewed_at < ?
`, PartitionLeasesTableName), consumerGroup, topic, threshold)

if err != nil {
return fmt.Errorf("failed to purge stale leases: %w", err)
}

// RowsAffected error is swallowed because the DELETE itself succeeded;
// the count is for observability only.
if deleted, err := result.RowsAffected(); err == nil && deleted > 0 {
metrics.NamedCounter(s.scope, "purge_stale", "rows_deleted", deleted, metrics.NewTag("topic", topic))
s.logger.Debugw("purged stale leases",
logTopic, topic,
"deleted", deleted,
)
}

return nil
}

// DiscoverAndAcquirePartitions discovers partitions from messages table and tries to acquire leases.
// Returns the number of new leases acquired and the full list of discovered partitions.
// maxPartitions limits how many total partitions this subscriber can own (0 = unlimited)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@ package mysql
import (
"context"
"database/sql"
"fmt"
"testing"
"time"

Expand Down Expand Up @@ -412,3 +413,54 @@ func TestPartitionLeaseStore_DiscoverAndAcquirePartitions(t *testing.T) {
})
}
}

func TestPartitionLeaseStore_PurgeStale(t *testing.T) {
tests := []struct {
name string
setup func(mock sqlmock.Sqlmock)
wantErr bool
}{
{
name: "deletes rows older than threshold",
setup: func(mock sqlmock.Sqlmock) {
mock.ExpectExec("DELETE FROM queue_partition_leases").
WithArgs(testConsumerGroup, "test_topic", sqlmock.AnyArg()).
WillReturnResult(sqlmock.NewResult(0, 2))
},
},
{
name: "no stale rows is a no-op",
setup: func(mock sqlmock.Sqlmock) {
mock.ExpectExec("DELETE FROM queue_partition_leases").
WithArgs(testConsumerGroup, "test_topic", sqlmock.AnyArg()).
WillReturnResult(sqlmock.NewResult(0, 0))
},
},
{
name: "database error",
setup: func(mock sqlmock.Sqlmock) {
mock.ExpectExec("DELETE FROM queue_partition_leases").
WithArgs(testConsumerGroup, "test_topic", sqlmock.AnyArg()).
WillReturnError(fmt.Errorf("db error"))
},
wantErr: true,
},
}

for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
db, mock, store := setuppartitionLeaseStoreTest(t)
defer db.Close()

tt.setup(mock)

err := store.PurgeStale(context.Background(), "test_topic", testConsumerGroup, 300_000)
if tt.wantErr {
require.Error(t, err)
} else {
require.NoError(t, err)
}
require.NoError(t, mock.ExpectationsWereMet())
})
}
}
14 changes: 14 additions & 0 deletions platform/extension/messagequeue/mysql/stores.go
Original file line number Diff line number Diff line change
Expand Up @@ -101,6 +101,12 @@ type offsetStore interface {
// Used by the subscriber to compute the GC threshold without messageStore
// needing to query the offsets table.
GetMinAckedOffset(ctx context.Context, topic string, partitionKey string) (offset int64, found bool, err error)

// DeleteOffset removes one consumer group's offset row for a partition.
// Callers use this when retiring a fully-drained partition; Initialize
// recreates the row if the partition ever receives messages again.
// Idempotent: no-op if the row is already gone.
DeleteOffset(ctx context.Context, topic string, partitionKey string, consumerGroup string) error
}

// leaseInfo describes one partition's current lease row (internal use only)
Expand Down Expand Up @@ -137,6 +143,14 @@ type partitionLeaseStore interface {
// instead of write-probing every lease row each discovery tick.
GetAllLeases(ctx context.Context, topic string, consumerGroup string) ([]leaseInfo, error)

// PurgeStale deletes lease rows not renewed within olderThanMs. Backstop
// for holders that crashed while owning a drained partition: acquisition
// only probes discovered partitions, so a stale lease on a partition
// with no messages is otherwise never refreshed or removed. Deleting a
// stale row is equivalent to lease expiry — a concurrent renewal makes
// the row fresh and the age predicate skips it.
PurgeStale(ctx context.Context, topic string, consumerGroup string, olderThanMs int64) error

// DiscoverAndAcquirePartitions discovers partitions from messages table and tries to acquire leases.
// Returns the number of new leases acquired and the full list of discovered partitions.
// leaseDurationMs is how long the lease is valid (in milliseconds)
Expand Down
Loading
Loading