Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
110 changes: 70 additions & 40 deletions manager/manager.go
Original file line number Diff line number Diff line change
Expand Up @@ -170,6 +170,12 @@ type Manager struct {
dekRotator *RaftDEKManager
roleManager *roleManager

// becomeFollowerCleanups holds the teardown for each leader-only
// component started by becomeLeader, registered with onBecomeFollower
// next to the code that starts it. becomeFollower runs and clears it.
// Guarded by mu.
becomeFollowerCleanups []func()

cancelFunc context.CancelFunc

// mu is a general mutex used to coordinate starting/stopping and
Expand Down Expand Up @@ -1062,6 +1068,10 @@ func (m *Manager) becomeLeader(ctx context.Context) {
log.G(ctx).WithError(err).Error("keymanager failed with an error")
}
}(m.keyManager)
m.onBecomeFollower(func() {
m.keyManager.Stop()
m.keyManager = nil
})
}

go func(d *dispatcher.Dispatcher) {
Expand All @@ -1080,16 +1090,23 @@ func (m *Manager) becomeLeader(ctx context.Context) {
log.G(ctx).WithError(err).Error("Dispatcher exited with an error")
}
}(m.dispatcher)
// The dispatcher, logbroker and CA server are gRPC services that are
// registered when creating the manager and would need to be re-registered
// if they were recreated. For simplicity, they are stopped but not nilled
// out.
m.onBecomeFollower(func() { m.dispatcher.Stop() })

if err := m.logbroker.Start(ctx); err != nil {
log.G(ctx).WithError(err).Error("LogBroker failed to start")
}
m.onBecomeFollower(func() { m.logbroker.Stop() })

go func(server *ca.Server) {
if err := server.Run(ctx); err != nil {
log.G(ctx).WithError(err).Error("CA signer exited with an error")
}
}(m.caserver)
m.onBecomeFollower(func() { m.caserver.Stop() })

// Start all sub-components in separate goroutines.
// TODO(aluzzardi): This should have some kind of error handling so that
Expand All @@ -1100,95 +1117,108 @@ func (m *Manager) becomeLeader(ctx context.Context) {
log.G(ctx).WithError(err).Error("allocator exited with an error")
}
}(m.allocator)
m.onBecomeFollower(func() {
m.allocator.Stop()
m.allocator = nil
})
}

go func(scheduler *scheduler.Scheduler) {
if err := scheduler.Run(ctx); err != nil {
log.G(ctx).WithError(err).Error("scheduler exited with an error")
}
}(m.scheduler)
m.onBecomeFollower(func() {
m.scheduler.Stop()
m.scheduler = nil
})

go func(constraintEnforcer *constraintenforcer.ConstraintEnforcer) {
constraintEnforcer.Run()
}(m.constraintEnforcer)
m.onBecomeFollower(func() {
m.constraintEnforcer.Stop()
m.constraintEnforcer = nil
})

go func(volumeEnforcer *volumeenforcer.VolumeEnforcer) {
volumeEnforcer.Run()
}(m.volumeEnforcer)
m.onBecomeFollower(func() {
m.volumeEnforcer.Stop()
m.volumeEnforcer = nil
})

go func(taskReaper *taskreaper.TaskReaper) {
taskReaper.Run(ctx)
}(m.taskReaper)
m.onBecomeFollower(func() {
m.taskReaper.Stop()
m.taskReaper = nil
})

go func(orchestrator *replicated.Orchestrator) {
if err := orchestrator.Run(ctx); err != nil {
log.G(ctx).WithError(err).Error("replicated orchestrator exited with an error")
}
}(m.replicatedOrchestrator)
m.onBecomeFollower(func() {
m.replicatedOrchestrator.Stop()
m.replicatedOrchestrator = nil
})

go func(orchestrator *jobs.Orchestrator) {
// jobs orchestrator does not return errors.
orchestrator.Run(ctx)
}(m.jobsOrchestrator)
m.onBecomeFollower(func() {
m.jobsOrchestrator.Stop()
m.jobsOrchestrator = nil
})

