From 593cf0112fd62da261c361dc0db4e154a4cf7655 Mon Sep 17 00:00:00 2001 From: Victor Lyuboslavsky <2685025+getvictor@users.noreply.github.com> Date: Fri, 27 Feb 2026 16:21:23 -0600 Subject: [PATCH] Moved cleanup activities logic to activity bounded context. (#40663) **Related issue:** Resolves #38536 Split the activities cleanup job from the queries cleanup job. # Checklist for submitter - [ ] Changes file added for user-visible changes in `changes/`, `orbit/changes/` or `ee/fleetd-chrome/changes`. - Present in previous PR ## Testing - [x] Added/updated automated tests - [x] QA'd all new/changed functionality manually ## Summary by CodeRabbit ## Release Notes * **New Features** * Added automated cleanup job for expired live queries based on activity expiration settings. * **Improvements** * Refactored activity data cleanup to use a dedicated service for better reliability and maintainability. * Enhanced scheduled cleanup operations with improved separation of concerns for activity and live query management. --- cmd/fleet/cron.go | 17 +- cmd/fleet/serve.go | 2 +- .../api/cleanup_expired_activities.go | 11 ++ server/activity/api/service.go | 1 + server/activity/internal/mysql/activity.go | 31 ++++ .../internal/mysql/activity_cleanup_test.go | 134 ++++++++++++++ .../service/cleanup_expired_activities.go | 16 ++ .../activity/internal/service/handler_test.go | 4 + .../activity/internal/service/service_test.go | 8 + server/activity/internal/types/activity.go | 3 + server/datastore/mysql/activities.go | 29 +-- server/datastore/mysql/activities_test.go | 171 ++++-------------- server/fleet/datastore.go | 8 +- server/mock/datastore_mock.go | 12 +- 14 files changed, 271 insertions(+), 176 deletions(-) create mode 100644 server/activity/api/cleanup_expired_activities.go create mode 100644 server/activity/internal/mysql/activity_cleanup_test.go create mode 100644 server/activity/internal/service/cleanup_expired_activities.go diff --git a/cmd/fleet/cron.go b/cmd/fleet/cron.go index 0022f38f66..d68995a28d 100644 --- a/cmd/fleet/cron.go +++ b/cmd/fleet/cron.go @@ -906,6 +906,7 @@ func newCleanupsAndAggregationSchedule( bootstrapPackageStore fleet.MDMBootstrapPackageStore, softwareTitleIconStore fleet.SoftwareTitleIconStore, androidSvc android.Service, + activitySvc activity_api.Service, ) (*schedule.Schedule, error) { const ( name = string(fleet.CronCleanupsThenAggregation) @@ -1080,10 +1081,20 @@ func newCleanupsAndAggregationSchedule( if !appConfig.ActivityExpirySettings.ActivityExpiryEnabled { return nil } - // A maxCount of 5,000 means that the cron job will keep the activities (and associated tables) - // sizes in control for deployments that generate (5k x 24 hours) ~120,000 activities per day. + // A maxCount of 5,000 means that the cron job will keep the activities table + // size in control for deployments that generate (5k x 24 hours) ~120,000 activities per day. const maxCount = 5000 - return ds.CleanupActivitiesAndAssociatedData(ctx, maxCount, appConfig.ActivityExpirySettings.ActivityExpiryWindow) + return activitySvc.CleanupExpiredActivities(ctx, maxCount, appConfig.ActivityExpirySettings.ActivityExpiryWindow) + }), + schedule.WithJob("cleanup_live_queries", func(ctx context.Context) error { + appConfig, err := ds.AppConfig(ctx) + if err != nil { + return err + } + if !appConfig.ActivityExpirySettings.ActivityExpiryEnabled { + return nil + } + return ds.CleanupExpiredLiveQueries(ctx, appConfig.ActivityExpirySettings.ActivityExpiryWindow) }), schedule.WithJob("cleanup_unused_software_installers", func(ctx context.Context) error { // remove only those unused created more than a minute ago to avoid a diff --git a/cmd/fleet/serve.go b/cmd/fleet/serve.go index 9cdbc17087..6eb8a75c02 100644 --- a/cmd/fleet/serve.go +++ b/cmd/fleet/serve.go @@ -1102,7 +1102,7 @@ the way that the Fleet server works. func() (fleet.CronSchedule, error) { commander := apple_mdm.NewMDMAppleCommander(mdmStorage, mdmPushService) return newCleanupsAndAggregationSchedule( - ctx, instanceID, ds, svc, logger, redisWrapperDS, &config, commander, softwareInstallStore, bootstrapPackageStore, softwareTitleIconStore, androidSvc, + ctx, instanceID, ds, svc, logger, redisWrapperDS, &config, commander, softwareInstallStore, bootstrapPackageStore, softwareTitleIconStore, androidSvc, activitySvc, ) }, ); err != nil { diff --git a/server/activity/api/cleanup_expired_activities.go b/server/activity/api/cleanup_expired_activities.go new file mode 100644 index 0000000000..c81879c8ac --- /dev/null +++ b/server/activity/api/cleanup_expired_activities.go @@ -0,0 +1,11 @@ +package api + +import "context" + +// CleanupExpiredActivitiesService cleans up expired activities. +type CleanupExpiredActivitiesService interface { + // CleanupExpiredActivities deletes up to maxCount activities older than expiryWindowDays + // that are not linked to any host. Host-linked activities are preserved + // (they are cleaned up when the host activity is processed). + CleanupExpiredActivities(ctx context.Context, maxCount int, expiryWindowDays int) error +} diff --git a/server/activity/api/service.go b/server/activity/api/service.go index 60fa246939..a52f26aaa1 100644 --- a/server/activity/api/service.go +++ b/server/activity/api/service.go @@ -9,4 +9,5 @@ type Service interface { ListHostPastActivitiesService StreamActivitiesService NewActivityService + CleanupExpiredActivitiesService } diff --git a/server/activity/internal/mysql/activity.go b/server/activity/internal/mysql/activity.go index b709841609..ae63e796d1 100644 --- a/server/activity/internal/mysql/activity.go +++ b/server/activity/internal/mysql/activity.go @@ -195,6 +195,37 @@ func (ds *Datastore) ListHostPastActivities(ctx context.Context, hostID uint, op return activities, metaData, nil } +// CleanupExpiredActivities deletes up to maxCount activities older than expiryWindowDays +// that are not linked to any host. Host-linked activities are preserved. +func (ds *Datastore) CleanupExpiredActivities(ctx context.Context, maxCount int, expiryWindowDays int) error { + ctx, span := tracer.Start(ctx, "activity.mysql.CleanupExpiredActivities") + defer span.End() + + const selectQuery = ` + SELECT a.id FROM activities a + LEFT JOIN host_activities ha ON (a.id=ha.activity_id) + WHERE ha.activity_id IS NULL AND a.created_at < DATE_SUB(NOW(), INTERVAL ? DAY) + ORDER BY a.id ASC + LIMIT ?` + + var activityIDs []uint + if err := sqlx.SelectContext(ctx, ds.primary, &activityIDs, selectQuery, expiryWindowDays, maxCount); err != nil { + return ctxerr.Wrap(ctx, err, "select expired activities for deletion") + } + if len(activityIDs) == 0 { + return nil + } + + deleteQuery, args, err := sqlx.In(`DELETE FROM activities WHERE id IN (?)`, activityIDs) + if err != nil { + return ctxerr.Wrap(ctx, err, "build expired activities IN query") + } + if _, err := ds.primary.ExecContext(ctx, deleteQuery, args...); err != nil { + return ctxerr.Wrap(ctx, err, "delete expired activities") + } + return nil +} + // fetchActivityDetails fetches details for activities in a separate query // to avoid MySQL sort buffer issues with large JSON entries. func (ds *Datastore) fetchActivityDetails(ctx context.Context, activities []*api.Activity) error { diff --git a/server/activity/internal/mysql/activity_cleanup_test.go b/server/activity/internal/mysql/activity_cleanup_test.go new file mode 100644 index 0000000000..844d4616f5 --- /dev/null +++ b/server/activity/internal/mysql/activity_cleanup_test.go @@ -0,0 +1,134 @@ +package mysql + +import ( + "fmt" + "testing" + "time" + + "github.com/fleetdm/fleet/v4/server/activity/internal/testutils" + "github.com/fleetdm/fleet/v4/server/ptr" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func TestCleanupExpiredActivities(t *testing.T) { + tdb := testutils.SetupTestDB(t, "activity_cleanup") + ds := NewDatastore(tdb.Conns(), tdb.Logger) + env := &testEnv{TestDB: tdb, ds: ds} + + cases := []struct { + name string + fn func(t *testing.T, env *testEnv) + }{ + {"NothingToDelete", testCleanupExpiredActivitiesNoop}, + {"DeletesExpiredNonHostActivities", testCleanupExpiredActivitiesBasic}, + {"RespectsMaxCount", testCleanupExpiredActivitiesBatch}, + } + for _, c := range cases { + t.Run(c.name, func(t *testing.T) { + defer env.TruncateTables(t) + c.fn(t, env) + }) + } +} + +func testCleanupExpiredActivitiesNoop(t *testing.T, env *testEnv) { + ctx := t.Context() + + // No activities exist -should be a no-op. + err := env.ds.CleanupExpiredActivities(ctx, 500, 1) + require.NoError(t, err) + + // Create a recent activity -should not be deleted. + userID := env.InsertUser(t, "user", "user@example.com") + env.InsertActivity(t, ptr.Uint(userID), "recent_activity", map[string]any{}) + + err = env.ds.CleanupExpiredActivities(ctx, 500, 1) + require.NoError(t, err) + + activities, _, err := env.ds.ListActivities(ctx, listOpts()) + require.NoError(t, err) + assert.Len(t, activities, 1) +} + +func testCleanupExpiredActivitiesBasic(t *testing.T, env *testEnv) { + ctx := t.Context() + userID := env.InsertUser(t, "user", "user@example.com") + hostID := env.InsertHost(t, "h1.local", nil) + + expiredTime := time.Now().Add(-48 * time.Hour) + recentTime := time.Now() + + // Create activities with different states: + // 1. Expired, no host link → should be deleted + expiredNoHost := env.InsertActivityWithTime(t, ptr.Uint(userID), "expired_no_host", map[string]any{}, expiredTime) + // 2. Expired, linked to host → should be preserved + expiredWithHost := env.InsertActivityWithTime(t, ptr.Uint(userID), "expired_with_host", map[string]any{}, expiredTime) + env.InsertHostActivity(t, hostID, expiredWithHost) + // 3. Recent, no host link → should be preserved + recentNoHost := env.InsertActivityWithTime(t, ptr.Uint(userID), "recent_no_host", map[string]any{}, recentTime) + + err := env.ds.CleanupExpiredActivities(ctx, 500, 1) + require.NoError(t, err) + + activities, _, err := env.ds.ListActivities(ctx, listOpts()) + require.NoError(t, err) + require.Len(t, activities, 2) + + activityIDs := make([]uint, len(activities)) + for i, a := range activities { + activityIDs[i] = a.ID + } + assert.NotContains(t, activityIDs, expiredNoHost, "expired non-host activity should be deleted") + assert.Contains(t, activityIDs, expiredWithHost, "expired host-linked activity should be preserved") + assert.Contains(t, activityIDs, recentNoHost, "recent activity should be preserved") + + // Verify host_activities link still exists for the preserved activity. + var hostActivityCount int + err = env.DB.GetContext(ctx, &hostActivityCount, "SELECT COUNT(*) FROM host_activities WHERE activity_id = ?", expiredWithHost) + require.NoError(t, err) + assert.Equal(t, 1, hostActivityCount) +} + +func testCleanupExpiredActivitiesBatch(t *testing.T, env *testEnv) { + ctx := t.Context() + userID := env.InsertUser(t, "user", "user@example.com") + expiredTime := time.Now().Add(-48 * time.Hour) + + // Create 10 expired activities (no host links). + for i := range 10 { + env.InsertActivityWithTime(t, ptr.Uint(userID), fmt.Sprintf("expired_%d", i), map[string]any{}, expiredTime) + } + + // Cleanup with maxCount=3 -only 3 should be deleted per call. + err := env.ds.CleanupExpiredActivities(ctx, 3, 1) + require.NoError(t, err) + + activities, _, err := env.ds.ListActivities(ctx, listOpts()) + require.NoError(t, err) + assert.Len(t, activities, 7, "only 3 of 10 expired activities should be deleted") + + // Run again -another 3 deleted. + err = env.ds.CleanupExpiredActivities(ctx, 3, 1) + require.NoError(t, err) + + activities, _, err = env.ds.ListActivities(ctx, listOpts()) + require.NoError(t, err) + assert.Len(t, activities, 4) + + // Run again -another 3 deleted. + err = env.ds.CleanupExpiredActivities(ctx, 3, 1) + require.NoError(t, err) + + activities, _, err = env.ds.ListActivities(ctx, listOpts()) + require.NoError(t, err) + assert.Len(t, activities, 1) + + // Run again -last one deleted. + err = env.ds.CleanupExpiredActivities(ctx, 3, 1) + require.NoError(t, err) + + activities, _, err = env.ds.ListActivities(ctx, listOpts()) + require.NoError(t, err) + assert.Len(t, activities, 0) +} diff --git a/server/activity/internal/service/cleanup_expired_activities.go b/server/activity/internal/service/cleanup_expired_activities.go new file mode 100644 index 0000000000..1e849f8ffa --- /dev/null +++ b/server/activity/internal/service/cleanup_expired_activities.go @@ -0,0 +1,16 @@ +package service + +import ( + "context" + + "github.com/fleetdm/fleet/v4/server/contexts/ctxerr" +) + +// CleanupExpiredActivities deletes up to maxCount activities older than expiryWindowDays +// that are not linked to any host. +func (s *Service) CleanupExpiredActivities(ctx context.Context, maxCount int, expiryWindowDays int) error { + if err := s.store.CleanupExpiredActivities(ctx, maxCount, expiryWindowDays); err != nil { + return ctxerr.Wrap(ctx, err, "cleanup expired activities") + } + return nil +} diff --git a/server/activity/internal/service/handler_test.go b/server/activity/internal/service/handler_test.go index 428ed88036..2b90f61c49 100644 --- a/server/activity/internal/service/handler_test.go +++ b/server/activity/internal/service/handler_test.go @@ -141,3 +141,7 @@ func (m *mockService) StreamActivities(_ context.Context, _ api.JSONLogger) erro func (m *mockService) NewActivity(_ context.Context, _ *api.User, _ api.ActivityDetails) error { panic("mockService.NewActivity should not be called in validation tests") } + +func (m *mockService) CleanupExpiredActivities(_ context.Context, _ int, _ int) error { + panic("mockService.CleanupExpiredActivities should not be called in validation tests") +} diff --git a/server/activity/internal/service/service_test.go b/server/activity/internal/service/service_test.go index 97781f6225..2f94d33d8b 100644 --- a/server/activity/internal/service/service_test.go +++ b/server/activity/internal/service/service_test.go @@ -64,6 +64,10 @@ func (m *mockDatastore) NewActivity(ctx context.Context, user *api.User, activit return nil } +func (m *mockDatastore) CleanupExpiredActivities(ctx context.Context, maxCount int, expiryWindowDays int) error { + return nil +} + type mockUserProvider struct { users []*activity.User listUsersErr error @@ -517,6 +521,10 @@ func (m *mockStreamingDatastore) NewActivity(ctx context.Context, user *api.User panic("not implemented") } +func (m *mockStreamingDatastore) CleanupExpiredActivities(ctx context.Context, maxCount int, expiryWindowDays int) error { + panic("not implemented") +} + func newTestActivity(id uint, actorName string, actorID uint, actType, details string) *api.Activity { jsonDetails := json.RawMessage(details) return &api.Activity{ diff --git a/server/activity/internal/types/activity.go b/server/activity/internal/types/activity.go index 6f94d44706..9420b7c0bb 100644 --- a/server/activity/internal/types/activity.go +++ b/server/activity/internal/types/activity.go @@ -83,4 +83,7 @@ type Datastore interface { // NewActivity stores a new activity record in the database. // The webhook context key must be set in the context before calling this method. NewActivity(ctx context.Context, user *api.User, activity api.ActivityDetails, details []byte, createdAt time.Time) error + // CleanupExpiredActivities deletes up to maxCount activities older than expiryWindowDays + // that are not linked to any host. Host-linked activities are preserved. + CleanupExpiredActivities(ctx context.Context, maxCount int, expiryWindowDays int) error } diff --git a/server/datastore/mysql/activities.go b/server/datastore/mysql/activities.go index 1261b7a2b1..b1298b36c9 100644 --- a/server/datastore/mysql/activities.go +++ b/server/datastore/mysql/activities.go @@ -297,35 +297,10 @@ func (ds *Datastore) ListHostUpcomingActivities(ctx context.Context, hostID uint return activities, metaData, nil } -func (ds *Datastore) CleanupActivitiesAndAssociatedData(ctx context.Context, maxCount int, expiredWindowDays int) error { - const selectActivitiesQuery = ` - SELECT a.id FROM activities a - LEFT JOIN host_activities ha ON (a.id=ha.activity_id) - WHERE ha.activity_id IS NULL AND a.created_at < DATE_SUB(NOW(), INTERVAL ? DAY) - ORDER BY a.id ASC - LIMIT ?;` - var activityIDs []uint - if err := sqlx.SelectContext(ctx, ds.writer(ctx), &activityIDs, selectActivitiesQuery, expiredWindowDays, maxCount); err != nil { - return ctxerr.Wrap(ctx, err, "select activities for deletion") - } - if len(activityIDs) > 0 { - deleteActivitiesQuery, args, err := sqlx.In(`DELETE FROM activities WHERE id IN (?);`, activityIDs) - if err != nil { - return ctxerr.Wrap(ctx, err, "build activities IN query") - } - if _, err := ds.writer(ctx).ExecContext(ctx, deleteActivitiesQuery, args...); err != nil { - return ctxerr.Wrap(ctx, err, "delete expired activities") - } - } - - // `activities` and `queries` are not tied because the activity itself holds - // the query SQL so they don't need to be executed on the same transaction. - // +func (ds *Datastore) CleanupExpiredLiveQueries(ctx context.Context, expiredWindowDays int) error { // All expired live queries are deleted in batch sizes of // `deleteIDsBatchSize` to ensure the table size is kept in check - // with high volumes of live queries (zero-trust workflows). This differs - // from the `activities` cleanup which uses maxCount as a limit to the - // number of activities to delete. + // with high volumes of live queries (zero-trust workflows). const selectUnsavedQueryIDs = ` SELECT id diff --git a/server/datastore/mysql/activities_test.go b/server/datastore/mysql/activities_test.go index 04b7e1eac1..51266031d0 100644 --- a/server/datastore/mysql/activities_test.go +++ b/server/datastore/mysql/activities_test.go @@ -32,8 +32,8 @@ func TestActivity(t *testing.T) { }{ {"UsernameChange", testActivityUsernameChange}, {"ListHostUpcomingActivities", testListHostUpcomingActivities}, - {"CleanupActivitiesAndAssociatedData", testCleanupActivitiesAndAssociatedData}, - {"CleanupActivitiesAndAssociatedDataBatch", testCleanupActivitiesAndAssociatedDataBatch}, + {"CleanupExpiredLiveQueries", testCleanupExpiredLiveQueries}, + {"CleanupExpiredLiveQueriesBatch", testCleanupExpiredLiveQueriesBatch}, {"ActivateNextActivity", testActivateNextActivity}, {"ActivateItselfOnEmptyQueue", testActivateItselfOnEmptyQueue}, {"CancelNonActivatedUpcomingActivity", testCancelNonActivatedUpcomingActivity}, @@ -536,8 +536,7 @@ func testListHostUpcomingActivities(t *testing.T, ds *Datastore) { } } -func testCleanupActivitiesAndAssociatedData(t *testing.T, ds *Datastore) { - activitySvc := NewTestActivityService(t, ds) +func testCleanupExpiredLiveQueries(t *testing.T, ds *Datastore) { ctx := context.Background() user1 := &fleet.User{ Password: []byte("p4ssw0rd.123"), @@ -549,219 +548,123 @@ func testCleanupActivitiesAndAssociatedData(t *testing.T, ds *Datastore) { require.NoError(t, err) // Nothing to delete. - err = ds.CleanupActivitiesAndAssociatedData(ctx, 500, 1) + err = ds.CleanupExpiredLiveQueries(ctx, 1) require.NoError(t, err) - nonSavedQuery1, err := ds.NewQuery(ctx, &fleet.Query{ - Name: "nonSavedQuery1", + nonSavedQuery, err := ds.NewQuery(ctx, &fleet.Query{ + Name: "nonSavedQuery", Saved: false, Query: "SELECT 1;", Logging: fleet.LoggingSnapshot, }) require.NoError(t, err) - savedQuery1, err := ds.NewQuery(ctx, &fleet.Query{ - Name: "savedQuery1", + savedQuery, err := ds.NewQuery(ctx, &fleet.Query{ + Name: "savedQuery", Saved: true, Query: "SELECT 2;", Logging: fleet.LoggingSnapshot, }) require.NoError(t, err) - distributedQueryCampaign1, err := ds.NewDistributedQueryCampaign(ctx, &fleet.DistributedQueryCampaign{ - QueryID: nonSavedQuery1.ID, + campaign, err := ds.NewDistributedQueryCampaign(ctx, &fleet.DistributedQueryCampaign{ + QueryID: nonSavedQuery.ID, Status: fleet.QueryComplete, UserID: user1.ID, }) require.NoError(t, err) _, err = ds.NewDistributedQueryCampaignTarget(ctx, &fleet.DistributedQueryCampaignTarget{ - DistributedQueryCampaignID: distributedQueryCampaign1.ID, + DistributedQueryCampaignID: campaign.ID, TargetID: 1, Type: fleet.TargetHost, }) require.NoError(t, err) - apiUser := &activity_api.User{ID: user1.ID, Name: user1.Name, Email: user1.Email} - err = activitySvc.NewActivity(ctx, apiUser, dummyActivity{ - name: "other activity", - details: map[string]interface{}{"detail": 0, "foo": "zoo"}, - }) - require.NoError(t, err) - err = activitySvc.NewActivity(ctx, apiUser, dummyActivity{ - name: "live query", - details: map[string]interface{}{"detail": 1, "foo": "bar"}, - }) - require.NoError(t, err) - err = activitySvc.NewActivity(ctx, apiUser, dummyActivity{ - name: "some host activity", - details: map[string]interface{}{"detail": 0, "foo": "zoo"}, - hostIDs: []uint{1}, - }) - require.NoError(t, err) - err = activitySvc.NewActivity(ctx, apiUser, dummyActivity{ - name: "some host activity 2", - details: map[string]interface{}{"detail": 0, "foo": "bar"}, - hostIDs: []uint{2}, - }) + + // Nothing is deleted because the data is recent. + err = ds.CleanupExpiredLiveQueries(ctx, 1) require.NoError(t, err) - // Nothing is deleted, as the activities and associated data is recent. - const maxCount = 500 - err = ds.CleanupActivitiesAndAssociatedData(ctx, maxCount, 1) + _, err = ds.Query(ctx, nonSavedQuery.ID) require.NoError(t, err) - - activities := ListActivitiesAPI(t, ctx, activitySvc, activity_api.ListOptions{}) - require.Len(t, activities, 4) - nonExpiredActivityID := activities[0].ID - expiredActivityID := activities[1].ID - nonExpiredHostActivityID := activities[2].ID - expiredHostActivityID := activities[3].ID - _, err = ds.Query(ctx, nonSavedQuery1.ID) + _, err = ds.DistributedQueryCampaign(ctx, campaign.ID) require.NoError(t, err) - _, err = ds.DistributedQueryCampaign(ctx, distributedQueryCampaign1.ID) - require.NoError(t, err) - targets, err := ds.DistributedQueryCampaignTargetIDs(ctx, distributedQueryCampaign1.ID) + targets, err := ds.DistributedQueryCampaignTargetIDs(ctx, campaign.ID) require.NoError(t, err) require.Len(t, targets.HostIDs, 1) - // Make some of the activity and associated data older. - _, err = ds.writer(context.Background()).Exec(` - UPDATE activities SET created_at = ? WHERE id = ? OR id = ?`, - time.Now().Add(-48*time.Hour), expiredActivityID, expiredHostActivityID, - ) - require.NoError(t, err) + // Make the queries older. _, err = ds.writer(context.Background()).Exec(` UPDATE queries SET created_at = ? WHERE id = ? OR id = ?`, - time.Now().Add(-48*time.Hour), nonSavedQuery1.ID, savedQuery1.ID, + time.Now().Add(-48*time.Hour), nonSavedQuery.ID, savedQuery.ID, ) require.NoError(t, err) - // Expired activity and associated data should be cleaned up. - err = ds.CleanupActivitiesAndAssociatedData(ctx, maxCount, 1) + // Expired unsaved query, its campaign, and campaign targets should be cleaned up. + err = ds.CleanupExpiredLiveQueries(ctx, 1) require.NoError(t, err) - activities = ListActivitiesAPI(t, ctx, activitySvc, activity_api.ListOptions{}) - require.Len(t, activities, 3) - require.Equal(t, nonExpiredActivityID, activities[0].ID) - require.Equal(t, nonExpiredHostActivityID, activities[1].ID) - require.Equal(t, expiredHostActivityID, activities[2].ID) - _, err = ds.Query(ctx, nonSavedQuery1.ID) + _, err = ds.Query(ctx, nonSavedQuery.ID) require.ErrorIs(t, err, sql.ErrNoRows) - _, err = ds.DistributedQueryCampaign(ctx, distributedQueryCampaign1.ID) + _, err = ds.DistributedQueryCampaign(ctx, campaign.ID) require.ErrorIs(t, err, sql.ErrNoRows) - targets, err = ds.DistributedQueryCampaignTargetIDs(ctx, distributedQueryCampaign1.ID) + targets, err = ds.DistributedQueryCampaignTargetIDs(ctx, campaign.ID) require.NoError(t, err) require.Empty(t, targets.HostIDs) require.Empty(t, targets.LabelIDs) require.Empty(t, targets.TeamIDs) // Saved query should not be cleaned up. - savedQuery1, err = ds.Query(ctx, savedQuery1.ID) + savedQuery, err = ds.Query(ctx, savedQuery.ID) require.NoError(t, err) - require.NotNil(t, savedQuery1) + require.NotNil(t, savedQuery) } -func testCleanupActivitiesAndAssociatedDataBatch(t *testing.T, ds *Datastore) { - activitySvc := NewTestActivityService(t, ds) +func testCleanupExpiredLiveQueriesBatch(t *testing.T, ds *Datastore) { ctx := context.Background() - user1 := &fleet.User{ - Password: []byte("p4ssw0rd.123"), - Name: "user1", - Email: "user1@example.com", - GlobalRole: ptr.String(fleet.RoleAdmin), - } - user1, err := ds.NewUser(ctx, user1) - require.NoError(t, err) - - const maxCount = 500 - - // Create 1500 activities. - insertActivitiesStmt := ` - INSERT INTO activities - (user_id, user_name, activity_type, details, user_email) - VALUES ` - var insertActivitiesArgs []interface{} - for i := 0; i < 1500; i++ { - insertActivitiesArgs = append(insertActivitiesArgs, - user1.ID, user1.Name, "foobar", `{"foo": "bar"}`, user1.Email, - ) - } - insertActivitiesStmt += strings.TrimSuffix(strings.Repeat("(?, ?, ?, ?, ?),", 1500), ",") - _, err = ds.writer(ctx).ExecContext(ctx, insertActivitiesStmt, insertActivitiesArgs...) - require.NoError(t, err) // Create 1500 non-saved queries. insertQueriesStmt := ` INSERT INTO queries (name, description, query) VALUES ` - var insertQueriesArgs []interface{} - for i := 0; i < 1500; i++ { + var insertQueriesArgs []any + for i := range 1500 { insertQueriesArgs = append(insertQueriesArgs, fmt.Sprintf("foobar%d", i), "foobar", "SELECT 1;", ) } insertQueriesStmt += strings.TrimSuffix(strings.Repeat("(?, ?, ?),", 1500), ",") - _, err = ds.writer(ctx).ExecContext(ctx, insertQueriesStmt, insertQueriesArgs...) + _, err := ds.writer(ctx).ExecContext(ctx, insertQueriesStmt, insertQueriesArgs...) require.NoError(t, err) - err = ds.CleanupActivitiesAndAssociatedData(ctx, maxCount, 1) + // Nothing deleted; all recent. + err = ds.CleanupExpiredLiveQueries(ctx, 1) require.NoError(t, err) - activities := ListActivitiesAPI(t, ctx, activitySvc, activity_api.ListOptions{}) - require.Len(t, activities, 1500) var queriesLen int ExecAdhocSQL(t, ds, func(q sqlx.ExtContext) error { return sqlx.GetContext(ctx, q, &queriesLen, `SELECT COUNT(*) FROM queries WHERE NOT saved;`) }) require.Equal(t, 1500, queriesLen) - // Make 1250 activities as expired. - _, err = ds.writer(context.Background()).Exec(` - UPDATE activities SET created_at = ? WHERE id <= 1250`, - time.Now().Add(-48*time.Hour), - ) - require.NoError(t, err) - - // Make 1250 queries as expired. + // Make 1250 queries expired. _, err = ds.writer(context.Background()).Exec(` UPDATE queries SET created_at = ? WHERE id <= 1250`, time.Now().Add(-48*time.Hour), ) require.NoError(t, err) - err = ds.CleanupActivitiesAndAssociatedData(ctx, maxCount, 1) + // All 1250 expired queries should be cleaned up in one call (batched internally). + err = ds.CleanupExpiredLiveQueries(ctx, 1) require.NoError(t, err) - activities = ListActivitiesAPI(t, ctx, activitySvc, activity_api.ListOptions{}) - require.Len(t, activities, 1000) - ExecAdhocSQL(t, ds, func(q sqlx.ExtContext) error { - return sqlx.GetContext(ctx, q, &queriesLen, `SELECT COUNT(*) FROM queries WHERE NOT saved;`) - }) - require.Equal(t, 250, queriesLen) // All expired queries should be cleaned up. - - err = ds.CleanupActivitiesAndAssociatedData(ctx, maxCount, 1) - require.NoError(t, err) - - activities = ListActivitiesAPI(t, ctx, activitySvc, activity_api.ListOptions{}) - require.Len(t, activities, 500) ExecAdhocSQL(t, ds, func(q sqlx.ExtContext) error { return sqlx.GetContext(ctx, q, &queriesLen, `SELECT COUNT(*) FROM queries WHERE NOT saved;`) }) require.Equal(t, 250, queriesLen) - err = ds.CleanupActivitiesAndAssociatedData(ctx, maxCount, 1) + // Running again should be a no-op (remaining 250 are not expired). + err = ds.CleanupExpiredLiveQueries(ctx, 1) require.NoError(t, err) - activities = ListActivitiesAPI(t, ctx, activitySvc, activity_api.ListOptions{}) - require.Len(t, activities, 250) - ExecAdhocSQL(t, ds, func(q sqlx.ExtContext) error { - return sqlx.GetContext(ctx, q, &queriesLen, `SELECT COUNT(*) FROM queries WHERE NOT saved;`) - }) - require.Equal(t, 250, queriesLen) - - err = ds.CleanupActivitiesAndAssociatedData(ctx, maxCount, 1) - require.NoError(t, err) - - activities = ListActivitiesAPI(t, ctx, activitySvc, activity_api.ListOptions{}) - require.Len(t, activities, 250) ExecAdhocSQL(t, ds, func(q sqlx.ExtContext) error { return sqlx.GetContext(ctx, q, &queriesLen, `SELECT COUNT(*) FROM queries WHERE NOT saved;`) }) diff --git a/server/fleet/datastore.go b/server/fleet/datastore.go index 844f86b98b..d9d08e200b 100644 --- a/server/fleet/datastore.go +++ b/server/fleet/datastore.go @@ -2004,11 +2004,9 @@ type Datastore interface { // CleanupUnusedScriptContents will remove script contents that have no references to them from // the scripts or host_script_results tables. CleanupUnusedScriptContents(ctx context.Context) error - // CleanupActivitiesAndAssociatedData will cleanup (up to maxCount) activities and their associated data - // that are older than the given expiration window. - // - // The argument maxCount is used to not lock the database for long periods of time. - CleanupActivitiesAndAssociatedData(ctx context.Context, maxCount int, expiryWindowDays int) error + // CleanupExpiredLiveQueries cleans up unsaved queries older than the given expiration window (in days), + // orphaned distributed query campaigns that reference non-existing queries, and orphaned campaign targets that reference non-existing campaigns. + CleanupExpiredLiveQueries(ctx context.Context, expiryWindowDays int) error // WipeHostViaScript sends a script to wipe a host and updates the // states in host_mdm_actions. WipeHostViaScript(ctx context.Context, request *HostScriptRequestPayload, hostFleetPlatform string) error diff --git a/server/mock/datastore_mock.go b/server/mock/datastore_mock.go index f5e4153bb4..0050c6a48b 100644 --- a/server/mock/datastore_mock.go +++ b/server/mock/datastore_mock.go @@ -1305,7 +1305,7 @@ type DeleteHostLocationDataFunc func(ctx context.Context, hostID uint) error type CleanupUnusedScriptContentsFunc func(ctx context.Context) error -type CleanupActivitiesAndAssociatedDataFunc func(ctx context.Context, maxCount int, expiryWindowDays int) error +type CleanupExpiredLiveQueriesFunc func(ctx context.Context, expiryWindowDays int) error type WipeHostViaScriptFunc func(ctx context.Context, request *fleet.HostScriptRequestPayload, hostFleetPlatform string) error @@ -3707,8 +3707,8 @@ type DataStore struct { CleanupUnusedScriptContentsFunc CleanupUnusedScriptContentsFunc CleanupUnusedScriptContentsFuncInvoked bool - CleanupActivitiesAndAssociatedDataFunc CleanupActivitiesAndAssociatedDataFunc - CleanupActivitiesAndAssociatedDataFuncInvoked bool + CleanupExpiredLiveQueriesFunc CleanupExpiredLiveQueriesFunc + CleanupExpiredLiveQueriesFuncInvoked bool WipeHostViaScriptFunc WipeHostViaScriptFunc WipeHostViaScriptFuncInvoked bool @@ -8914,11 +8914,11 @@ func (s *DataStore) CleanupUnusedScriptContents(ctx context.Context) error { return s.CleanupUnusedScriptContentsFunc(ctx) } -func (s *DataStore) CleanupActivitiesAndAssociatedData(ctx context.Context, maxCount int, expiryWindowDays int) error { +func (s *DataStore) CleanupExpiredLiveQueries(ctx context.Context, expiryWindowDays int) error { s.mu.Lock() - s.CleanupActivitiesAndAssociatedDataFuncInvoked = true + s.CleanupExpiredLiveQueriesFuncInvoked = true s.mu.Unlock() - return s.CleanupActivitiesAndAssociatedDataFunc(ctx, maxCount, expiryWindowDays) + return s.CleanupExpiredLiveQueriesFunc(ctx, expiryWindowDays) } func (s *DataStore) WipeHostViaScript(ctx context.Context, request *fleet.HostScriptRequestPayload, hostFleetPlatform string) error {