From af00a3dbf693790a3278c262767beed8b82be693 Mon Sep 17 00:00:00 2001 From: Richard Davenport Date: Wed, 16 Sep 2026 10:01:54 -0500 Subject: [PATCH 1/4] manager: pair leader-only startup with its teardown becomeLeader() starts the components that only run on the leader, and becomeFollower() stops them again. The two lists sit two hundred lines apart and are kept in step by hand, so a component can be added to one and forgotten in the other. Register each component's teardown with onBecomeFollower() immediately after the code that starts it, and have becomeFollower() run and clear the registered teardowns. Both functions are already called under m.mu, so the slice needs no lock of its own. No functional change: the same thirteen components are stopped, and the same ones nilled out. They are now stopped in the order they were started, where becomeFollower() previously kept its own hand-maintained order. Each component watches the store independently and none calls into another, so the order has no effect. Manager.Stop() is left alone. It stops each component directly under a nil guard, so it works whether or not becomeFollower() has run. Signed-off-by: Richard Davenport --- manager/manager.go | 106 +++++++++++++++++++++++++--------------- manager/manager_test.go | 16 ++++++ 2 files changed, 82 insertions(+), 40 deletions(-) diff --git a/manager/manager.go b/manager/manager.go index 8b14005641..cbb0b7e817 100644 --- a/manager/manager.go +++ b/manager/manager.go @@ -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 @@ -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) { @@ -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 @@ -1100,6 +1117,10 @@ 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) { @@ -1107,24 +1128,44 @@ func (m *Manager) becomeLeader(ctx context.Context) { 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. @@ -1136,59 +1177,44 @@ func (m *Manager) becomeLeader(ctx context.Context) { 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. diff --git a/manager/manager_test.go b/manager/manager_test.go index 50872b1fca..008e45c30c 100644 --- a/manager/manager_test.go +++ b/manager/manager_test.go @@ -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") }) + + m.becomeFollower() + require.Equal(t, []string{"first", "second"}, ran) + + m.becomeFollower() + require.Equal(t, []string{"first", "second"}, ran) +} From 24e25ad801f7da16d876217ef5a7104cdaed2fd6 Mon Sep 17 00:00:00 2001 From: Richard Davenport Date: Wed, 16 Sep 2026 10:02:04 -0500 Subject: [PATCH 2/4] manager: stop the jobs orchestrator when a manager is demoted becomeLeader() starts the jobs orchestrator alongside the other leader-only components, but becomeFollower() never stopped it. Every other leader-only component -- the replicated and global orchestrators, the task reaper, the scheduler, the constraint and volume enforcers, the role and key managers -- is stopped and nilled out there. The jobs orchestrator was the only one that survived demotion. A manager that has held leadership and then lost it therefore kept a live jobs orchestrator for the lifetime of the daemon. Only Manager.Stop() shut it down, which in practice means a dockerd restart. Cancelling the context does not help: in this package the context carries logging information only, and Stop() is the sole shutdown path. In a cluster where several managers have held leadership at some point, this leaves multiple jobs orchestrators reconciling the same service concurrently. Each one creates its own task for the same job iteration, so a replicated job declared with MaxConcurrent: 1 and TotalCompletions: 1 starts one task per live orchestrator, in the same slot, milliseconds apart. Observed on a five-manager swarm running 27.3.1, where two demoted managers produced two tasks for iteration 0 slot 0, both of which ran the job's payload concurrently. Orchestrator.Stop() closes stopChan inside a sync.Once, so it is safe against repeated calls and cannot hit the double-close problem fixed for the task reaper in 2491. Signed-off-by: Richard Davenport --- manager/manager.go | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/manager/manager.go b/manager/manager.go index cbb0b7e817..d2ab18e56e 100644 --- a/manager/manager.go +++ b/manager/manager.go @@ -1171,6 +1171,10 @@ func (m *Manager) becomeLeader(ctx context.Context) { // 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 { From 163b7b92b28b7821a120b15b16915c603d3c57d2 Mon Sep 17 00:00:00 2001 From: Richard Davenport Date: Wed, 16 Sep 2026 10:02:52 -0500 Subject: [PATCH 3/4] manager/orchestrator/jobs: plumb a context into ReconcileService Neither reconciler had a context to log or trace with, and both called restart.Restart() with context.Background() under a TODO asking for the real one. Take a context.Context in ReconcileService and pass it down. Every caller already has one to hand. This clears both TODOs, and lets the next commit log with the reconciliation's own context attached. No functional change. Signed-off-by: Richard Davenport --- manager/orchestrator/jobs/fakes_test.go | 2 +- manager/orchestrator/jobs/global/reconciler.go | 6 ++---- manager/orchestrator/jobs/global/reconciler_test.go | 4 ++-- manager/orchestrator/jobs/orchestrator.go | 10 +++++----- manager/orchestrator/jobs/replicated/reconciler.go | 5 ++--- .../orchestrator/jobs/replicated/reconciler_test.go | 8 ++++---- 6 files changed, 16 insertions(+), 19 deletions(-) diff --git a/manager/orchestrator/jobs/fakes_test.go b/manager/orchestrator/jobs/fakes_test.go index a2202248d5..bf39a05ff8 100644 --- a/manager/orchestrator/jobs/fakes_test.go +++ b/manager/orchestrator/jobs/fakes_test.go @@ -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) diff --git a/manager/orchestrator/jobs/global/reconciler.go b/manager/orchestrator/jobs/global/reconciler.go index c519507b81..bfbb06372d 100644 --- a/manager/orchestrator/jobs/global/reconciler.go +++ b/manager/orchestrator/jobs/global/reconciler.go @@ -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 @@ -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. diff --git a/manager/orchestrator/jobs/global/reconciler_test.go b/manager/orchestrator/jobs/global/reconciler_test.go index 1c9cf4641d..cd9dc09060 100644 --- a/manager/orchestrator/jobs/global/reconciler_test.go +++ b/manager/orchestrator/jobs/global/reconciler_test.go @@ -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()) }) @@ -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) { diff --git a/manager/orchestrator/jobs/orchestrator.go b/manager/orchestrator/jobs/orchestrator.go index 5d53e7019c..bebbe26b21 100644 --- a/manager/orchestrator/jobs/orchestrator.go +++ b/manager/orchestrator/jobs/orchestrator.go @@ -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 @@ -134,7 +134,7 @@ 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") @@ -142,7 +142,7 @@ func (o *Orchestrator) init(ctx context.Context) { } 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") @@ -226,7 +226,7 @@ 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") @@ -234,7 +234,7 @@ func (o *Orchestrator) handleEvent(ctx context.Context, event events.Event) { } 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") diff --git a/manager/orchestrator/jobs/replicated/reconciler.go b/manager/orchestrator/jobs/replicated/reconciler.go index e3b0d5dc69..aad65f9b78 100644 --- a/manager/orchestrator/jobs/replicated/reconciler.go +++ b/manager/orchestrator/jobs/replicated/reconciler.go @@ -39,7 +39,7 @@ func NewReconciler(store *store.MemoryStore, restart restartSupervisor) *Reconci // checking to see if new replicas should be created. reconcileService returns // an error if there is some case prevent it from correctly reconciling the // service. -func (r *Reconciler) ReconcileService(id string) error { +func (r *Reconciler) ReconcileService(ctx context.Context, id string) error { var ( service *api.Service tasks []*api.Task @@ -235,8 +235,7 @@ func (r *Reconciler) ReconcileService(id string) error { return nil } - // TODO(dperny): pass in context from above - return r.restart.Restart(context.Background(), tx, cluster, service, *t) + return r.restart.Restart(ctx, tx, cluster, service, *t) }); err != nil { return err } diff --git a/manager/orchestrator/jobs/replicated/reconciler_test.go b/manager/orchestrator/jobs/replicated/reconciler_test.go index e4d5e03c16..9f3db18a4e 100644 --- a/manager/orchestrator/jobs/replicated/reconciler_test.go +++ b/manager/orchestrator/jobs/replicated/reconciler_test.go @@ -153,7 +153,7 @@ var _ = Describe("Replicated Job reconciler", func() { }) Expect(err).ToNot(HaveOccurred()) - err = r.ReconcileService(serviceID) + err = r.ReconcileService(context.Background(), serviceID) Expect(err).ToNot(HaveOccurred()) // verify there are maxConcurrent tasks @@ -177,7 +177,7 @@ var _ = Describe("Replicated Job reconciler", func() { return store.UpdateService(tx, service) }) Expect(err).ToNot(HaveOccurred()) - err = r.ReconcileService(serviceID) + err = r.ReconcileService(context.Background(), serviceID) Expect(err).ToNot(HaveOccurred()) // fetch the tasks before we get to the test case itself, @@ -231,7 +231,7 @@ var _ = Describe("Replicated Job reconciler", func() { }) Expect(err).ToNot(HaveOccurred()) - reconcileErr = r.ReconcileService(serviceID) + reconcileErr = r.ReconcileService(context.Background(), serviceID) }) When("the job has no tasks yet created", func() { @@ -516,7 +516,7 @@ var _ = Describe("Replicated Job reconciler", func() { }) Expect(err).ToNot(HaveOccurred()) - reconcileErr := r.ReconcileService("someService") + reconcileErr := r.ReconcileService(context.Background(), "someService") Expect(reconcileErr).To(HaveOccurred()) Expect(reconcileErr.Error()).To(ContainSubstring("underflow")) }) From 4b1b562078505c93764ae4dff384cb0074268802 Mon Sep 17 00:00:00 2001 From: Richard Davenport Date: Wed, 16 Sep 2026 10:03:11 -0500 Subject: [PATCH 4/4] manager/orchestrator/jobs: recover from a completion overshoot ReconcileService computes how many new tasks to create with unsigned subtraction: possibleNewTasks := rj.MaxConcurrent - runningTasks allowedNewTasks := rj.TotalCompletions - completeTasks - runningTasks If more tasks exist for the current job iteration than the service asked for, these underflow, and the guard below them returns an error rather than creating a preposterous number of tasks. Refusing to create the tasks is right; returning an error is not, because it abandons the rest of the reconciliation. The removal of tasks belonging to previous job iterations and the restarting of failed tasks both happen after that point, so a service that has overshot TotalCompletions once is never reconciled again -- the count cannot come back down, so every subsequent pass takes the same early return. Make both subtractions saturating and downgrade the guard to a warning. An overshot job now creates zero new tasks, which is what the arithmetic was trying to express, while the rest of the reconcile continues to run. Dropping the early return leaves the restart loop reachable for a job that has already met TotalCompletions, where it would restart surplus failed tasks and overshoot further. Clear the restart list in that case. This does not prevent the overshoot; it stops the overshoot from being permanent. The existing spec asserted the error. It now asserts the new behaviour, and that the reconcile really does continue: a task belonging to a previous job iteration is still marked for removal. A second spec covers the restart list being cleared once the job has met its goal. Signed-off-by: Richard Davenport --- .../jobs/replicated/reconciler.go | 49 +++++++---- .../jobs/replicated/reconciler_test.go | 84 +++++++++++++++++-- 2 files changed, 112 insertions(+), 21 deletions(-) diff --git a/manager/orchestrator/jobs/replicated/reconciler.go b/manager/orchestrator/jobs/replicated/reconciler.go index aad65f9b78..e50b4b80bd 100644 --- a/manager/orchestrator/jobs/replicated/reconciler.go +++ b/manager/orchestrator/jobs/replicated/reconciler.go @@ -2,9 +2,9 @@ package replicated import ( "context" - "fmt" "github.com/moby/swarmkit/v2/api" + "github.com/moby/swarmkit/v2/log" "github.com/moby/swarmkit/v2/manager/orchestrator" "github.com/moby/swarmkit/v2/manager/state/store" ) @@ -157,31 +157,39 @@ func (r *Reconciler) ReconcileService(ctx context.Context, id string) error { rj := service.Spec.GetReplicatedJob() // possibleNewTasks gives us the upper bound for how many tasks we'll - // create. also, ugh, subtracting uints. there's no way this can ever go - // wrong. - possibleNewTasks := rj.MaxConcurrent - runningTasks + // create. subtractions here are saturating: if more tasks exist than the + // service asks for, we want to create none, not to underflow. + possibleNewTasks := subOrZero(rj.MaxConcurrent, runningTasks) // allowedNewTasks is how many tasks we could create, if there were no // restriction on maximum concurrency. This is the total number of tasks // we want completed, minus the tasks that are already completed, minus // the tasks that are in progress. - // - // seriously, ugh, subtracting unsigned ints. totally a fine and not at all - // risky operation, with no possibility for catastrophe - allowedNewTasks := rj.TotalCompletions - completeTasks - runningTasks + allowedNewTasks := subOrZero(subOrZero(rj.TotalCompletions, completeTasks), runningTasks) // the lower number of allowedNewTasks and possibleNewTasks is how many we // can create. actualNewTasks := min(possibleNewTasks, allowedNewTasks) - // this check might seem odd, but it protects us from an underflow of the - // above subtractions, which, again, is a totally impossible thing that can - // never happen, ever, obviously. - if actualNewTasks > rj.TotalCompletions { - return fmt.Errorf( - "uint64 underflow, we're not going to create %v tasks", - actualNewTasks, - ) + // a job that has overshot its TotalCompletions is not something we can + // undo, but it is also not a reason to stop reconciling: the removal of + // tasks belonging to older job iterations, and the restarting of failed + // tasks, both still need to happen. log it and carry on creating zero new + // tasks. + if completeTasks+runningTasks > rj.TotalCompletions { + log.G(ctx).WithFields(log.Fields{ + "service.id": service.ID, + "job.iteration": jobVersion, + "tasks.complete": completeTasks, + "tasks.running": runningTasks, + "totalCompletions": rj.TotalCompletions, + }).Warn("replicated job has more tasks than TotalCompletions; creating no new tasks") + } + + if completeTasks >= rj.TotalCompletions { + // The job has already reached its goal. Restarting any failed tasks + // now would risk overshooting the desired number of completions. + restartTasks = nil } // finally, we can create these tasks. do this in a batch operation, to @@ -289,3 +297,12 @@ func (r *Reconciler) SlotTuple(t *api.Task) orchestrator.SlotTuple { Slot: t.Slot, } } + +// subOrZero subtracts b from a, returning 0 rather than underflowing when b is +// greater than a. +func subOrZero(a, b uint64) uint64 { + if b > a { + return 0 + } + return a - b +} diff --git a/manager/orchestrator/jobs/replicated/reconciler_test.go b/manager/orchestrator/jobs/replicated/reconciler_test.go index 9f3db18a4e..ba4b0ef9a5 100644 --- a/manager/orchestrator/jobs/replicated/reconciler_test.go +++ b/manager/orchestrator/jobs/replicated/reconciler_test.go @@ -481,12 +481,15 @@ var _ = Describe("Replicated Job reconciler", func() { }) }) - It("should return an underflow error if there are more running tasks than TotalCompletions", func() { + It("should create no tasks, and no error, if there are more running tasks than TotalCompletions", func() { // this is an error condition which should not happen in real life, // but i want to make sure that we can't accidentally start - // creating nearly the maximum 64-bit unsigned int number of tasks. + // creating nearly the maximum 64-bit unsigned int number of tasks, + // and that overshooting TotalCompletions does not wedge + // reconciliation for the service permanently. maxConcurrent := uint64(10) totalCompletions := uint64(20) + var oldTaskID string err := s.Update(func(tx store.Tx) error { service := &api.Service{ ID: "someService", @@ -498,6 +501,9 @@ var _ = Describe("Replicated Job reconciler", func() { }, }, }, + JobStatus: &api.JobStatus{ + JobIteration: api.Version{Index: 1}, + }, } if err := store.CreateService(tx, service); err != nil { return err @@ -505,20 +511,88 @@ var _ = Describe("Replicated Job reconciler", func() { for range totalCompletions + 10 { task := orchestrator.NewTask(nil, service, 0, "") + task.JobIteration = &api.Version{Index: 1} + task.DesiredState = api.TaskStateCompleted + + if err := store.CreateTask(tx, task); err != nil { + return err + } + } + + // a task left over from the previous iteration of the job. the + // reconcile pass is expected to mark it for removal, which + // happens after the point where an overshoot used to return + // early. + oldTask := orchestrator.NewTask(nil, service, 0, "") + oldTask.JobIteration = &api.Version{Index: 0} + oldTask.DesiredState = api.TaskStateCompleted + oldTaskID = oldTask.ID + return store.CreateTask(tx, oldTask) + }) + Expect(err).ToNot(HaveOccurred()) + + reconcileErr := r.ReconcileService(context.Background(), "someService") + Expect(reconcileErr).ToNot(HaveOccurred()) + + s.View(func(tx store.ReadTx) { + tasks, err := store.FindTasks(tx, store.ByServiceID("someService")) + Expect(err).ToNot(HaveOccurred()) + Expect(tasks).To(HaveLen(int(totalCompletions + 10 + 1))) + + // the rest of the reconcile pass still runs: the task from the + // previous iteration is marked for removal. + oldTask := store.GetTask(tx, oldTaskID) + Expect(oldTask).ToNot(BeNil()) + Expect(oldTask.DesiredState).To(Equal(api.TaskStateRemove)) + }) + }) + + It("should not restart failed tasks once TotalCompletions has been reached", func() { + // restarting a failed task after the job has met its goal would + // push the number of completions past what was asked for. + maxConcurrent := uint64(1) + totalCompletions := uint64(2) + err := s.Update(func(tx store.Tx) error { + service := &api.Service{ + ID: "someService", + Spec: api.ServiceSpec{ + Mode: &api.ServiceSpec_ReplicatedJob{ + ReplicatedJob: &api.ReplicatedJob{ + MaxConcurrent: maxConcurrent, + TotalCompletions: totalCompletions, + }, + }, + }, + } + if err := store.CreateService(tx, service); err != nil { + return err + } + + for i := range totalCompletions { + task := orchestrator.NewTask(nil, service, i, "") task.JobIteration = &api.Version{} task.DesiredState = api.TaskStateCompleted + task.Status.State = api.TaskStateCompleted if err := store.CreateTask(tx, task); err != nil { return err } } - return nil + + // a surplus task that failed. the job is already done, so it + // must not be handed to the restart supervisor. + failed := orchestrator.NewTask(nil, service, totalCompletions, "") + failed.JobIteration = &api.Version{} + failed.DesiredState = api.TaskStateCompleted + failed.Status.State = api.TaskStateFailed + return store.CreateTask(tx, failed) }) Expect(err).ToNot(HaveOccurred()) reconcileErr := r.ReconcileService(context.Background(), "someService") - Expect(reconcileErr).To(HaveOccurred()) - Expect(reconcileErr.Error()).To(ContainSubstring("underflow")) + Expect(reconcileErr).ToNot(HaveOccurred()) + + Expect(f.tasks).To(BeEmpty()) }) }) })