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
14 changes: 14 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.

12 changes: 11 additions & 1 deletion platform/extension/messagequeue/mysql/stores.go
Original file line number Diff line number Diff line change
Expand Up @@ -154,8 +154,18 @@ type subscriberHeartbeatStore interface {
// within this duration are considered dead.
ActiveSubscribers(ctx context.Context, topic string, consumerGroup string, staleDurationMs int64) ([]string, error)

// Deregister removes a subscriber's heartbeat entry
// Deregister removes a subscriber's heartbeat row. Hard delete: the row
// is not needed once the subscriber is gone, and subscriber names are
// unique per process (hostname-pid), so rows would otherwise accumulate
// forever across deploys. Re-subscribing re-inserts via Heartbeat.
Deregister(ctx context.Context, topic string, subscriberName string, consumerGroup string) error

// PurgeStale deletes heartbeat rows whose last heartbeat is older than
// olderThanMs. Backstop for subscribers that never deregistered
// (crashes, SIGKILL): without it the table grows monotonically since
// every process registers under a fresh name. Deleting a live-but-stalled
// subscriber's row is harmless — its next heartbeat re-inserts it.
PurgeStale(ctx context.Context, topic string, consumerGroup string, olderThanMs int64) error
}

// DeliveryState represents the full per-message delivery tracking state.
Expand Down
15 changes: 15 additions & 0 deletions platform/extension/messagequeue/mysql/subscriber.go
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,14 @@ const (
// queries when many partitions are idle (e.g., 50 idle partitions at 100ms
// poll interval = 500 GC queries/sec without throttling).
gcIdleTickInterval = 100

// heartbeatPurgeAfterLeaseDurations sets the age threshold for purging
// abandoned heartbeat rows, as a multiple of LeaseDurationMs (10x = 5min
// at defaults). Well past every transient window in the protocol — a row
// that stale belongs to a subscriber that crashed without deregistering.
// Purging a live-but-stalled subscriber's row is harmless: its next
// heartbeat re-inserts it.
heartbeatPurgeAfterLeaseDurations = 10
)

