diff --git a/platform/extension/messagequeue/mysql/mock_stores.go b/platform/extension/messagequeue/mysql/mock_stores.go index 28c446bf0..87417b0cc 100644 --- a/platform/extension/messagequeue/mysql/mock_stores.go +++ b/platform/extension/messagequeue/mysql/mock_stores.go @@ -391,6 +391,20 @@ func (mr *MocksubscriberHeartbeatStoreMockRecorder) Heartbeat(ctx, topic, subscr return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "Heartbeat", reflect.TypeOf((*MocksubscriberHeartbeatStore)(nil).Heartbeat), ctx, topic, subscriberName, consumerGroup) } +// PurgeStale mocks base method. +func (m *MocksubscriberHeartbeatStore) PurgeStale(ctx context.Context, topic, consumerGroup string, olderThanMs int64) error { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "PurgeStale", ctx, topic, consumerGroup, olderThanMs) + ret0, _ := ret[0].(error) + return ret0 +} + +// PurgeStale indicates an expected call of PurgeStale. +func (mr *MocksubscriberHeartbeatStoreMockRecorder) PurgeStale(ctx, topic, consumerGroup, olderThanMs any) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "PurgeStale", reflect.TypeOf((*MocksubscriberHeartbeatStore)(nil).PurgeStale), ctx, topic, consumerGroup, olderThanMs) +} + // MockdeliveryStateStore is a mock of deliveryStateStore interface. type MockdeliveryStateStore struct { ctrl *gomock.Controller diff --git a/platform/extension/messagequeue/mysql/stores.go b/platform/extension/messagequeue/mysql/stores.go index 9b28c26d3..c739f2fb5 100644 --- a/platform/extension/messagequeue/mysql/stores.go +++ b/platform/extension/messagequeue/mysql/stores.go @@ -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. diff --git a/platform/extension/messagequeue/mysql/subscriber.go b/platform/extension/messagequeue/mysql/subscriber.go index 6277c7fe4..12791588d 100644 --- a/platform/extension/messagequeue/mysql/subscriber.go +++ b/platform/extension/messagequeue/mysql/subscriber.go @@ -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. @@ -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: diff --git a/platform/extension/messagequeue/mysql/subscriber_heartbeat_store.go b/platform/extension/messagequeue/mysql/subscriber_heartbeat_store.go index 7a2c18fff..f12a2f7c7 100644 --- a/platform/extension/messagequeue/mysql/subscriber_heartbeat_store.go +++ b/platform/extension/messagequeue/mysql/subscriber_heartbeat_store.go @@ -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) @@ -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 +} diff --git a/platform/extension/messagequeue/mysql/subscriber_heartbeat_store_test.go b/platform/extension/messagequeue/mysql/subscriber_heartbeat_store_test.go index 55dc7e226..8aa8fa88b 100644 --- a/platform/extension/messagequeue/mysql/subscriber_heartbeat_store_test.go +++ b/platform/extension/messagequeue/mysql/subscriber_heartbeat_store_test.go @@ -172,15 +172,16 @@ 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) @@ -188,6 +189,57 @@ func TestSubscriberHeartbeatStore_Deregister_SoftDelete(t *testing.T) { 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() @@ -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) @@ -230,8 +282,8 @@ 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, @@ -239,8 +291,8 @@ func TestSubscriberHeartbeatStore_Deregister(t *testing.T) { { 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, @@ -248,8 +300,8 @@ func TestSubscriberHeartbeatStore_Deregister(t *testing.T) { { 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, diff --git a/platform/extension/messagequeue/mysql/subscriber_test.go b/platform/extension/messagequeue/mysql/subscriber_test.go index 02112b2e0..1949a6453 100644 --- a/platform/extension/messagequeue/mysql/subscriber_test.go +++ b/platform/extension/messagequeue/mysql/subscriber_test.go @@ -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 } diff --git a/test/integration/extension/messagequeue/mysql/queue_test.go b/test/integration/extension/messagequeue/mysql/queue_test.go index 310c18055..5e9c15d6f 100644 --- a/test/integration/extension/messagequeue/mysql/queue_test.go +++ b/test/integration/extension/messagequeue/mysql/queue_test.go @@ -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") }