Moved cleanup activities logic to activity bounded context. (#40663)
<!-- Add the related story/sub-task/bug number, like Resolves #123, or remove if NA --> **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 <!-- This is an auto-generated comment: release notes by coderabbit.ai --> ## 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. <!-- end of auto-generated comment: release notes by coderabbit.ai -->
This commit is contained in:
+14
-3
@@ -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
|
||||
|
||||
+1
-1
@@ -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 {
|
||||
|
||||
@@ -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
|
||||
}
|
||||
@@ -9,4 +9,5 @@ type Service interface {
|
||||
ListHostPastActivitiesService
|
||||
StreamActivitiesService
|
||||
NewActivityService
|
||||
CleanupExpiredActivitiesService
|
||||
}
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
@@ -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")
|
||||
}
|
||||
|
||||
@@ -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{
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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;`)
|
||||
})
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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 {
|
||||
|
||||
Reference in New Issue
Block a user