diff --git a/Makefile b/Makefile index 682548127b7d..4aaaac92cf9d 100644 --- a/Makefile +++ b/Makefile @@ -30,6 +30,7 @@ OCM_CLIENT_ID := ${OCM_CLIENT_ID} OCM_CLIENT_SECRET := ${OCM_CLIENT_SECRET} ENABLE_AUTH := $(or ${ENABLE_AUTH},False) DELETE_PVC := $(or ${DELETE_PVC},False) +REPLICAS_COUNT = $(shell if ! [ "${TARGET}" = "minikube" ];then echo 3; else echo $(or ${SERVICE_REPLICAS_COUNT},3);fi) ifdef INSTALLATION_TIMEOUT INSTALLATION_TIMEOUT_FLAG = --installation-timeout $(INSTALLATION_TIMEOUT) @@ -172,7 +173,7 @@ deploy-service-requirements: deploy-namespace deploy-inventory-service-file deploy-service: deploy-namespace deploy-service-requirements deploy-role python3 ./tools/deploy_assisted_installer.py $(DEPLOY_TAG_OPTION) --namespace "$(NAMESPACE)" \ - --profile "$(PROFILE)" $(TEST_FLAGS) --target "$(TARGET)" + --profile "$(PROFILE)" $(TEST_FLAGS) --target "$(TARGET)" --replicas-count $(REPLICAS_COUNT) python3 ./tools/wait_for_assisted_service.py --target $(TARGET) --namespace "$(NAMESPACE)" \ --profile "$(PROFILE)" --domain "$(INGRESS_DOMAIN)" diff --git a/cmd/main.go b/cmd/main.go index e66e7a5b52a3..29fe4d642f0f 100644 --- a/cmd/main.go +++ b/cmd/main.go @@ -1,6 +1,7 @@ package main import ( + "context" "flag" "fmt" "log" @@ -8,6 +9,13 @@ import ( "strings" "time" + "github.com/openshift/assisted-service/pkg/thread" + + "github.com/openshift/assisted-service/internal/imgexpirer" + + "k8s.io/client-go/kubernetes" + "k8s.io/client-go/tools/clientcmd" + "github.com/go-openapi/strfmt" "github.com/go-openapi/swag" "github.com/jinzhu/gorm" @@ -21,7 +29,6 @@ import ( "github.com/openshift/assisted-service/internal/events" "github.com/openshift/assisted-service/internal/hardware" "github.com/openshift/assisted-service/internal/host" - "github.com/openshift/assisted-service/internal/imgexpirer" "github.com/openshift/assisted-service/internal/metrics" "github.com/openshift/assisted-service/internal/versions" "github.com/openshift/assisted-service/models" @@ -30,10 +37,10 @@ import ( "github.com/openshift/assisted-service/pkg/db" "github.com/openshift/assisted-service/pkg/generator" "github.com/openshift/assisted-service/pkg/job" + "github.com/openshift/assisted-service/pkg/leader" "github.com/openshift/assisted-service/pkg/ocm" "github.com/openshift/assisted-service/pkg/requestid" "github.com/openshift/assisted-service/pkg/s3wrapper" - "github.com/openshift/assisted-service/pkg/thread" "github.com/openshift/assisted-service/restapi" "github.com/prometheus/client_golang/prometheus" "github.com/sirupsen/logrus" @@ -67,6 +74,7 @@ var Options struct { OCMConfig ocm.Config HostConfig host.Config LogLevel string `envconfig:"LOG_LEVEL" default:"info"` + LeaderConfig leader.Config } func main() { @@ -115,6 +123,7 @@ func main() { } } + var lead leader.ElectorInterface authHandler := auth.NewAuthHandler(Options.Auth, ocmClient, log.WithField("pkg", "auth")) authzHandler := auth.NewAuthzHandler(Options.Auth, ocmClient, log.WithField("pkg", "authz")) versionHandler := versions.NewHandler(Options.Versions) @@ -125,19 +134,6 @@ func main() { instructionApi := host.NewInstructionManager(log.WithField("pkg", "instructions"), db, hwValidator, Options.InstructionConfig, connectivityValidator) prometheusRegistry := prometheus.DefaultRegisterer metricsManager := metrics.NewMetricsManager(prometheusRegistry) - hostApi := host.NewManager(log.WithField("pkg", "host-state"), db, eventsHandler, hwValidator, instructionApi, &Options.HWValidatorConfig, metricsManager, &Options.HostConfig) - clusterApi := cluster.NewManager(Options.ClusterConfig, log.WithField("pkg", "cluster-state"), db, - eventsHandler, hostApi, metricsManager) - - clusterStateMonitor := thread.New( - log.WithField("pkg", "cluster-monitor"), "Cluster State Monitor", Options.ClusterStateMonitorInterval, clusterApi.ClusterMonitoring) - clusterStateMonitor.Start() - defer clusterStateMonitor.Stop() - - hostStateMonitor := thread.New( - log.WithField("pkg", "host-monitor"), "Host State Monitor", Options.HostStateMonitorInterval, hostApi.HostMonitoring) - hostStateMonitor.Start() - defer hostStateMonitor.Stop() log.Println("DeployTarget: " + Options.DeployTarget) @@ -171,7 +167,22 @@ func main() { log.Fatal("failed to create client:", err) } generator = job.New(log.WithField("pkg", "k8s-job-wrapper"), kclient, Options.JobConfig) + + cfg, cerr := clientcmd.BuildConfigFromFlags("", "") + if cerr != nil { + log.WithError(cerr).Fatalf("Failed to create kubernetes cluster config") + } + k8sClient := kubernetes.NewForConfigOrDie(cfg) + lead = leader.NewElector(k8sClient, Options.LeaderConfig, "assisted-service-leader-election-helper", + log.WithField("pkg", "monitor-runner")) + err = lead.StartLeaderElection(context.Background()) + if err != nil { + log.WithError(cerr).Fatalf("Failed to start leader") + } + case "onprem": + + lead = &leader.DummyElector{} // in on-prem mode, setup file system s3 driver and use localjob implementation objectHandler = s3wrapper.NewFSClient("/data", log) if objectHandler == nil { @@ -183,6 +194,21 @@ func main() { log.Fatalf("not supported deploy target %s", Options.DeployTarget) } + hostApi := host.NewManager(log.WithField("pkg", "host-state"), db, eventsHandler, hwValidator, + instructionApi, &Options.HWValidatorConfig, metricsManager, &Options.HostConfig, lead) + clusterApi := cluster.NewManager(Options.ClusterConfig, log.WithField("pkg", "cluster-state"), db, + eventsHandler, hostApi, metricsManager, lead) + + clusterStateMonitor := thread.New( + log.WithField("pkg", "cluster-monitor"), "Cluster State Monitor", Options.ClusterStateMonitorInterval, clusterApi.ClusterMonitoring) + clusterStateMonitor.StartWithCondition(lead.IsLeader) + defer clusterStateMonitor.Stop() + + hostStateMonitor := thread.New( + log.WithField("pkg", "host-monitor"), "Host State Monitor", Options.HostStateMonitorInterval, hostApi.HostMonitoring) + hostStateMonitor.StartWithCondition(lead.IsLeader) + defer hostStateMonitor.Stop() + if newUrl, err = s3wrapper.FixEndpointURL(Options.BMConfig.S3EndpointURL); err != nil { log.WithError(err).Fatalf("failed to create valid bm config S3 endpoint URL from %s", Options.BMConfig.S3EndpointURL) } else { @@ -193,10 +219,10 @@ func main() { events := events.NewApi(eventsHandler, logrus.WithField("pkg", "eventsApi")) - expirer := imgexpirer.NewManager(objectHandler, eventsHandler, Options.BMConfig.ImageExpirationTime) + expirer := imgexpirer.NewManager(objectHandler, eventsHandler, Options.BMConfig.ImageExpirationTime, lead) imageExpirationMonitor := thread.New( log.WithField("pkg", "image-expiration-monitor"), "Image Expiration Monitor", Options.ImageExpirationInterval, expirer.ExpirationTask) - imageExpirationMonitor.Start() + imageExpirationMonitor.StartWithCondition(lead.IsLeader) defer imageExpirationMonitor.Stop() h, err := restapi.Handler(restapi.Config{ diff --git a/deploy/assisted-service.yaml b/deploy/assisted-service.yaml index b486a988c970..c5c6df36dfc6 100644 --- a/deploy/assisted-service.yaml +++ b/deploy/assisted-service.yaml @@ -7,7 +7,7 @@ spec: selector: matchLabels: app: assisted-service - replicas: 1 + replicas: 3 template: metadata: labels: diff --git a/deploy/roles/default_role.yaml b/deploy/roles/default_role.yaml index 9f5577a45409..a8f7a6eecc06 100644 --- a/deploy/roles/default_role.yaml +++ b/deploy/roles/default_role.yaml @@ -30,6 +30,19 @@ rules: - batch resources: - jobs + - verbs: + - '*' + apiGroups: + - '' + resources: + - configmaps + - verbs: + - '*' + apiGroups: + - coordination.k8s.io + resources: + - leases + --- kind: RoleBinding apiVersion: rbac.authorization.k8s.io/v1 diff --git a/internal/bminventory/inventory_test.go b/internal/bminventory/inventory_test.go index 16af68a9aba0..6cf8ca0215fd 100644 --- a/internal/bminventory/inventory_test.go +++ b/internal/bminventory/inventory_test.go @@ -2149,7 +2149,7 @@ var _ = Describe("KubeConfig download", func() { mockS3Client = s3wrapper.NewMockAPI(ctrl) mockJob = job.NewMockAPI(ctrl) clusterApi = cluster.NewManager(cluster.Config{}, getTestLog().WithField("pkg", "cluster-monitor"), - db, nil, nil, nil) + db, nil, nil, nil, nil) bm = NewBareMetalInventory(db, getTestLog(), nil, clusterApi, cfg, mockJob, nil, mockS3Client, nil) c = common.Cluster{Cluster: models.Cluster{ @@ -2278,7 +2278,7 @@ var _ = Describe("UploadClusterIngressCert test", func() { mockS3Client = s3wrapper.NewMockAPI(ctrl) mockJob = job.NewMockAPI(ctrl) clusterApi = cluster.NewManager(cluster.Config{}, getTestLog().WithField("pkg", "cluster-monitor"), - db, nil, nil, nil) + db, nil, nil, nil, nil) bm = NewBareMetalInventory(db, getTestLog(), nil, clusterApi, cfg, mockJob, nil, mockS3Client, nil) c = common.Cluster{Cluster: models.Cluster{ ID: &clusterID, diff --git a/internal/cluster/cluster.go b/internal/cluster/cluster.go index 149e9693eae3..521267afaf50 100644 --- a/internal/cluster/cluster.go +++ b/internal/cluster/cluster.go @@ -6,6 +6,8 @@ import ( "net/http" "time" + "github.com/openshift/assisted-service/pkg/leader" + "github.com/openshift/assisted-service/pkg/s3wrapper" "github.com/filanov/stateswitch" @@ -82,9 +84,11 @@ type Manager struct { metricAPI metrics.API hostAPI host.API rp *refreshPreprocessor + leaderElector leader.Leader } -func NewManager(cfg Config, log logrus.FieldLogger, db *gorm.DB, eventsHandler events.Handler, hostAPI host.API, metricApi metrics.API) *Manager { +func NewManager(cfg Config, log logrus.FieldLogger, db *gorm.DB, eventsHandler events.Handler, hostAPI host.API, metricApi metrics.API, + leaderElector leader.Leader) *Manager { th := &transitionHandler{ log: log, db: db, @@ -100,6 +104,7 @@ func NewManager(cfg Config, log logrus.FieldLogger, db *gorm.DB, eventsHandler e metricAPI: metricApi, rp: newRefreshPreprocessor(log, hostAPI), hostAPI: hostAPI, + leaderElector: leaderElector, } } @@ -170,6 +175,12 @@ func (m *Manager) GetMasterNodesIds(ctx context.Context, c *common.Cluster, db * } func (m *Manager) ClusterMonitoring() { + if !m.leaderElector.IsLeader() { + m.log.Debugf("Not a leader, exiting ClusterMonitoring") + return + } + + m.log.Debugf("Running ClusterMonitoring") var ( clusters []*common.Cluster clusterAfterRefresh *common.Cluster @@ -184,6 +195,11 @@ func (m *Manager) ClusterMonitoring() { return } for _, cluster := range clusters { + + if !m.leaderElector.IsLeader() { + m.log.Debugf("Not a leader, exiting ClusterMonitoring") + return + } if clusterAfterRefresh, err = m.RefreshStatus(ctx, cluster, m.db); err != nil { log.WithError(err).Errorf("failed to refresh cluster %s state", cluster.ID) continue diff --git a/internal/cluster/cluster_test.go b/internal/cluster/cluster_test.go index e4f50e6a84a8..6b712059ec8a 100644 --- a/internal/cluster/cluster_test.go +++ b/internal/cluster/cluster_test.go @@ -9,6 +9,8 @@ import ( "net/http" "time" + "github.com/openshift/assisted-service/pkg/leader" + "github.com/openshift/assisted-service/pkg/s3wrapper" "github.com/kelseyhightower/envconfig" @@ -51,7 +53,8 @@ var _ = Describe("stateMachine", func() { BeforeEach(func() { db = common.PrepareTestDB(dbName, &events.Event{}) - state = NewManager(defaultTestConfig, getTestLog(), db, nil, nil, nil) + dummy := &leader.DummyElector{} + state = NewManager(defaultTestConfig, getTestLog(), db, nil, nil, nil, dummy) id := strfmt.UUID(uuid.New().String()) cluster = &common.Cluster{Cluster: models.Cluster{ ID: &id, @@ -117,8 +120,9 @@ var _ = Describe("cluster monitor", func() { mockHostAPI = host.NewMockAPI(ctrl) mockMetric = metrics.NewMockAPI(ctrl) mockEvents = events.NewMockHandler(ctrl) + dummy := &leader.DummyElector{} clusterApi = NewManager(defaultTestConfig, getTestLog().WithField("pkg", "cluster-monitor"), db, - mockEvents, mockHostAPI, mockMetric) + mockEvents, mockHostAPI, mockMetric, dummy) expectedState = "" shouldHaveUpdated = false }) @@ -415,8 +419,9 @@ var _ = Describe("VerifyRegisterHost", func() { BeforeEach(func() { db = common.PrepareTestDB(dbName, &events.Event{}) id = strfmt.UUID(uuid.New().String()) + dummy := &leader.DummyElector{} clusterApi = NewManager(defaultTestConfig, getTestLog().WithField("pkg", "cluster-monitor"), db, - nil, nil, nil) + nil, nil, nil, dummy) }) checkVerifyRegisterHost := func(clusterStatus string, expectErr bool) { @@ -466,8 +471,9 @@ var _ = Describe("VerifyClusterUpdatability", func() { BeforeEach(func() { db = common.PrepareTestDB(dbName, &events.Event{}) id = strfmt.UUID(uuid.New().String()) + dummy := &leader.DummyElector{} clusterApi = NewManager(defaultTestConfig, getTestLog().WithField("pkg", "cluster-monitor"), db, - nil, nil, nil) + nil, nil, nil, dummy) }) checkVerifyClusterUpdatability := func(clusterStatus string, expectErr bool) { @@ -513,8 +519,9 @@ var _ = Describe("SetGeneratorVersion", func() { It("set generator version", func() { db = common.PrepareTestDB(dbName, &events.Event{}) id = strfmt.UUID(uuid.New().String()) + dummy := &leader.DummyElector{} clusterApi = NewManager(defaultTestConfig, getTestLog().WithField("pkg", "cluster-monitor"), db, - nil, nil, nil) + nil, nil, nil, dummy) cluster := common.Cluster{Cluster: models.Cluster{ID: &id, Status: swag.String(models.ClusterStatusReady)}} Expect(db.Create(&cluster).Error).ShouldNot(HaveOccurred()) cluster = geCluster(id, db) @@ -544,7 +551,8 @@ var _ = Describe("CancelInstallation", func() { eventsHandler = events.New(db, logrus.New()) ctrl = gomock.NewController(GinkgoT()) mockMetric = metrics.NewMockAPI(ctrl) - state = NewManager(defaultTestConfig, getTestLog(), db, eventsHandler, nil, mockMetric) + dummy := &leader.DummyElector{} + state = NewManager(defaultTestConfig, getTestLog(), db, eventsHandler, nil, mockMetric, dummy) id := strfmt.UUID(uuid.New().String()) c = common.Cluster{Cluster: models.Cluster{ ID: &id, @@ -615,7 +623,8 @@ var _ = Describe("ResetCluster", func() { BeforeEach(func() { db = common.PrepareTestDB(dbName, &events.Event{}) eventsHandler = events.New(db, logrus.New()) - state = NewManager(defaultTestConfig, getTestLog(), db, eventsHandler, nil, nil) + dummy := &leader.DummyElector{} + state = NewManager(defaultTestConfig, getTestLog(), db, eventsHandler, nil, nil, dummy) }) It("reset_cluster", func() { @@ -725,7 +734,8 @@ var _ = Describe("PrepareForInstallation", func() { BeforeEach(func() { db = common.PrepareTestDB(dbName, &events.Event{}) - capi = NewManager(defaultTestConfig, getTestLog(), db, nil, nil, nil) + dummy := &leader.DummyElector{} + capi = NewManager(defaultTestConfig, getTestLog(), db, nil, nil, nil, dummy) clusterId = strfmt.UUID(uuid.New().String()) }) @@ -806,7 +816,8 @@ var _ = Describe("HandlePreInstallationError", func() { BeforeEach(func() { db = common.PrepareTestDB(dbName, &events.Event{}) - capi = NewManager(defaultTestConfig, getTestLog(), db, nil, nil, nil) + dummy := &leader.DummyElector{} + capi = NewManager(defaultTestConfig, getTestLog(), db, nil, nil, nil, dummy) clusterId = strfmt.UUID(uuid.New().String()) }) @@ -891,7 +902,8 @@ var _ = Describe("SetVips", func() { db = common.PrepareTestDB(dbName, &events.Event{}) ctrl = gomock.NewController(GinkgoT()) mockEvents = events.NewMockHandler(ctrl) - capi = NewManager(defaultTestConfig, getTestLog(), db, mockEvents, nil, nil) + dummy := &leader.DummyElector{} + capi = NewManager(defaultTestConfig, getTestLog(), db, mockEvents, nil, nil, dummy) clusterId = strfmt.UUID(uuid.New().String()) }) AfterEach(func() { @@ -1024,8 +1036,9 @@ var _ = Describe("ready_state", func() { db = common.PrepareTestDB(dbName, &events.Event{}) ctrl = gomock.NewController(GinkgoT()) mockEvents = events.NewMockHandler(ctrl) + dummy := &leader.DummyElector{} clusterApi = NewManager(defaultTestConfig, getTestLog().WithField("pkg", "cluster-monitor"), db, - mockEvents, nil, nil) + mockEvents, nil, nil, dummy) id = strfmt.UUID(uuid.New().String()) cluster = common.Cluster{Cluster: models.Cluster{ @@ -1088,8 +1101,9 @@ var _ = Describe("insufficient_state", func() { mockHostAPI = host.NewMockAPI(ctrl) mockEvents := events.NewMockHandler(ctrl) db = common.PrepareTestDB(dbName, &events.Event{}) + dummy := &leader.DummyElector{} clusterApi = NewManager(defaultTestConfig, getTestLog().WithField("pkg", "cluster-monitor"), db, - mockEvents, mockHostAPI, nil) + mockEvents, mockHostAPI, nil, dummy) id = strfmt.UUID(uuid.New().String()) cluster = common.Cluster{Cluster: models.Cluster{ @@ -1129,7 +1143,8 @@ var _ = Describe("prepare-for-installation refresh status", func() { ctrl = gomock.NewController(GinkgoT()) mockHostAPI = host.NewMockAPI(ctrl) mockEvents := events.NewMockHandler(ctrl) - capi = NewManager(cfg, getTestLog(), db, mockEvents, mockHostAPI, nil) + dummy := &leader.DummyElector{} + capi = NewManager(cfg, getTestLog(), db, mockEvents, mockHostAPI, nil, dummy) clusterId = strfmt.UUID(uuid.New().String()) cl = common.Cluster{ Cluster: models.Cluster{ @@ -1189,7 +1204,8 @@ var _ = Describe("Cluster tarred files", func() { mockS3Client = s3wrapper.NewMockAPI(ctrl) mockHostAPI = host.NewMockAPI(ctrl) mockEvents := events.NewMockHandler(ctrl) - capi = NewManager(cfg, getTestLog(), db, mockEvents, mockHostAPI, nil) + dummy := &leader.DummyElector{} + capi = NewManager(cfg, getTestLog(), db, mockEvents, mockHostAPI, nil, dummy) clusterId = strfmt.UUID(uuid.New().String()) cl = common.Cluster{ Cluster: models.Cluster{ diff --git a/internal/cluster/mock_cluster_api.go b/internal/cluster/mock_cluster_api.go index a386d0a39573..2fb41f0a9956 100644 --- a/internal/cluster/mock_cluster_api.go +++ b/internal/cluster/mock_cluster_api.go @@ -12,7 +12,6 @@ import ( gomock "github.com/golang/mock/gomock" gorm "github.com/jinzhu/gorm" common "github.com/openshift/assisted-service/internal/common" - s3wrapper "github.com/openshift/assisted-service/pkg/s3wrapper" ) diff --git a/internal/cluster/transition_test.go b/internal/cluster/transition_test.go index 2cfa9828209d..e881318a4813 100644 --- a/internal/cluster/transition_test.go +++ b/internal/cluster/transition_test.go @@ -39,7 +39,7 @@ var _ = Describe("Transition tests", func() { eventsHandler = events.New(db, logrus.New()) ctrl = gomock.NewController(GinkgoT()) mockMetric = metrics.NewMockAPI(ctrl) - capi = NewManager(defaultTestConfig, getTestLog(), db, eventsHandler, nil, mockMetric) + capi = NewManager(defaultTestConfig, getTestLog(), db, eventsHandler, nil, mockMetric, nil) clusterId = strfmt.UUID(uuid.New().String()) }) @@ -162,7 +162,7 @@ var _ = Describe("Cancel cluster installation", func() { ctrl = gomock.NewController(GinkgoT()) mockEventsHandler = events.NewMockHandler(ctrl) mockMetric = metrics.NewMockAPI(ctrl) - capi = NewManager(defaultTestConfig, getTestLog(), db, mockEventsHandler, nil, mockMetric) + capi = NewManager(defaultTestConfig, getTestLog(), db, mockEventsHandler, nil, mockMetric, nil) }) acceptNewEvents := func(times int) { @@ -230,7 +230,7 @@ var _ = Describe("Reset cluster", func() { db = common.PrepareTestDB(dbName, &events.Event{}) ctrl = gomock.NewController(GinkgoT()) mockEventsHandler = events.NewMockHandler(ctrl) - capi = NewManager(defaultTestConfig, getTestLog(), db, mockEventsHandler, nil, nil) + capi = NewManager(defaultTestConfig, getTestLog(), db, mockEventsHandler, nil, nil, nil) }) acceptNewEvents := func(times int) { @@ -357,7 +357,7 @@ var _ = Describe("Refresh Cluster - No DHCP", func() { mockHostAPI = host.NewMockAPI(ctrl) mockMetric = metrics.NewMockAPI(ctrl) clusterApi = NewManager(defaultTestConfig, getTestLog().WithField("pkg", "cluster-monitor"), db, - mockEvents, mockHostAPI, mockMetric) + mockEvents, mockHostAPI, mockMetric, nil) hid1 = strfmt.UUID(uuid.New().String()) hid2 = strfmt.UUID(uuid.New().String()) @@ -819,7 +819,7 @@ var _ = Describe("Refresh Cluster - Advanced networking validations", func() { mockHostAPI = host.NewMockAPI(ctrl) mockMetric = metrics.NewMockAPI(ctrl) clusterApi = NewManager(defaultTestConfig, getTestLog().WithField("pkg", "cluster-monitor"), db, - mockEvents, mockHostAPI, mockMetric) + mockEvents, mockHostAPI, mockMetric, nil) hid1 = strfmt.UUID(uuid.New().String()) hid2 = strfmt.UUID(uuid.New().String()) @@ -1017,7 +1017,7 @@ var _ = Describe("Refresh Cluster - With DHCP", func() { mockHostAPI = host.NewMockAPI(ctrl) mockMetric = metrics.NewMockAPI(ctrl) clusterApi = NewManager(defaultTestConfig, getTestLog().WithField("pkg", "cluster-monitor"), db, - mockEvents, mockHostAPI, mockMetric) + mockEvents, mockHostAPI, mockMetric, nil) hid1 = strfmt.UUID(uuid.New().String()) hid2 = strfmt.UUID(uuid.New().String()) diff --git a/internal/host/host.go b/internal/host/host.go index 6f32f380f9c9..0685a8726fa7 100644 --- a/internal/host/host.go +++ b/internal/host/host.go @@ -7,6 +7,8 @@ import ( "strconv" "time" + "github.com/openshift/assisted-service/pkg/leader" + "github.com/openshift/assisted-service/internal/hostutil" "github.com/go-openapi/strfmt" @@ -98,10 +100,11 @@ type Manager struct { rp *refreshPreprocessor metricApi metrics.API Config Config + leaderElector leader.Leader } func NewManager(log logrus.FieldLogger, db *gorm.DB, eventsHandler events.Handler, hwValidator hardware.Validator, instructionApi InstructionApi, - hwValidatorCfg *hardware.ValidatorCfg, metricApi metrics.API, config *Config) *Manager { + hwValidatorCfg *hardware.ValidatorCfg, metricApi metrics.API, config *Config, leaderElector leader.ElectorInterface) *Manager { th := &transitionHandler{ db: db, log: log, @@ -117,6 +120,7 @@ func NewManager(log logrus.FieldLogger, db *gorm.DB, eventsHandler events.Handle rp: newRefreshPreprocessor(log, hwValidatorCfg), metricApi: metricApi, Config: *config, + leaderElector: leaderElector, } } diff --git a/internal/host/host_test.go b/internal/host/host_test.go index 9045ed873b66..4fdc91a7ddd0 100644 --- a/internal/host/host_test.go +++ b/internal/host/host_test.go @@ -9,6 +9,8 @@ import ( "strconv" "time" + "github.com/openshift/assisted-service/pkg/leader" + "github.com/openshift/assisted-service/internal/hostutil" "github.com/go-openapi/strfmt" @@ -48,8 +50,9 @@ var _ = Describe("update_role", func() { ) BeforeEach(func() { + dummy := &leader.DummyElector{} db = common.PrepareTestDB(dbName, &events.Event{}) - state = NewManager(getTestLog(), db, nil, nil, nil, createValidatorCfg(), nil, defaultConfig) + state = NewManager(getTestLog(), db, nil, nil, nil, createValidatorCfg(), nil, defaultConfig, dummy) id = strfmt.UUID(uuid.New().String()) clusterID = strfmt.UUID(uuid.New().String()) }) @@ -189,7 +192,8 @@ var _ = Describe("update_progress", func() { ctrl = gomock.NewController(GinkgoT()) mockEvents = events.NewMockHandler(ctrl) mockMetric = metrics.NewMockAPI(ctrl) - state = NewManager(getTestLog(), db, mockEvents, nil, nil, createValidatorCfg(), mockMetric, defaultConfig) + dummy := &leader.DummyElector{} + state = NewManager(getTestLog(), db, mockEvents, nil, nil, createValidatorCfg(), mockMetric, defaultConfig, dummy) id := strfmt.UUID(uuid.New().String()) clusterId := strfmt.UUID(uuid.New().String()) host = getTestHost(id, clusterId, "") @@ -376,7 +380,8 @@ var _ = Describe("monitor_disconnection", func() { db = common.PrepareTestDB(dbName, &events.Event{}) ctrl = gomock.NewController(GinkgoT()) mockEvents = events.NewMockHandler(ctrl) - state = NewManager(getTestLog(), db, mockEvents, nil, nil, createValidatorCfg(), nil, defaultConfig) + dummy := &leader.DummyElector{} + state = NewManager(getTestLog(), db, mockEvents, nil, nil, createValidatorCfg(), nil, defaultConfig, dummy) clusterID := strfmt.UUID(uuid.New().String()) host = getTestHost(strfmt.UUID(uuid.New().String()), clusterID, models.HostStatusDiscovering) cluster := getTestCluster(clusterID, "1.1.0.0/16") @@ -458,7 +463,8 @@ var _ = Describe("cancel installation", func() { BeforeEach(func() { db = common.PrepareTestDB(dbName, &events.Event{}) eventsHandler = events.New(db, logrus.New()) - state = NewManager(getTestLog(), db, eventsHandler, nil, nil, nil, nil, defaultConfig) + dummy := &leader.DummyElector{} + state = NewManager(getTestLog(), db, eventsHandler, nil, nil, nil, nil, defaultConfig, dummy) id := strfmt.UUID(uuid.New().String()) clusterId := strfmt.UUID(uuid.New().String()) h = getTestHost(id, clusterId, models.HostStatusDiscovering) @@ -540,7 +546,8 @@ var _ = Describe("reset host", func() { db = common.PrepareTestDB(dbName, &events.Event{}) eventsHandler = events.New(db, logrus.New()) config = *defaultConfig - state = NewManager(getTestLog(), db, eventsHandler, nil, nil, nil, nil, &config) + dummy := &leader.DummyElector{} + state = NewManager(getTestLog(), db, eventsHandler, nil, nil, nil, nil, &config, dummy) }) AfterEach(func() { common.DeleteTestDB(db, dbName) @@ -798,7 +805,8 @@ var _ = Describe("UpdateInventory", func() { BeforeEach(func() { db = common.PrepareTestDB(dbName, &events.Event{}) - hapi = NewManager(getTestLog(), db, nil, nil, nil, createValidatorCfg(), nil, defaultConfig) + dummy := &leader.DummyElector{} + hapi = NewManager(getTestLog(), db, nil, nil, nil, createValidatorCfg(), nil, defaultConfig, dummy) hostId = strfmt.UUID(uuid.New().String()) clusterId = strfmt.UUID(uuid.New().String()) }) @@ -913,7 +921,8 @@ var _ = Describe("Update hostname", func() { BeforeEach(func() { db = common.PrepareTestDB(dbName, &events.Event{}) - hapi = NewManager(getTestLog(), db, nil, nil, nil, createValidatorCfg(), nil, defaultConfig) + dummy := &leader.DummyElector{} + hapi = NewManager(getTestLog(), db, nil, nil, nil, createValidatorCfg(), nil, defaultConfig, dummy) hostId = strfmt.UUID(uuid.New().String()) clusterId = strfmt.UUID(uuid.New().String()) }) @@ -1029,7 +1038,8 @@ var _ = Describe("SetBootstrap", func() { db = common.PrepareTestDB(dbName, &events.Event{}) ctrl = gomock.NewController(GinkgoT()) mockEvents = events.NewMockHandler(ctrl) - hapi = NewManager(getTestLog(), db, mockEvents, nil, nil, createValidatorCfg(), nil, defaultConfig) + dummy := &leader.DummyElector{} + hapi = NewManager(getTestLog(), db, mockEvents, nil, nil, createValidatorCfg(), nil, defaultConfig, dummy) hostId = strfmt.UUID(uuid.New().String()) clusterId = strfmt.UUID(uuid.New().String()) @@ -1082,7 +1092,8 @@ var _ = Describe("PrepareForInstallation", func() { db = common.PrepareTestDB(dbName, &events.Event{}) ctrl = gomock.NewController(GinkgoT()) mockEvents = events.NewMockHandler(ctrl) - hapi = NewManager(getTestLog(), db, mockEvents, nil, nil, createValidatorCfg(), nil, defaultConfig) + dummy := &leader.DummyElector{} + hapi = NewManager(getTestLog(), db, mockEvents, nil, nil, createValidatorCfg(), nil, defaultConfig, dummy) hostId = strfmt.UUID(uuid.New().String()) clusterId = strfmt.UUID(uuid.New().String()) }) @@ -1148,6 +1159,7 @@ var _ = Describe("AutoAssignRole", func() { BeforeEach(func() { db = common.PrepareTestDB(dbName, &events.Event{}) clusterId = strfmt.UUID(uuid.New().String()) + dummy := &leader.DummyElector{} hapi = NewManager( getTestLog(), db, @@ -1157,6 +1169,7 @@ var _ = Describe("AutoAssignRole", func() { createValidatorCfg(), nil, defaultConfig, + dummy, ) Expect(db.Create(&common.Cluster{Cluster: models.Cluster{ID: &clusterId}}).Error).ShouldNot(HaveOccurred()) }) @@ -1254,6 +1267,7 @@ var _ = Describe("IsValidMasterCandidate", func() { BeforeEach(func() { db = common.PrepareTestDB(dbName, &events.Event{}) clusterId = strfmt.UUID(uuid.New().String()) + dummy := &leader.DummyElector{} hapi = NewManager( getTestLog(), db, @@ -1263,6 +1277,7 @@ var _ = Describe("IsValidMasterCandidate", func() { createValidatorCfg(), nil, defaultConfig, + dummy, ) Expect(db.Create(&common.Cluster{Cluster: models.Cluster{ID: &clusterId}}).Error).ShouldNot(HaveOccurred()) }) diff --git a/internal/host/monitor.go b/internal/host/monitor.go index 48b21e98d30f..ea37b43a66b0 100644 --- a/internal/host/monitor.go +++ b/internal/host/monitor.go @@ -8,6 +8,12 @@ import ( ) func (m *Manager) HostMonitoring() { + if !m.leaderElector.IsLeader() { + m.log.Debugf("Not a leader, exiting HostMonitoring") + return + } + + m.log.Debugf("Running HostMonitoring") var ( hosts []*models.Host requestID = requestid.NewID() diff --git a/internal/host/transition_test.go b/internal/host/transition_test.go index 9b4565f69e24..d651f690a082 100644 --- a/internal/host/transition_test.go +++ b/internal/host/transition_test.go @@ -51,7 +51,7 @@ var _ = Describe("RegisterHost", func() { ctrl = gomock.NewController(GinkgoT()) db = common.PrepareTestDB(dbName, &events.Event{}) mockEvents = events.NewMockHandler(ctrl) - hapi = NewManager(getTestLog(), db, mockEvents, nil, nil, createValidatorCfg(), nil, defaultConfig) + hapi = NewManager(getTestLog(), db, mockEvents, nil, nil, createValidatorCfg(), nil, defaultConfig, nil) hostId = strfmt.UUID(uuid.New().String()) clusterId = strfmt.UUID(uuid.New().String()) }) @@ -360,7 +360,7 @@ var _ = Describe("HostInstallationFailed", func() { ctrl = gomock.NewController(GinkgoT()) mockMetric = metrics.NewMockAPI(ctrl) mockEvents = events.NewMockHandler(ctrl) - hapi = NewManager(getTestLog(), db, mockEvents, nil, nil, createValidatorCfg(), mockMetric, defaultConfig) + hapi = NewManager(getTestLog(), db, mockEvents, nil, nil, createValidatorCfg(), mockMetric, defaultConfig, nil) hostId = strfmt.UUID(uuid.New().String()) clusterId = strfmt.UUID(uuid.New().String()) host = getTestHost(hostId, clusterId, "") @@ -400,7 +400,7 @@ var _ = Describe("Cancel host installation", func() { db = common.PrepareTestDB(dbName, &events.Event{}) ctrl = gomock.NewController(GinkgoT()) mockEventsHandler = events.NewMockHandler(ctrl) - hapi = NewManager(getTestLog(), db, mockEventsHandler, nil, nil, createValidatorCfg(), nil, defaultConfig) + hapi = NewManager(getTestLog(), db, mockEventsHandler, nil, nil, createValidatorCfg(), nil, defaultConfig, nil) }) tests := []struct { @@ -473,7 +473,7 @@ var _ = Describe("Reset host", func() { db = common.PrepareTestDB(dbName, &events.Event{}) ctrl = gomock.NewController(GinkgoT()) mockEventsHandler = events.NewMockHandler(ctrl) - hapi = NewManager(getTestLog(), db, mockEventsHandler, nil, nil, createValidatorCfg(), nil, defaultConfig) + hapi = NewManager(getTestLog(), db, mockEventsHandler, nil, nil, createValidatorCfg(), nil, defaultConfig, nil) }) tests := []struct { @@ -546,7 +546,7 @@ var _ = Describe("Install", func() { db = common.PrepareTestDB(dbName, &events.Event{}) ctrl = gomock.NewController(GinkgoT()) mockEvents = events.NewMockHandler(ctrl) - hapi = NewManager(getTestLog(), db, mockEvents, nil, nil, createValidatorCfg(), nil, defaultConfig) + hapi = NewManager(getTestLog(), db, mockEvents, nil, nil, createValidatorCfg(), nil, defaultConfig, nil) hostId = strfmt.UUID(uuid.New().String()) clusterId = strfmt.UUID(uuid.New().String()) }) @@ -699,7 +699,7 @@ var _ = Describe("Disable", func() { db = common.PrepareTestDB(dbName, &events.Event{}) ctrl = gomock.NewController(GinkgoT()) mockEvents = events.NewMockHandler(ctrl) - hapi = NewManager(getTestLog(), db, mockEvents, nil, nil, createValidatorCfg(), nil, defaultConfig) + hapi = NewManager(getTestLog(), db, mockEvents, nil, nil, createValidatorCfg(), nil, defaultConfig, nil) hostId = strfmt.UUID(uuid.New().String()) clusterId = strfmt.UUID(uuid.New().String()) }) @@ -830,7 +830,7 @@ var _ = Describe("Enable", func() { db = common.PrepareTestDB(dbName, &events.Event{}) ctrl = gomock.NewController(GinkgoT()) mockEvents = events.NewMockHandler(ctrl) - hapi = NewManager(getTestLog(), db, mockEvents, nil, nil, createValidatorCfg(), nil, defaultConfig) + hapi = NewManager(getTestLog(), db, mockEvents, nil, nil, createValidatorCfg(), nil, defaultConfig, nil) hostId = strfmt.UUID(uuid.New().String()) clusterId = strfmt.UUID(uuid.New().String()) }) @@ -1013,7 +1013,7 @@ var _ = Describe("Refresh Host", func() { db = common.PrepareTestDB(dbName, &events.Event{}) ctrl = gomock.NewController(GinkgoT()) mockEvents = events.NewMockHandler(ctrl) - hapi = NewManager(getTestLog(), db, mockEvents, nil, nil, createValidatorCfg(), nil, defaultConfig) + hapi = NewManager(getTestLog(), db, mockEvents, nil, nil, createValidatorCfg(), nil, defaultConfig, nil) hostId = strfmt.UUID(uuid.New().String()) clusterId = strfmt.UUID(uuid.New().String()) }) diff --git a/internal/imgexpirer/imgexpirer.go b/internal/imgexpirer/imgexpirer.go index 0c4c1d4d7707..ea8508a391e8 100644 --- a/internal/imgexpirer/imgexpirer.go +++ b/internal/imgexpirer/imgexpirer.go @@ -5,6 +5,8 @@ import ( "regexp" "time" + "github.com/openshift/assisted-service/pkg/leader" + "github.com/go-openapi/strfmt" "github.com/openshift/assisted-service/internal/events" "github.com/openshift/assisted-service/models" @@ -25,17 +27,22 @@ type Manager struct { objectHandler s3wrapper.API eventsHandler events.Handler deleteTime time.Duration + leaderElector leader.Leader } -func NewManager(objectHandler s3wrapper.API, eventsHandler events.Handler, deleteTime time.Duration) *Manager { +func NewManager(objectHandler s3wrapper.API, eventsHandler events.Handler, deleteTime time.Duration, leaderElector leader.ElectorInterface) *Manager { return &Manager{ objectHandler: objectHandler, eventsHandler: eventsHandler, deleteTime: deleteTime, + leaderElector: leaderElector, } } func (m *Manager) ExpirationTask() { + if !m.leaderElector.IsLeader() { + return + } ctx := requestid.ToContext(context.Background(), requestid.NewID()) m.objectHandler.ExpireObjects(ctx, imagePrefix, m.deleteTime, m.DeletedImageCallback) } diff --git a/internal/imgexpirer/imgexpirer_test.go b/internal/imgexpirer/imgexpirer_test.go index 86ffdc1e1a6f..b86bc1d08529 100644 --- a/internal/imgexpirer/imgexpirer_test.go +++ b/internal/imgexpirer/imgexpirer_test.go @@ -35,7 +35,7 @@ var _ = Describe("imgexpirer", func() { ctrl = gomock.NewController(GinkgoT()) mockEvents = events.NewMockHandler(ctrl) deleteTime, _ := time.ParseDuration("60m") - imgExp = NewManager(nil, mockEvents, deleteTime) + imgExp = NewManager(nil, mockEvents, deleteTime, nil) }) It("callback_valid_objname", func() { clusterId := "53116787-3eb0-4211-93ac-611d5cedaa30" diff --git a/pkg/leader/leaderelector.go b/pkg/leader/leaderelector.go new file mode 100644 index 000000000000..93d25452bd5c --- /dev/null +++ b/pkg/leader/leaderelector.go @@ -0,0 +1,122 @@ +package leader + +import ( + "context" + "os" + "time" + + "github.com/sirupsen/logrus" + + "k8s.io/apimachinery/pkg/util/uuid" + "k8s.io/client-go/kubernetes" + "k8s.io/client-go/tools/leaderelection" + "k8s.io/client-go/tools/leaderelection/resourcelock" +) + +type Config struct { + LeaseDuration time.Duration `envconfig:"LEADER_LEASE_DURATION" default:"15s"` + RetryInterval time.Duration `envconfig:"LEADER_RETRY_INTERVAL" default:"2s"` + RenewDeadline time.Duration `envconfig:"LEADER_RENEW_DEADLINE" default:"10s"` + Namespace string `envconfig:"NAMESPACE" default:"assisted-installer"` +} + +//go:generate mockgen -source=leaderelector.go -package=leader -destination=mock_leader_elector.go + +type Leader interface { + IsLeader() bool +} +type ElectorInterface interface { + Leader + StartLeaderElection(ctx context.Context) error +} + +type DummyElector struct{} + +func (f *DummyElector) StartLeaderElection(ctx context.Context) error { + return nil +} + +func (f *DummyElector) IsLeader() bool { + return true +} + +var _ ElectorInterface = &Elector{} + +type Elector struct { + log logrus.FieldLogger + config Config + kube *kubernetes.Clientset + isLeader bool + configMapName string +} + +func NewElector(kubeClient *kubernetes.Clientset, config Config, configMapName string, logger logrus.FieldLogger) *Elector { + return &Elector{log: logger, config: config, kube: kubeClient, configMapName: configMapName, isLeader: false} +} + +func (l *Elector) IsLeader() bool { + return l.isLeader +} + +func (l *Elector) StartLeaderElection(ctx context.Context) error { + + resourceLock, err := l.createResourceLock(l.configMapName) + if err != nil { + return err + } + + leaderElector, err := l.createLeaderElector(resourceLock) + if err != nil { + return err + } + + l.log.Info("Attempting to acquire leader lease") + // Running loop cause leaderElector.Run is blocking + // and needs to be restarted while leader is lost + go func() { + for { + l.log.Infof("Starting leader elections process") + leaderElector.Run(ctx) + } + }() + + return nil +} + +func (l *Elector) createResourceLock(name string) (resourcelock.Interface, error) { + // Leader id, needs to be unique + id, err := os.Hostname() + if err != nil { + return nil, err + } + + id = id + "_" + string(uuid.NewUUID()) + + return resourcelock.New(resourcelock.ConfigMapsResourceLock, + l.config.Namespace, + name, + l.kube.CoreV1(), + l.kube.CoordinationV1(), + resourcelock.ResourceLockConfig{ + Identity: id, + }) +} + +func (l *Elector) createLeaderElector(resourceLock resourcelock.Interface) (*leaderelection.LeaderElector, error) { + return leaderelection.NewLeaderElector(leaderelection.LeaderElectionConfig{ + Lock: resourceLock, + LeaseDuration: l.config.LeaseDuration, + RenewDeadline: l.config.RenewDeadline, + RetryPeriod: l.config.RetryInterval, + Callbacks: leaderelection.LeaderCallbacks{ + OnStartedLeading: func(_ context.Context) { + l.log.Info("Successfully acquired leadership lease") + l.isLeader = true + }, + OnStoppedLeading: func() { + l.log.Infof("NO LONGER LEADER.") + l.isLeader = false + }, + }, + }) +} diff --git a/pkg/leader/mock_leader_elector.go b/pkg/leader/mock_leader_elector.go new file mode 100644 index 000000000000..e0197df926a3 --- /dev/null +++ b/pkg/leader/mock_leader_elector.go @@ -0,0 +1,63 @@ +// Code generated by MockGen. DO NOT EDIT. +// Source: leaderelector.go + +// Package leader is a generated GoMock package. +package leader + +import ( + context "context" + reflect "reflect" + + gomock "github.com/golang/mock/gomock" +) + +// MockElectorInterface is a mock of ElectorInterface interface +type MockElectorInterface struct { + ctrl *gomock.Controller + recorder *MockElectorInterfaceMockRecorder +} + +// MockElectorInterfaceMockRecorder is the mock recorder for MockElectorInterface +type MockElectorInterfaceMockRecorder struct { + mock *MockElectorInterface +} + +// NewMockElectorInterface creates a new mock instance +func NewMockElectorInterface(ctrl *gomock.Controller) *MockElectorInterface { + mock := &MockElectorInterface{ctrl: ctrl} + mock.recorder = &MockElectorInterfaceMockRecorder{mock} + return mock +} + +// EXPECT returns an object that allows the caller to indicate expected use +func (m *MockElectorInterface) EXPECT() *MockElectorInterfaceMockRecorder { + return m.recorder +} + +// StartLeaderElection mocks base method +func (m *MockElectorInterface) StartLeaderElection(ctx context.Context) error { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "StartLeaderElection", ctx) + ret0, _ := ret[0].(error) + return ret0 +} + +// StartLeaderElection indicates an expected call of StartLeaderElection +func (mr *MockElectorInterfaceMockRecorder) StartLeaderElection(ctx interface{}) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "StartLeaderElection", reflect.TypeOf((*MockElectorInterface)(nil).StartLeaderElection), ctx) +} + +// IsLeader mocks base method +func (m *MockElectorInterface) IsLeader() bool { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "IsLeader") + ret0, _ := ret[0].(bool) + return ret0 +} + +// IsLeader indicates an expected call of IsLeader +func (mr *MockElectorInterfaceMockRecorder) IsLeader() *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "IsLeader", reflect.TypeOf((*MockElectorInterface)(nil).IsLeader)) +} diff --git a/pkg/thread/thread.go b/pkg/thread/thread.go index a8a48417d1f8..6b2a3b83019f 100644 --- a/pkg/thread/thread.go +++ b/pkg/thread/thread.go @@ -38,7 +38,15 @@ func New(log logrus.FieldLogger, name string, interval time.Duration, exec func( // Start thread func (t *Thread) Start() { t.log.Infof("Started %s", t.name) - go t.loop() + go t.loop(func() bool { + return true + }) +} + +// Start thread with condition +func (t *Thread) StartWithCondition(condition func() bool) { + t.log.Infof("Started %s with given condition", t.name) + go t.loop(condition) } // Stop thread @@ -49,19 +57,19 @@ func (t *Thread) Stop() { t.log.Infof("Stopped %s", t.name) } -func (t *Thread) loop() { - intervalTimer := time.NewTimer(0) - +func (t *Thread) loop(condition func() bool) { defer close(t.done) - defer intervalTimer.Stop() + ticker := time.NewTicker(t.interval) + defer ticker.Stop() for { select { case <-t.done: return - case <-intervalTimer.C: - t.exec() - intervalTimer.Reset(t.interval) + case <-ticker.C: + if condition() { + t.exec() + } } } } diff --git a/tools/deploy_assisted_installer.py b/tools/deploy_assisted_installer.py index 80c03eb1d66b..00e2627a4e82 100644 --- a/tools/deploy_assisted_installer.py +++ b/tools/deploy_assisted_installer.py @@ -33,10 +33,10 @@ def main(): with open(SRC_FILE, "r") as src: raw_data = src.read() raw_data = raw_data.replace('REPLACE_NAMESPACE', deploy_options.namespace) - data = yaml.safe_load(raw_data) image_fqdn = deployment_options.get_image_override(deploy_options, "assisted-service", "SERVICE") + data["spec"]["replicas"] = deploy_options.replicas_count data["spec"]["template"]["spec"]["containers"][0]["image"] = image_fqdn if deploy_options.subsystem_test: if data["spec"]["template"]["spec"]["containers"][0].get("env", None) is None: diff --git a/tools/deployment_options.py b/tools/deployment_options.py index eeb80bdddc64..e3b409bae73a 100644 --- a/tools/deployment_options.py +++ b/tools/deployment_options.py @@ -36,6 +36,13 @@ def load_deployment_options(parser=None): type=str ) + parser.add_argument( + '--replicas-count', + help='Replicas count of assisted-service', + type=int, + default=3 + ) + deploy_options = parser.add_mutually_exclusive_group() deploy_options.add_argument("--deploy-tag", help='Tag for all deployment images', type=str) deploy_options.add_argument("--deploy-manifest-tag", help='Tag of the assisted-installer-deployment repo to get the deployment images manifest from', type=str)