// HookSignal identifies the type of subscriber lifecycle event.
Expand Down Expand Up @@ -565,6 +573,13 @@ func (s *subscriber) managePartitions(ctx context.Context, sub *subscription) {
if err := s.sendHeartbeat(ctx, sub); err != nil {
s.logger.Errorw("periodic heartbeat failed", append(logFields, "error", err)...)
}
// Purge heartbeat rows abandoned by subscribers that never
// deregistered (crashes) — without this the table grows
// monotonically, since every process registers under a fresh
// hostname-pid name.
if err := s.heartbeatStore.PurgeStale(ctx, sub.topic, cfg.ConsumerGroup, heartbeatPurgeAfterLeaseDurations*cfg.LeaseDurationMs); err != nil {
s.logger.Errorw("stale heartbeat purge failed", append(logFields, "error", err)...)
}
s.emitSignal(SignalPartitionUpdate)

case <-discoveryTicker.C:
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -101,18 +101,16 @@ func (s *sqlSubscriberHeartbeatStore) ActiveSubscribers(ctx context.Context, top
return names, nil
}

// Deregister soft-deletes a subscriber's heartbeat entry by setting deregistered_at.
// Idempotent: no-op if already deregistered.
// Deregister removes a subscriber's heartbeat row (hard delete — see the
// subscriberHeartbeatStore interface doc). Idempotent: no-op if already gone.
func (s *sqlSubscriberHeartbeatStore) Deregister(ctx context.Context, topic string, subscriberName string, consumerGroup string) (retErr error) {
op := metrics.Begin(s.scope, "deregister", metrics.StorageLatencyBuckets)
defer func() { op.Complete(retErr) }()

now := s.nowFunc().UnixMilli()

_, err := s.db.ExecContext(ctx, fmt.Sprintf(`
UPDATE %s SET deregistered_at = ?
WHERE consumer_group = ? AND topic = ? AND subscriber_name = ? AND deregistered_at = 0
`, SubscriberHeartbeatsTableName), now, consumerGroup, topic, subscriberName)
DELETE FROM %s
WHERE consumer_group = ? AND topic = ? AND subscriber_name = ?
`, SubscriberHeartbeatsTableName), consumerGroup, topic, subscriberName)

if err != nil {
return fmt.Errorf("failed to deregister subscriber: %w", err)
Expand All @@ -125,3 +123,33 @@ func (s *sqlSubscriberHeartbeatStore) Deregister(ctx context.Context, topic stri

return nil
}

// PurgeStale deletes heartbeat rows older than olderThanMs for the topic and
// consumer group. See the subscriberHeartbeatStore interface doc.
func (s *sqlSubscriberHeartbeatStore) PurgeStale(ctx context.Context, topic string, consumerGroup string, olderThanMs int64) (retErr error) {
op := metrics.Begin(s.scope, "purge_stale", metrics.StorageLatencyBuckets)
defer func() { op.Complete(retErr) }()

threshold := s.nowFunc().UnixMilli() - olderThanMs

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

if err != nil {
return fmt.Errorf("failed to purge stale heartbeats: %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 heartbeats",
logTopic, topic,
"deleted", deleted,
)
}

return nil
}
Original file line number Diff line number Diff line change
Expand Up @@ -172,22 +172,74 @@ func TestSubscriberHeartbeatStore_ActiveSubscribers_ExcludesDeregistered(t *test
require.NoError(t, mock.ExpectationsWereMet())
}

func TestSubscriberHeartbeatStore_Deregister_SoftDelete(t *testing.T) {
func TestSubscriberHeartbeatStore_Deregister_HardDelete(t *testing.T) {
db, mock, store := setupSubscriberHeartbeatStoreTest(t)
defer db.Close()

ctx := context.Background()

// Verify deregister uses UPDATE (not DELETE) and targets only active rows (deregistered_at = 0)
mock.ExpectExec(`UPDATE queue_subscriber_heartbeats SET deregistered_at.*AND deregistered_at = 0`).
WithArgs(sqlmock.AnyArg(), testConsumerGroup, "test_topic", testSubscriberName).
// Verify deregister deletes the row outright — subscriber names are
// unique per process, so soft-deleted rows would accumulate forever.
mock.ExpectExec(`DELETE FROM queue_subscriber_heartbeats`).
WithArgs(testConsumerGroup, "test_topic", testSubscriberName).
WillReturnResult(sqlmock.NewResult(0, 1))

err := store.Deregister(ctx, "test_topic", testSubscriberName, testConsumerGroup)
require.NoError(t, err)
require.NoError(t, mock.ExpectationsWereMet())
}

func TestSubscriberHeartbeatStore_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_subscriber_heartbeats`).
WithArgs(testConsumerGroup, "test_topic", sqlmock.AnyArg()).
WillReturnResult(sqlmock.NewResult(0, 3))
},
},
{
name: "no stale rows is a no-op",
setup: func(mock sqlmock.Sqlmock) {
mock.ExpectExec(`DELETE FROM queue_subscriber_heartbeats`).
WithArgs(testConsumerGroup, "test_topic", sqlmock.AnyArg()).
WillReturnResult(sqlmock.NewResult(0, 0))
},
},
{
name: "database error",
setup: func(mock sqlmock.Sqlmock) {
mock.ExpectExec(`DELETE FROM queue_subscriber_heartbeats`).
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 := setupSubscriberHeartbeatStoreTest(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())
})
}
}

