修复合并错误

This commit is contained in:
JeshuaRen 2024-12-31 09:35:06 +08:00
parent 8d49c8684b
commit a03436cc10
7 changed files with 29 additions and 13 deletions

View File

@ -30,7 +30,7 @@ func NewServer(svc Service, cfg *mymq.Config) (*Server, error) {
func(msg *mq.Message) (*mq.Message, error) {
return msgDispatcher.Handle(srv.service, msg)
},
cfg.RabbitMQParam,
//cfg.RabbitMQParam,
)
if err != nil {
return nil, err

View File

@ -36,7 +36,7 @@ func NewServer(svc Service, cfg *mymq.Config) (*Server, error) {
func(msg *mq.Message) (*mq.Message, error) {
return msgDispatcher.Handle(srv.service, msg)
},
cfg.RabbitMQParam,
//cfg.RabbitMQParam,
)
if err != nil {
return nil, err

View File

@ -1,7 +1,6 @@
package mq
import (
"fmt"
"gitlink.org.cn/cloudream/common/pkgs/mq"
)
@ -13,6 +12,14 @@ type Config struct {
RabbitMQParam mq.RabbitMQParam `json:"param"`
}
func (cfg *Config) MakeConnectingURL() string {
return fmt.Sprintf("amqp://%s:%s@%s%s", cfg.Account, cfg.Password, cfg.Address, cfg.VHost)
func (cfg *Config) MakeConnectingURL() mq.Config {
//return fmt.Sprintf("amqp://%s:%s@%s%s", cfg.Account, cfg.Password, cfg.Address, cfg.VHost)
conf := mq.Config{
Address: cfg.Address,
Account: cfg.Account,
Password: cfg.Password,
VHost: cfg.VHost,
Param: cfg.RabbitMQParam,
}
return conf
}

View File

@ -33,7 +33,7 @@ func NewServer(svc Service, cfg *mymq.Config) (*Server, error) {
func(msg *mq.Message) (*mq.Message, error) {
return msgDispatcher.Handle(srv.service, msg)
},
cfg.RabbitMQParam,
//cfg.RabbitMQParam,
)
if err != nil {
return nil, err

View File

@ -40,7 +40,7 @@ func NewServer(svc Service, cfg *mymq.Config) (*Server, error) {
func(msg *mq.Message) (*mq.Message, error) {
return msgDispatcher.Handle(srv.service, msg)
},
cfg.RabbitMQParam,
//cfg.RabbitMQParam,
)
if err != nil {
return nil, err

View File

@ -8,6 +8,7 @@ import (
"gitlink.org.cn/cloudream/scheduler/schedulerMiddleware/internal/manager/jobmgr/event"
"gitlink.org.cn/cloudream/scheduler/schedulerMiddleware/internal/manager/jobmgr/job"
"path/filepath"
"strconv"
"time"
"github.com/samber/lo"
@ -148,16 +149,19 @@ func loadDatasetPackage(userID cdssdk.UserID, packageID cdssdk.PackageID, storag
}
defer schglb.CloudreamStoragePool.Release(stgCli)
loadPackageResp, err := stgCli.StorageLoadPackage(cdsapi.StorageLoadPackageReq{
rootPath := "dataset_" + strconv.FormatInt(time.Now().Unix(), 10)
_, err = stgCli.StorageLoadPackage(cdsapi.StorageLoadPackageReq{
UserID: userID,
PackageID: packageID,
StorageID: storageID,
RootPath: rootPath,
})
if err != nil {
return "", err
}
logger.Info("load pacakge path: " + loadPackageResp.FullPath)
return loadPackageResp.FullPath, nil
logger.Info("load pacakge path: " + rootPath)
return rootPath, nil
}
func (s *JobExecuting) submitNormalTask(rtx jobmgr.JobStateRunContext, cmd string, envs []schsdk.KVPair, ccInfo schmod.ComputingCenter, pcmImgInfo schmod.PCMImage, resourceID pcmsdk.ResourceID) error {

View File

@ -6,6 +6,8 @@ import (
"gitlink.org.cn/cloudream/scheduler/schedulerMiddleware/internal/manager/jobmgr"
"gitlink.org.cn/cloudream/scheduler/schedulerMiddleware/internal/manager/jobmgr/event"
"gitlink.org.cn/cloudream/scheduler/schedulerMiddleware/internal/manager/jobmgr/job"
"strconv"
"time"
"gitlink.org.cn/cloudream/common/pkgs/future"
"gitlink.org.cn/cloudream/common/pkgs/logger"
@ -84,16 +86,19 @@ func (s *MultiInstanceUpdate) do(rtx jobmgr.JobStateRunContext, jo *jobmgr.Job)
StorageID: ccInfo.CDSStorageID,
})
loadPackageResp, err := stgCli.StorageLoadPackage(cdsapi.StorageLoadPackageReq{
rootPath := "model_update_" + strconv.FormatInt(time.Now().Unix(), 10)
_, err = stgCli.StorageLoadPackage(cdsapi.StorageLoadPackageReq{
UserID: userID,
PackageID: dtrJob.DataReturnPackageID,
StorageID: getStg.StorageID,
RootPath: rootPath,
})
if err != nil {
return fmt.Errorf("loading package: %w", err)
}
logger.Info("load pacakge path: " + loadPackageResp.FullPath)
fullPath = loadPackageResp.FullPath
logger.Info("load pacakge path: " + rootPath)
fullPath = rootPath
}
// 发送事件更新各个instance