diff --git a/manager/dispatcher/dispatcher.go b/manager/dispatcher/dispatcher.go index 5bc230fec3..06a01e1c71 100644 --- a/manager/dispatcher/dispatcher.go +++ b/manager/dispatcher/dispatcher.go @@ -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. diff --git a/manager/dispatcher/dispatcher_test.go b/manager/dispatcher/dispatcher_test.go index 6f593ed792..673a64493a 100644 --- a/manager/dispatcher/dispatcher_test.go +++ b/manager/dispatcher/dispatcher_test.go @@ -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()