diff --git a/cmd/fleet/serve.go b/cmd/fleet/serve.go index d645f419b8..44bdd5fe83 100644 --- a/cmd/fleet/serve.go +++ b/cmd/fleet/serve.go @@ -50,6 +50,7 @@ import ( "github.com/fleetdm/fleet/v4/server/pubsub" "github.com/fleetdm/fleet/v4/server/service" "github.com/fleetdm/fleet/v4/server/service/async" + "github.com/fleetdm/fleet/v4/server/service/redis_lock" "github.com/fleetdm/fleet/v4/server/service/redis_policy_set" "github.com/fleetdm/fleet/v4/server/sso" "github.com/fleetdm/fleet/v4/server/version" @@ -691,6 +692,7 @@ the way that the Fleet server works. } var softwareInstallStore fleet.SoftwareInstallerStore + var distributedLock fleet.Lock if license.IsPremium() { profileMatcher := apple_mdm.NewProfileMatcher(redisPool) if config.S3.SoftwareInstallersBucket != "" { @@ -718,6 +720,7 @@ the way that the Fleet server works. } } + distributedLock = redis_lock.NewLock(redisPool) svc, err = eeservice.NewService( svc, ds, @@ -730,6 +733,7 @@ the way that the Fleet server works. ssoSessionStore, profileMatcher, softwareInstallStore, + distributedLock, ) if err != nil { initFatal(err, "initial Fleet Premium service") @@ -870,7 +874,7 @@ the way that the Fleet server works. } else { config.Calendar.Periodicity = 5 * time.Minute } - return cron.NewCalendarSchedule(ctx, instanceID, ds, config.Calendar, logger) + return cron.NewCalendarSchedule(ctx, instanceID, ds, distributedLock, config.Calendar, logger) }, ); err != nil { initFatal(err, "failed to register calendar schedule") diff --git a/ee/server/service/calendar.go b/ee/server/service/calendar.go index eee0e84fa4..34b4e5daba 100644 --- a/ee/server/service/calendar.go +++ b/ee/server/service/calendar.go @@ -3,16 +3,24 @@ package service import ( "context" "fmt" + "sync" + "github.com/fleetdm/fleet/v4/server/authz" "github.com/fleetdm/fleet/v4/server/contexts/ctxerr" "github.com/fleetdm/fleet/v4/server/fleet" "github.com/fleetdm/fleet/v4/server/service/calendar" "github.com/go-kit/log/level" - "sync" + "github.com/google/uuid" ) +var asyncCalendarProcessing bool +var asyncMutex sync.Mutex + func (svc *Service) CalendarWebhook(ctx context.Context, eventUUID string, channelID string, resourceState string) error { + // We don't want the sender to cancel the context since we want to make sure we process the webhook. + ctx = context.WithoutCancel(ctx) + appConfig, err := svc.ds.AppConfig(ctx) if err != nil { return fmt.Errorf("load app config: %w", err) @@ -36,9 +44,9 @@ func (svc *Service) CalendarWebhook(ctx context.Context, eventUUID string, chann svc.authz.SkipAuthorization(ctx) if fleet.IsNotFound(err) { // We could try to stop the channel callbacks here, but that may not be secure since we don't know if the request is legitimate - level.Warn(svc.logger).Log("msg", "Received calendar callback, but did not find corresponding event in database", "event_uuid", + level.Info(svc.logger).Log("msg", "Received calendar callback, but did not find corresponding event in database", "event_uuid", eventUUID, "channel_id", channelID) - return err + return nil } return err } @@ -47,7 +55,7 @@ func (svc *Service) CalendarWebhook(ctx context.Context, eventUUID string, chann return fmt.Errorf("calendar event %s has no team ID", eventUUID) } - localConfig := &calendar.CalendarConfig{ + localConfig := &calendar.Config{ GoogleCalendarIntegration: *googleCalendarIntegrationConfig, ServerURL: appConfig.ServerSettings.ServerURL, } @@ -63,6 +71,60 @@ func (svc *Service) CalendarWebhook(ctx context.Context, eventUUID string, chann return authz.ForbiddenWithInternal(fmt.Sprintf("calendar channel ID mismatch: %s != %s", savedChannelID, channelID), nil, nil, nil) } + lockValue, reserved, err := svc.getCalendarLock(ctx, eventUUID, true) + if err != nil { + return err + } + // If lock has been reserved by cron, we will need to re-process this event in case the calendar event was changed after the cron job read it. + if lockValue == "" && !reserved { + // We did not get a lock, so there is nothing to do here + return nil + } + + if !reserved { + unlocked := false + defer func() { + if !unlocked { + svc.releaseCalendarLock(ctx, eventUUID, lockValue) + } + }() + + // Remove event from the queue so that we don't process this event again. + // Note: This item can be added back to the queue while we are processing it. + err = svc.distributedLock.RemoveFromSet(ctx, calendar.QueueKey, eventUUID) + if err != nil { + return ctxerr.Wrap(ctx, err, "remove calendar event from queue") + } + + err = svc.processCalendarEvent(ctx, eventDetails, googleCalendarIntegrationConfig, userCalendar) + if err != nil { + return err + } + svc.releaseCalendarLock(ctx, eventUUID, lockValue) + unlocked = true + } + + // Now, we need to check if there are any events in the queue that need to be re-processed. + asyncMutex.Lock() + defer asyncMutex.Unlock() + if !asyncCalendarProcessing { + eventIDs, err := svc.distributedLock.GetSet(ctx, calendar.QueueKey) + if err != nil { + return ctxerr.Wrap(ctx, err, "get calendar event queue") + } + if len(eventIDs) > 0 { + asyncCalendarProcessing = true + go svc.processCalendarAsync(ctx, eventIDs) + } + return nil + } + + return nil +} + +func (svc *Service) processCalendarEvent(ctx context.Context, eventDetails *fleet.CalendarEventDetails, + googleCalendarIntegrationConfig *fleet.GoogleCalendarIntegration, userCalendar fleet.UserCalendar) error { + genBodyFn := func(conflict bool) (body string, ok bool, err error) { // This function is called when a new event is being created. @@ -113,7 +175,7 @@ func (svc *Service) CalendarWebhook(ctx context.Context, eventUUID string, chann return calendar.GenerateCalendarEventBody(ctx, svc.ds, team.Name, host, &sync.Map{}, conflict, svc.logger), true, nil } - err = userCalendar.Configure(eventDetails.Email) + err := userCalendar.Configure(eventDetails.Email) if err != nil { return ctxerr.Wrap(ctx, err, "configure calendar") } @@ -124,10 +186,156 @@ func (svc *Service) CalendarWebhook(ctx context.Context, eventUUID string, chann if updated && event != nil { // Event was updated, so we need to save it _, err = svc.ds.CreateOrUpdateCalendarEvent(ctx, event.UUID, event.Email, event.StartTime, event.EndTime, event.Data, - event.TimeZone, eventDetails.ID, fleet.CalendarWebhookStatusNone) + event.TimeZone, eventDetails.HostID, fleet.CalendarWebhookStatusNone) if err != nil { return ctxerr.Wrap(ctx, err, "create or update calendar event") } } + return nil } + +func (svc *Service) releaseCalendarLock(ctx context.Context, eventUUID string, lockValue string) { + ok, err := svc.distributedLock.ReleaseLock(ctx, calendar.LockKeyPrefix+eventUUID, lockValue) + if err != nil { + level.Error(svc.logger).Log("msg", "Failed to release calendar lock", "err", err) + } + if !ok { + // If the lock was not released, it will expire on its own. + level.Warn(svc.logger).Log("msg", "Failed to release calendar lock") + } +} + +func (svc *Service) getCalendarLock(ctx context.Context, eventUUID string, addToQueue bool) (lockValue string, reserved bool, err error) { + // Check if lock has been reserved, which means we can't have it. + reservedValue, err := svc.distributedLock.Get(ctx, calendar.ReservedLockKeyPrefix+eventUUID) + if err != nil { + return "", false, ctxerr.Wrap(ctx, err, "get calendar reserved lock") + } + reserved = reservedValue != nil + if reserved && !addToQueue { + // We flag the lock as reserved. + return "", reserved, nil + } + var lockAcquired bool + if !reserved { + // Try to acquire the lock + lockValue = uuid.New().String() + lockAcquired, err = svc.distributedLock.AcquireLock(ctx, calendar.LockKeyPrefix+eventUUID, lockValue, 0) + if err != nil { + return "", false, ctxerr.Wrap(ctx, err, "acquire calendar lock") + } + } + if (!lockAcquired || reserved) && addToQueue { + // Could not acquire lock, so we are already processing this event. In this case, we add the event to + // the queue (actually a set) to indicate that we need to re-process the event. + err = svc.distributedLock.AddToSet(ctx, calendar.QueueKey, eventUUID) + if err != nil { + return "", false, ctxerr.Wrap(ctx, err, "add calendar event to queue") + } + + if reserved { + // We flag the lock as reserved. + return "", reserved, nil + } + + // Try to acquire the lock again in case it was released while we were adding the event to the queue. + lockAcquired, err = svc.distributedLock.AcquireLock(ctx, calendar.LockKeyPrefix+eventUUID, lockValue, 0) + if err != nil { + return "", false, ctxerr.Wrap(ctx, err, "acquire calendar lock again") + } + + if !lockAcquired { + // We could not acquire the lock, so we are done here. + return "", reserved, nil + } + } + return lockValue, false, nil +} + +func (svc *Service) processCalendarAsync(ctx context.Context, eventIDs []string) { + defer func() { + asyncMutex.Lock() + asyncCalendarProcessing = false + asyncMutex.Unlock() + }() + for { + if len(eventIDs) == 0 { + return + } + for _, eventUUID := range eventIDs { + if ok := svc.processCalendarEventAsync(ctx, eventUUID); !ok { + return + } + } + + // Now we check whether there are any more events in the queue. + var err error + eventIDs, err = svc.distributedLock.GetSet(ctx, calendar.QueueKey) + if err != nil { + level.Error(svc.logger).Log("msg", "Failed to get calendar event queue", "err", err) + return + } + } +} + +func (svc *Service) processCalendarEventAsync(ctx context.Context, eventUUID string) bool { + lockValue, _, err := svc.getCalendarLock(ctx, eventUUID, false) + if err != nil { + level.Error(svc.logger).Log("msg", "Failed to get calendar lock", "err", err) + return false + } + if lockValue == "" { + // We did not get a lock, so there is nothing to do here + return true + } + defer svc.releaseCalendarLock(ctx, eventUUID, lockValue) + + // Remove event from the queue so that we don't process this event again. + // Note: This item can be added back to the queue while we are processing it. + err = svc.distributedLock.RemoveFromSet(ctx, calendar.QueueKey, eventUUID) + if err != nil { + level.Error(svc.logger).Log("msg", "Failed to remove calendar event from queue", "err", err) + return false + } + + appConfig, err := svc.ds.AppConfig(ctx) + if err != nil { + level.Error(svc.logger).Log("msg", "Failed to load app config", "err", err) + return false + } + + if len(appConfig.Integrations.GoogleCalendar) == 0 { + // Google Calendar integration is not configured + return true + } + googleCalendarIntegrationConfig := appConfig.Integrations.GoogleCalendar[0] + + eventDetails, err := svc.ds.GetCalendarEventDetailsByUUID(ctx, eventUUID) + if err != nil { + if fleet.IsNotFound(err) { + // We found this event when the callback initially came in. So the event may have been removed or re-created since then. + return true + } + level.Error(svc.logger).Log("msg", "Failed to get calendar event details", "err", err) + return false + } + if eventDetails.TeamID == nil { + // Should not happen + level.Error(svc.logger).Log("msg", "Calendar event has no team ID", "uuid", eventUUID) + return false + } + + localConfig := &calendar.Config{ + GoogleCalendarIntegration: *googleCalendarIntegrationConfig, + ServerURL: appConfig.ServerSettings.ServerURL, + } + userCalendar := calendar.CreateUserCalendarFromConfig(ctx, localConfig, svc.logger) + + err = svc.processCalendarEvent(ctx, eventDetails, googleCalendarIntegrationConfig, userCalendar) + if err != nil { + level.Error(svc.logger).Log("msg", "Failed to process calendar event", "err", err) + return false + } + return true +} diff --git a/ee/server/service/mdm_external_test.go b/ee/server/service/mdm_external_test.go index 424d417ee6..9c3c3fbeb0 100644 --- a/ee/server/service/mdm_external_test.go +++ b/ee/server/service/mdm_external_test.go @@ -82,6 +82,7 @@ func setupMockDatastorePremiumService() (*mock.Store, *eeservice.Service, contex nil, nil, nil, + nil, ) if err != nil { panic(err) diff --git a/ee/server/service/service.go b/ee/server/service/service.go index 0a0eab6861..3a1ac66a16 100644 --- a/ee/server/service/service.go +++ b/ee/server/service/service.go @@ -28,6 +28,7 @@ type Service struct { depService *apple_mdm.DEPService profileMatcher fleet.ProfileMatcher softwareInstallStore fleet.SoftwareInstallerStore + distributedLock fleet.Lock } func NewService( @@ -42,6 +43,7 @@ func NewService( sso sso.SessionStore, profileMatcher fleet.ProfileMatcher, softwareInstallStore fleet.SoftwareInstallerStore, + distributedLock fleet.Lock, ) (*Service, error) { authorizer, err := authz.NewAuthorizer() if err != nil { @@ -61,6 +63,7 @@ func NewService( depService: apple_mdm.NewDEPService(ds, depStorage, logger), profileMatcher: profileMatcher, softwareInstallStore: softwareInstallStore, + distributedLock: distributedLock, } // Override methods that can't be easily overriden via diff --git a/server/cron/calendar_cron.go b/server/cron/calendar_cron.go index c1aaab2217..908c9a6c0f 100644 --- a/server/cron/calendar_cron.go +++ b/server/cron/calendar_cron.go @@ -4,6 +4,7 @@ import ( "context" "errors" "fmt" + "github.com/google/uuid" "slices" "sync" "time" @@ -27,6 +28,7 @@ func NewCalendarSchedule( ctx context.Context, instanceID string, ds fleet.Datastore, + distributedLock fleet.Lock, serverConfig config.CalendarConfig, logger kitlog.Logger, ) (*schedule.Schedule, error) { @@ -47,7 +49,7 @@ func NewCalendarSchedule( schedule.WithJob( "calendar_events", func(ctx context.Context) error { - return cronCalendarEvents(ctx, ds, serverConfig, logger) + return cronCalendarEvents(ctx, ds, distributedLock, serverConfig, logger) }, ), ) @@ -55,7 +57,8 @@ func NewCalendarSchedule( return s, nil } -func cronCalendarEvents(ctx context.Context, ds fleet.Datastore, serverConfig config.CalendarConfig, logger kitlog.Logger) error { +func cronCalendarEvents(ctx context.Context, ds fleet.Datastore, distributedLock fleet.Lock, serverConfig config.CalendarConfig, + logger kitlog.Logger) error { appConfig, err := ds.AppConfig(ctx) if err != nil { return fmt.Errorf("load app config: %w", err) @@ -78,14 +81,14 @@ func cronCalendarEvents(ctx context.Context, ds fleet.Datastore, serverConfig co return fmt.Errorf("list teams: %w", err) } - localConfig := &calendar.CalendarConfig{ + localConfig := &calendar.Config{ CalendarConfig: serverConfig, GoogleCalendarIntegration: *googleCalendarIntegrationConfig, ServerURL: appConfig.ServerSettings.ServerURL, } for _, team := range teams { if err := cronCalendarEventsForTeam( - ctx, ds, localConfig, *team, appConfig.OrgInfo.OrgName, domain, logger, + ctx, ds, distributedLock, localConfig, *team, appConfig.OrgInfo.OrgName, domain, logger, ); err != nil { level.Info(logger).Log("msg", "events calendar cron", "team_id", team.ID, "err", err) } @@ -97,7 +100,8 @@ func cronCalendarEvents(ctx context.Context, ds fleet.Datastore, serverConfig co func cronCalendarEventsForTeam( ctx context.Context, ds fleet.Datastore, - calendarConfig *calendar.CalendarConfig, + distributedLock fleet.Lock, + calendarConfig *calendar.Config, team fleet.Team, orgName string, domain string, @@ -179,7 +183,7 @@ func cronCalendarEventsForTeam( // Process hosts that are failing calendar policies. start = time.Now() - processCalendarFailingHosts(ctx, ds, calendarConfig, orgName, failingHosts, logger) + processCalendarFailingHosts(ctx, ds, distributedLock, calendarConfig, orgName, failingHosts, logger) level.Debug(logger).Log( "msg", "failing_hosts", "took", time.Since(start), ) @@ -197,7 +201,8 @@ func cronCalendarEventsForTeam( func processCalendarFailingHosts( ctx context.Context, ds fleet.Datastore, - calendarConfig *calendar.CalendarConfig, + distributedLock fleet.Lock, + calendarConfig *calendar.Config, orgName string, hosts []fleet.HostPolicyMembershipData, logger kitlog.Logger, @@ -253,7 +258,8 @@ func processCalendarFailingHosts( switch { case err == nil && !expiredEvent: if err := processFailingHostExistingCalendarEvent( - ctx, ds, userCalendar, orgName, hostCalendarEvent, calendarEvent, host, &policyIDtoPolicy, calendarConfig, logger, + ctx, ds, distributedLock, userCalendar, orgName, hostCalendarEvent, calendarEvent, host, &policyIDtoPolicy, + calendarConfig, logger, ); err != nil { level.Info(logger).Log("msg", "process failing host existing calendar event", "err", err) continue // continue with next host @@ -303,15 +309,89 @@ func filterHostsWithSameEmail(hosts []fleet.HostPolicyMembershipData) []fleet.Ho func processFailingHostExistingCalendarEvent( ctx context.Context, ds fleet.Datastore, + distributedLock fleet.Lock, userCalendar fleet.UserCalendar, orgName string, hostCalendarEvent *fleet.HostCalendarEvent, calendarEvent *fleet.CalendarEvent, host fleet.HostPolicyMembershipData, policyIDtoPolicy *sync.Map, - calendarConfig *calendar.CalendarConfig, + calendarConfig *calendar.Config, logger kitlog.Logger, ) error { + + // Try to acquire the lock. Lock is needed to ensure calendar callback is not processed for this event at the same time. + eventUUID := calendarEvent.UUID + lockValue := uuid.New().String() + lockAcquired, err := distributedLock.AcquireLock(ctx, calendar.LockKeyPrefix+eventUUID, lockValue, 0) + if err != nil { + return fmt.Errorf("acquire calendar lock: %w", err) + } + lockReserved := false + if !lockAcquired { + // Lock was not acquired. We reserve the lock and try to acquire it until we do. + var timeoutMs uint64 = 2 * 60 * 1000 + lockAcquired, err = distributedLock.AcquireLock(ctx, calendar.ReservedLockKeyPrefix+eventUUID, lockValue, timeoutMs) + if err != nil { + return fmt.Errorf("reserve calendar lock: %w", err) + } + if !lockAcquired { + // Lock was not reserved. Another cron job is processing this event. This is not expected. + return errors.New("could not reserve calendar lock") + } + lockReserved = true + done := make(chan struct{}) + go func() { + for { + // Keep trying to get the lock. + lockAcquired, err = distributedLock.AcquireLock(ctx, calendar.LockKeyPrefix+eventUUID, lockValue, 0) + if err != nil || lockAcquired { + done <- struct{}{} + return + } + time.Sleep(100 * time.Millisecond) + } + }() + select { + case <-done: + // Lock was acquired. + if err != nil { + return fmt.Errorf("try to acquire calendar lock: %w", err) + } + case <-time.After(time.Duration(timeoutMs) * time.Millisecond): + // We couldn't acquire the lock in time. + return errors.New("could not acquire calendar lock in time") + } + } + defer func() { + // Release locks. + if lockReserved { + ok, err := distributedLock.ReleaseLock(ctx, calendar.ReservedLockKeyPrefix+eventUUID, lockValue) + if err != nil { + level.Error(logger).Log("msg", "Failed to release calendar reserve lock", "err", err) + } + if !ok { + // If the lock was not released, it will expire on its own. + level.Warn(logger).Log("msg", "Failed to release calendar reserve lock") + } + } + ok, err := distributedLock.ReleaseLock(ctx, calendar.LockKeyPrefix+eventUUID, lockValue) + if err != nil { + level.Error(logger).Log("msg", "Failed to release calendar lock", "err", err) + } + if !ok { + // If the lock was not released, it will expire on its own. + level.Warn(logger).Log("msg", "Failed to release calendar lock") + } + }() + + // Remove event from the queue so that we don't process this event again. + // Note: This item can be added back to the queue while we are processing it. + err = distributedLock.RemoveFromSet(ctx, calendar.QueueKey, eventUUID) + if err != nil { + return fmt.Errorf("remove calendar event from queue: %w", err) + } + updatedEvent := calendarEvent updated := false now := time.Now() @@ -492,7 +572,7 @@ func addBusinessDay(date time.Time) time.Time { func removeCalendarEventsFromPassingHosts( ctx context.Context, ds fleet.Datastore, - calendarConfig *calendar.CalendarConfig, + calendarConfig *calendar.Config, hosts []fleet.HostPolicyMembershipData, logger kitlog.Logger, ) { @@ -601,9 +681,9 @@ func cronCalendarEventsCleanup(ctx context.Context, ds fleet.Datastore, logger k } var userCalendar fleet.UserCalendar - var calConfig *calendar.CalendarConfig + var calConfig *calendar.Config if len(appConfig.Integrations.GoogleCalendar) > 0 { - calConfig = &calendar.CalendarConfig{ + calConfig = &calendar.Config{ GoogleCalendarIntegration: *appConfig.Integrations.GoogleCalendar[0], ServerURL: appConfig.ServerSettings.ServerURL, } @@ -657,7 +737,7 @@ func cronCalendarEventsCleanup(ctx context.Context, ds fleet.Datastore, logger k func deleteAllCalendarEvents( ctx context.Context, ds fleet.Datastore, - calendarConfig *calendar.CalendarConfig, + calendarConfig *calendar.Config, teamID *uint, logger kitlog.Logger, ) error { @@ -670,7 +750,7 @@ func deleteAllCalendarEvents( } func deleteCalendarEventsInParallel( - ctx context.Context, ds fleet.Datastore, calendarConfig *calendar.CalendarConfig, calendarEvents []*fleet.CalendarEvent, + ctx context.Context, ds fleet.Datastore, calendarConfig *calendar.Config, calendarEvents []*fleet.CalendarEvent, logger kitlog.Logger, ) { if len(calendarEvents) > 0 { @@ -703,7 +783,7 @@ func deleteCalendarEventsInParallel( func cleanupTeamCalendarEvents( ctx context.Context, ds fleet.Datastore, - calendarConfig *calendar.CalendarConfig, + calendarConfig *calendar.Config, team fleet.Team, logger kitlog.Logger, ) error { diff --git a/server/cron/calendar_cron_test.go b/server/cron/calendar_cron_test.go index bafb937422..0111f72291 100644 --- a/server/cron/calendar_cron_test.go +++ b/server/cron/calendar_cron_test.go @@ -11,14 +11,15 @@ import ( "testing" "time" - "github.com/fleetdm/fleet/v4/server/config" - "github.com/fleetdm/fleet/v4/server/ptr" - "github.com/stretchr/testify/assert" - "github.com/fleetdm/fleet/v4/ee/server/calendar" + "github.com/fleetdm/fleet/v4/server/config" + "github.com/fleetdm/fleet/v4/server/datastore/redis/redistest" "github.com/fleetdm/fleet/v4/server/fleet" "github.com/fleetdm/fleet/v4/server/mock" + "github.com/fleetdm/fleet/v4/server/ptr" + "github.com/fleetdm/fleet/v4/server/service/redis_lock" kitlog "github.com/go-kit/log" + "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" ) @@ -199,7 +200,8 @@ func TestEventForDifferentHost(t *testing.T) { return hcEvent, calEvent, nil } - err := cronCalendarEvents(ctx, ds, defaultCalendarConfig, logger) + pool := redistest.SetupRedis(t, t.Name(), false, false, false) + err := cronCalendarEvents(ctx, ds, redis_lock.NewLock(pool), defaultCalendarConfig, logger) require.NoError(t, err) } @@ -373,7 +375,8 @@ func TestCalendarEventsMultipleHosts(t *testing.T) { return nil, nil } - err := cronCalendarEvents(ctx, ds, defaultCalendarConfig, logger) + pool := redistest.SetupRedis(t, t.Name(), false, false, false) + err := cronCalendarEvents(ctx, ds, redis_lock.NewLock(pool), defaultCalendarConfig, logger) require.NoError(t, err) eventsMu.Lock() @@ -662,7 +665,9 @@ func TestCalendarEvents1KHosts(t *testing.T) { return nil, nil } - err := cronCalendarEvents(ctx, ds, defaultCalendarConfig, logger) + pool := redistest.SetupRedis(t, t.Name(), false, false, false) + distributedLock := redis_lock.NewLock(pool) + err := cronCalendarEvents(ctx, ds, distributedLock, defaultCalendarConfig, logger) require.NoError(t, err) createdCalendarEvents := calendar.ListGoogleMockEvents() @@ -699,7 +704,7 @@ func TestCalendarEvents1KHosts(t *testing.T) { return nil } - err = cronCalendarEvents(ctx, ds, defaultCalendarConfig, logger) + err = cronCalendarEvents(ctx, ds, distributedLock, defaultCalendarConfig, logger) require.NoError(t, err) createdCalendarEvents = calendar.ListGoogleMockEvents() @@ -948,7 +953,8 @@ func TestEventBody(t *testing.T) { return nil, nil } - err := cronCalendarEvents(ctx, ds, defaultCalendarConfig, logger) + pool := redistest.SetupRedis(t, t.Name(), false, false, false) + err := cronCalendarEvents(ctx, ds, redis_lock.NewLock(pool), defaultCalendarConfig, logger) require.NoError(t, err) numberOfEvents := 7 diff --git a/server/datastore/mysql/calendar_events.go b/server/datastore/mysql/calendar_events.go index fa58b8f3af..433e15a16a 100644 --- a/server/datastore/mysql/calendar_events.go +++ b/server/datastore/mysql/calendar_events.go @@ -5,6 +5,7 @@ import ( "database/sql" "errors" "fmt" + "github.com/google/uuid" "time" "github.com/fleetdm/fleet/v4/server/contexts/ctxerr" @@ -12,9 +13,11 @@ import ( "github.com/jmoiron/sqlx" ) +const calendarEventCols = `ce.id, ce.uuid, ce.email, ce.start_time, ce.end_time, ce.event, ce.timezone, ce.created_at, ce.updated_at` + func (ds *Datastore) CreateOrUpdateCalendarEvent( ctx context.Context, - uuid string, + uuidStr string, email string, startTime time.Time, endTime time.Time, @@ -23,11 +26,15 @@ func (ds *Datastore) CreateOrUpdateCalendarEvent( hostID uint, webhookStatus fleet.CalendarWebhookStatus, ) (*fleet.CalendarEvent, error) { + UUID, err := uuid.Parse(uuidStr) + if err != nil { + return nil, ctxerr.Wrap(ctx, err, "invalid uuid") + } var id int64 if err := ds.withRetryTxx(ctx, func(tx sqlx.ExtContext) error { const calendarEventsQuery = ` INSERT INTO calendar_events ( - uuid, + uuid_bin, email, start_time, end_time, @@ -35,7 +42,7 @@ func (ds *Datastore) CreateOrUpdateCalendarEvent( timezone ) VALUES (?, ?, ?, ?, ?, ?) ON DUPLICATE KEY UPDATE - uuid = VALUES(uuid), + uuid_bin = VALUES(uuid_bin), start_time = VALUES(start_time), end_time = VALUES(end_time), event = VALUES(event), @@ -45,7 +52,7 @@ func (ds *Datastore) CreateOrUpdateCalendarEvent( result, err := tx.ExecContext( ctx, calendarEventsQuery, - uuid, + UUID[:], email, startTime, endTime, @@ -98,9 +105,7 @@ func (ds *Datastore) CreateOrUpdateCalendarEvent( } func getCalendarEventByID(ctx context.Context, q sqlx.QueryerContext, id uint) (*fleet.CalendarEvent, error) { - const calendarEventsQuery = ` - SELECT * FROM calendar_events WHERE id = ?; - ` + const calendarEventsQuery = "SELECT " + calendarEventCols + " FROM calendar_events ce WHERE id = ?" var calendarEvent fleet.CalendarEvent err := sqlx.GetContext(ctx, q, &calendarEvent, calendarEventsQuery, id) if err != nil { @@ -113,9 +118,7 @@ func getCalendarEventByID(ctx context.Context, q sqlx.QueryerContext, id uint) ( } func (ds *Datastore) GetCalendarEvent(ctx context.Context, email string) (*fleet.CalendarEvent, error) { - const calendarEventsQuery = ` - SELECT * FROM calendar_events WHERE email = ?; - ` + const calendarEventsQuery = "SELECT " + calendarEventCols + " FROM calendar_events ce WHERE email = ?" var calendarEvent fleet.CalendarEvent err := sqlx.GetContext(ctx, ds.reader(ctx), &calendarEvent, calendarEventsQuery, email) if err != nil { @@ -127,29 +130,37 @@ func (ds *Datastore) GetCalendarEvent(ctx context.Context, email string) (*fleet return &calendarEvent, nil } -func (ds *Datastore) GetCalendarEventDetailsByUUID(ctx context.Context, uuid string) (*fleet.CalendarEventDetails, error) { +func (ds *Datastore) GetCalendarEventDetailsByUUID(ctx context.Context, uuidStr string) (*fleet.CalendarEventDetails, error) { + UUID, err := uuid.Parse(uuidStr) + if err != nil { + return nil, ctxerr.Wrap(ctx, err, "invalid uuid") + } const calendarEventsByUUIDQuery = ` - SELECT ce.*, h.team_id as team_id, h.id as host_id FROM calendar_events ce + SELECT ` + calendarEventCols + `, h.team_id as team_id, h.id as host_id FROM calendar_events ce LEFT JOIN host_calendar_events hce ON hce.calendar_event_id = ce.id LEFT JOIN hosts h ON h.id = hce.host_id - WHERE ce.uuid = ?; + WHERE ce.uuid_bin = ?; ` var calendarEvent fleet.CalendarEventDetails - err := sqlx.GetContext(ctx, ds.reader(ctx), &calendarEvent, calendarEventsByUUIDQuery, uuid) + err = sqlx.GetContext(ctx, ds.reader(ctx), &calendarEvent, calendarEventsByUUIDQuery, UUID[:]) if err != nil { if errors.Is(err, sql.ErrNoRows) { - return nil, ctxerr.Wrap(ctx, notFound("CalendarEvent").WithMessage(fmt.Sprintf("uuid: %s", uuid))) + return nil, ctxerr.Wrap(ctx, notFound("CalendarEvent").WithMessage(fmt.Sprintf("uuid: %s", UUID.String()))) } return nil, ctxerr.Wrap(ctx, err, "get calendar event") } return &calendarEvent, nil } -func (ds *Datastore) UpdateCalendarEvent(ctx context.Context, calendarEventID uint, uuid string, startTime time.Time, endTime time.Time, +func (ds *Datastore) UpdateCalendarEvent(ctx context.Context, calendarEventID uint, uuidStr string, startTime time.Time, endTime time.Time, data []byte, timeZone string) error { + UUID, err := uuid.Parse(uuidStr) + if err != nil { + return ctxerr.Wrap(ctx, err, "invalid uuid") + } const calendarEventsQuery = ` UPDATE calendar_events SET - uuid = ?, + uuid_bin = ?, start_time = ?, end_time = ?, event = ?, @@ -157,7 +168,7 @@ func (ds *Datastore) UpdateCalendarEvent(ctx context.Context, calendarEventID ui updated_at = CURRENT_TIMESTAMP WHERE id = ?; ` - if _, err := ds.writer(ctx).ExecContext(ctx, calendarEventsQuery, uuid, startTime, endTime, data, timeZone, + if _, err := ds.writer(ctx).ExecContext(ctx, calendarEventsQuery, UUID[:], startTime, endTime, data, timeZone, calendarEventID); err != nil { return ctxerr.Wrap(ctx, err, "update calendar event") } @@ -186,7 +197,7 @@ func (ds *Datastore) GetHostCalendarEvent(ctx context.Context, hostID uint) (*fl return nil, nil, ctxerr.Wrap(ctx, err, "get host calendar event") } const calendarEventsQuery = ` - SELECT * FROM calendar_events WHERE id = ? + SELECT ` + calendarEventCols + ` FROM calendar_events ce WHERE id = ? ` var calendarEvent fleet.CalendarEvent if err := sqlx.GetContext(ctx, ds.reader(ctx), &calendarEvent, calendarEventsQuery, hostCalendarEvent.CalendarEventID); err != nil { @@ -200,7 +211,7 @@ func (ds *Datastore) GetHostCalendarEvent(ctx context.Context, hostID uint) (*fl func (ds *Datastore) GetHostCalendarEventByEmail(ctx context.Context, email string) (*fleet.HostCalendarEvent, *fleet.CalendarEvent, error) { const calendarEventsQuery = ` - SELECT * FROM calendar_events WHERE email = ? + SELECT ` + calendarEventCols + ` FROM calendar_events ce WHERE email = ? ` var calendarEvent fleet.CalendarEvent if err := sqlx.GetContext(ctx, ds.reader(ctx), &calendarEvent, calendarEventsQuery, email); err != nil { @@ -236,7 +247,7 @@ func (ds *Datastore) UpdateHostCalendarWebhookStatus(ctx context.Context, hostID func (ds *Datastore) ListCalendarEvents(ctx context.Context, teamID *uint) ([]*fleet.CalendarEvent, error) { calendarEventsQuery := ` - SELECT ce.* FROM calendar_events ce + SELECT ` + calendarEventCols + ` FROM calendar_events ce ` var args []interface{} @@ -258,7 +269,7 @@ func (ds *Datastore) ListCalendarEvents(ctx context.Context, teamID *uint) ([]*f func (ds *Datastore) ListOutOfDateCalendarEvents(ctx context.Context, t time.Time) ([]*fleet.CalendarEvent, error) { calendarEventsQuery := ` - SELECT ce.* FROM calendar_events ce WHERE updated_at < ? + SELECT ` + calendarEventCols + ` FROM calendar_events ce WHERE updated_at < ? ` var calendarEvents []*fleet.CalendarEvent if err := sqlx.SelectContext(ctx, ds.reader(ctx), &calendarEvents, calendarEventsQuery, t); err != nil { diff --git a/server/datastore/mysql/calendar_events_test.go b/server/datastore/mysql/calendar_events_test.go index b4ec9d510a..bee3ef6ed9 100644 --- a/server/datastore/mysql/calendar_events_test.go +++ b/server/datastore/mysql/calendar_events_test.go @@ -4,6 +4,7 @@ import ( "context" "github.com/google/uuid" "github.com/stretchr/testify/assert" + "strings" "testing" "time" @@ -79,7 +80,7 @@ func testUpdateCalendarEvent(t *testing.T, ds *Datastore) { eventDetails, err := ds.GetCalendarEventDetailsByUUID(ctx, eventUUIDNew) require.NoError(t, err) - assert.Equal(t, eventUUIDNew, eventDetails.UUID) + assert.Equal(t, strings.ToUpper(eventUUIDNew), eventDetails.UUID) assert.Equal(t, *calendarEvent, eventDetails.CalendarEvent) assert.Equal(t, host.ID, eventDetails.HostID) assert.Nil(t, eventDetails.TeamID) diff --git a/server/datastore/mysql/migrations/tables/20240707134035_AddUUIDToCalendarEvents_test.go b/server/datastore/mysql/migrations/tables/20240707134035_AddUUIDToCalendarEvents_test.go index 851997f251..4372053c53 100644 --- a/server/datastore/mysql/migrations/tables/20240707134035_AddUUIDToCalendarEvents_test.go +++ b/server/datastore/mysql/migrations/tables/20240707134035_AddUUIDToCalendarEvents_test.go @@ -21,7 +21,7 @@ func TestUp_20240707134035(t *testing.T) { // Apply current migration. applyNext(t, db) - // check that it's NULL + // check that UUID is not NULL const selectUUIDStmt = `SELECT uuid FROM calendar_events WHERE id = ?` var uuid1, uuid2 string err := db.Get(&uuid1, selectUUIDStmt, event1ID) diff --git a/server/datastore/mysql/migrations/tables/20240709132642_StoreCalendarEventsUUIDAsBinary.go b/server/datastore/mysql/migrations/tables/20240709132642_StoreCalendarEventsUUIDAsBinary.go new file mode 100644 index 0000000000..185fcf3546 --- /dev/null +++ b/server/datastore/mysql/migrations/tables/20240709132642_StoreCalendarEventsUUIDAsBinary.go @@ -0,0 +1,53 @@ +package tables + +import ( + "database/sql" + "fmt" +) + +func init() { + MigrationClient.AddMigration(Up_20240709132642, Down_20240709132642) +} + +func Up_20240709132642(tx *sql.Tx) error { + // Implementation based on: https://dev.mysql.com/blog-archive/storing-uuid-values-in-mysql-tables/ + + if _, err := tx.Exec(`ALTER TABLE calendar_events ADD COLUMN uuid_bin BINARY(16) NOT NULL`); err != nil { + return fmt.Errorf("failed to add `uuid_bin` column to `calendar_events` table: %w", err) + } + + // Convert existing UUIDs to binary format + if _, err := tx.Exec(`UPDATE calendar_events SET uuid_bin = UNHEX(REPLACE(uuid,'-','')), updated_at = updated_at`); err != nil { + return fmt.Errorf("failed to convert UUIDs to binary form for existing calendar events: %w", err) + } + + // Add unique constraint to uuid_bin column + if _, err := tx.Exec(`ALTER TABLE calendar_events ADD CONSTRAINT idx_calendar_events_uuid_bin_unique UNIQUE (uuid_bin)`); err != nil { + return fmt.Errorf("failed to add unique constraint to `uuid_bin` column in `calendar_events` table: %w", err) + } + + // Drop existing uuid column + if _, err := tx.Exec(`ALTER TABLE calendar_events DROP COLUMN uuid`); err != nil { + return fmt.Errorf("failed to drop `uuid` column from `calendar_events` table: %w", err) + } + + // Add a new GENERATED uuid column + if _, err := tx.Exec( + `ALTER TABLE calendar_events ADD COLUMN uuid VARCHAR(36) COLLATE utf8mb4_unicode_ci GENERATED ALWAYS AS ( + (INSERT( + INSERT( + INSERT( + INSERT(hex(uuid_bin),9,0,'-'), + 14,0,'-'), + 19,0,'-'), + 24,0,'-') + )) VIRTUAL`); err != nil { + return fmt.Errorf("failed to add `uuid` column to `calendar_events` table: %w", err) + } + + return nil +} + +func Down_20240709132642(_ *sql.Tx) error { + return nil +} diff --git a/server/datastore/mysql/migrations/tables/20240709132642_StoreCalendarEventsUUIDAsBinary_test.go b/server/datastore/mysql/migrations/tables/20240709132642_StoreCalendarEventsUUIDAsBinary_test.go new file mode 100644 index 0000000000..5257df21fb --- /dev/null +++ b/server/datastore/mysql/migrations/tables/20240709132642_StoreCalendarEventsUUIDAsBinary_test.go @@ -0,0 +1,49 @@ +package tables + +import ( + "github.com/google/uuid" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + "strings" + "testing" + "time" +) + +func TestUp_20240709132642(t *testing.T) { + db := applyUpToPrev(t) + + testUUID := strings.ToUpper(uuid.New().String()) + startTime := time.Now().UTC() + endTime := time.Now().UTC().Add(30 * time.Minute) + data := []byte("{\"foo\": \"bar\"}") + const insertStmtUUID = `INSERT INTO calendar_events (email, start_time, end_time, event, uuid) VALUES (?, ?, ?, ?, ?)` + eventID := execNoErrLastID(t, db, insertStmtUUID, "bob@example.com", startTime, endTime, data, testUUID) + + applyNext(t, db) + // Check that uuid and uuid_bin are correct + const selectUUIDStmt = `SELECT uuid, uuid_bin FROM calendar_events WHERE id = ?` + type event struct { + UUID string `db:"uuid"` + UUIDBin []byte `db:"uuid_bin"` + } + var e event + err := db.Get(&e, selectUUIDStmt, eventID) + require.NoError(t, err) + assert.Equal(t, testUUID, e.UUID) + uuidFromBytes, err := uuid.FromBytes(e.UUIDBin) + require.NoError(t, err) + assert.Equal(t, uuid.MustParse(testUUID), uuidFromBytes) + + // Try to use the same uuid again + const insertStmtUUIDBin = `INSERT INTO calendar_events (email, start_time, end_time, event, uuid_bin) VALUES (?, ?, ?, ?, ?)` + _, err = db.Exec(insertStmtUUIDBin, "alice@example.com", startTime, endTime, data, e.UUIDBin) + assert.Error(t, err) + + // Insert a new event with a new UUID + uuidBin := uuid.New() + eventID = execNoErrLastID(t, db, insertStmtUUIDBin, "jane@example.com", startTime, endTime, data, uuidBin[:]) + err = db.Get(&e, selectUUIDStmt, eventID) + require.NoError(t, err) + assert.Equal(t, uuidBin[:], e.UUIDBin) + assert.Equal(t, strings.ToUpper(uuidBin.String()), e.UUID) +} diff --git a/server/datastore/mysql/policies_test.go b/server/datastore/mysql/policies_test.go index f505ebc3bb..fdf9012509 100644 --- a/server/datastore/mysql/policies_test.go +++ b/server/datastore/mysql/policies_test.go @@ -3650,11 +3650,11 @@ func testGetTeamHostsPolicyMemberships(t *testing.T, ds *Datastore) { // tZ := "America/Argentina/Buenos_Aires" now := time.Now() - eventUUID1 := "event-uuid" + eventUUID1 := uuid.New().String() _, err = ds.CreateOrUpdateCalendarEvent(ctx, eventUUID1, "foo@example.com", now, now.Add(30*time.Minute), []byte(`{"foo": "bar"}`), tZ, host1.ID, fleet.CalendarWebhookStatusPending) require.NoError(t, err) - eventUUID2 := "event-uuid2" + eventUUID2 := uuid.New().String() _, err = ds.CreateOrUpdateCalendarEvent(ctx, eventUUID2, "bar@example.com", now, now.Add(30*time.Minute), []byte(`{"foo": "bar"}`), tZ, host6.ID, fleet.CalendarWebhookStatusPending) require.NoError(t, err) @@ -3740,7 +3740,7 @@ func testGetTeamHostsPolicyMemberships(t *testing.T, ds *Datastore) { _, err = ds.CreateOrUpdateCalendarEvent(ctx, eventUUID1, "foo@example.com", now, now.Add(30*time.Minute), []byte(`{"foo": "bar"}`), tZ, host2.ID, fleet.CalendarWebhookStatusPending) require.NoError(t, err) - eventUUID3 := "event-uuid3" + eventUUID3 := uuid.New().String() calendarEventHost3, err := ds.CreateOrUpdateCalendarEvent(ctx, eventUUID3, "zoo@example.com", now, now.Add(30*time.Minute), []byte(`{"foo": "bar"}`), tZ, host3.ID, fleet.CalendarWebhookStatusPending) require.NoError(t, err) diff --git a/server/datastore/mysql/schema.sql b/server/datastore/mysql/schema.sql index b4cf5151e7..3ee0c37f54 100644 --- a/server/datastore/mysql/schema.sql +++ b/server/datastore/mysql/schema.sql @@ -53,10 +53,11 @@ CREATE TABLE `calendar_events` ( `created_at` timestamp NOT NULL DEFAULT CURRENT_TIMESTAMP, `updated_at` timestamp NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP, `timezone` varchar(64) COLLATE utf8mb4_unicode_ci DEFAULT NULL, - `uuid` varchar(36) COLLATE utf8mb4_unicode_ci NOT NULL, + `uuid_bin` binary(16) NOT NULL, + `uuid` varchar(36) COLLATE utf8mb4_unicode_ci GENERATED ALWAYS AS (insert(insert(insert(insert(hex(`uuid_bin`),9,0,'-'),14,0,'-'),19,0,'-'),24,0,'-')) VIRTUAL, PRIMARY KEY (`id`), UNIQUE KEY `idx_one_calendar_event_per_email` (`email`), - UNIQUE KEY `idx_calendar_events_uuid_unique` (`uuid`) + UNIQUE KEY `idx_calendar_events_uuid_bin_unique` (`uuid_bin`) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci; /*!40101 SET character_set_client = @saved_cs_client */; /*!40101 SET @saved_cs_client = @@character_set_client */; @@ -945,9 +946,9 @@ CREATE TABLE `migration_status_tables` ( `tstamp` timestamp NULL DEFAULT CURRENT_TIMESTAMP, PRIMARY KEY (`id`), UNIQUE KEY `id` (`id`) -) ENGINE=InnoDB AUTO_INCREMENT=281 DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci; +) ENGINE=InnoDB AUTO_INCREMENT=282 DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci; /*!40101 SET character_set_client = @saved_cs_client */; -INSERT INTO `migration_status_tables` VALUES (1,0,1,'2020-01-01 01:01:01'),(2,20161118193812,1,'2020-01-01 01:01:01'),(3,20161118211713,1,'2020-01-01 01:01:01'),(4,20161118212436,1,'2020-01-01 01:01:01'),(5,20161118212515,1,'2020-01-01 01:01:01'),(6,20161118212528,1,'2020-01-01 01:01:01'),(7,20161118212538,1,'2020-01-01 01:01:01'),(8,20161118212549,1,'2020-01-01 01:01:01'),(9,20161118212557,1,'2020-01-01 01:01:01'),(10,20161118212604,1,'2020-01-01 01:01:01'),(11,20161118212613,1,'2020-01-01 01:01:01'),(12,20161118212621,1,'2020-01-01 01:01:01'),(13,20161118212630,1,'2020-01-01 01:01:01'),(14,20161118212641,1,'2020-01-01 01:01:01'),(15,20161118212649,1,'2020-01-01 01:01:01'),(16,20161118212656,1,'2020-01-01 01:01:01'),(17,20161118212758,1,'2020-01-01 01:01:01'),(18,20161128234849,1,'2020-01-01 01:01:01'),(19,20161230162221,1,'2020-01-01 01:01:01'),(20,20170104113816,1,'2020-01-01 01:01:01'),(21,20170105151732,1,'2020-01-01 01:01:01'),(22,20170108191242,1,'2020-01-01 01:01:01'),(23,20170109094020,1,'2020-01-01 01:01:01'),(24,20170109130438,1,'2020-01-01 01:01:01'),(25,20170110202752,1,'2020-01-01 01:01:01'),(26,20170111133013,1,'2020-01-01 01:01:01'),(27,20170117025759,1,'2020-01-01 01:01:01'),(28,20170118191001,1,'2020-01-01 01:01:01'),(29,20170119234632,1,'2020-01-01 01:01:01'),(30,20170124230432,1,'2020-01-01 01:01:01'),(31,20170127014618,1,'2020-01-01 01:01:01'),(32,20170131232841,1,'2020-01-01 01:01:01'),(33,20170223094154,1,'2020-01-01 01:01:01'),(34,20170306075207,1,'2020-01-01 01:01:01'),(35,20170309100733,1,'2020-01-01 01:01:01'),(36,20170331111922,1,'2020-01-01 01:01:01'),(37,20170502143928,1,'2020-01-01 01:01:01'),(38,20170504130602,1,'2020-01-01 01:01:01'),(39,20170509132100,1,'2020-01-01 01:01:01'),(40,20170519105647,1,'2020-01-01 01:01:01'),(41,20170519105648,1,'2020-01-01 01:01:01'),(42,20170831234300,1,'2020-01-01 01:01:01'),(43,20170831234301,1,'2020-01-01 01:01:01'),(44,20170831234303,1,'2020-01-01 01:01:01'),(45,20171116163618,1,'2020-01-01 01:01:01'),(46,20171219164727,1,'2020-01-01 01:01:01'),(47,20180620164811,1,'2020-01-01 01:01:01'),(48,20180620175054,1,'2020-01-01 01:01:01'),(49,20180620175055,1,'2020-01-01 01:01:01'),(50,20191010101639,1,'2020-01-01 01:01:01'),(51,20191010155147,1,'2020-01-01 01:01:01'),(52,20191220130734,1,'2020-01-01 01:01:01'),(53,20200311140000,1,'2020-01-01 01:01:01'),(54,20200405120000,1,'2020-01-01 01:01:01'),(55,20200407120000,1,'2020-01-01 01:01:01'),(56,20200420120000,1,'2020-01-01 01:01:01'),(57,20200504120000,1,'2020-01-01 01:01:01'),(58,20200512120000,1,'2020-01-01 01:01:01'),(59,20200707120000,1,'2020-01-01 01:01:01'),(60,20201011162341,1,'2020-01-01 01:01:01'),(61,20201021104586,1,'2020-01-01 01:01:01'),(62,20201102112520,1,'2020-01-01 01:01:01'),(63,20201208121729,1,'2020-01-01 01:01:01'),(64,20201215091637,1,'2020-01-01 01:01:01'),(65,20210119174155,1,'2020-01-01 01:01:01'),(66,20210326182902,1,'2020-01-01 01:01:01'),(67,20210421112652,1,'2020-01-01 01:01:01'),(68,20210506095025,1,'2020-01-01 01:01:01'),(69,20210513115729,1,'2020-01-01 01:01:01'),(70,20210526113559,1,'2020-01-01 01:01:01'),(71,20210601000001,1,'2020-01-01 01:01:01'),(72,20210601000002,1,'2020-01-01 01:01:01'),(73,20210601000003,1,'2020-01-01 01:01:01'),(74,20210601000004,1,'2020-01-01 01:01:01'),(75,20210601000005,1,'2020-01-01 01:01:01'),(76,20210601000006,1,'2020-01-01 01:01:01'),(77,20210601000007,1,'2020-01-01 01:01:01'),(78,20210601000008,1,'2020-01-01 01:01:01'),(79,20210606151329,1,'2020-01-01 01:01:01'),(80,20210616163757,1,'2020-01-01 01:01:01'),(81,20210617174723,1,'2020-01-01 01:01:01'),(82,20210622160235,1,'2020-01-01 01:01:01'),(83,20210623100031,1,'2020-01-01 01:01:01'),(84,20210623133615,1,'2020-01-01 01:01:01'),(85,20210708143152,1,'2020-01-01 01:01:01'),(86,20210709124443,1,'2020-01-01 01:01:01'),(87,20210712155608,1,'2020-01-01 01:01:01'),(88,20210714102108,1,'2020-01-01 01:01:01'),(89,20210719153709,1,'2020-01-01 01:01:01'),(90,20210721171531,1,'2020-01-01 01:01:01'),(91,20210723135713,1,'2020-01-01 01:01:01'),(92,20210802135933,1,'2020-01-01 01:01:01'),(93,20210806112844,1,'2020-01-01 01:01:01'),(94,20210810095603,1,'2020-01-01 01:01:01'),(95,20210811150223,1,'2020-01-01 01:01:01'),(96,20210818151827,1,'2020-01-01 01:01:01'),(97,20210818151828,1,'2020-01-01 01:01:01'),(98,20210818182258,1,'2020-01-01 01:01:01'),(99,20210819131107,1,'2020-01-01 01:01:01'),(100,20210819143446,1,'2020-01-01 01:01:01'),(101,20210903132338,1,'2020-01-01 01:01:01'),(102,20210915144307,1,'2020-01-01 01:01:01'),(103,20210920155130,1,'2020-01-01 01:01:01'),(104,20210927143115,1,'2020-01-01 01:01:01'),(105,20210927143116,1,'2020-01-01 01:01:01'),(106,20211013133706,1,'2020-01-01 01:01:01'),(107,20211013133707,1,'2020-01-01 01:01:01'),(108,20211102135149,1,'2020-01-01 01:01:01'),(109,20211109121546,1,'2020-01-01 01:01:01'),(110,20211110163320,1,'2020-01-01 01:01:01'),(111,20211116184029,1,'2020-01-01 01:01:01'),(112,20211116184030,1,'2020-01-01 01:01:01'),(113,20211202092042,1,'2020-01-01 01:01:01'),(114,20211202181033,1,'2020-01-01 01:01:01'),(115,20211207161856,1,'2020-01-01 01:01:01'),(116,20211216131203,1,'2020-01-01 01:01:01'),(117,20211221110132,1,'2020-01-01 01:01:01'),(118,20220107155700,1,'2020-01-01 01:01:01'),(119,20220125105650,1,'2020-01-01 01:01:01'),(120,20220201084510,1,'2020-01-01 01:01:01'),(121,20220208144830,1,'2020-01-01 01:01:01'),(122,20220208144831,1,'2020-01-01 01:01:01'),(123,20220215152203,1,'2020-01-01 01:01:01'),(124,20220223113157,1,'2020-01-01 01:01:01'),(125,20220307104655,1,'2020-01-01 01:01:01'),(126,20220309133956,1,'2020-01-01 01:01:01'),(127,20220316155700,1,'2020-01-01 01:01:01'),(128,20220323152301,1,'2020-01-01 01:01:01'),(129,20220330100659,1,'2020-01-01 01:01:01'),(130,20220404091216,1,'2020-01-01 01:01:01'),(131,20220419140750,1,'2020-01-01 01:01:01'),(132,20220428140039,1,'2020-01-01 01:01:01'),(133,20220503134048,1,'2020-01-01 01:01:01'),(134,20220524102918,1,'2020-01-01 01:01:01'),(135,20220526123327,1,'2020-01-01 01:01:01'),(136,20220526123328,1,'2020-01-01 01:01:01'),(137,20220526123329,1,'2020-01-01 01:01:01'),(138,20220608113128,1,'2020-01-01 01:01:01'),(139,20220627104817,1,'2020-01-01 01:01:01'),(140,20220704101843,1,'2020-01-01 01:01:01'),(141,20220708095046,1,'2020-01-01 01:01:01'),(142,20220713091130,1,'2020-01-01 01:01:01'),(143,20220802135510,1,'2020-01-01 01:01:01'),(144,20220818101352,1,'2020-01-01 01:01:01'),(145,20220822161445,1,'2020-01-01 01:01:01'),(146,20220831100036,1,'2020-01-01 01:01:01'),(147,20220831100151,1,'2020-01-01 01:01:01'),(148,20220908181826,1,'2020-01-01 01:01:01'),(149,20220914154915,1,'2020-01-01 01:01:01'),(150,20220915165115,1,'2020-01-01 01:01:01'),(151,20220915165116,1,'2020-01-01 01:01:01'),(152,20220928100158,1,'2020-01-01 01:01:01'),(153,20221014084130,1,'2020-01-01 01:01:01'),(154,20221027085019,1,'2020-01-01 01:01:01'),(155,20221101103952,1,'2020-01-01 01:01:01'),(156,20221104144401,1,'2020-01-01 01:01:01'),(157,20221109100749,1,'2020-01-01 01:01:01'),(158,20221115104546,1,'2020-01-01 01:01:01'),(159,20221130114928,1,'2020-01-01 01:01:01'),(160,20221205112142,1,'2020-01-01 01:01:01'),(161,20221216115820,1,'2020-01-01 01:01:01'),(162,20221220195934,1,'2020-01-01 01:01:01'),(163,20221220195935,1,'2020-01-01 01:01:01'),(164,20221223174807,1,'2020-01-01 01:01:01'),(165,20221227163855,1,'2020-01-01 01:01:01'),(166,20221227163856,1,'2020-01-01 01:01:01'),(167,20230202224725,1,'2020-01-01 01:01:01'),(168,20230206163608,1,'2020-01-01 01:01:01'),(169,20230214131519,1,'2020-01-01 01:01:01'),(170,20230303135738,1,'2020-01-01 01:01:01'),(171,20230313135301,1,'2020-01-01 01:01:01'),(172,20230313141819,1,'2020-01-01 01:01:01'),(173,20230315104937,1,'2020-01-01 01:01:01'),(174,20230317173844,1,'2020-01-01 01:01:01'),(175,20230320133602,1,'2020-01-01 01:01:01'),(176,20230330100011,1,'2020-01-01 01:01:01'),(177,20230330134823,1,'2020-01-01 01:01:01'),(178,20230405232025,1,'2020-01-01 01:01:01'),(179,20230408084104,1,'2020-01-01 01:01:01'),(180,20230411102858,1,'2020-01-01 01:01:01'),(181,20230421155932,1,'2020-01-01 01:01:01'),(182,20230425082126,1,'2020-01-01 01:01:01'),(183,20230425105727,1,'2020-01-01 01:01:01'),(184,20230501154913,1,'2020-01-01 01:01:01'),(185,20230503101418,1,'2020-01-01 01:01:01'),(186,20230515144206,1,'2020-01-01 01:01:01'),(187,20230517140952,1,'2020-01-01 01:01:01'),(188,20230517152807,1,'2020-01-01 01:01:01'),(189,20230518114155,1,'2020-01-01 01:01:01'),(190,20230520153236,1,'2020-01-01 01:01:01'),(191,20230525151159,1,'2020-01-01 01:01:01'),(192,20230530122103,1,'2020-01-01 01:01:01'),(193,20230602111827,1,'2020-01-01 01:01:01'),(194,20230608103123,1,'2020-01-01 01:01:01'),(195,20230629140529,1,'2020-01-01 01:01:01'),(196,20230629140530,1,'2020-01-01 01:01:01'),(197,20230711144622,1,'2020-01-01 01:01:01'),(198,20230721135421,1,'2020-01-01 01:01:01'),(199,20230721161508,1,'2020-01-01 01:01:01'),(200,20230726115701,1,'2020-01-01 01:01:01'),(201,20230807100822,1,'2020-01-01 01:01:01'),(202,20230814150442,1,'2020-01-01 01:01:01'),(203,20230823122728,1,'2020-01-01 01:01:01'),(204,20230906152143,1,'2020-01-01 01:01:01'),(205,20230911163618,1,'2020-01-01 01:01:01'),(206,20230912101759,1,'2020-01-01 01:01:01'),(207,20230915101341,1,'2020-01-01 01:01:01'),(208,20230918132351,1,'2020-01-01 01:01:01'),(209,20231004144339,1,'2020-01-01 01:01:01'),(210,20231009094541,1,'2020-01-01 01:01:01'),(211,20231009094542,1,'2020-01-01 01:01:01'),(212,20231009094543,1,'2020-01-01 01:01:01'),(213,20231009094544,1,'2020-01-01 01:01:01'),(214,20231016091915,1,'2020-01-01 01:01:01'),(215,20231024174135,1,'2020-01-01 01:01:01'),(216,20231025120016,1,'2020-01-01 01:01:01'),(217,20231025160156,1,'2020-01-01 01:01:01'),(218,20231031165350,1,'2020-01-01 01:01:01'),(219,20231106144110,1,'2020-01-01 01:01:01'),(220,20231107130934,1,'2020-01-01 01:01:01'),(221,20231109115838,1,'2020-01-01 01:01:01'),(222,20231121054530,1,'2020-01-01 01:01:01'),(223,20231122101320,1,'2020-01-01 01:01:01'),(224,20231130132828,1,'2020-01-01 01:01:01'),(225,20231130132931,1,'2020-01-01 01:01:01'),(226,20231204155427,1,'2020-01-01 01:01:01'),(227,20231206142340,1,'2020-01-01 01:01:01'),(228,20231207102320,1,'2020-01-01 01:01:01'),(229,20231207102321,1,'2020-01-01 01:01:01'),(230,20231207133731,1,'2020-01-01 01:01:01'),(231,20231212094238,1,'2020-01-01 01:01:01'),(232,20231212095734,1,'2020-01-01 01:01:01'),(233,20231212161121,1,'2020-01-01 01:01:01'),(234,20231215122713,1,'2020-01-01 01:01:01'),(235,20231219143041,1,'2020-01-01 01:01:01'),(236,20231224070653,1,'2020-01-01 01:01:01'),(237,20240110134315,1,'2020-01-01 01:01:01'),(238,20240119091637,1,'2020-01-01 01:01:01'),(239,20240126020642,1,'2020-01-01 01:01:01'),(240,20240126020643,1,'2020-01-01 01:01:01'),(241,20240129162819,1,'2020-01-01 01:01:01'),(242,20240130115133,1,'2020-01-01 01:01:01'),(243,20240131083822,1,'2020-01-01 01:01:01'),(244,20240205095928,1,'2020-01-01 01:01:01'),(245,20240205121956,1,'2020-01-01 01:01:01'),(246,20240209110212,1,'2020-01-01 01:01:01'),(247,20240212111533,1,'2020-01-01 01:01:01'),(248,20240221112844,1,'2020-01-01 01:01:01'),(249,20240222073518,1,'2020-01-01 01:01:01'),(250,20240222135115,1,'2020-01-01 01:01:01'),(251,20240226082255,1,'2020-01-01 01:01:01'),(252,20240228082706,1,'2020-01-01 01:01:01'),(253,20240301173035,1,'2020-01-01 01:01:01'),(254,20240302111134,1,'2020-01-01 01:01:01'),(255,20240312103753,1,'2020-01-01 01:01:01'),(256,20240313143416,1,'2020-01-01 01:01:01'),(257,20240314085226,1,'2020-01-01 01:01:01'),(258,20240314151747,1,'2020-01-01 01:01:01'),(259,20240320145650,1,'2020-01-01 01:01:01'),(260,20240327115530,1,'2020-01-01 01:01:01'),(261,20240327115617,1,'2020-01-01 01:01:01'),(262,20240408085837,1,'2020-01-01 01:01:01'),(263,20240415104633,1,'2020-01-01 01:01:01'),(264,20240430111727,1,'2020-01-01 01:01:01'),(265,20240515200020,1,'2020-01-01 01:01:01'),(266,20240521143023,1,'2020-01-01 01:01:01'),(267,20240521143024,1,'2020-01-01 01:01:01'),(268,20240601174138,1,'2020-01-01 01:01:01'),(269,20240607133721,1,'2020-01-01 01:01:01'),(270,20240612150059,1,'2020-01-01 01:01:01'),(271,20240613162201,1,'2020-01-01 01:01:01'),(272,20240613172616,1,'2020-01-01 01:01:01'),(273,20240618142419,1,'2020-01-01 01:01:01'),(274,20240625093543,1,'2020-01-01 01:01:01'),(275,20240626195531,1,'2020-01-01 01:01:01'),(276,20240702123921,1,'2020-01-01 01:01:01'),(277,20240703154849,1,'2020-01-01 01:01:01'),(278,20240707134035,1,'2020-01-01 01:01:01'),(279,20240707134036,1,'2020-01-01 01:01:01'),(280,20240709124958,1,'2020-01-01 01:01:01'); +INSERT INTO `migration_status_tables` VALUES (1,0,1,'2020-01-01 01:01:01'),(2,20161118193812,1,'2020-01-01 01:01:01'),(3,20161118211713,1,'2020-01-01 01:01:01'),(4,20161118212436,1,'2020-01-01 01:01:01'),(5,20161118212515,1,'2020-01-01 01:01:01'),(6,20161118212528,1,'2020-01-01 01:01:01'),(7,20161118212538,1,'2020-01-01 01:01:01'),(8,20161118212549,1,'2020-01-01 01:01:01'),(9,20161118212557,1,'2020-01-01 01:01:01'),(10,20161118212604,1,'2020-01-01 01:01:01'),(11,20161118212613,1,'2020-01-01 01:01:01'),(12,20161118212621,1,'2020-01-01 01:01:01'),(13,20161118212630,1,'2020-01-01 01:01:01'),(14,20161118212641,1,'2020-01-01 01:01:01'),(15,20161118212649,1,'2020-01-01 01:01:01'),(16,20161118212656,1,'2020-01-01 01:01:01'),(17,20161118212758,1,'2020-01-01 01:01:01'),(18,20161128234849,1,'2020-01-01 01:01:01'),(19,20161230162221,1,'2020-01-01 01:01:01'),(20,20170104113816,1,'2020-01-01 01:01:01'),(21,20170105151732,1,'2020-01-01 01:01:01'),(22,20170108191242,1,'2020-01-01 01:01:01'),(23,20170109094020,1,'2020-01-01 01:01:01'),(24,20170109130438,1,'2020-01-01 01:01:01'),(25,20170110202752,1,'2020-01-01 01:01:01'),(26,20170111133013,1,'2020-01-01 01:01:01'),(27,20170117025759,1,'2020-01-01 01:01:01'),(28,20170118191001,1,'2020-01-01 01:01:01'),(29,20170119234632,1,'2020-01-01 01:01:01'),(30,20170124230432,1,'2020-01-01 01:01:01'),(31,20170127014618,1,'2020-01-01 01:01:01'),(32,20170131232841,1,'2020-01-01 01:01:01'),(33,20170223094154,1,'2020-01-01 01:01:01'),(34,20170306075207,1,'2020-01-01 01:01:01'),(35,20170309100733,1,'2020-01-01 01:01:01'),(36,20170331111922,1,'2020-01-01 01:01:01'),(37,20170502143928,1,'2020-01-01 01:01:01'),(38,20170504130602,1,'2020-01-01 01:01:01'),(39,20170509132100,1,'2020-01-01 01:01:01'),(40,20170519105647,1,'2020-01-01 01:01:01'),(41,20170519105648,1,'2020-01-01 01:01:01'),(42,20170831234300,1,'2020-01-01 01:01:01'),(43,20170831234301,1,'2020-01-01 01:01:01'),(44,20170831234303,1,'2020-01-01 01:01:01'),(45,20171116163618,1,'2020-01-01 01:01:01'),(46,20171219164727,1,'2020-01-01 01:01:01'),(47,20180620164811,1,'2020-01-01 01:01:01'),(48,20180620175054,1,'2020-01-01 01:01:01'),(49,20180620175055,1,'2020-01-01 01:01:01'),(50,20191010101639,1,'2020-01-01 01:01:01'),(51,20191010155147,1,'2020-01-01 01:01:01'),(52,20191220130734,1,'2020-01-01 01:01:01'),(53,20200311140000,1,'2020-01-01 01:01:01'),(54,20200405120000,1,'2020-01-01 01:01:01'),(55,20200407120000,1,'2020-01-01 01:01:01'),(56,20200420120000,1,'2020-01-01 01:01:01'),(57,20200504120000,1,'2020-01-01 01:01:01'),(58,20200512120000,1,'2020-01-01 01:01:01'),(59,20200707120000,1,'2020-01-01 01:01:01'),(60,20201011162341,1,'2020-01-01 01:01:01'),(61,20201021104586,1,'2020-01-01 01:01:01'),(62,20201102112520,1,'2020-01-01 01:01:01'),(63,20201208121729,1,'2020-01-01 01:01:01'),(64,20201215091637,1,'2020-01-01 01:01:01'),(65,20210119174155,1,'2020-01-01 01:01:01'),(66,20210326182902,1,'2020-01-01 01:01:01'),(67,20210421112652,1,'2020-01-01 01:01:01'),(68,20210506095025,1,'2020-01-01 01:01:01'),(69,20210513115729,1,'2020-01-01 01:01:01'),(70,20210526113559,1,'2020-01-01 01:01:01'),(71,20210601000001,1,'2020-01-01 01:01:01'),(72,20210601000002,1,'2020-01-01 01:01:01'),(73,20210601000003,1,'2020-01-01 01:01:01'),(74,20210601000004,1,'2020-01-01 01:01:01'),(75,20210601000005,1,'2020-01-01 01:01:01'),(76,20210601000006,1,'2020-01-01 01:01:01'),(77,20210601000007,1,'2020-01-01 01:01:01'),(78,20210601000008,1,'2020-01-01 01:01:01'),(79,20210606151329,1,'2020-01-01 01:01:01'),(80,20210616163757,1,'2020-01-01 01:01:01'),(81,20210617174723,1,'2020-01-01 01:01:01'),(82,20210622160235,1,'2020-01-01 01:01:01'),(83,20210623100031,1,'2020-01-01 01:01:01'),(84,20210623133615,1,'2020-01-01 01:01:01'),(85,20210708143152,1,'2020-01-01 01:01:01'),(86,20210709124443,1,'2020-01-01 01:01:01'),(87,20210712155608,1,'2020-01-01 01:01:01'),(88,20210714102108,1,'2020-01-01 01:01:01'),(89,20210719153709,1,'2020-01-01 01:01:01'),(90,20210721171531,1,'2020-01-01 01:01:01'),(91,20210723135713,1,'2020-01-01 01:01:01'),(92,20210802135933,1,'2020-01-01 01:01:01'),(93,20210806112844,1,'2020-01-01 01:01:01'),(94,20210810095603,1,'2020-01-01 01:01:01'),(95,20210811150223,1,'2020-01-01 01:01:01'),(96,20210818151827,1,'2020-01-01 01:01:01'),(97,20210818151828,1,'2020-01-01 01:01:01'),(98,20210818182258,1,'2020-01-01 01:01:01'),(99,20210819131107,1,'2020-01-01 01:01:01'),(100,20210819143446,1,'2020-01-01 01:01:01'),(101,20210903132338,1,'2020-01-01 01:01:01'),(102,20210915144307,1,'2020-01-01 01:01:01'),(103,20210920155130,1,'2020-01-01 01:01:01'),(104,20210927143115,1,'2020-01-01 01:01:01'),(105,20210927143116,1,'2020-01-01 01:01:01'),(106,20211013133706,1,'2020-01-01 01:01:01'),(107,20211013133707,1,'2020-01-01 01:01:01'),(108,20211102135149,1,'2020-01-01 01:01:01'),(109,20211109121546,1,'2020-01-01 01:01:01'),(110,20211110163320,1,'2020-01-01 01:01:01'),(111,20211116184029,1,'2020-01-01 01:01:01'),(112,20211116184030,1,'2020-01-01 01:01:01'),(113,20211202092042,1,'2020-01-01 01:01:01'),(114,20211202181033,1,'2020-01-01 01:01:01'),(115,20211207161856,1,'2020-01-01 01:01:01'),(116,20211216131203,1,'2020-01-01 01:01:01'),(117,20211221110132,1,'2020-01-01 01:01:01'),(118,20220107155700,1,'2020-01-01 01:01:01'),(119,20220125105650,1,'2020-01-01 01:01:01'),(120,20220201084510,1,'2020-01-01 01:01:01'),(121,20220208144830,1,'2020-01-01 01:01:01'),(122,20220208144831,1,'2020-01-01 01:01:01'),(123,20220215152203,1,'2020-01-01 01:01:01'),(124,20220223113157,1,'2020-01-01 01:01:01'),(125,20220307104655,1,'2020-01-01 01:01:01'),(126,20220309133956,1,'2020-01-01 01:01:01'),(127,20220316155700,1,'2020-01-01 01:01:01'),(128,20220323152301,1,'2020-01-01 01:01:01'),(129,20220330100659,1,'2020-01-01 01:01:01'),(130,20220404091216,1,'2020-01-01 01:01:01'),(131,20220419140750,1,'2020-01-01 01:01:01'),(132,20220428140039,1,'2020-01-01 01:01:01'),(133,20220503134048,1,'2020-01-01 01:01:01'),(134,20220524102918,1,'2020-01-01 01:01:01'),(135,20220526123327,1,'2020-01-01 01:01:01'),(136,20220526123328,1,'2020-01-01 01:01:01'),(137,20220526123329,1,'2020-01-01 01:01:01'),(138,20220608113128,1,'2020-01-01 01:01:01'),(139,20220627104817,1,'2020-01-01 01:01:01'),(140,20220704101843,1,'2020-01-01 01:01:01'),(141,20220708095046,1,'2020-01-01 01:01:01'),(142,20220713091130,1,'2020-01-01 01:01:01'),(143,20220802135510,1,'2020-01-01 01:01:01'),(144,20220818101352,1,'2020-01-01 01:01:01'),(145,20220822161445,1,'2020-01-01 01:01:01'),(146,20220831100036,1,'2020-01-01 01:01:01'),(147,20220831100151,1,'2020-01-01 01:01:01'),(148,20220908181826,1,'2020-01-01 01:01:01'),(149,20220914154915,1,'2020-01-01 01:01:01'),(150,20220915165115,1,'2020-01-01 01:01:01'),(151,20220915165116,1,'2020-01-01 01:01:01'),(152,20220928100158,1,'2020-01-01 01:01:01'),(153,20221014084130,1,'2020-01-01 01:01:01'),(154,20221027085019,1,'2020-01-01 01:01:01'),(155,20221101103952,1,'2020-01-01 01:01:01'),(156,20221104144401,1,'2020-01-01 01:01:01'),(157,20221109100749,1,'2020-01-01 01:01:01'),(158,20221115104546,1,'2020-01-01 01:01:01'),(159,20221130114928,1,'2020-01-01 01:01:01'),(160,20221205112142,1,'2020-01-01 01:01:01'),(161,20221216115820,1,'2020-01-01 01:01:01'),(162,20221220195934,1,'2020-01-01 01:01:01'),(163,20221220195935,1,'2020-01-01 01:01:01'),(164,20221223174807,1,'2020-01-01 01:01:01'),(165,20221227163855,1,'2020-01-01 01:01:01'),(166,20221227163856,1,'2020-01-01 01:01:01'),(167,20230202224725,1,'2020-01-01 01:01:01'),(168,20230206163608,1,'2020-01-01 01:01:01'),(169,20230214131519,1,'2020-01-01 01:01:01'),(170,20230303135738,1,'2020-01-01 01:01:01'),(171,20230313135301,1,'2020-01-01 01:01:01'),(172,20230313141819,1,'2020-01-01 01:01:01'),(173,20230315104937,1,'2020-01-01 01:01:01'),(174,20230317173844,1,'2020-01-01 01:01:01'),(175,20230320133602,1,'2020-01-01 01:01:01'),(176,20230330100011,1,'2020-01-01 01:01:01'),(177,20230330134823,1,'2020-01-01 01:01:01'),(178,20230405232025,1,'2020-01-01 01:01:01'),(179,20230408084104,1,'2020-01-01 01:01:01'),(180,20230411102858,1,'2020-01-01 01:01:01'),(181,20230421155932,1,'2020-01-01 01:01:01'),(182,20230425082126,1,'2020-01-01 01:01:01'),(183,20230425105727,1,'2020-01-01 01:01:01'),(184,20230501154913,1,'2020-01-01 01:01:01'),(185,20230503101418,1,'2020-01-01 01:01:01'),(186,20230515144206,1,'2020-01-01 01:01:01'),(187,20230517140952,1,'2020-01-01 01:01:01'),(188,20230517152807,1,'2020-01-01 01:01:01'),(189,20230518114155,1,'2020-01-01 01:01:01'),(190,20230520153236,1,'2020-01-01 01:01:01'),(191,20230525151159,1,'2020-01-01 01:01:01'),(192,20230530122103,1,'2020-01-01 01:01:01'),(193,20230602111827,1,'2020-01-01 01:01:01'),(194,20230608103123,1,'2020-01-01 01:01:01'),(195,20230629140529,1,'2020-01-01 01:01:01'),(196,20230629140530,1,'2020-01-01 01:01:01'),(197,20230711144622,1,'2020-01-01 01:01:01'),(198,20230721135421,1,'2020-01-01 01:01:01'),(199,20230721161508,1,'2020-01-01 01:01:01'),(200,20230726115701,1,'2020-01-01 01:01:01'),(201,20230807100822,1,'2020-01-01 01:01:01'),(202,20230814150442,1,'2020-01-01 01:01:01'),(203,20230823122728,1,'2020-01-01 01:01:01'),(204,20230906152143,1,'2020-01-01 01:01:01'),(205,20230911163618,1,'2020-01-01 01:01:01'),(206,20230912101759,1,'2020-01-01 01:01:01'),(207,20230915101341,1,'2020-01-01 01:01:01'),(208,20230918132351,1,'2020-01-01 01:01:01'),(209,20231004144339,1,'2020-01-01 01:01:01'),(210,20231009094541,1,'2020-01-01 01:01:01'),(211,20231009094542,1,'2020-01-01 01:01:01'),(212,20231009094543,1,'2020-01-01 01:01:01'),(213,20231009094544,1,'2020-01-01 01:01:01'),(214,20231016091915,1,'2020-01-01 01:01:01'),(215,20231024174135,1,'2020-01-01 01:01:01'),(216,20231025120016,1,'2020-01-01 01:01:01'),(217,20231025160156,1,'2020-01-01 01:01:01'),(218,20231031165350,1,'2020-01-01 01:01:01'),(219,20231106144110,1,'2020-01-01 01:01:01'),(220,20231107130934,1,'2020-01-01 01:01:01'),(221,20231109115838,1,'2020-01-01 01:01:01'),(222,20231121054530,1,'2020-01-01 01:01:01'),(223,20231122101320,1,'2020-01-01 01:01:01'),(224,20231130132828,1,'2020-01-01 01:01:01'),(225,20231130132931,1,'2020-01-01 01:01:01'),(226,20231204155427,1,'2020-01-01 01:01:01'),(227,20231206142340,1,'2020-01-01 01:01:01'),(228,20231207102320,1,'2020-01-01 01:01:01'),(229,20231207102321,1,'2020-01-01 01:01:01'),(230,20231207133731,1,'2020-01-01 01:01:01'),(231,20231212094238,1,'2020-01-01 01:01:01'),(232,20231212095734,1,'2020-01-01 01:01:01'),(233,20231212161121,1,'2020-01-01 01:01:01'),(234,20231215122713,1,'2020-01-01 01:01:01'),(235,20231219143041,1,'2020-01-01 01:01:01'),(236,20231224070653,1,'2020-01-01 01:01:01'),(237,20240110134315,1,'2020-01-01 01:01:01'),(238,20240119091637,1,'2020-01-01 01:01:01'),(239,20240126020642,1,'2020-01-01 01:01:01'),(240,20240126020643,1,'2020-01-01 01:01:01'),(241,20240129162819,1,'2020-01-01 01:01:01'),(242,20240130115133,1,'2020-01-01 01:01:01'),(243,20240131083822,1,'2020-01-01 01:01:01'),(244,20240205095928,1,'2020-01-01 01:01:01'),(245,20240205121956,1,'2020-01-01 01:01:01'),(246,20240209110212,1,'2020-01-01 01:01:01'),(247,20240212111533,1,'2020-01-01 01:01:01'),(248,20240221112844,1,'2020-01-01 01:01:01'),(249,20240222073518,1,'2020-01-01 01:01:01'),(250,20240222135115,1,'2020-01-01 01:01:01'),(251,20240226082255,1,'2020-01-01 01:01:01'),(252,20240228082706,1,'2020-01-01 01:01:01'),(253,20240301173035,1,'2020-01-01 01:01:01'),(254,20240302111134,1,'2020-01-01 01:01:01'),(255,20240312103753,1,'2020-01-01 01:01:01'),(256,20240313143416,1,'2020-01-01 01:01:01'),(257,20240314085226,1,'2020-01-01 01:01:01'),(258,20240314151747,1,'2020-01-01 01:01:01'),(259,20240320145650,1,'2020-01-01 01:01:01'),(260,20240327115530,1,'2020-01-01 01:01:01'),(261,20240327115617,1,'2020-01-01 01:01:01'),(262,20240408085837,1,'2020-01-01 01:01:01'),(263,20240415104633,1,'2020-01-01 01:01:01'),(264,20240430111727,1,'2020-01-01 01:01:01'),(265,20240515200020,1,'2020-01-01 01:01:01'),(266,20240521143023,1,'2020-01-01 01:01:01'),(267,20240521143024,1,'2020-01-01 01:01:01'),(268,20240601174138,1,'2020-01-01 01:01:01'),(269,20240607133721,1,'2020-01-01 01:01:01'),(270,20240612150059,1,'2020-01-01 01:01:01'),(271,20240613162201,1,'2020-01-01 01:01:01'),(272,20240613172616,1,'2020-01-01 01:01:01'),(273,20240618142419,1,'2020-01-01 01:01:01'),(274,20240625093543,1,'2020-01-01 01:01:01'),(275,20240626195531,1,'2020-01-01 01:01:01'),(276,20240702123921,1,'2020-01-01 01:01:01'),(277,20240703154849,1,'2020-01-01 01:01:01'),(278,20240707134035,1,'2020-01-01 01:01:01'),(279,20240707134036,1,'2020-01-01 01:01:01'),(280,20240709124958,1,'2020-01-01 01:01:01'),(281,20240709132642,1,'2020-01-01 01:01:01'); /*!40101 SET @saved_cs_client = @@character_set_client */; /*!40101 SET character_set_client = utf8 */; CREATE TABLE `mobile_device_management_solutions` ( diff --git a/server/fleet/calendar.go b/server/fleet/calendar.go index 6f34b20051..e4a76354d8 100644 --- a/server/fleet/calendar.go +++ b/server/fleet/calendar.go @@ -42,6 +42,25 @@ type UserCalendar interface { Get(event *CalendarEvent, key string) (interface{}, error) } +// Lock interface for managing distributed locks. +type Lock interface { + // AcquireLock attempts to acquire a lock with the given key. value is the value to set for the key, which is used to release the lock. + // expireMs is the time in milliseconds after which the lock is automatically released. expireMs=0 means a default expiration time is used. + // Returns true if the lock was acquired, false otherwise. + AcquireLock(ctx context.Context, key string, value string, expireMs uint64) (ok bool, err error) + // ReleaseLock attempts to release a lock with the given key and value. If key does not exist or value does not match, the lock is not released. + // Returns true if the lock was released, false otherwise. + ReleaseLock(ctx context.Context, key string, value string) (ok bool, err error) + // Get retrieves the value of the given key. If the key does not exist, nil is returned. + Get(ctx context.Context, key string) (*string, error) + // AddToSet adds the value to the set identified by the given key. + AddToSet(ctx context.Context, key string, value string) error + // RemoveFromSet removes the value from the set identified by the given key. + RemoveFromSet(ctx context.Context, key string, value string) error + // GetSet retrieves a slice of string values from the set identified by the given key. + GetSet(ctx context.Context, key string) ([]string, error) +} + type CalendarWebhookPayload struct { Timestamp time.Time `json:"timestamp"` HostID uint `json:"host_id"` diff --git a/server/service/calendar/calendar.go b/server/service/calendar/calendar.go index 59c1784dd8..a125ddb7a3 100644 --- a/server/service/calendar/calendar.go +++ b/server/service/calendar/calendar.go @@ -16,13 +16,19 @@ import ( "github.com/go-kit/log/level" ) -type CalendarConfig struct { +const ( + LockKeyPrefix = "calendar:lock:" + ReservedLockKeyPrefix = "calendar:reserved:" + QueueKey = "calendar:queue" +) + +type Config struct { config.CalendarConfig fleet.GoogleCalendarIntegration ServerURL string } -func CreateUserCalendarFromConfig(ctx context.Context, config *CalendarConfig, logger kitlog.Logger) fleet.UserCalendar { +func CreateUserCalendarFromConfig(ctx context.Context, config *Config, logger kitlog.Logger) fleet.UserCalendar { googleCalendarConfig := calendar.GoogleCalendarConfig{ Context: ctx, IntegrationConfig: &config.GoogleCalendarIntegration, diff --git a/server/service/integration_enterprise_test.go b/server/service/integration_enterprise_test.go index c393bf45b1..b48cb7d714 100644 --- a/server/service/integration_enterprise_test.go +++ b/server/service/integration_enterprise_test.go @@ -35,6 +35,8 @@ import ( "github.com/fleetdm/fleet/v4/server/mdm" "github.com/fleetdm/fleet/v4/server/ptr" "github.com/fleetdm/fleet/v4/server/pubsub" + commonCalendar "github.com/fleetdm/fleet/v4/server/service/calendar" + "github.com/fleetdm/fleet/v4/server/service/redis_lock" "github.com/fleetdm/fleet/v4/server/service/schedule" "github.com/fleetdm/fleet/v4/server/test" "github.com/go-kit/log" @@ -87,7 +89,8 @@ func (s *integrationEnterpriseTestSuite) SetupSuite() { cronLog = kitlog.NewNopLogger() } calendarSchedule, err = cron.NewCalendarSchedule( - ctx, s.T().Name(), s.ds, config.CalendarConfig{Periodicity: 24 * time.Hour}, cronLog, + ctx, s.T().Name(), s.ds, redis_lock.NewLock(s.redisPool), config.CalendarConfig{Periodicity: 24 * time.Hour}, + cronLog, ) return calendarSchedule, err } @@ -10959,9 +10962,9 @@ func (s *integrationEnterpriseTestSuite) TestCalendarCallback() { }, ) require.NoError(t, err) - team1Policy2, err := s.ds.NewTeamPolicy( + team1Policy2Calendar, err := s.ds.NewTeamPolicy( ctx, team1.ID, nil, fleet.PolicyPayload{ - Name: "team1Policy2", + Name: "team1Policy2Calendar", Query: "SELECT 2;", CalendarEventsEnabled: true, }, @@ -11012,7 +11015,7 @@ func (s *integrationEnterpriseTestSuite) TestCalendarCallback() { host1Team1, map[uint]*bool{ team1Policy1Calendar.ID: ptr.Bool(false), - team1Policy2.ID: ptr.Bool(true), + team1Policy2Calendar.ID: ptr.Bool(true), globalPolicy.ID: nil, }, ), http.StatusOK, &distributedResp) @@ -11022,7 +11025,7 @@ func (s *integrationEnterpriseTestSuite) TestCalendarCallback() { host2Team1, map[uint]*bool{ team1Policy1Calendar.ID: ptr.Bool(true), - team1Policy2.ID: ptr.Bool(false), + team1Policy2Calendar.ID: ptr.Bool(false), globalPolicy.ID: nil, }, ), http.StatusOK, &distributedResp) @@ -11106,28 +11109,95 @@ func (s *integrationEnterpriseTestSuite) TestCalendarCallback() { // Delete the event on the calendar calendar.ClearMockEvents() - // This callback should recreate the event + // Grab the distributed lock for this event + distributedLock := redis_lock.NewLock(s.redisPool) + lockValue := uuid.New().String() + result, err := distributedLock.AcquireLock(ctx, commonCalendar.LockKeyPrefix+event.UUID, lockValue, 0) + require.NoError(t, err) + assert.NotEmpty(t, result) + + // This callback should put the event processing in a queue for async processing. It does not start async + // processing because it assumes another server is handling this webhook, and that server will start + // async processing. _ = s.DoRawWithHeaders("POST", "/api/v1/fleet/calendar/webhook/"+event.UUID, []byte(""), http.StatusOK, map[string]string{ "X-Goog-Channel-Id": details.ChannelID, "X-Goog-Resource-State": "exists", }) - team1CalendarEvents, err = s.ds.ListCalendarEvents(ctx, &team1.ID) + uuids, err := distributedLock.GetSet(ctx, commonCalendar.QueueKey) require.NoError(t, err) - require.Len(t, team1CalendarEvents, 1) + assert.ElementsMatch(t, []string{event.UUID}, uuids) + // The calendar should still be empty since event hasn't processed yet + assert.Zero(t, len(calendar.ListGoogleMockEvents())) + // We clear the queue + assert.NoError(t, distributedLock.RemoveFromSet(ctx, commonCalendar.QueueKey, event.UUID)) + + // We release the normal lock, but grab the reserve lock instead + ok, err := distributedLock.ReleaseLock(ctx, commonCalendar.LockKeyPrefix+event.UUID, lockValue) + require.NoError(t, err) + assert.True(t, ok) + result, err = distributedLock.AcquireLock(ctx, commonCalendar.ReservedLockKeyPrefix+event.UUID, lockValue, 0) + require.NoError(t, err) + assert.NotEmpty(t, result) + + // This callback should put the event processing in a queue for async processing, AND start the async processing + _ = s.DoRawWithHeaders("POST", "/api/v1/fleet/calendar/webhook/"+event.UUID, []byte(""), http.StatusOK, map[string]string{ + "X-Goog-Channel-Id": details.ChannelID, + "X-Goog-Resource-State": "exists", + }) + + uuids, err = distributedLock.GetSet(ctx, commonCalendar.QueueKey) + require.NoError(t, err) + assert.ElementsMatch(t, []string{event.UUID}, uuids) + // The calendar should still be empty since event hasn't processed yet + assert.Zero(t, len(calendar.ListGoogleMockEvents())) + + // We grab the normal lock again. + lockValue2 := uuid.New().String() + result, err = distributedLock.AcquireLock(ctx, commonCalendar.LockKeyPrefix+event.UUID, lockValue2, 0) + require.NoError(t, err) + assert.NotEmpty(t, result) + // We release the reserve lock. + ok, err = distributedLock.ReleaseLock(ctx, commonCalendar.ReservedLockKeyPrefix+event.UUID, lockValue) + require.NoError(t, err) + assert.True(t, ok) + // We release the normal lock. + ok, err = distributedLock.ReleaseLock(ctx, commonCalendar.LockKeyPrefix+event.UUID, lockValue2) + require.NoError(t, err) + assert.True(t, ok) + + done := make(chan struct{}) + go func() { + for { + time.Sleep(100 * time.Millisecond) + team1CalendarEvents, err = s.ds.ListCalendarEvents(ctx, &team1.ID) + require.NoError(t, err) + require.Len(t, team1CalendarEvents, 1) + if event.UUID != team1CalendarEvents[0].UUID { + done <- struct{}{} + return + } + } + }() + select { + case <-done: // All good + case <-time.After(5 * time.Second): + t.Fatal("timeout waiting for calendar event processing") + } + eventRecreated := team1CalendarEvents[0] assert.NotZero(t, eventRecreated.ID) assert.Equal(t, user1Email, eventRecreated.Email) assert.NotZero(t, eventRecreated.StartTime) assert.NotZero(t, eventRecreated.EndTime) assert.NotEmpty(t, eventRecreated.UUID) - assert.NotEqual(t, event.UUID, eventRecreated.UUID) assert.NotEqual(t, event.StartTime, eventRecreated.StartTime) assert.NotEqual(t, event.EndTime, eventRecreated.EndTime) assert.Equal(t, 1, calendar.MockChannelsCount()) + assert.Equal(t, 1, len(calendar.ListGoogleMockEvents())) - // The previous event UUID should not work anymore - _ = s.DoRawWithHeaders("POST", "/api/v1/fleet/calendar/webhook/"+event.UUID, []byte(""), http.StatusNotFound, map[string]string{ + // The previous event UUID should not work anymore, but API returns OK because this is a common occurrence. + _ = s.DoRawWithHeaders("POST", "/api/v1/fleet/calendar/webhook/"+event.UUID, []byte(""), http.StatusOK, map[string]string{ "X-Goog-Channel-Id": details.ChannelID, "X-Goog-Resource-State": "exists", }) @@ -11171,6 +11241,75 @@ func (s *integrationEnterpriseTestSuite) TestCalendarCallback() { assert.Equal(t, eventRecreated.EndTime, eventUpdated.EndTime) assert.Equal(t, 1, calendar.MockChannelsCount()) + // Update the time of the event again + events = calendar.ListGoogleMockEvents() + require.Len(t, events, 1) + for _, e := range events { + st, err := time.Parse(time.RFC3339, e.Start.DateTime) + require.NoError(t, err) + newStartTime := st.Add(5 * time.Minute).Format(time.RFC3339) + e.Start.DateTime = newStartTime + } + + // Grab the lock + event = eventUpdated + lockValue = uuid.New().String() + result, err = distributedLock.AcquireLock(ctx, commonCalendar.LockKeyPrefix+event.UUID, lockValue, 0) + require.NoError(t, err) + assert.NotEmpty(t, result) + + mysql.ExecAdhocSQL(t, s.ds, func(db sqlx.ExtContext) error { + // Update updated_at so the event gets updated (the event is updated regularly) + _, err := db.ExecContext(ctx, + `UPDATE calendar_events SET updated_at = DATE_SUB(CURRENT_TIMESTAMP, INTERVAL 25 HOUR) WHERE id = ?`, event.ID) + return err + }) + + // Trigger the calendar cron async. It should wait for the lock and set reserve lock. + go triggerAndWait(ctx, t, s.ds, s.calendarSchedule, 10*time.Second) + done = make(chan struct{}) + go func() { + for { + time.Sleep(100 * time.Millisecond) + reserveLock, err := distributedLock.Get(ctx, commonCalendar.ReservedLockKeyPrefix+event.UUID) + require.NoError(t, err) + if reserveLock != nil { + done <- struct{}{} + return + } + } + }() + select { + case <-done: // All good + case <-time.After(5 * time.Second): + t.Fatal("timeout waiting for cron to set reserve lock") + } + + // Release the normal lock + ok, err = distributedLock.ReleaseLock(ctx, commonCalendar.LockKeyPrefix+event.UUID, lockValue) + require.NoError(t, err) + assert.True(t, ok) + + // Wait for the event to update + done = make(chan struct{}) + go func() { + for { + time.Sleep(100 * time.Millisecond) + team1CalendarEvents, err = s.ds.ListCalendarEvents(ctx, &team1.ID) + require.NoError(t, err) + if len(team1CalendarEvents) == 1 && team1CalendarEvents[0].UUID == event.UUID && + team1CalendarEvents[0].StartTime.After(event.StartTime) { + done <- struct{}{} + return + } + } + }() + select { + case <-done: // All good + case <-time.After(5 * time.Second): + t.Fatal("timeout waiting for event to update during cron") + } + // Delete the event on the calendar calendar.ClearMockEvents() @@ -11179,7 +11318,7 @@ func (s *integrationEnterpriseTestSuite) TestCalendarCallback() { host1Team1, map[uint]*bool{ team1Policy1Calendar.ID: ptr.Bool(true), - team1Policy2.ID: ptr.Bool(true), + team1Policy2Calendar.ID: ptr.Bool(true), globalPolicy.ID: nil, }, ), http.StatusOK, &distributedResp) @@ -11192,10 +11331,11 @@ func (s *integrationEnterpriseTestSuite) TestCalendarCallback() { }) assert.Equal(t, 0, calendar.MockChannelsCount()) + previousEvent := team1CalendarEvents[0] team1CalendarEvents, err = s.ds.ListCalendarEvents(ctx, &team1.ID) require.NoError(t, err) require.Len(t, team1CalendarEvents, 1) - assert.Equal(t, eventUpdated, team1CalendarEvents[0]) + assert.Equal(t, previousEvent, team1CalendarEvents[0]) // Trigger calendar should cleanup the events triggerAndWait(ctx, t, s.ds, s.calendarSchedule, 5*time.Second) diff --git a/server/service/redis_lock/redis_lock.go b/server/service/redis_lock/redis_lock.go new file mode 100644 index 0000000000..a5443f1df8 --- /dev/null +++ b/server/service/redis_lock/redis_lock.go @@ -0,0 +1,118 @@ +package redis_lock + +import ( + "context" + "errors" + + "github.com/fleetdm/fleet/v4/server/contexts/ctxerr" + "github.com/fleetdm/fleet/v4/server/datastore/redis" + "github.com/fleetdm/fleet/v4/server/fleet" + redigo "github.com/gomodule/redigo/redis" +) + +// This package implements a distributed lock using Redis. The lock can be used +// to prevent multiple Fleet servers from accessing a shared resource. + +const ( + defaultExpireMs = 60 * 1000 +) + +type redisLock struct { + pool fleet.RedisPool + testPrefix string // for tests, the key prefix to use to avoid conflicts +} + +func NewLock(pool fleet.RedisPool) fleet.Lock { + lock := &redisLock{ + pool: pool, + } + return fleet.Lock(lock) +} + +func (r *redisLock) AcquireLock(ctx context.Context, key string, value string, expireMs uint64) (ok bool, err error) { + conn := redis.ConfigureDoer(r.pool, r.pool.Get()) + defer conn.Close() + + if expireMs == 0 { + expireMs = defaultExpireMs + } + + // Reference: https://redis.io/docs/latest/commands/set/ + // NX -- Only set the key if it does not already exist. + result, err := redigo.String(conn.Do("SET", r.testPrefix+key, value, "NX", "PX", expireMs)) + if err != nil && !errors.Is(err, redigo.ErrNil) { + return false, ctxerr.Wrap(ctx, err, "redis acquire lock") + } + return result != "", nil +} + +func (r *redisLock) ReleaseLock(ctx context.Context, key string, value string) (ok bool, err error) { + conn := redis.ConfigureDoer(r.pool, r.pool.Get()) + defer conn.Close() + + const unlockScript = ` + if redis.call("get", KEYS[1]) == ARGV[1] then + return redis.call("del", KEYS[1]) + else + return 0 + end + ` + + // Reference: https://redis.io/docs/latest/commands/set/ + // Only release the lock if the value matches. + res, err := redigo.Int64(conn.Do("EVAL", unlockScript, 1, r.testPrefix+key, value)) + if err != nil && !errors.Is(err, redigo.ErrNil) { + return false, ctxerr.Wrap(ctx, err, "redis release lock") + } + return res > 0, nil +} + +func (r *redisLock) AddToSet(ctx context.Context, key string, value string) error { + conn := redis.ConfigureDoer(r.pool, r.pool.Get()) + defer conn.Close() + + // Reference: https://redis.io/docs/latest/commands/sadd/ + _, err := conn.Do("SADD", r.testPrefix+key, value) + if err != nil { + return ctxerr.Wrap(ctx, err, "redis add to set") + } + return nil +} + +func (r *redisLock) RemoveFromSet(ctx context.Context, key string, value string) error { + conn := redis.ConfigureDoer(r.pool, r.pool.Get()) + defer conn.Close() + + // Reference: https://redis.io/docs/latest/commands/srem/ + _, err := conn.Do("SREM", r.testPrefix+key, value) + if err != nil { + return ctxerr.Wrap(ctx, err, "redis add to set") + } + return nil +} + +func (r *redisLock) GetSet(ctx context.Context, key string) ([]string, error) { + conn := redis.ConfigureDoer(r.pool, r.pool.Get()) + defer conn.Close() + + // Reference: https://redis.io/docs/latest/commands/smembers/ + members, err := redigo.Strings(conn.Do("SMEMBERS", r.testPrefix+key)) + if err != nil && !errors.Is(err, redigo.ErrNil) { + return nil, ctxerr.Wrap(ctx, err, "redis get set members") + } + return members, nil +} + +func (r *redisLock) Get(ctx context.Context, key string) (*string, error) { + conn := redis.ConfigureDoer(r.pool, r.pool.Get()) + defer conn.Close() + + res, err := redigo.String(conn.Do("GET", r.testPrefix+key)) + if errors.Is(err, redigo.ErrNil) { + return nil, nil + } + if err != nil { + return nil, ctxerr.Wrap(ctx, err, "redis get") + } + return &res, nil +} diff --git a/server/service/redis_lock/redis_lock_test.go b/server/service/redis_lock/redis_lock_test.go new file mode 100644 index 0000000000..74df2a8334 --- /dev/null +++ b/server/service/redis_lock/redis_lock_test.go @@ -0,0 +1,145 @@ +package redis_lock + +import ( + "context" + "github.com/fleetdm/fleet/v4/server/datastore/redis/redistest" + "github.com/fleetdm/fleet/v4/server/fleet" + "github.com/fleetdm/fleet/v4/server/test" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + "testing" + "time" +) + +func TestRedisLock(t *testing.T) { + for _, f := range []func(*testing.T, fleet.Lock){ + testRedisAcquireLock, + testRedisSet, + } { + t.Run(test.FunctionName(f), func(t *testing.T) { + t.Run("standalone", func(t *testing.T) { + lock := setupRedis(t, false, false) + f(t, lock) + }) + t.Run("cluster", func(t *testing.T) { + lock := setupRedis(t, true, true) + f(t, lock) + }) + }) + } +} + +func setupRedis(t testing.TB, cluster, redir bool) fleet.Lock { + pool := redistest.SetupRedis(t, t.Name(), cluster, redir, true) + return NewLockTest(t, pool) +} + +type TestName interface { + Name() string +} + +// NewFailingTest creates a redis policy set for failing policies to be used +// only in tests. +func NewLockTest(t TestName, pool fleet.RedisPool) fleet.Lock { + lock := &redisLock{ + pool: pool, + testPrefix: t.Name() + ":", + } + return fleet.Lock(lock) +} + +func testRedisAcquireLock(t *testing.T, lock fleet.Lock) { + ctx := context.Background() + result, err := lock.AcquireLock(ctx, "test", "1", 0) + require.NoError(t, err) + assert.True(t, result) + + // Try to acquire the same lock + result, err = lock.AcquireLock(ctx, "test", "1", 0) + assert.NoError(t, err) + assert.False(t, result) + + // Try to release the lock with a wrong value + ok, err := lock.ReleaseLock(ctx, "test", "2") + require.NoError(t, err) + assert.False(t, ok) + + // Try to release the lock with the wrong key + ok, err = lock.ReleaseLock(ctx, "bad", "1") + require.NoError(t, err) + assert.False(t, ok) + + // Try to release the lock with the correct key/value + ok, err = lock.ReleaseLock(ctx, "test", "1") + require.NoError(t, err) + assert.True(t, ok) + + // Acquire the lock again + result, err = lock.AcquireLock(ctx, "test", "1", 0) + require.NoError(t, err) + assert.True(t, result) + + // Get lock + getResult, err := lock.Get(ctx, "test") + assert.NoError(t, err) + require.NotNil(t, getResult) + assert.Equal(t, "1", *getResult) + + // Try to set lock with expiration + var expire uint64 = 10 + result, err = lock.AcquireLock(ctx, "testE", "1", expire) + require.NoError(t, err) + assert.True(t, result) + + // Try to acquire the same lock after waiting + duration := time.Duration(expire+1) * time.Millisecond + time.Sleep(duration) + result, err = lock.AcquireLock(ctx, "testE", "1", 0) + require.NoError(t, err) + assert.True(t, result) + + // Get non-existent key + getResult, err = lock.Get(ctx, "testNonExistent") + assert.NoError(t, err) + assert.Nil(t, getResult) + +} + +func testRedisSet(t *testing.T, lock fleet.Lock) { + ctx := context.Background() + + // Get a non-existent set + result, err := lock.GetSet(ctx, "missingSet") + assert.NoError(t, err) + assert.Empty(t, result) + + // Add to a set + values := []string{"foo", "bar"} + err = lock.AddToSet(ctx, "testSet", values[0]) + assert.NoError(t, err) + err = lock.AddToSet(ctx, "testSet", values[1]) + assert.NoError(t, err) + + // Get the set + result, err = lock.GetSet(ctx, "testSet") + assert.NoError(t, err) + assert.ElementsMatch(t, values, result) + + // Remove from set + err = lock.RemoveFromSet(ctx, "testSet", values[0]) + assert.NoError(t, err) + + // Get the set + result, err = lock.GetSet(ctx, "testSet") + assert.NoError(t, err) + assert.Equal(t, []string{values[1]}, result) + + // Remove from set + err = lock.RemoveFromSet(ctx, "testSet", values[1]) + assert.NoError(t, err) + + // Get the set + result, err = lock.GetSet(ctx, "testSet") + assert.NoError(t, err) + assert.Empty(t, result) +} diff --git a/server/service/testing_utils.go b/server/service/testing_utils.go index 73c0d495c0..df5fa5c219 100644 --- a/server/service/testing_utils.go +++ b/server/service/testing_utils.go @@ -34,6 +34,7 @@ import ( "github.com/fleetdm/fleet/v4/server/ptr" "github.com/fleetdm/fleet/v4/server/service/async" "github.com/fleetdm/fleet/v4/server/service/mock" + "github.com/fleetdm/fleet/v4/server/service/redis_lock" "github.com/fleetdm/fleet/v4/server/sso" "github.com/fleetdm/fleet/v4/server/test" kitlog "github.com/go-kit/log" @@ -69,6 +70,7 @@ func newTestServiceWithConfig(t *testing.T, ds fleet.Datastore, fleetConfig conf ssoStore sso.SessionStore profMatcher fleet.ProfileMatcher softwareInstallStore fleet.SoftwareInstallerStore + distributedLock fleet.Lock ) if len(opts) > 0 { if opts[0].Clock != nil { @@ -95,6 +97,7 @@ func newTestServiceWithConfig(t *testing.T, ds fleet.Datastore, fleetConfig conf if opts[0].Pool != nil { ssoStore = sso.NewSessionStore(opts[0].Pool) profMatcher = apple_mdm.NewProfileMatcher(opts[0].Pool) + distributedLock = redis_lock.NewLock(opts[0].Pool) } if opts[0].ProfileMatcher != nil { profMatcher = opts[0].ProfileMatcher @@ -194,6 +197,7 @@ func newTestServiceWithConfig(t *testing.T, ds fleet.Datastore, fleetConfig conf ssoStore, profMatcher, softwareInstallStore, + distributedLock, ) if err != nil { panic(err) diff --git a/tools/calendar/delete-events/delete-events.go b/tools/calendar/delete-events/delete-events.go index 676da9d88f..f3f1355548 100644 --- a/tools/calendar/delete-events/delete-events.go +++ b/tools/calendar/delete-events/delete-events.go @@ -70,16 +70,20 @@ func main() { log.Fatalf("Unable to create Calendar service: %v", err) } numberDeleted := 0 + var maxResults int64 = 1000 + pageToken := "" + now := time.Now() for { list, err := withRetry( func() (any, error) { return service.Events.List("primary"). EventTypes("default"). - MaxResults(1000). + MaxResults(maxResults). OrderBy("startTime"). SingleEvents(true). ShowDeleted(false). Q(eventTitle). + PageToken(pageToken). Do() }, ) @@ -89,7 +93,17 @@ func main() { if len(list.(*calendar.Events).Items) == 0 { break } + foundNewEvents := false for _, item := range list.(*calendar.Events).Items { + created, err := time.Parse(time.RFC3339, item.Created) + if err != nil { + log.Fatalf("Unable to parse event created time: %v", err) + } + if created.After(now) { + // Found events created after we started deleting events, so we should stop + foundNewEvents = true + continue // Skip this event but finish the loop to make sure we don't miss something + } if item.Summary == eventTitle { _, err := withRetry( func() (any, error) { @@ -105,6 +119,10 @@ func main() { } } } + pageToken = list.(*calendar.Events).NextPageToken + if pageToken == "" || foundNewEvents { + break + } } log.Printf("DONE. Deleted %d events total for %s", numberDeleted, userEmail) }(userEmail) diff --git a/tools/calendar/move-events/move-events.go b/tools/calendar/move-events/move-events.go index 3e22032c99..616418b196 100644 --- a/tools/calendar/move-events/move-events.go +++ b/tools/calendar/move-events/move-events.go @@ -81,16 +81,20 @@ func main() { } numberMoved := 0 + var maxResults int64 = 1000 + pageToken := "" + now := time.Now() for { list, err := withRetry( func() (any, error) { return service.Events.List("primary").EventTypes("default"). - MaxResults(1000). + MaxResults(maxResults). OrderBy("startTime"). SingleEvents(true). ShowDeleted(false). TimeMin(dateTimeEndStr). Q(eventTitle). + PageToken(pageToken). Do() }, ) @@ -101,7 +105,17 @@ func main() { if len(list.(*calendar.Events).Items) == 0 { break } + foundNewEvents := false for _, item := range list.(*calendar.Events).Items { + created, err := time.Parse(time.RFC3339, item.Created) + if err != nil { + log.Fatalf("Unable to parse event created time: %v", err) + } + if created.After(now) { + // Found events created after we started moving events, so we should stop + foundNewEvents = true + continue // Skip this event but finish the loop to make sure we don't miss something + } if item.Summary == eventTitle { item.Start.DateTime = dateTime.Format(time.RFC3339) item.End.DateTime = dateTime.Add(30 * time.Minute).Format(time.RFC3339) @@ -120,6 +134,10 @@ func main() { } } + pageToken = list.(*calendar.Events).NextPageToken + if pageToken == "" || foundNewEvents { + break + } } log.Printf("DONE. Moved total %d events for %s", numberMoved, userEmail) }(userEmail)