Bugfix: catch-all cron job to avoid blocked upcoming activities queue (#29477)
This commit is contained in:
+39
-19
@@ -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
|
||||
}
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user