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
47 changes: 47 additions & 0 deletions manager/scheduler/revision_recovery_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,47 @@
package scheduler

import (
"context"
"testing"

"github.com/moby/swarmkit/v2/api"
"github.com/moby/swarmkit/v2/manager/state/store"
"github.com/stretchr/testify/require"
)

func TestPendingOldRevisionRecoversWhenNodeReturns(t *testing.T) {
ctx := context.Background()
state := store.NewMemoryStore(nil)
defer state.Close()
service := &api.Service{ID: "service", SpecVersion: &api.Version{Index: 2}}
task := &api.Task{
ID: "pending", ServiceID: service.ID, SpecVersion: &api.Version{Index: 1},
DesiredState: api.TaskStateRunning, Status: api.TaskStatus{State: api.TaskStatePending},
Spec: api.TaskSpec{Runtime: &api.TaskSpec_Container{Container: &api.ContainerSpec{Image: "busybox:1.37"}}},
}
// A service-level metadata update changes the revision, not the task spec.
service.Spec.Task = task.Spec
node := &api.Node{
ID: "worker", Spec: api.NodeSpec{Availability: api.NodeAvailabilityActive},
Status: api.NodeStatus{State: api.NodeStatus_DOWN}, Description: &api.NodeDescription{},
}
require.NoError(t, state.Update(func(tx store.Tx) error {
require.NoError(t, store.CreateService(tx, service))
require.NoError(t, store.CreateTask(tx, task))
return store.CreateNode(tx, node)
}))
scheduler := New(state)
state.View(func(tx store.ReadTx) { require.NoError(t, scheduler.setupTasksList(tx)) })
scheduler.tick(ctx)
require.Contains(t, scheduler.unassignedTasks, task.ID,
"an older spec revision must not discard a task that is still desired running")
node.Status.State = api.NodeStatus_READY
require.NoError(t, state.Update(func(tx store.Tx) error { return store.UpdateNode(tx, node) }))
scheduler.createOrUpdateNode(node)
scheduler.tick(ctx)
state.View(func(tx store.ReadTx) {
recovered := store.GetTask(tx, task.ID)
require.Equal(t, api.TaskStateAssigned, recovered.Status.State)
require.Equal(t, node.ID, recovered.NodeID)
})
}
4 changes: 3 additions & 1 deletion manager/scheduler/scheduler.go
Original file line number Diff line number Diff line change
Expand Up @@ -943,7 +943,9 @@ func (s *Scheduler) noSuitableNode(ctx context.Context, taskGroup map[string]*ap
newT.Status.Timestamp = ptypes.MustTimestampProto(time.Now())
sv := service.SpecVersion
tv := newT.SpecVersion
if sv != nil && tv != nil && sv.Index > tv.Index {
// A metadata-only service update can leave a valid task on an older
// revision. Keep retrying it until the orchestrator requests shutdown.
if sv != nil && tv != nil && sv.Index > tv.Index && t.DesiredState >= api.TaskStateShutdown {
log.G(ctx).WithField("task.id", t.ID).Debug(
"task belongs to old revision of service",
)
Expand Down