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
1 change: 1 addition & 0 deletions manager/dispatcher/dispatcher.go
Original file line number Diff line number Diff line change
Expand Up @@ -847,6 +847,7 @@ func (d *Dispatcher) processUpdates(ctx context.Context) {
volume := store.GetVolume(tx, volumeID)
if volume == nil {
logger.Error("volume unavailable")
return nil
}

// buckle your seatbelts, we're going quadratic.
Expand Down
93 changes: 93 additions & 0 deletions manager/dispatcher/dispatcher_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -1554,6 +1554,99 @@ func TestTaskUpdate(t *testing.T) {

}

// TestVolumeUnpublishedDeletedVolume tests that if a volume is removed from the
// store after a node has reported it as unpublished, but before the dispatcher
// has processed that report (as happens when the volume is force-removed while
// the node is still unpublishing it), the dispatcher does not crash, and
// updates to other volumes are still processed.
func TestVolumeUnpublishedDeletedVolume(t *testing.T) {
// the dispatcher is deliberately not started, because the background loop
// in Run could process the unpublish reports before the volume is removed.
// instead, processUpdates is called explicitly below.
s := store.NewMemoryStore(nil)
defer s.Close()
d := New()
d.store = s

nodeID := "nodeID"

// both volumes are no longer in use on the node, so the scheduler has
// moved them to PENDING_NODE_UNPUBLISH, and the dispatcher has told the
// node to unpublish them.
volumes := []*api.Volume{
{
ID: "volumeID0",
Spec: api.VolumeSpec{
Annotations: api.Annotations{
Name: "volumeName0",
},
Driver: &api.Driver{
Name: "someDriver",
},
},
PublishStatus: []*api.VolumePublishStatus{
{
NodeID: nodeID,
State: api.VolumePublishStatus_PENDING_NODE_UNPUBLISH,
},
},
}, {
ID: "volumeID1",
Spec: api.VolumeSpec{
Annotations: api.Annotations{
Name: "volumeName1",
},
Driver: &api.Driver{
Name: "someDriver",
},
},
PublishStatus: []*api.VolumePublishStatus{
{
NodeID: nodeID,
State: api.VolumePublishStatus_PENDING_NODE_UNPUBLISH,
},
},
},
}
err := s.Update(func(tx store.Tx) error {
for _, v := range volumes {
if err := store.CreateVolume(tx, v); err != nil {
return err
}
}
return nil
})
require.NoError(t, err)

// the node reports that it has unpublished both volumes. this is the
// state UpdateVolumeStatus leaves behind.
d.unpublishedVolumes = map[string][]string{
volumes[0].ID: {nodeID},
volumes[1].ID: {nodeID},
}

// before the dispatcher processes the report, the first volume is
// force-removed, which deletes it from the store regardless of its
// PublishStatus (see (*controlapi.Server).RemoveVolume).
err = s.Update(func(tx store.Tx) error {
return store.DeleteVolume(tx, volumes[0].ID)
})
require.NoError(t, err)

require.NotPanics(t, func() {
d.processUpdates(context.Background())
})

s.View(func(readTx store.ReadTx) {
assert.Nil(t, store.GetVolume(readTx, volumes[0].ID))

storeVolume := store.GetVolume(readTx, volumes[1].ID)
require.NotNil(t, storeVolume)
require.Len(t, storeVolume.PublishStatus, 1)
assert.Equal(t, api.VolumePublishStatus_PENDING_UNPUBLISH, storeVolume.PublishStatus[0].State)
})
}

func TestTaskUpdateNoCert(t *testing.T) {
gd := startDispatcher(t, DefaultConfig())
defer gd.Close()
Expand Down