diff --git a/log/schedulerclient.log b/log/schedulerclient.log index d544c91..2b7d8f0 100644 --- a/log/schedulerclient.log +++ b/log/schedulerclient.log @@ -9971,3 +9971,4 @@ Key: 'PackageDelete.PackageID' Error:Field validation for 'PackageID' failed on }|... 2025-05-14 09:03:12 [DEBU] JobStatusReport job 2025-05-14 09:03:16 [DEBU] JobStatusReport taskID: 1920782156997595179, taskName: task0038, messages: &{{} {} Train task0038 1920782156997595179 true 1790300942428540928 /nudt-cloudream2/cds/result/5/2205/1745568810} +2025-05-14 15:37:51 [INFO] start serving http at: :7891 diff --git a/manager/internal/mq/executor.go b/manager/internal/mq/executor.go index bf2f6e3..357ff39 100644 --- a/manager/internal/mq/executor.go +++ b/manager/internal/mq/executor.go @@ -1,6 +1,6 @@ package mq //func (svc *Service) ReportExecutorTaskStatus(msg *mgrmq.ReportExecutorTaskStatus) (*mgrmq.ReportExecutorTaskStatusResp, *mq.CodeMessage) { -// svc.exeMgr.Report(msg.ExecutorID, msg.TaskStatus) +// svc.exeMgr.ReportMessage(msg.ExecutorID, msg.TaskStatus) // return mq.ReplyOK(mgrmq.NewReportExecutorTaskStatusResp()) //} diff --git a/schedulerMiddleware/internal/http/jobflow_test.go b/schedulerMiddleware/internal/http/jobflow_test.go new file mode 100644 index 0000000..cdac8ad --- /dev/null +++ b/schedulerMiddleware/internal/http/jobflow_test.go @@ -0,0 +1,245 @@ +package http + +// +//import ( +// "bytes" +// "encoding/json" +// "net/http/httptest" +// "testing" +// +// "github.com/gin-gonic/gin" +// . "github.com/smartystreets/goconvey/convey" +// "github.com/stretchr/testify/mock" +// +// sch "gitlink.org.cn/cloudream/common/sdks/pcmscheduler" +// schsdk "gitlink.org.cn/cloudream/common/sdks/scheduler" +// cdssdk "gitlink.org.cn/cloudream/common/sdks/storage" +// jobTask "gitlink.org.cn/cloudream/scheduler/schedulerMiddleware/internal/manager/task" +//) +// +//// 模拟接口定义 +//type TestJobFlowService interface { +// Submit(userID cdssdk.UserID, id schsdk.DataID, token string, task *jobTask.JobTask[sch.TaskMessage]) error +// Create(jobFlow schsdk.JobFlow) (int64, error) +// Update(jobFlow schsdk.JobFlow) error +// SaveFlowData(id int64, userID cdssdk.UserID, data schsdk.FlowData) error +// Query(id schsdk.DataID, userID cdssdk.UserID, currentPage, pageSize int) ([]schsdk.JobFlow, int, int, error) +// DeleteByID(id int64, userID cdssdk.UserID) error +// QueryRunJob(id schsdk.JobSetID, userID cdssdk.UserID, currentPage, pageSize int) ([]schsdk.JobFlowRun, int, int, error) +// QueryStatus(runID schsdk.JobSetID, userID cdssdk.UserID) ([]schsdk.JobFlowRunStatus, error) +//} +// +//// 实现 mock 结构体 +//type MockJobFlowSvc struct { +// mock.Mock +//} +// +//func (m *MockJobFlowSvc) Submit(userID cdssdk.UserID, id schsdk.DataID, token string, task *jobTask.JobTask[sch.TaskMessage]) error { +// args := m.Called(userID, id, token, task) +// return args.Error(0) +//} +// +//func (m *MockJobFlowSvc) Create(jobFlow schsdk.JobFlow) (int64, error) { +// args := m.Called(jobFlow) +// return args.Get(0).(int64), args.Error(1) +//} +// +//func (m *MockJobFlowSvc) Update(jobFlow schsdk.JobFlow) error { +// args := m.Called(jobFlow) +// return args.Error(0) +//} +// +//func (m *MockJobFlowSvc) SaveFlowData(id int64, userID cdssdk.UserID, data schsdk.FlowData) error { +// args := m.Called(id, userID, data) +// return args.Error(0) +//} +// +//func (m *MockJobFlowSvc) Query(id schsdk.DataID, userID cdssdk.UserID, currentPage, pageSize int) ([]schsdk.JobFlow, int, int, error) { +// args := m.Called(id, userID, currentPage, pageSize) +// return args.Get(0).([]schsdk.JobFlow), args.Int(1), args.Int(2), args.Error(3) +//} +// +//func (m *MockJobFlowSvc) DeleteByID(id int64, userID cdssdk.UserID) error { +// args := m.Called(id, userID) +// return args.Error(0) +//} +// +//func (m *MockJobFlowSvc) QueryRunJob(id schsdk.JobSetID, userID cdssdk.UserID, currentPage, pageSize int) ([]schsdk.JobFlowRun, int, int, error) { +// args := m.Called(id, userID, currentPage, pageSize) +// return args.Get(0).([]schsdk.JobFlowRun), args.Int(1), args.Int(2), args.Error(3) +//} +// +//func (m *MockJobFlowSvc) QueryStatus(runID schsdk.JobSetID, userID cdssdk.UserID) ([]schsdk.JobFlowRunStatus, error) { +// args := m.Called(runID, userID) +// return args.Get(0).([]schsdk.JobFlowRunStatus), args.Error(1) +//} +// +//// 模拟 services.Service 接口 +//type MockService struct { +// MockJobFlowSvc *MockJobFlowSvc +//} +// +//func (s *MockService) JobFlowSvc() TestJobFlowService { +// return s.MockJobFlowSvc +//} +// +//// 模拟 Server 和 TestJobFlowService +//type TestServer struct { +// mockService *MockService +//} +// +//func (ts *TestServer) JobFlowSvc() *TestJobFlowService { +// return &JobFlowService{ +// Server: &Server{ +// svc: ts.mockService, +// }, +// } +//} +// +//func TestJobFlowHandlers(t *testing.T) { +// gin.SetMode(gin.TestMode) +// +// Convey("Given a job flow service", t, func() { +// mockSvc := new(MockJobFlowSvc) +// mockService := &MockService{MockJobFlowSvc: mockSvc} +// testServer := &TestServer{mockService: mockService} +// jobFlowService := testServer.JobFlowSvc() +// +// Convey("When handling Submit request", func() { +// reqBody, _ := json.Marshal(JobFlowSubmitReq{ +// UserID: 1, +// ID: 100, +// }) +// req := httptest.NewRequest("POST", "/submit", bytes.NewBuffer(reqBody)) +// req.Header.Set("Authorization", "test-token") +// w := httptest.NewRecorder() +// c, _ := gin.CreateTestContext(w) +// c.Request = req +// +// task := jobTask.NewJobTask[sch.TaskMessage]() +// mockSvc.On("Submit", cdssdk.UserID(1), schsdk.DataID(100), "test-token", task).Return(nil).Once() +// +// go func() { +// task.Send(sch.TaskMessage{Status: sch.SuccessStatus}) +// }() +// +// jobFlowService.Submit(c) +// So(w.Code, ShouldEqual, 200) +// So(w.Body.String(), ShouldContainSubstring, `"code":0`) +// }) +// +// Convey("When handling Create request", func() { +// reqBody, _ := json.Marshal(CreateReq{ +// JobFlow: schsdk.JobFlow{ +// Name: "test-flow", +// UserID: 1, +// }, +// }) +// req := httptest.NewRequest("POST", "/create", bytes.NewBuffer(reqBody)) +// w := httptest.NewRecorder() +// c, _ := gin.CreateTestContext(w) +// c.Request = req +// +// mockSvc.On("Create", mock.Anything).Return(int64(1), nil).Once() +// +// jobFlowService.Create(c) +// So(w.Code, ShouldEqual, 200) +// So(w.Body.String(), ShouldContainSubstring, `"id":1`) +// }) +// +// Convey("When handling Update request", func() { +// reqBody, _ := json.Marshal(UpdateReq{ +// JobFlow: schsdk.JobFlow{ +// ID: 1, +// Name: "updated-flow", +// UserID: 1, +// }, +// }) +// req := httptest.NewRequest("POST", "/update", bytes.NewBuffer(reqBody)) +// w := httptest.NewRecorder() +// c, _ := gin.CreateTestContext(w) +// c.Request = req +// +// mockSvc.On("Update", mock.Anything).Return(nil).Once() +// +// jobFlowService.Update(c) +// So(w.Code, ShouldEqual, 200) +// So(w.Body.String(), ShouldContainSubstring, `"code":0`) +// }) +// +// Convey("When handling SaveFlowData request", func() { +// reqBody, _ := json.Marshal(SaveFlowDataReq{ +// ID: 1, +// UserID: 1, +// Data: schsdk.FlowData{ +// Nodes: []schsdk.Node{}, +// }, +// }) +// req := httptest.NewRequest("POST", "/save", bytes.NewBuffer(reqBody)) +// w := httptest.NewRecorder() +// c, _ := gin.CreateTestContext(w) +// c.Request = req +// +// mockSvc.On("SaveFlowData", int64(1), cdssdk.UserID(1), mock.Anything).Return(nil).Once() +// +// jobFlowService.SaveFlowData(c) +// So(w.Code, ShouldEqual, 200) +// So(w.Body.String(), ShouldContainSubstring, `"code":0`) +// }) +// +// Convey("When handling Query request", func() { +// req := httptest.NewRequest("GET", "/query?id=1&userID=1¤tPage=1&pageSize=10", nil) +// w := httptest.NewRecorder() +// c, _ := gin.CreateTestContext(w) +// c.Request = req +// +// mockSvc.On("Query", schsdk.DataID(1), cdssdk.UserID(1), 1, 10). +// Return([]schsdk.JobFlow{{ID: 1}}, 1, 1, nil).Once() +// +// jobFlowService.Query(c) +// So(w.Code, ShouldEqual, 200) +// So(w.Body.String(), ShouldContainSubstring, `"totalPages":1`) +// }) +// +// Convey("When handling Delete request", func() { +// req := httptest.NewRequest("GET", "/delete?id=1&userID=1", nil) +// w := httptest.NewRecorder() +// c, _ := gin.CreateTestContext(w) +// c.Request = req +// +// mockSvc.On("DeleteByID", int64(1), cdssdk.UserID(1)).Return(nil).Once() +// +// jobFlowService.Delete(c) +// So(w.Code, ShouldEqual, 200) +// So(w.Body.String(), ShouldContainSubstring, `"code":0`) +// }) +// +// Convey("When handling QueryRunJob request", func() { +// req := httptest.NewRequest("GET", "/run-job?id=1&userID=1¤tPage=1&pageSize=10", nil) +// w := httptest.NewRecorder() +// c, _ := gin.CreateTestContext(w) +// c.Request = req +// +// mockSvc.On("QueryRunJob", schsdk.JobSetID("1"), cdssdk.UserID(1), 1, 10). +// Return([]schsdk.JobFlowRun{{ID: "1"}}, 1, 1, nil).Once() +// +// jobFlowService.QueryRunJob(c) +// So(w.Code, ShouldEqual, 200) +// So(w.Body.String(), ShouldContainSubstring, `"totalPages":1`) +// }) +// +// Convey("When handling QueryStatus request", func() { +// req := httptest.NewRequest("GET", "/status?runID=1&userID=1", nil) +// w := httptest.NewRecorder() +// c, _ := gin.CreateTestContext(w) +// c.Request = req +// +// mockSvc.On("QueryStatus", schsdk.JobSetID("1"), cdssdk.UserID(1)). +// Return([]schsdk.JobFlowRunStatus{{}}, nil).Once() +// +// jobFlowService.QueryStatus(c) +// So(w.Code, ShouldEqual, 200) +// So(w.Body.String(), ShouldContainSubstring, `"code":0`) +// }) +// }) +//} diff --git a/schedulerMiddleware/internal/manager/jobmgr/job/state2/data_return.go b/schedulerMiddleware/internal/manager/jobmgr/job/state2/data_return.go index 25ca6b4..09bee16 100644 --- a/schedulerMiddleware/internal/manager/jobmgr/job/state2/data_return.go +++ b/schedulerMiddleware/internal/manager/jobmgr/job/state2/data_return.go @@ -77,7 +77,7 @@ func (s *DataReturn) do(rtx jobmgr.JobStateRunContext, jo *jobmgr.Job) (cdssdk.P reportMsg := reJob.Info.Report jo.SetOutput(&schsdk.DataReturnJobOutput{ - Report: reportMsg, + ReportMessage: reportMsg, }) cdsCli, err := schglb.CloudreamStoragePool.Acquire() diff --git a/schedulerMiddleware/internal/manager/jobmgr/job/state2/data_schedule_binding.go b/schedulerMiddleware/internal/manager/jobmgr/job/state2/data_schedule_binding.go index 7e9ab7f..622fa40 100644 --- a/schedulerMiddleware/internal/manager/jobmgr/job/state2/data_schedule_binding.go +++ b/schedulerMiddleware/internal/manager/jobmgr/job/state2/data_schedule_binding.go @@ -46,9 +46,8 @@ func NewDataScheduleAndBinding(data DataScheduleAndBinding) *DataScheduleAndBind } type dataScheduleParam struct { - UserID cdssdk.UserID - PackageID cdssdk.PackageID - //BindingClusterIDs []schsdk.ClusterID + UserID cdssdk.UserID + PackageID cdssdk.PackageID UploadedClusters []uploadersdk.Cluster ClusterMap map[schsdk.ClusterID]uploadersdk.ClusterMapping RootPath string @@ -121,7 +120,7 @@ func (s *DataScheduleAndBinding) hasDependence(rtx jobmgr.JobStateRunContext, jo } var bindingInfos []sch.DataBinding - message := output.Report + message := output.ReportMessage dataRetureData, err := rtx.Mgr.DB.PCMJob().GetDataReturnByID(rtx.Mgr.DB.DefCtx(), message.TaskID, message.ClusterID) if err != nil { diff --git a/schedulerMiddleware/internal/services/jobset.go b/schedulerMiddleware/internal/services/jobset.go index 830cc4a..da3021a 100644 --- a/schedulerMiddleware/internal/services/jobset.go +++ b/schedulerMiddleware/internal/services/jobset.go @@ -334,8 +334,8 @@ func (svc *JobSetService) JobStatusReport(reportInfo schsdk.JobStatusReport) err } drJob.SetOutput(&schsdk.DataReturnJobOutput{ - Type: schsdk.JobTypeDataReturn, - Report: *report, + Type: schsdk.JobTypeDataReturn, + ReportMessage: *report, }) logger.Infof("pcmJob info: %v", drJob)