From 2ccfab253f55c983f0829be41bd2acf619412d34 Mon Sep 17 00:00:00 2001 From: Martin Angers Date: Tue, 27 May 2025 16:38:39 -0400 Subject: [PATCH] Bugfix: catch-all cron job to avoid blocked upcoming activities queue (#29477) --- cmd/fleet/cron.go | 58 ++++++--- cmd/fleet/serve.go | 8 ++ server/datastore/mysql/activities.go | 27 ++++ server/datastore/mysql/activities_test.go | 149 ++++++++++++++++++++++ server/fleet/cron_schedules.go | 5 +- server/fleet/datastore.go | 1 + server/mock/datastore_mock.go | 12 ++ 7 files changed, 239 insertions(+), 21 deletions(-) diff --git a/cmd/fleet/cron.go b/cmd/fleet/cron.go index 394aa6bf29..4131c43fc7 100644 --- a/cmd/fleet/cron.go +++ b/cmd/fleet/cron.go @@ -1013,26 +1013,23 @@ func newFrequentCleanupsSchedule( schedule.WithAltLockID("leader_frequent_cleanups"), schedule.WithLogger(kitlog.With(logger, "cron", name)), // Run cleanup jobs first. - schedule.WithJob( - "redis_live_queries", - func(ctx context.Context) error { - // It's necessary to avoid lingering live queries in case of: - // - (Unknown) bug in the implementation, or, - // - Redis is so overloaded already that the lq.StopQuery in svc.CompleteCampaign fails to execute, or, - // - MySQL is so overloaded that ds.SaveDistributedQueryCampaign in svc.CompleteCampaign fails to execute. - names, err := lq.LoadActiveQueryNames() - if err != nil { - return err - } - ids := stringSliceToUintSlice(names, logger) - completed, err := ds.GetCompletedCampaigns(ctx, ids) - if err != nil { - return err - } - err = lq.CleanupInactiveQueries(ctx, completed) + schedule.WithJob("redis_live_queries", func(ctx context.Context) error { + // It's necessary to avoid lingering live queries in case of: + // - (Unknown) bug in the implementation, or, + // - Redis is so overloaded already that the lq.StopQuery in svc.CompleteCampaign fails to execute, or, + // - MySQL is so overloaded that ds.SaveDistributedQueryCampaign in svc.CompleteCampaign fails to execute. + names, err := lq.LoadActiveQueryNames() + if err != nil { return err - }, - ), + } + ids := stringSliceToUintSlice(names, logger) + completed, err := ds.GetCompletedCampaigns(ctx, ids) + if err != nil { + return err + } + err = lq.CleanupInactiveQueries(ctx, completed) + return err + }), ) return s, nil @@ -1545,3 +1542,26 @@ func newIPhoneIPadReviver( return s, nil } + +func newUpcomingActivitiesSchedule( + ctx context.Context, + instanceID string, + ds fleet.Datastore, + logger kitlog.Logger, +) (*schedule.Schedule, error) { + const ( + name = string(fleet.CronUpcomingActivitiesMaintenance) + defaultInterval = 10 * time.Minute + ) + s := schedule.New( + ctx, name, instanceID, defaultInterval, ds, ds, + schedule.WithLogger(kitlog.With(logger, "cron", name)), + schedule.WithJob("unblock_hosts_upcoming_activity_queue", func(ctx context.Context) error { + const maxUnblockHosts = 500 + _, err := ds.UnblockHostsUpcomingActivityQueue(ctx, maxUnblockHosts) + return err + }), + ) + + return s, nil +} diff --git a/cmd/fleet/serve.go b/cmd/fleet/serve.go index d5a95cd3fa..22f222b713 100644 --- a/cmd/fleet/serve.go +++ b/cmd/fleet/serve.go @@ -897,6 +897,14 @@ the way that the Fleet server works. initFatal(err, "failed to register cleanups_then_aggregations schedule") } + if err := cronSchedules.StartCronSchedule( + func() (fleet.CronSchedule, error) { + return newUpcomingActivitiesSchedule(ctx, instanceID, ds, logger) + }, + ); err != nil { + initFatal(err, "failed to register upcoming_activities_maintenance schedule") + } + if err := cronSchedules.StartCronSchedule(func() (fleet.CronSchedule, error) { return newUsageStatisticsSchedule(ctx, instanceID, ds, config, license, logger) }); err != nil { diff --git a/server/datastore/mysql/activities.go b/server/datastore/mysql/activities.go index b2b6250de2..58cbe26cb1 100644 --- a/server/datastore/mysql/activities.go +++ b/server/datastore/mysql/activities.go @@ -1053,6 +1053,33 @@ func (ds *Datastore) GetHostUpcomingActivityMeta(ctx context.Context, hostID uin return &actMeta, nil } +// UnblockHostsUpcomingActivityQueue checks for hosts that have upcoming +// activities but none is "activated", meaning that the queue is blocked +// (cannot make progress anymore), possibly due to a failure when activating +// the next activity, or to a missing call to activateNextUpcomingActivity. It +// unblocks up to maxHosts found in this situation (by activating the next +// activity for each host). +func (ds *Datastore) UnblockHostsUpcomingActivityQueue(ctx context.Context, maxHosts int) (int, error) { + const findBlockedHostsStmt = ` + SELECT + DISTINCT inactive_ua.host_id + FROM + upcoming_activities inactive_ua + LEFT OUTER JOIN upcoming_activities active_ua ON + active_ua.host_id = inactive_ua.host_id AND + active_ua.activated_at IS NOT NULL + WHERE + active_ua.host_id IS NULL AND + inactive_ua.activated_at IS NULL + LIMIT ?` + + var blockedHostIDs []uint + if err := sqlx.SelectContext(ctx, ds.reader(ctx), &blockedHostIDs, findBlockedHostsStmt, maxHosts); err != nil { + return 0, ctxerr.Wrap(ctx, err, "select blocked hosts") + } + return len(blockedHostIDs), ds.activateNextUpcomingActivityForBatchOfHosts(ctx, blockedHostIDs) +} + func (ds *Datastore) activateNextUpcomingActivityForBatchOfHosts(ctx context.Context, hostIDs []uint) error { const maxHostIDsPerBatch = 500 diff --git a/server/datastore/mysql/activities_test.go b/server/datastore/mysql/activities_test.go index 2f222364d1..c568555e92 100644 --- a/server/datastore/mysql/activities_test.go +++ b/server/datastore/mysql/activities_test.go @@ -45,6 +45,7 @@ func TestActivity(t *testing.T) { {"CancelActivatedUpcomingActivity", testCancelActivatedUpcomingActivity}, {"SetResultAfterCancelUpcomingActivity", testSetResultAfterCancelUpcomingActivity}, {"GetHostUpcomingActivityMeta", testGetHostUpcomingActivityMeta}, + {"UnblockHostsUpcomingActivityQueue", testUnblockHostsUpcomingActivityQueue}, } for _, c := range cases { t.Run(c.name, func(t *testing.T) { @@ -2009,3 +2010,151 @@ func testGetHostUpcomingActivityMeta(t *testing.T, ds *Datastore) { _, err = ds.GetHostUpcomingActivityMeta(ctx, host2.ID, wipeExecID) require.ErrorAs(t, err, &nfe) } + +func testUnblockHostsUpcomingActivityQueue(t *testing.T, ds *Datastore) { + ctx := t.Context() + u := test.NewUser(t, ds, "user1", "user1@example.com", false) + + // create a few hosts + hosts := make([]*fleet.Host, 5) + for i := range hosts { + host := test.NewHost(t, ds, fmt.Sprintf("h%d.local", i+1), fmt.Sprintf("10.10.10.%d", i+1), fmt.Sprint(i+1), fmt.Sprint(i+1), time.Now()) + hosts[i] = host + } + + // run without anything in any host queue + n, err := ds.UnblockHostsUpcomingActivityQueue(ctx, 10) + require.NoError(t, err) + require.Equal(t, 0, n) + + deleteUpcomingActivityToBlockQueue := func(execID string) { + ExecAdhocSQL(t, ds, func(q sqlx.ExtContext) error { + _, err := q.ExecContext(ctx, `DELETE FROM upcoming_activities WHERE execution_id = ?`, execID) + return err + }) + } + + // enqueue some activities on some hosts (the nature of the activity is not relevant) + host0ScriptA, err := ds.NewHostScriptExecutionRequest(ctx, &fleet.HostScriptRequestPayload{HostID: hosts[0].ID, ScriptContents: "A", UserID: &u.ID, SyncRequest: true}) + require.NoError(t, err) + host0ScriptB, err := ds.NewHostScriptExecutionRequest(ctx, &fleet.HostScriptRequestPayload{HostID: hosts[0].ID, ScriptContents: "B", UserID: &u.ID, SyncRequest: true}) + require.NoError(t, err) + host1ScriptA, err := ds.NewHostScriptExecutionRequest(ctx, &fleet.HostScriptRequestPayload{HostID: hosts[1].ID, ScriptContents: "A", UserID: &u.ID, SyncRequest: true}) + require.NoError(t, err) + host2ScriptA, err := ds.NewHostScriptExecutionRequest(ctx, &fleet.HostScriptRequestPayload{HostID: hosts[2].ID, ScriptContents: "A", UserID: &u.ID, SyncRequest: true}) + require.NoError(t, err) + + checkUpcomingActivities(t, ds, hosts[0], host0ScriptA.ExecutionID, host0ScriptB.ExecutionID) + checkUpcomingActivities(t, ds, hosts[1], host1ScriptA.ExecutionID) + checkUpcomingActivities(t, ds, hosts[2], host2ScriptA.ExecutionID) + checkUpcomingActivities(t, ds, hosts[3]) + checkUpcomingActivities(t, ds, hosts[4]) + + // nothing to unblock + n, err = ds.UnblockHostsUpcomingActivityQueue(ctx, 10) + require.NoError(t, err) + require.Equal(t, 0, n) + + // block queue for host 0 + deleteUpcomingActivityToBlockQueue(host0ScriptA.ExecutionID) + + n, err = ds.UnblockHostsUpcomingActivityQueue(ctx, 10) + require.NoError(t, err) + require.Equal(t, 1, n) + + checkUpcomingActivities(t, ds, hosts[0], host0ScriptB.ExecutionID) + checkUpcomingActivities(t, ds, hosts[1], host1ScriptA.ExecutionID) + checkUpcomingActivities(t, ds, hosts[2], host2ScriptA.ExecutionID) + checkUpcomingActivities(t, ds, hosts[3]) + checkUpcomingActivities(t, ds, hosts[4]) + + // enqueue script C for all hosts + host0ScriptC, err := ds.NewHostScriptExecutionRequest(ctx, &fleet.HostScriptRequestPayload{HostID: hosts[0].ID, ScriptContents: "C", UserID: &u.ID, SyncRequest: true}) + require.NoError(t, err) + host1ScriptC, err := ds.NewHostScriptExecutionRequest(ctx, &fleet.HostScriptRequestPayload{HostID: hosts[1].ID, ScriptContents: "C", UserID: &u.ID, SyncRequest: true}) + require.NoError(t, err) + host2ScriptC, err := ds.NewHostScriptExecutionRequest(ctx, &fleet.HostScriptRequestPayload{HostID: hosts[2].ID, ScriptContents: "C", UserID: &u.ID, SyncRequest: true}) + require.NoError(t, err) + host3ScriptC, err := ds.NewHostScriptExecutionRequest(ctx, &fleet.HostScriptRequestPayload{HostID: hosts[3].ID, ScriptContents: "C", UserID: &u.ID, SyncRequest: true}) + require.NoError(t, err) + host4ScriptC, err := ds.NewHostScriptExecutionRequest(ctx, &fleet.HostScriptRequestPayload{HostID: hosts[4].ID, ScriptContents: "C", UserID: &u.ID, SyncRequest: true}) + require.NoError(t, err) + + checkUpcomingActivities(t, ds, hosts[0], host0ScriptB.ExecutionID, host0ScriptC.ExecutionID) + checkUpcomingActivities(t, ds, hosts[1], host1ScriptA.ExecutionID, host1ScriptC.ExecutionID) + checkUpcomingActivities(t, ds, hosts[2], host2ScriptA.ExecutionID, host2ScriptC.ExecutionID) + checkUpcomingActivities(t, ds, hosts[3], host3ScriptC.ExecutionID) + checkUpcomingActivities(t, ds, hosts[4], host4ScriptC.ExecutionID) + + // block queue for all hosts, but since hosts 3 and 4 are now empty, no need + // to unblock + deleteUpcomingActivityToBlockQueue(host0ScriptB.ExecutionID) + deleteUpcomingActivityToBlockQueue(host1ScriptA.ExecutionID) + deleteUpcomingActivityToBlockQueue(host2ScriptA.ExecutionID) + deleteUpcomingActivityToBlockQueue(host3ScriptC.ExecutionID) + deleteUpcomingActivityToBlockQueue(host4ScriptC.ExecutionID) + + n, err = ds.UnblockHostsUpcomingActivityQueue(ctx, 10) + require.NoError(t, err) + require.Equal(t, 3, n) + + checkUpcomingActivities(t, ds, hosts[0], host0ScriptC.ExecutionID) + checkUpcomingActivities(t, ds, hosts[1], host1ScriptC.ExecutionID) + checkUpcomingActivities(t, ds, hosts[2], host2ScriptC.ExecutionID) + checkUpcomingActivities(t, ds, hosts[3]) + checkUpcomingActivities(t, ds, hosts[4]) + + // enqueue script D and E for all hosts + host0ScriptD, err := ds.NewHostScriptExecutionRequest(ctx, &fleet.HostScriptRequestPayload{HostID: hosts[0].ID, ScriptContents: "D", UserID: &u.ID, SyncRequest: true}) + require.NoError(t, err) + host1ScriptD, err := ds.NewHostScriptExecutionRequest(ctx, &fleet.HostScriptRequestPayload{HostID: hosts[1].ID, ScriptContents: "D", UserID: &u.ID, SyncRequest: true}) + require.NoError(t, err) + host2ScriptD, err := ds.NewHostScriptExecutionRequest(ctx, &fleet.HostScriptRequestPayload{HostID: hosts[2].ID, ScriptContents: "D", UserID: &u.ID, SyncRequest: true}) + require.NoError(t, err) + host3ScriptD, err := ds.NewHostScriptExecutionRequest(ctx, &fleet.HostScriptRequestPayload{HostID: hosts[3].ID, ScriptContents: "D", UserID: &u.ID, SyncRequest: true}) + require.NoError(t, err) + host4ScriptD, err := ds.NewHostScriptExecutionRequest(ctx, &fleet.HostScriptRequestPayload{HostID: hosts[4].ID, ScriptContents: "D", UserID: &u.ID, SyncRequest: true}) + require.NoError(t, err) + host0ScriptE, err := ds.NewHostScriptExecutionRequest(ctx, &fleet.HostScriptRequestPayload{HostID: hosts[0].ID, ScriptContents: "E", UserID: &u.ID, SyncRequest: true}) + require.NoError(t, err) + host1ScriptE, err := ds.NewHostScriptExecutionRequest(ctx, &fleet.HostScriptRequestPayload{HostID: hosts[1].ID, ScriptContents: "E", UserID: &u.ID, SyncRequest: true}) + require.NoError(t, err) + host2ScriptE, err := ds.NewHostScriptExecutionRequest(ctx, &fleet.HostScriptRequestPayload{HostID: hosts[2].ID, ScriptContents: "E", UserID: &u.ID, SyncRequest: true}) + require.NoError(t, err) + host3ScriptE, err := ds.NewHostScriptExecutionRequest(ctx, &fleet.HostScriptRequestPayload{HostID: hosts[3].ID, ScriptContents: "E", UserID: &u.ID, SyncRequest: true}) + require.NoError(t, err) + host4ScriptE, err := ds.NewHostScriptExecutionRequest(ctx, &fleet.HostScriptRequestPayload{HostID: hosts[4].ID, ScriptContents: "E", UserID: &u.ID, SyncRequest: true}) + require.NoError(t, err) + + checkUpcomingActivities(t, ds, hosts[0], host0ScriptC.ExecutionID, host0ScriptD.ExecutionID, host0ScriptE.ExecutionID) + checkUpcomingActivities(t, ds, hosts[1], host1ScriptC.ExecutionID, host1ScriptD.ExecutionID, host1ScriptE.ExecutionID) + checkUpcomingActivities(t, ds, hosts[2], host2ScriptC.ExecutionID, host2ScriptD.ExecutionID, host2ScriptE.ExecutionID) + checkUpcomingActivities(t, ds, hosts[3], host3ScriptD.ExecutionID, host3ScriptE.ExecutionID) + checkUpcomingActivities(t, ds, hosts[4], host4ScriptD.ExecutionID, host4ScriptE.ExecutionID) + + // block queue for all hosts + deleteUpcomingActivityToBlockQueue(host0ScriptC.ExecutionID) + deleteUpcomingActivityToBlockQueue(host1ScriptC.ExecutionID) + deleteUpcomingActivityToBlockQueue(host2ScriptC.ExecutionID) + deleteUpcomingActivityToBlockQueue(host3ScriptD.ExecutionID) + deleteUpcomingActivityToBlockQueue(host4ScriptD.ExecutionID) + + // process max 3 hosts + n, err = ds.UnblockHostsUpcomingActivityQueue(ctx, 3) + require.NoError(t, err) + require.Equal(t, 3, n) + // run again, should process the next 2 hosts + n, err = ds.UnblockHostsUpcomingActivityQueue(ctx, 3) + require.NoError(t, err) + require.Equal(t, 2, n) + // run again, nothing to unblock + n, err = ds.UnblockHostsUpcomingActivityQueue(ctx, 3) + require.NoError(t, err) + require.Equal(t, 0, n) + + checkUpcomingActivities(t, ds, hosts[0], host0ScriptD.ExecutionID, host0ScriptE.ExecutionID) + checkUpcomingActivities(t, ds, hosts[1], host1ScriptD.ExecutionID, host1ScriptE.ExecutionID) + checkUpcomingActivities(t, ds, hosts[2], host2ScriptD.ExecutionID, host2ScriptE.ExecutionID) + checkUpcomingActivities(t, ds, hosts[3], host3ScriptE.ExecutionID) + checkUpcomingActivities(t, ds, hosts[4], host4ScriptE.ExecutionID) +} diff --git a/server/fleet/cron_schedules.go b/server/fleet/cron_schedules.go index 7c7dc82f40..9ccc3f3932 100644 --- a/server/fleet/cron_schedules.go +++ b/server/fleet/cron_schedules.go @@ -30,8 +30,9 @@ const ( CronMaintainedApps CronScheduleName = "maintained_apps" // CronRefreshVPPAppVersions updates the versions of VPP apps in Fleet to the latest value. Runs // every 1h. - CronRefreshVPPAppVersions CronScheduleName = "refresh_vpp_app_versions" - CronAppleMDMIPhoneIPadReviver CronScheduleName = "apple_mdm_iphone_ipad_reviver" + CronRefreshVPPAppVersions CronScheduleName = "refresh_vpp_app_versions" + CronAppleMDMIPhoneIPadReviver CronScheduleName = "apple_mdm_iphone_ipad_reviver" + CronUpcomingActivitiesMaintenance CronScheduleName = "upcoming_activities_maintenance" ) type CronSchedulesService interface { diff --git a/server/fleet/datastore.go b/server/fleet/datastore.go index c2376290bc..9a816c044d 100644 --- a/server/fleet/datastore.go +++ b/server/fleet/datastore.go @@ -693,6 +693,7 @@ type Datastore interface { ListHostPastActivities(ctx context.Context, hostID uint, opt ListOptions) ([]*Activity, *PaginationMetadata, error) IsExecutionPendingForHost(ctx context.Context, hostID uint, scriptID uint) (bool, error) GetHostUpcomingActivityMeta(ctx context.Context, hostID uint, executionID string) (*UpcomingActivityMeta, error) + UnblockHostsUpcomingActivityQueue(ctx context.Context, maxHosts int) (int, error) /////////////////////////////////////////////////////////////////////////////// // StatisticsStore diff --git a/server/mock/datastore_mock.go b/server/mock/datastore_mock.go index af2726e028..1877ccd31a 100644 --- a/server/mock/datastore_mock.go +++ b/server/mock/datastore_mock.go @@ -514,6 +514,8 @@ type IsExecutionPendingForHostFunc func(ctx context.Context, hostID uint, script type GetHostUpcomingActivityMetaFunc func(ctx context.Context, hostID uint, executionID string) (*fleet.UpcomingActivityMeta, error) +type UnblockHostsUpcomingActivityQueueFunc func(ctx context.Context, maxHosts int) (int, error) + type ShouldSendStatisticsFunc func(ctx context.Context, frequency time.Duration, config config.FleetConfig) (fleet.StatisticsPayload, bool, error) type RecordStatisticsSentFunc func(ctx context.Context) error @@ -2099,6 +2101,9 @@ type DataStore struct { GetHostUpcomingActivityMetaFunc GetHostUpcomingActivityMetaFunc GetHostUpcomingActivityMetaFuncInvoked bool + UnblockHostsUpcomingActivityQueueFunc UnblockHostsUpcomingActivityQueueFunc + UnblockHostsUpcomingActivityQueueFuncInvoked bool + ShouldSendStatisticsFunc ShouldSendStatisticsFunc ShouldSendStatisticsFuncInvoked bool @@ -5093,6 +5098,13 @@ func (s *DataStore) GetHostUpcomingActivityMeta(ctx context.Context, hostID uint return s.GetHostUpcomingActivityMetaFunc(ctx, hostID, executionID) } +func (s *DataStore) UnblockHostsUpcomingActivityQueue(ctx context.Context, maxHosts int) (int, error) { + s.mu.Lock() + s.UnblockHostsUpcomingActivityQueueFuncInvoked = true + s.mu.Unlock() + return s.UnblockHostsUpcomingActivityQueueFunc(ctx, maxHosts) +} + func (s *DataStore) ShouldSendStatistics(ctx context.Context, frequency time.Duration, config config.FleetConfig) (fleet.StatisticsPayload, bool, error) { s.mu.Lock() s.ShouldSendStatisticsFuncInvoked = true