mirror of
https://github.com/yunionio/cloudpods.git
synced 2026-09-19 10:46:58 +08:00
refactor(region,host): container status management (#22634)
This commit is contained in:
@@ -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"`
|
||||
}
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -0,0 +1 @@
|
||||
package statusman // import "yunion.io/x/onecloud/pkg/hostman/guestman/pod/statusman"
|
||||
@@ -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
|
||||
}
|
||||
@@ -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 {
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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))
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user