From 780d55f9460a72c1ba27ef658bbf60a50df38ef5 Mon Sep 17 00:00:00 2001 From: Zexi Li Date: Fri, 30 May 2025 17:48:35 +0800 Subject: [PATCH] refactor(region,host): container status management (#22634) --- pkg/apis/compute/guests.go | 5 + pkg/apis/compute/host.go | 6 +- pkg/compute/models/guest_actions.go | 21 ++- .../baremetal_server_sync_status_task.go | 11 +- .../tasks/guest/guest_syncstatus_task.go | 2 +- pkg/compute/tasks/guest/pod_stop_task.go | 5 + .../container/status/status_manager.go | 39 +++++ pkg/hostman/guestman/guestman.go | 10 +- pkg/hostman/guestman/pod.go | 43 ++++- pkg/hostman/guestman/pod/statusman/doc.go | 1 + .../guestman/pod/statusman/pod_state.go | 157 ++++++++++++++++++ pkg/hostman/guestman/qemu-kvm.go | 8 +- pkg/hostman/guestman/runtime.go | 6 +- pkg/hostman/hostutils/hostutils.go | 6 +- 14 files changed, 288 insertions(+), 32 deletions(-) create mode 100644 pkg/hostman/guestman/pod/statusman/doc.go create mode 100644 pkg/hostman/guestman/pod/statusman/pod_state.go diff --git a/pkg/apis/compute/guests.go b/pkg/apis/compute/guests.go index 2bb0c430f3..9142bfbda2 100644 --- a/pkg/apis/compute/guests.go +++ b/pkg/apis/compute/guests.go @@ -1427,3 +1427,8 @@ type ServerChangeBillingTypeInput struct { // required: true BillingType string `json:"billing_type"` } + +type ServerPerformStatusInput struct { + apis.PerformStatusInput + Containers map[string]*ContainerPerformStatusInput `json:"containers"` +} diff --git a/pkg/apis/compute/host.go b/pkg/apis/compute/host.go index 9f6d636488..c0cceed22d 100644 --- a/pkg/apis/compute/host.go +++ b/pkg/apis/compute/host.go @@ -714,13 +714,13 @@ type HostUploadGuestsStatusRequest struct { GuestIds []string `json:"guest_ids"` } -type HostUploadGuestStatusResponse struct { +type HostUploadGuestStatusInput struct { apis.PerformStatusInput Containers map[string]*ContainerPerformStatusInput `json:"containers"` } -type HostUploadGuestsStatusResponse struct { - Guests map[string]*HostUploadGuestStatusResponse `json:"guests"` +type HostUploadGuestsStatusInput struct { + Guests map[string]*HostUploadGuestStatusInput `json:"guests"` } type GuestUploadContainerStatusResponse struct { diff --git a/pkg/compute/models/guest_actions.go b/pkg/compute/models/guest_actions.go index 290e19e72d..4f640a4539 100644 --- a/pkg/compute/models/guest_actions.go +++ b/pkg/compute/models/guest_actions.go @@ -3466,7 +3466,7 @@ func (self *SGuest) SetBackupGuestStatus(userCred mcclient.TokenCredential, stat return nil } -func (g *SGuest) SetStatusFromHost(ctx context.Context, userCred mcclient.TokenCredential, resp api.HostUploadGuestStatusResponse, hasParentTask bool, originStatus string) error { +func (g *SGuest) SetStatusFromHost(ctx context.Context, userCred mcclient.TokenCredential, resp api.HostUploadGuestStatusInput, hasParentTask bool, originStatus string) error { statusStr := resp.Status switch statusStr { case cloudprovider.CloudVMStatusRunning: @@ -3490,7 +3490,9 @@ func (g *SGuest) SetStatusFromHost(ctx context.Context, userCred mcclient.TokenC statusStr = originStatus } } - input := resp.PerformStatusInput + input := api.ServerPerformStatusInput{ + PerformStatusInput: resp.PerformStatusInput, + } input.Status = statusStr if _, err := g.PerformStatus(ctx, userCred, nil, input); err != nil { return errors.Wrapf(err, "perform status of %s", jsonutils.Marshal(resp)) @@ -3498,7 +3500,7 @@ func (g *SGuest) SetStatusFromHost(ctx context.Context, userCred mcclient.TokenC return nil } -func (m *SGuestManager) PerformUploadStatus(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, input *api.HostUploadGuestsStatusResponse) (*api.GuestUploadStatusesResponse, error) { +func (m *SGuestManager) PerformUploadStatus(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, input *api.HostUploadGuestsStatusInput) (*api.GuestUploadStatusesResponse, error) { out := &api.GuestUploadStatusesResponse{ Guests: make(map[string]*api.GuestUploadStatusResponse), } @@ -3545,7 +3547,7 @@ func (m *SGuestManager) PerformUploadStatus(ctx context.Context, userCred mcclie } // 同步状态 -func (self *SGuest) PerformStatus(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, input apis.PerformStatusInput) (jsonutils.JSONObject, error) { +func (self *SGuest) PerformStatus(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, input api.ServerPerformStatusInput) (jsonutils.JSONObject, error) { if input.HostId != "" && self.BackupHostId != "" && input.HostId == self.BackupHostId { // perform status called from slave guest return nil, self.SetBackupGuestStatus(userCred, input.Status, input.Reason) @@ -3563,7 +3565,7 @@ func (self *SGuest) PerformStatus(ctx context.Context, userCred mcclient.TokenCr } preStatus := self.Status - _, err := self.SVirtualResourceBase.PerformStatus(ctx, userCred, query, input) + _, err := self.SVirtualResourceBase.PerformStatus(ctx, userCred, query, input.PerformStatusInput) if err != nil { return nil, errors.Wrap(err, "SVirtualResourceBase.PerformStatus") } @@ -3604,6 +3606,15 @@ func (self *SGuest) PerformStatus(ctx context.Context, userCred mcclient.TokenCr return nil, err } } + for cId, cStatus := range input.Containers { + ctr, err := GetContainerManager().FetchById(cId) + if err != nil { + return nil, errors.Wrapf(err, "GetContainerManager(%s)", cId) + } + if _, err := ctr.(*SContainer).PerformStatus(ctx, userCred, query, *cStatus); err != nil { + return nil, errors.Wrapf(err, "PerformStatus(%s) of container", cId) + } + } return nil, nil } diff --git a/pkg/compute/tasks/guest/baremetal_server_sync_status_task.go b/pkg/compute/tasks/guest/baremetal_server_sync_status_task.go index cd3a7c4971..4fe84a3c53 100644 --- a/pkg/compute/tasks/guest/baremetal_server_sync_status_task.go +++ b/pkg/compute/tasks/guest/baremetal_server_sync_status_task.go @@ -88,8 +88,9 @@ func (self *BaremetalServerSyncStatusTask) OnGuestStatusTaskComplete(ctx context } func (self *BaremetalServerSyncStatusTask) OnGuestStatusTaskCompleteFailed(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { - input := apis.PerformStatusInput{ - Status: api.VM_UNKNOWN, + input := api.ServerPerformStatusInput{ + PerformStatusInput: apis.PerformStatusInput{ + Status: api.VM_UNKNOWN}, } guest.PerformStatus(ctx, self.UserCred, nil, input) } @@ -97,8 +98,10 @@ func (self *BaremetalServerSyncStatusTask) OnGuestStatusTaskCompleteFailed(ctx c func (self *BaremetalServerSyncStatusTask) OnGetStatusFail(ctx context.Context, guest *models.SGuest) { kwargs := jsonutils.NewDict() kwargs.Set("status", jsonutils.NewString(api.VM_UNKNOWN)) - input := apis.PerformStatusInput{ - Status: api.VM_UNKNOWN, + input := api.ServerPerformStatusInput{ + PerformStatusInput: apis.PerformStatusInput{ + Status: api.VM_UNKNOWN, + }, } guest.PerformStatus(ctx, self.UserCred, nil, input) self.SetStageComplete(ctx, nil) diff --git a/pkg/compute/tasks/guest/guest_syncstatus_task.go b/pkg/compute/tasks/guest/guest_syncstatus_task.go index 041a50bada..ed50b2c586 100644 --- a/pkg/compute/tasks/guest/guest_syncstatus_task.go +++ b/pkg/compute/tasks/guest/guest_syncstatus_task.go @@ -70,7 +70,7 @@ func (self *GuestSyncstatusTask) getOriginStatus() string { func (self *GuestSyncstatusTask) OnGetStatusComplete(ctx context.Context, obj db.IStandaloneModel, body jsonutils.JSONObject) { guest := obj.(*models.SGuest) log.Debugf("OnGetStatusSucc guest %s(%s) status %s", guest.Name, guest.Id, body) - resp := new(api.HostUploadGuestStatusResponse) + resp := new(api.HostUploadGuestStatusInput) body.Unmarshal(resp) if err := guest.SetStatusFromHost(ctx, self.GetUserCred(), *resp, self.HasParentTask(), self.getOriginStatus()); err != nil { log.Warningf("SetStatusFromHost for guest %s error: %v", guest.GetId(), err) diff --git a/pkg/compute/tasks/guest/pod_stop_task.go b/pkg/compute/tasks/guest/pod_stop_task.go index 4244495e82..79593e1dc9 100644 --- a/pkg/compute/tasks/guest/pod_stop_task.go +++ b/pkg/compute/tasks/guest/pod_stop_task.go @@ -81,6 +81,11 @@ func (t *PodStopTask) OnContainerStopped(ctx context.Context, pod *models.SGuest task.ScheduleRun(nil) } +func (t *PodStopTask) OnContainerStoppedFailed(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) OnPodStopped(ctx context.Context, pod *models.SGuest, data jsonutils.JSONObject) { pod.SetStatus(ctx, t.GetUserCred(), api.VM_READY, "") t.SetStageComplete(ctx, nil) diff --git a/pkg/hostman/container/status/status_manager.go b/pkg/hostman/container/status/status_manager.go index 27506049f6..9b221bba4f 100644 --- a/pkg/hostman/container/status/status_manager.go +++ b/pkg/hostman/container/status/status_manager.go @@ -25,6 +25,7 @@ import ( "yunion.io/x/onecloud/pkg/apis" computeapi "yunion.io/x/onecloud/pkg/apis/compute" "yunion.io/x/onecloud/pkg/hostman/container/prober/results" + "yunion.io/x/onecloud/pkg/hostman/guestman/pod/statusman" "yunion.io/x/onecloud/pkg/hostman/hostutils" ) @@ -49,6 +50,44 @@ func (m *manager) SetContainerStartup(podId string, containerId string, started status = computeapi.CONTAINER_STATUS_NET_FAILED } } + + input := &statusman.PodStatusUpdateRequest{ + Id: podId, + Pod: pod.(statusman.IPod), + Status: computeapi.VM_RUNNING, + Reason: result.Reason, + ContainerStatuses: map[string]*statusman.ContainerStatus{ + containerId: {Status: status}, + }, + } + + if err := statusman.GetManager().UpdateStatus(input); err != nil { + err = errors.Wrapf(err, "set container(%s/%s) status failed, input: %s", podId, containerId, jsonutils.Marshal(input)) + log.Warningf(err.Error()) + errMsg := []string{ + "can't set container status", + } + for _, msg := range errMsg { + if strings.Contains(err.Error(), msg) { + return nil + } + } + return errors.Wrap(err, "update container status") + } else { + log.Infof("set container(%s/%s) status to %s", podId, containerId, jsonutils.Marshal(input).String()) + } + return nil +} + +func (m *manager) SetContainerStartupOld(podId string, containerId string, started bool, result results.ProbeResult, pod results.IPod) error { + status := computeapi.CONTAINER_STATUS_PROBE_FAILED + if started { + status = computeapi.CONTAINER_STATUS_RUNNING + } else { + if result.IsNetFailedError() && pod.IsRunning() { + status = computeapi.CONTAINER_STATUS_NET_FAILED + } + } input := &computeapi.ContainerPerformStatusInput{ PerformStatusInput: apis.PerformStatusInput{ Status: status, diff --git a/pkg/hostman/guestman/guestman.go b/pkg/hostman/guestman/guestman.go index 9d30566976..470208883c 100644 --- a/pkg/hostman/guestman/guestman.go +++ b/pkg/hostman/guestman/guestman.go @@ -47,6 +47,7 @@ import ( fwdpb "yunion.io/x/onecloud/pkg/hostman/guestman/forwarder/api" "yunion.io/x/onecloud/pkg/hostman/guestman/pod/pleg" "yunion.io/x/onecloud/pkg/hostman/guestman/pod/runtime" + "yunion.io/x/onecloud/pkg/hostman/guestman/pod/statusman" "yunion.io/x/onecloud/pkg/hostman/guestman/types" deployapi "yunion.io/x/onecloud/pkg/hostman/hostdeployer/apis" "yunion.io/x/onecloud/pkg/hostman/hostinfo/hostconsts" @@ -144,6 +145,7 @@ func NewGuestManager(host hostutils.IHost, serversPath string, workerCnt int) (* return nil, errors.Wrap(err, "mkdir qemu log dir") } if manager.host.IsContainerHost() { + statusman.GetManager().Start() manager.startContainerProbeManager() runtimeMan, err := runtime.NewRuntimeManager(manager.GetCRI()) if err != nil { @@ -1783,8 +1785,8 @@ func (m *SGuestManager) GetGuestTrafficRecord(sid string) (map[string]compute.SN func (m *SGuestManager) UploadGuestsStatus(ctx context.Context, i interface{}) (jsonutils.JSONObject, error) { input := i.(*compute.HostUploadGuestsStatusRequest) errs := []error{} - resp := &compute.HostUploadGuestsStatusResponse{ - Guests: make(map[string]*compute.HostUploadGuestStatusResponse, 0), + resp := &compute.HostUploadGuestsStatusInput{ + Guests: make(map[string]*compute.HostUploadGuestStatusInput, 0), } reason := "upload guest status by host" for _, id := range input.GuestIds { @@ -1814,10 +1816,10 @@ func (m *SGuestManager) UploadGuestsStatus(ctx context.Context, i interface{}) ( return ret, err } -func (m *SGuestManager) ProbeGuestInitStatus(sid string) *compute.HostUploadGuestStatusResponse { +func (m *SGuestManager) ProbeGuestInitStatus(sid string) *compute.HostUploadGuestStatusInput { guest, _ := m.GetServer(sid) status := m.getStatus(sid) - resp := &compute.HostUploadGuestStatusResponse{ + resp := &compute.HostUploadGuestStatusInput{ PerformStatusInput: apis.PerformStatusInput{ Status: status, BlockJobsCount: -1, diff --git a/pkg/hostman/guestman/pod.go b/pkg/hostman/guestman/pod.go index 6aa4f50347..17e41d60da 100644 --- a/pkg/hostman/guestman/pod.go +++ b/pkg/hostman/guestman/pod.go @@ -50,6 +50,7 @@ import ( _ "yunion.io/x/onecloud/pkg/hostman/container/volume_mount/disk" "yunion.io/x/onecloud/pkg/hostman/guestman/desc" "yunion.io/x/onecloud/pkg/hostman/guestman/pod/runtime" + "yunion.io/x/onecloud/pkg/hostman/guestman/pod/statusman" deployapi "yunion.io/x/onecloud/pkg/hostman/hostdeployer/apis" "yunion.io/x/onecloud/pkg/hostman/hostinfo" "yunion.io/x/onecloud/pkg/hostman/hostutils" @@ -351,7 +352,7 @@ func (s *sPodGuestInstance) getStatus(ctx context.Context, defaultStatus string) return status } -func (s *sPodGuestInstance) GetUploadStatus(ctx context.Context, reason string) (*computeapi.HostUploadGuestStatusResponse, error) { +func (s *sPodGuestInstance) GetUploadStatus(ctx context.Context, reason string) (*computeapi.HostUploadGuestStatusInput, error) { // sync pod status var status = computeapi.VM_READY if s.IsRunning() { @@ -421,14 +422,14 @@ func (s *sPodGuestInstance) GetUploadStatus(ctx context.Context, reason string) HostId: hostinfo.Instance().HostId, } - return &computeapi.HostUploadGuestStatusResponse{ + return &computeapi.HostUploadGuestStatusInput{ PerformStatusInput: *statusInput, Containers: cStatuss, }, nil } func (s *sPodGuestInstance) UploadStatus(ctx context.Context, reason string) error { - resp, err := s.GetUploadStatus(ctx, reason) + /*resp, err := s.GetUploadStatus(ctx, reason) if err != nil { return errors.Wrapf(err, "get upload status of pod: %s", reason) } @@ -444,11 +445,35 @@ func (s *sPodGuestInstance) UploadStatus(ctx context.Context, reason string) err if _, err := hostutils.UpdateServerStatus(ctx, s.Id, &resp.PerformStatusInput); err != nil { errs = append(errs, errors.Wrapf(err, "failed update guest status")) + }*/ + // return errors.NewAggregate(errs) + + resp, err := s.GetUploadStatus(ctx, reason) + if err != nil { + return errors.Wrapf(err, "get upload status of pod: %s", reason) } - return errors.NewAggregate(errs) + containerStatuses := make(map[string]*statusman.ContainerStatus) + for ctrId, cStatus := range resp.Containers { + containerStatuses[ctrId] = &statusman.ContainerStatus{ + Status: cStatus.Status, + RestartCount: cStatus.RestartCount, + StartedAt: cStatus.StartedAt, + LastFinishedAt: cStatus.LastFinishedAt, + } + } + if err := statusman.GetManager().UpdateStatus(&statusman.PodStatusUpdateRequest{ + Id: s.Id, + Pod: s, + Status: resp.Status, + ContainerStatuses: containerStatuses, + Reason: reason, + }); err != nil { + return errors.Wrapf(err, "update status of pod: %s", reason) + } + return nil } -func (s *sPodGuestInstance) PostUploadStatus(resp *computeapi.HostUploadGuestStatusResponse, reason string) { +func (s *sPodGuestInstance) PostUploadStatus(resp *computeapi.HostUploadGuestStatusInput, reason string) { for ctrId, cStatus := range resp.Containers { s.markContainerProbeDirty(cStatus.Status, ctrId, reason) } @@ -534,11 +559,11 @@ func (s *sPodGuestInstance) IsContainerRunning(ctx context.Context, ctrId string return false, nil } -func (s *sPodGuestInstance) probeGuestStatus(ctx context.Context, resp *computeapi.HostUploadGuestStatusResponse) { +func (s *sPodGuestInstance) probeGuestStatus(ctx context.Context, resp *computeapi.HostUploadGuestStatusInput) { resp.Status = s.getStatus(ctx, resp.Status) } -func (s *sPodGuestInstance) HandleGuestStatus(ctx context.Context, resp *computeapi.HostUploadGuestStatusResponse) (jsonutils.JSONObject, error) { +func (s *sPodGuestInstance) HandleGuestStatus(ctx context.Context, resp *computeapi.HostUploadGuestStatusInput) (jsonutils.JSONObject, error) { s.probeGuestStatus(ctx, resp) hostutils.TaskComplete(ctx, jsonutils.Marshal(resp)) return nil, nil @@ -2075,6 +2100,10 @@ func (s *sPodGuestInstance) getContainerStatus(ctx context.Context, ctrId string return status, cs, nil } +func (s *sPodGuestInstance) MarkContainerProbeDirty(ctrStatus string, ctrId string, reason string) { + s.markContainerProbeDirty(ctrStatus, ctrId, reason) +} + func (s *sPodGuestInstance) markContainerProbeDirty(status, ctrId string, reason string) { if status == computeapi.CONTAINER_STATUS_PROBING { reason = fmt.Sprintf("status is probing: %s", reason) diff --git a/pkg/hostman/guestman/pod/statusman/doc.go b/pkg/hostman/guestman/pod/statusman/doc.go new file mode 100644 index 0000000000..69ef1214e1 --- /dev/null +++ b/pkg/hostman/guestman/pod/statusman/doc.go @@ -0,0 +1 @@ +package statusman // import "yunion.io/x/onecloud/pkg/hostman/guestman/pod/statusman" diff --git a/pkg/hostman/guestman/pod/statusman/pod_state.go b/pkg/hostman/guestman/pod/statusman/pod_state.go new file mode 100644 index 0000000000..99a3a3af49 --- /dev/null +++ b/pkg/hostman/guestman/pod/statusman/pod_state.go @@ -0,0 +1,157 @@ +package statusman + +import ( + "context" + "time" + + "yunion.io/x/pkg/errors" + + "yunion.io/x/onecloud/pkg/apis" + computeapi "yunion.io/x/onecloud/pkg/apis/compute" + "yunion.io/x/onecloud/pkg/hostman/hostutils" +) + +var ( + statusManager IPodStatusManager +) + +func init() { + statusManager = newPodStatusManager() +} + +func GetManager() IPodStatusManager { + return statusManager +} + +type IPodStatusManager interface { + UpdateStatus(req *PodStatusUpdateRequest) error + Start() + Stop() +} + +type ContainerStatus struct { + Status string + RestartCount int + StartedAt *time.Time + LastFinishedAt *time.Time +} + +type IPod interface { + MarkContainerProbeDirty(ctrStatus string, ctrId string, reason string) +} + +type PodStatusUpdateRequest struct { + Id string + Pod IPod + Status string + ContainerStatuses map[string]*ContainerStatus + Reason string + Result chan error +} + +func (r PodStatusUpdateRequest) ToServerPerformStatusInput() *computeapi.ServerPerformStatusInput { + powerState := computeapi.VM_POWER_STATES_OFF + if r.Status == computeapi.VM_RUNNING { + powerState = computeapi.VM_POWER_STATES_ON + } + guestStatus := &computeapi.ServerPerformStatusInput{ + PerformStatusInput: apis.PerformStatusInput{ + Status: r.Status, + PowerStates: powerState, + Reason: r.Reason, + }, + Containers: make(map[string]*computeapi.ContainerPerformStatusInput), + } + for ctrId, ctrStatus := range r.ContainerStatuses { + guestStatus.Containers[ctrId] = &computeapi.ContainerPerformStatusInput{ + PerformStatusInput: apis.PerformStatusInput{ + Status: ctrStatus.Status, + Reason: r.Reason, + }, + RestartCount: ctrStatus.RestartCount, + StartedAt: ctrStatus.StartedAt, + LastFinishedAt: ctrStatus.LastFinishedAt, + } + } + + return guestStatus +} + +func (r PodStatusUpdateRequest) ToHostUploadGuestsStatusInput() *computeapi.HostUploadGuestsStatusInput { + id := r.Id + guestStatus := &computeapi.HostUploadGuestStatusInput{ + PerformStatusInput: apis.PerformStatusInput{ + Status: r.Status, + Reason: r.Reason, + }, + Containers: make(map[string]*computeapi.ContainerPerformStatusInput), + } + for ctrId, ctrStatus := range r.ContainerStatuses { + guestStatus.Containers[ctrId] = &computeapi.ContainerPerformStatusInput{ + PerformStatusInput: apis.PerformStatusInput{ + Status: ctrStatus.Status, + Reason: r.Reason, + }, + RestartCount: ctrStatus.RestartCount, + StartedAt: ctrStatus.StartedAt, + LastFinishedAt: ctrStatus.LastFinishedAt, + } + } + + return &computeapi.HostUploadGuestsStatusInput{ + Guests: map[string]*computeapi.HostUploadGuestStatusInput{ + id: guestStatus, + }, + } +} + +type podStatusManager struct { + updateChan chan *PodStatusUpdateRequest + stopChan chan struct{} +} + +func newPodStatusManager() IPodStatusManager { + return &podStatusManager{ + updateChan: make(chan *PodStatusUpdateRequest), + stopChan: make(chan struct{}), + } +} + +func (m *podStatusManager) Start() { + go m.processLoop() +} + +func (m *podStatusManager) Stop() { + close(m.stopChan) +} + +func (m *podStatusManager) UpdateStatus(req *PodStatusUpdateRequest) error { + result := make(chan error, 1) + req.Result = result + m.updateChan <- req + return <-result +} + +func (m *podStatusManager) processLoop() { + for { + select { + case <-m.stopChan: + return + case req := <-m.updateChan: + err := m.handleUpdate(req) + req.Result <- err + } + } +} + +func (m *podStatusManager) handleUpdate(req *PodStatusUpdateRequest) error { + input := req.ToServerPerformStatusInput() + if _, err := hostutils.UpdateServerContainersStatus(context.Background(), req.Id, input); err != nil { + return errors.Wrapf(err, "update server containers status") + } + for ctrId, ctrStatus := range input.Containers { + // 同步容器状态可能会出现 probing 状态,所以需要 mark 成 dirty,等待 probe manager 重新探测容器状态 + req.Pod.MarkContainerProbeDirty(ctrStatus.Status, ctrId, input.Reason) + } + return nil +} diff --git a/pkg/hostman/guestman/qemu-kvm.go b/pkg/hostman/guestman/qemu-kvm.go index f04f46b055..fd686ce200 100644 --- a/pkg/hostman/guestman/qemu-kvm.go +++ b/pkg/hostman/guestman/qemu-kvm.go @@ -3610,7 +3610,7 @@ func (s *SKVMGuestInstance) CPUSetRemove(ctx context.Context) error { return nil } -func (s *SKVMGuestInstance) GetUploadStatus(ctx context.Context, reason string) (*api.HostUploadGuestStatusResponse, error) { +func (s *SKVMGuestInstance) GetUploadStatus(ctx context.Context, reason string) (*api.HostUploadGuestStatusInput, error) { var status = api.VM_READY if s.IsSuspend() { status = api.VM_SUSPEND @@ -3628,15 +3628,15 @@ func (s *SKVMGuestInstance) GetUploadStatus(ctx context.Context, reason string) statusInput.Status = api.VM_BLOCK_STREAM } } - return &api.HostUploadGuestStatusResponse{ + return &api.HostUploadGuestStatusInput{ PerformStatusInput: *statusInput, }, nil } -func (s *SKVMGuestInstance) PostUploadStatus(resp *api.HostUploadGuestStatusResponse, reason string) { +func (s *SKVMGuestInstance) PostUploadStatus(resp *api.HostUploadGuestStatusInput, reason string) { } -func (s *SKVMGuestInstance) HandleGuestStatus(ctx context.Context, resp *api.HostUploadGuestStatusResponse) (jsonutils.JSONObject, error) { +func (s *SKVMGuestInstance) HandleGuestStatus(ctx context.Context, resp *api.HostUploadGuestStatusInput) (jsonutils.JSONObject, error) { if resp.Status == GUEST_RUNNING && s.pciUninitialized { resp.Status = api.VM_UNSYNC } else if resp.Status == GUEST_RUNNING { diff --git a/pkg/hostman/guestman/runtime.go b/pkg/hostman/guestman/runtime.go index 6e2b029425..cc61a559f0 100644 --- a/pkg/hostman/guestman/runtime.go +++ b/pkg/hostman/guestman/runtime.go @@ -61,9 +61,9 @@ type GuestRuntimeInstance interface { ImportServer(pendingDelete bool) ExitCleanup(bool) - HandleGuestStatus(ctx context.Context, resp *computeapi.HostUploadGuestStatusResponse) (jsonutils.JSONObject, error) - GetUploadStatus(ctx context.Context, reason string) (*computeapi.HostUploadGuestStatusResponse, error) - PostUploadStatus(resp *computeapi.HostUploadGuestStatusResponse, reason string) + HandleGuestStatus(ctx context.Context, resp *computeapi.HostUploadGuestStatusInput) (jsonutils.JSONObject, error) + GetUploadStatus(ctx context.Context, reason string) (*computeapi.HostUploadGuestStatusInput, error) + PostUploadStatus(resp *computeapi.HostUploadGuestStatusInput, reason string) HandleGuestStart(ctx context.Context, userCred mcclient.TokenCredential, body jsonutils.JSONObject) (jsonutils.JSONObject, error) HandleStop(ctx context.Context, timeout int64) error diff --git a/pkg/hostman/hostutils/hostutils.go b/pkg/hostman/hostutils/hostutils.go index e86f110b96..c377de279d 100644 --- a/pkg/hostman/hostutils/hostutils.go +++ b/pkg/hostman/hostutils/hostutils.go @@ -196,6 +196,10 @@ func UpdateServerStatus(ctx context.Context, sid string, statusInput *apis.Perfo return UpdateResourceStatus(ctx, &modules.Servers, sid, statusInput) } +func UpdateServerContainersStatus(ctx context.Context, sid string, input *computeapi.ServerPerformStatusInput) (jsonutils.JSONObject, error) { + return modules.Servers.PerformAction(GetComputeSession(ctx), sid, "status", jsonutils.Marshal(input)) +} + func UpdateServerProgress(ctx context.Context, sid string, progress, progressMbps float64) (jsonutils.JSONObject, error) { params := map[string]float64{ "progress": progress, @@ -204,7 +208,7 @@ func UpdateServerProgress(ctx context.Context, sid string, progress, progressMbp return modules.Servers.Update(GetComputeSession(ctx), sid, jsonutils.Marshal(params)) } -func UploadGuestsStatus(ctx context.Context, resp *computeapi.HostUploadGuestsStatusResponse) (jsonutils.JSONObject, error) { +func UploadGuestsStatus(ctx context.Context, resp *computeapi.HostUploadGuestsStatusInput) (jsonutils.JSONObject, error) { return modules.Servers.PerformClassAction(GetComputeSession(ctx), "upload-status", jsonutils.Marshal(resp)) }