go func(globalOrchestrator *global.Orchestrator) {
if err := globalOrchestrator.Run(ctx); err != nil {
log.G(ctx).WithError(err).Error("global orchestrator exited with an error")
}
}(m.globalOrchestrator)
m.onBecomeFollower(func() {
m.globalOrchestrator.Stop()
m.globalOrchestrator = nil
})

go func(roleManager *roleManager) {
roleManager.Run(ctx)
}(m.roleManager)
m.onBecomeFollower(func() {
m.roleManager.Stop()
m.roleManager = nil
})

go func(volumeManager *csi.Manager) {
volumeManager.Run(ctx)
}(m.volumeManager)
m.onBecomeFollower(func() {
m.volumeManager.Stop()
m.volumeManager = nil
})
}

// becomeFollower shuts down the subsystems that are only run by the leader.
// becomeFollower shuts down the subsystems that are only run by the leader,
// by running the teardowns that becomeLeader registered with onBecomeFollower.
func (m *Manager) becomeFollower() {
// The following components are gRPC services that are
// registered when creating the manager and will need
// to be re-registered if they are recreated.
// For simplicity, they are not nilled out.
m.dispatcher.Stop()
m.logbroker.Stop()
m.caserver.Stop()

if m.allocator != nil {
m.allocator.Stop()
m.allocator = nil
}

m.constraintEnforcer.Stop()
m.constraintEnforcer = nil

m.volumeEnforcer.Stop()
m.volumeEnforcer = nil

m.replicatedOrchestrator.Stop()
m.replicatedOrchestrator = nil

m.globalOrchestrator.Stop()
m.globalOrchestrator = nil

m.taskReaper.Stop()
m.taskReaper = nil

m.scheduler.Stop()
m.scheduler = nil

m.roleManager.Stop()
m.roleManager = nil

if m.keyManager != nil {
m.keyManager.Stop()
m.keyManager = nil
cleanups := m.becomeFollowerCleanups
m.becomeFollowerCleanups = nil
for _, f := range cleanups {
f()
}
}

m.volumeManager.Stop()
m.volumeManager = nil
// onBecomeFollower registers f to run when this manager loses leadership.
// becomeLeader calls it immediately after the code that starts each
// leader-only component, so that a component cannot be started on the leader
// without its shutdown being written alongside.
func (m *Manager) onBecomeFollower(f func()) {
m.becomeFollowerCleanups = append(m.becomeFollowerCleanups, f)
}

// defaultClusterObject creates a default cluster.
Expand Down
16 changes: 16 additions & 0 deletions manager/manager_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -432,3 +432,19 @@ func TestManagerLockUnlock(t *testing.T) {
// error.
<-done
}

// TestBecomeFollowerRunsRegisteredCleanups checks the registry that pairs each
// leader-only component's start with its stop. Cleanups run once, in
// registration order, and a second demotion runs nothing.
func TestBecomeFollowerRunsRegisteredCleanups(t *testing.T) {
m := &Manager{}
var ran []string
m.onBecomeFollower(func() { ran = append(ran, "first") })
m.onBecomeFollower(func() { ran = append(ran, "second") })
Comment on lines +439 to +443

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Not adding one. A real demotion test needs a live raft cluster: becomeFollower stops thirteen components that block on a doneChan only closed by Run, and roleManager.Run dereferences its raft node immediately. No manager test drives a leadership transition today.

An earlier revision of this PR asserted the invariant by parsing the source instead. That was rejected in review in favour of pairing each start with its stop in one place, which is what the first commit does.


m.becomeFollower()
require.Equal(t, []string{"first", "second"}, ran)

m.becomeFollower()
require.Equal(t, []string{"first", "second"}, ran)
}
2 changes: 1 addition & 1 deletion manager/orchestrator/jobs/fakes_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -36,7 +36,7 @@ type fakeReconciler struct {
// ReconcileService implements the reconciler's ReconcileService method, but
// just records what arguments it has been passed, and maybe also returns an
// error if desired.
func (f *fakeReconciler) ReconcileService(id string) error {
func (f *fakeReconciler) ReconcileService(_ context.Context, id string) error {
f.Lock()
defer f.Unlock()
f.servicesReconciled = append(f.servicesReconciled, id)
Expand Down
6 changes: 2 additions & 4 deletions manager/orchestrator/jobs/global/reconciler.go
Original file line number Diff line number Diff line change
Expand Up @@ -35,7 +35,7 @@ func NewReconciler(store *store.MemoryStore, restart restartSupervisor) *Reconci
}

// ReconcileService reconciles one global job service.
func (r *Reconciler) ReconcileService(id string) error {
func (r *Reconciler) ReconcileService(ctx context.Context, id string) error {
var (
service *api.Service
cluster *api.Cluster
Expand Down Expand Up @@ -199,9 +199,7 @@ func (r *Reconciler) ReconcileService(id string) error {
}

// Finally, restart it
// TODO(dperny): pass in context to ReconcileService, so we can
// pass it in here.
return r.restart.Restart(context.Background(), tx, cluster, service, *t)
return r.restart.Restart(ctx, tx, cluster, service, *t)
}); err != nil {
// TODO(dperny): probably should log like in the other
// orchestrators instead of returning here.
Expand Down
4 changes: 2 additions & 2 deletions manager/orchestrator/jobs/global/reconciler_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -157,7 +157,7 @@ var _ = Describe("Global Job Reconciler", func() {

Expect(err).ToNot(HaveOccurred())

err = r.ReconcileService(serviceID)
err = r.ReconcileService(context.Background(), serviceID)
Expect(err).ToNot(HaveOccurred())
})

Expand All @@ -175,7 +175,7 @@ var _ = Describe("Global Job Reconciler", func() {
})
Expect(err).ToNot(HaveOccurred())

err = r.ReconcileService(serviceID)
err = r.ReconcileService(context.Background(), serviceID)
Expect(err).ToNot(HaveOccurred())

s.View(func(tx store.ReadTx) {
Expand Down
10 changes: 5 additions & 5 deletions manager/orchestrator/jobs/orchestrator.go
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,7 @@ import (
type Reconciler interface {
taskinit.InitHandler

ReconcileService(id string) error
ReconcileService(ctx context.Context, id string) error
}

// Orchestrator is the combined orchestrator controlling both Global and
Expand Down Expand Up @@ -134,15 +134,15 @@ func (o *Orchestrator) init(ctx context.Context) {

for _, service := range services {
if orchestrator.IsReplicatedJob(service) {
if err := o.replicatedReconciler.ReconcileService(service.ID); err != nil {
if err := o.replicatedReconciler.ReconcileService(ctx, service.ID); err != nil {
log.G(ctx).WithField(
"service.id", service.ID,
).WithError(err).Error("error reconciling replicated job")
}
}

if orchestrator.IsGlobalJob(service) {
if err := o.globalReconciler.ReconcileService(service.ID); err != nil {
if err := o.globalReconciler.ReconcileService(ctx, service.ID); err != nil {
log.G(ctx).WithField(
"service.id", service.ID,
).WithError(err).Error("error reconciling global job")
Expand Down Expand Up @@ -226,15 +226,15 @@ func (o *Orchestrator) handleEvent(ctx context.Context, event events.Event) {
}

if orchestrator.IsReplicatedJob(service) {
if err := o.replicatedReconciler.ReconcileService(service.ID); err != nil {
if err := o.replicatedReconciler.ReconcileService(ctx, service.ID); err != nil {
log.G(ctx).WithField(
"service.id", service.ID,
).WithError(err).Error("error reconciling replicated job")
}
}

if orchestrator.IsGlobalJob(service) {
if err := o.globalReconciler.ReconcileService(service.ID); err != nil {
if err := o.globalReconciler.ReconcileService(ctx, service.ID); err != nil {
log.G(ctx).WithField(
"service.id", service.ID,
).WithError(err).Error("error reconciling global job")
Expand Down
Loading