func TestSubscriberHeartbeatStore_ReRegistration(t *testing.T) {
db, mock, store := setupSubscriberHeartbeatStoreTest(t)
defer db.Close()
Expand All @@ -199,15 +251,15 @@ func TestSubscriberHeartbeatStore_ReRegistration(t *testing.T) {
WithArgs(testConsumerGroup, "test_topic", testSubscriberName, sqlmock.AnyArg()).
WillReturnResult(sqlmock.NewResult(1, 1))

// Step 2: Deregister soft-deletes the subscriber
mock.ExpectExec("UPDATE queue_subscriber_heartbeats").
WithArgs(sqlmock.AnyArg(), testConsumerGroup, "test_topic", testSubscriberName).
// Step 2: Deregister deletes the subscriber's row
mock.ExpectExec("DELETE FROM queue_subscriber_heartbeats").
WithArgs(testConsumerGroup, "test_topic", testSubscriberName).
WillReturnResult(sqlmock.NewResult(0, 1))

// Step 3: Heartbeat again re-registers (ON DUPLICATE KEY UPDATE resets deregistered_at = 0)
// Step 3: Heartbeat again re-registers with a fresh insert
mock.ExpectExec("INSERT INTO queue_subscriber_heartbeats").
WithArgs(testConsumerGroup, "test_topic", testSubscriberName, sqlmock.AnyArg()).
WillReturnResult(sqlmock.NewResult(0, 2)) // 2 = ON DUPLICATE KEY UPDATE
WillReturnResult(sqlmock.NewResult(1, 1))

err := store.Heartbeat(ctx, "test_topic", testSubscriberName, testConsumerGroup)
require.NoError(t, err)
Expand All @@ -230,26 +282,26 @@ func TestSubscriberHeartbeatStore_Deregister(t *testing.T) {
{
name: "successfully deregister",
setup: func(mock sqlmock.Sqlmock) {
mock.ExpectExec("UPDATE queue_subscriber_heartbeats").
WithArgs(sqlmock.AnyArg(), testConsumerGroup, "test_topic", testSubscriberName).
mock.ExpectExec("DELETE FROM queue_subscriber_heartbeats").
WithArgs(testConsumerGroup, "test_topic", testSubscriberName).
WillReturnResult(sqlmock.NewResult(0, 1))
},
wantErr: false,
},
{
name: "idempotent - already deregistered",
setup: func(mock sqlmock.Sqlmock) {
mock.ExpectExec("UPDATE queue_subscriber_heartbeats").
WithArgs(sqlmock.AnyArg(), testConsumerGroup, "test_topic", testSubscriberName).
mock.ExpectExec("DELETE FROM queue_subscriber_heartbeats").
WithArgs(testConsumerGroup, "test_topic", testSubscriberName).
WillReturnResult(sqlmock.NewResult(0, 0))
},
wantErr: false,
},
{
name: "database error",
setup: func(mock sqlmock.Sqlmock) {
mock.ExpectExec("UPDATE queue_subscriber_heartbeats").
WithArgs(sqlmock.AnyArg(), testConsumerGroup, "test_topic", testSubscriberName).
mock.ExpectExec("DELETE FROM queue_subscriber_heartbeats").
WithArgs(testConsumerGroup, "test_topic", testSubscriberName).
WillReturnError(fmt.Errorf("db error"))
},
wantErr: true,
Expand Down
1 change: 1 addition & 0 deletions platform/extension/messagequeue/mysql/subscriber_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,7 @@ func newTestHeartbeatStore(ctrl *gomock.Controller) *MocksubscriberHeartbeatStor
mockHB.EXPECT().Heartbeat(gomock.Any(), gomock.Any(), gomock.Any(), gomock.Any()).Return(nil).AnyTimes()
mockHB.EXPECT().ActiveSubscribers(gomock.Any(), gomock.Any(), gomock.Any(), gomock.Any()).Return([]string{"self"}, nil).AnyTimes()
mockHB.EXPECT().Deregister(gomock.Any(), gomock.Any(), gomock.Any(), gomock.Any()).Return(nil).AnyTimes()
mockHB.EXPECT().PurgeStale(gomock.Any(), gomock.Any(), gomock.Any(), gomock.Any()).Return(nil).AnyTimes()
return mockHB
}

Expand Down
9 changes: 9 additions & 0 deletions test/integration/extension/messagequeue/mysql/queue_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -1814,6 +1814,15 @@ func (s *SQLQueueIntegrationSuite) TestRebalance_SubscriberLeaves() {
return len(leases["s1"]) == 4
}, "S1 should reacquire all 4 partitions after S2 leaves")

// Deregistration hard-deletes the heartbeat row — subscriber names are
// unique per process, so rows would otherwise accumulate forever.
var s2Rows int
require.NoError(t, s.db.QueryRowContext(s.ctx, `
SELECT COUNT(*) FROM queue_subscriber_heartbeats
WHERE consumer_group = ? AND topic = ? AND subscriber_name = ?
`, consumerGroup, topic, "s2").Scan(&s2Rows))
assert.Equal(t, 0, s2Rows, "closed subscriber's heartbeat row must be deleted")

t.Logf("Subscriber leave verified: S1 owns all 4 partitions after S2 departed")
}

Expand Down