diff --git a/pkg/apis/compute/pod.go b/pkg/apis/compute/pod.go index 65346ef99a..3b6ae25ea7 100644 --- a/pkg/apis/compute/pod.go +++ b/pkg/apis/compute/pod.go @@ -17,6 +17,10 @@ package compute const ( POD_STATUS_CREATING_CONTAINER = "creating_container" POD_STATUS_CREATE_CONTAINER_FAILED = "create_container_failed" + POD_STATUS_STARTING_CONTAINER = "starting_container" + POD_STATUS_START_CONTAINER_FAILED = "start_container_failed" + POD_STATUS_STOPPING_CONTAINER = "stopping_container" + POD_STATUS_STOP_CONTAINER_FAILED = "stop_container_failed" POD_STATUS_DELETING_CONTAINER = "deleting_container" POD_STATUS_DELETE_CONTAINER_FAILED = "delete_container_failed" ) diff --git a/pkg/compute/guestdrivers/pod.go b/pkg/compute/guestdrivers/pod.go index 3e65dbc23a..fe19a50828 100644 --- a/pkg/compute/guestdrivers/pod.go +++ b/pkg/compute/guestdrivers/pod.go @@ -235,11 +235,19 @@ func (p *SPodDriver) RequestGuestHotAddIso(ctx context.Context, guest *models.SG return task.ScheduleRun(nil) } +func (p *SPodDriver) PerformStart(ctx context.Context, userCred mcclient.TokenCredential, guest *models.SGuest, data *jsonutils.JSONDict) error { + task, err := taskman.TaskManager.NewTask(ctx, "PodStartTask", guest, userCred, nil, "", "", nil) + if err != nil { + return errors.Wrap(err, "New PodStartTask") + } + return task.ScheduleRun(nil) +} + func (p *SPodDriver) RequestStartOnHost(ctx context.Context, guest *models.SGuest, host *models.SHost, userCred mcclient.TokenCredential, task taskman.ITask) error { header := p.getTaskRequestHeader(task) config := jsonutils.NewDict() - desc, err := guest.GetDriver().GetJsonDescAtHost(ctx, userCred, guest, host, nil) + desc, err := guest.GetDriver().GetJsonDescAtHost(ctx, task.GetUserCred(), guest, host, nil) if err != nil { return errors.Wrapf(err, "GetJsonDescAtHost") } @@ -281,6 +289,14 @@ func (p *SPodDriver) OnGuestDeployTaskDataReceived(ctx context.Context, guest *m return nil } +func (p *SPodDriver) StartGuestStopTask(guest *models.SGuest, ctx context.Context, userCred mcclient.TokenCredential, params *jsonutils.JSONDict, parentTaskId string) error { + task, err := taskman.TaskManager.NewTask(ctx, "PodStopTask", guest, userCred, nil, parentTaskId, "", nil) + if err != nil { + return errors.Wrap(err, "New PodStopTask") + } + return task.ScheduleRun(nil) +} + func (p *SPodDriver) RequestUndeployGuestOnHost(ctx context.Context, guest *models.SGuest, host *models.SHost, task taskman.ITask) error { task, err := taskman.TaskManager.NewTask(ctx, "PodDeleteTask", guest, task.GetUserCred(), nil, task.GetTaskId(), "", nil) if err != nil { diff --git a/pkg/compute/tasks/pod_start_task.go b/pkg/compute/tasks/pod_start_task.go new file mode 100644 index 0000000000..c066094ff3 --- /dev/null +++ b/pkg/compute/tasks/pod_start_task.go @@ -0,0 +1,77 @@ +// Copyright 2019 Yunion +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package tasks + +import ( + "context" + + "yunion.io/x/jsonutils" + "yunion.io/x/pkg/errors" + + api "yunion.io/x/onecloud/pkg/apis/compute" + "yunion.io/x/onecloud/pkg/cloudcommon/db" + "yunion.io/x/onecloud/pkg/cloudcommon/db/taskman" + "yunion.io/x/onecloud/pkg/compute/models" +) + +type PodStartTask struct { + SGuestBaseTask +} + +func init() { + taskman.RegisterTask(PodStartTask{}) +} + +func (t *PodStartTask) OnInit(ctx context.Context, obj db.IStandaloneModel, body jsonutils.JSONObject) { + t.SetStage("OnPodStarted", nil) + pod := obj.(*models.SGuest) + pod.StartGueststartTask(ctx, t.GetUserCred(), jsonutils.NewDict(), t.GetTaskId()) +} + +func (t *PodStartTask) OnPodStarted(ctx context.Context, pod *models.SGuest, _ jsonutils.JSONObject) { + t.SetStage("OnContainerStarted", nil) + pod.SetStatus(ctx, t.GetUserCred(), api.POD_STATUS_STARTING_CONTAINER, "") + ctrs, err := models.GetContainerManager().GetContainersByPod(pod.GetId()) + if err != nil { + t.OnContainerStartedFailed(ctx, pod, jsonutils.NewString(errors.Wrap(err, "GetContainersByPod").Error())) + return + } + isAllStarted := true + for i := range ctrs { + if ctrs[i].GetStatus() != api.CONTAINER_STATUS_RUNNING { + isAllStarted = false + ctrs[i].StartStartTask(ctx, t.GetUserCred(), t.GetTaskId()) + } + } + if isAllStarted { + t.OnContainerStarted(ctx, pod, nil) + return + } +} + +func (t *PodStartTask) OnPodStartedFailed(ctx context.Context, pod *models.SGuest, reason jsonutils.JSONObject) { + pod.SetStatus(ctx, t.GetUserCred(), api.VM_START_FAILED, reason.String()) + t.SetStageFailed(ctx, reason) +} + +func (t *PodStartTask) OnContainerStarted(ctx context.Context, pod *models.SGuest, data jsonutils.JSONObject) { + pod.SetStatus(ctx, t.GetUserCred(), api.VM_RUNNING, "") + t.SetStageComplete(ctx, nil) +} + +func (t *PodStartTask) OnContainerStartedFailed(ctx context.Context, pod *models.SGuest, data jsonutils.JSONObject) { + pod.SetStatus(ctx, t.GetUserCred(), api.POD_STATUS_START_CONTAINER_FAILED, data.String()) + t.SetStageFailed(ctx, data) +} diff --git a/pkg/compute/tasks/pod_stop_task.go b/pkg/compute/tasks/pod_stop_task.go new file mode 100644 index 0000000000..a116071a1b --- /dev/null +++ b/pkg/compute/tasks/pod_stop_task.go @@ -0,0 +1,88 @@ +// Copyright 2019 Yunion +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package tasks + +import ( + "context" + + "yunion.io/x/jsonutils" + "yunion.io/x/log" + "yunion.io/x/pkg/errors" + + api "yunion.io/x/onecloud/pkg/apis/compute" + "yunion.io/x/onecloud/pkg/cloudcommon/db" + "yunion.io/x/onecloud/pkg/cloudcommon/db/taskman" + "yunion.io/x/onecloud/pkg/compute/models" +) + +type PodStopTask struct { + SGuestBaseTask +} + +func init() { + taskman.RegisterTask(PodStopTask{}) +} + +func (t *PodStopTask) OnInit(ctx context.Context, obj db.IStandaloneModel, body jsonutils.JSONObject) { + t.SetStage("OnWaitContainerStopped", nil) + t.OnWaitContainerStopped(ctx, obj.(*models.SGuest), nil) +} + +func (t *PodStopTask) OnWaitContainerStopped(ctx context.Context, pod *models.SGuest, _ jsonutils.JSONObject) { + pod.SetStatus(ctx, t.GetUserCred(), api.POD_STATUS_STOPPING_CONTAINER, "") + ctrs, err := models.GetContainerManager().GetContainersByPod(pod.GetId()) + if err != nil { + t.OnWaitContainerStoppedFailed(ctx, pod, jsonutils.NewString(errors.Wrap(err, "GetContainersByPod").Error())) + return + } + isAllStopped := true + for i := range ctrs { + curCtr := ctrs[i] + log.Infof("========container status: %s", curCtr.GetStatus()) + if curCtr.GetStatus() != api.CONTAINER_STATUS_EXITED { + isAllStopped = false + curCtr.StartStopTask(ctx, t.GetUserCred(), &api.ContainerStopInput{Timeout: 15}, t.GetTaskId()) + } + } + if isAllStopped { + t.OnContainerStopped(ctx, pod) + return + } +} + +func (t *PodStopTask) OnWaitContainerStoppedFailed(ctx context.Context, pod *models.SGuest, data jsonutils.JSONObject) { + pod.SetStatus(ctx, t.GetUserCred(), api.POD_STATUS_STOP_CONTAINER_FAILED, data.String()) + t.SetStageFailed(ctx, data) +} + +func (t *PodStopTask) OnContainerStopped(ctx context.Context, pod *models.SGuest) { + t.SetStage("OnPodStopped", nil) + task, err := taskman.TaskManager.NewTask(ctx, "GuestStopTask", pod, t.GetUserCred(), nil, t.GetTaskId(), "", nil) + if err != nil { + t.OnPodStoppedFailed(ctx, pod, jsonutils.NewString(err.Error())) + return + } + task.ScheduleRun(nil) +} + +func (t *PodStopTask) OnPodStopped(ctx context.Context, pod *models.SGuest, data jsonutils.JSONObject) { + pod.SetStatus(ctx, t.GetUserCred(), api.VM_READY, "") + t.SetStageComplete(ctx, nil) +} + +func (t *PodStopTask) OnPodStoppedFailed(ctx context.Context, pod *models.SGuest, reason jsonutils.JSONObject) { + pod.SetStatus(ctx, t.GetUserCred(), api.VM_STOP_FAILED, reason.String()) + t.SetStageFailed(ctx, reason) +} diff --git a/pkg/hostman/guestman/pod.go b/pkg/hostman/guestman/pod.go index 9b31c7bf92..c238c27686 100644 --- a/pkg/hostman/guestman/pod.go +++ b/pkg/hostman/guestman/pod.go @@ -20,6 +20,7 @@ import ( "io/ioutil" "path" "path/filepath" + "strings" "time" runtimeapi "k8s.io/cri-api/pkg/apis/runtime/v1" @@ -508,12 +509,17 @@ func (s *sPodGuestInstance) StartContainer(ctx context.Context, userCred mcclien if hasCtr { status, err := s.getContainerStatus(ctx, ctrId) if err != nil { - return nil, errors.Wrap(err, "get container status") - } - if status == computeapi.CONTAINER_STATUS_EXITED { - needRecreate = true - } else if status != computeapi.CONTAINER_STATUS_CREATED { - return nil, errors.Wrapf(err, "can't start container when status is %s", status) + if strings.Contains(err.Error(), "not found") { + needRecreate = true + } else { + return nil, errors.Wrap(err, "get container status") + } + } else { + if status == computeapi.CONTAINER_STATUS_EXITED { + needRecreate = true + } else if status != computeapi.CONTAINER_STATUS_CREATED { + return nil, errors.Wrapf(err, "can't start container when status is %s", status) + } } } if !hasCtr || needRecreate { @@ -899,7 +905,7 @@ func (s *sPodGuestInstance) DeleteContainer(ctx context.Context, userCred mcclie return nil, errors.Wrap(err, "getContainerCRIId") } if criId != "" { - if err := s.getCRI().RemoveContainer(ctx, criId); err != nil { + if err := s.getCRI().RemoveContainer(ctx, criId); err != nil && !strings.Contains(err.Error(), "not found") { return nil, errors.Wrap(err, "cri.RemoveContainer") } }