Skip to content

Commit 2fa92f9

Browse files
zhaoxianhuagithubgxll
authored andcommitted
[fix][dingospeed] Task state repair in stop
1 parent d2d0299 commit 2fa92f9

4 files changed

Lines changed: 32 additions & 24 deletions

File tree

internal/service/cache_job_service.go

Lines changed: 16 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -40,18 +40,19 @@ func (p *CacheJobService) CreateCacheJob(c echo.Context, jobReq *query.CreateCac
4040
ctx, cancelFunc := context.WithCancel(appInfo.Ctx())
4141
var task common.Task
4242
cacheTask := task2.CacheTask{
43-
Ctx: ctx,
44-
Job: jobReq,
45-
CancelFunc: cancelFunc,
46-
SchedulerDao: p.schedulerDao,
43+
Ctx: ctx,
44+
Job: jobReq,
45+
CancelFunc: cancelFunc,
46+
SchedulerDao: p.schedulerDao,
47+
RunningStatus: consts.RunningStatusJobBreak,
4748
}
4849
req := &manager.CreateCacheJobReq{
4950
Type: jobReq.Type,
5051
InstanceId: jobReq.InstanceId,
5152
Datatype: jobReq.Datatype,
5253
Org: jobReq.Org,
5354
Repo: jobReq.Repo,
54-
Status: consts.StatusCacheJobIng,
55+
Status: consts.RunningStatusJobIng,
5556
}
5657
authorization := c.Request().Header.Get("Authorization")
5758
if jobReq.Type == consts.CacheTypePreheat {
@@ -81,7 +82,7 @@ func (p *CacheJobService) CreateCacheJob(c echo.Context, jobReq *query.CreateCac
8182
UsedStorage: uint64(sha.UsedStorage),
8283
}
8384
if err = p.cachePool.SubmitForTimeout(ctx, task); err != nil {
84-
p.schedulerDao.ExecUpdateCacheJobStatus(int(cacheJob.Id), consts.StatusCacheJobWait, jobReq.InstanceId, "", "", consts.TaskMoreErrMsg, 0)
85+
p.schedulerDao.ExecUpdateCacheJobStatus(int(cacheJob.Id), consts.RunningStatusJobWait, jobReq.InstanceId, "", "", consts.TaskMoreErrMsg, 0)
8586
return 0, err
8687
}
8788
} else if jobReq.Type == consts.CacheTypeMount {
@@ -91,7 +92,7 @@ func (p *CacheJobService) CreateCacheJob(c echo.Context, jobReq *query.CreateCac
9192
Authorization: authorization,
9293
}
9394
if err := p.cachePool.SubmitForTimeout(ctx, task); err != nil {
94-
p.schedulerDao.ExecUpdateRepositoryMountStatus(cacheTask.TaskNo, consts.StatusCacheJobWait, consts.TaskMoreErrMsg)
95+
p.schedulerDao.ExecUpdateRepositoryMountStatus(cacheTask.TaskNo, consts.RunningStatusJobWait, consts.TaskMoreErrMsg)
9596
return 0, err
9697
}
9798
} else {
@@ -103,10 +104,15 @@ func (p *CacheJobService) CreateCacheJob(c echo.Context, jobReq *query.CreateCac
103104

104105
func (p *CacheJobService) StopCacheJob(jobStatusReq *query.JobStatusReq) error {
105106
if task, ok := p.cachePool.GetTask(int(jobStatusReq.Id)); ok {
107+
if t, ok := task.(*task2.PreheatCacheTask); ok {
108+
t.RunningStatus = consts.RunningStatusJobStop
109+
} else if t, ok := task.(*task2.MountCacheTask); ok {
110+
t.RunningStatus = consts.RunningStatusJobStop
111+
}
106112
cancelFun := task.GetCancelFun()
107113
cancelFun()
108114
} else {
109-
p.schedulerDao.ExecUpdateCacheJobStatus(int(jobStatusReq.Id), consts.StatusCacheJobBreak, jobStatusReq.InstanceId, "", "", "speed未注册该任务,下载已中断。", 0)
115+
p.schedulerDao.ExecUpdateCacheJobStatus(int(jobStatusReq.Id), consts.RunningStatusJobStop, jobStatusReq.InstanceId, "", "", "speed未注册该任务,下载已中断。", 0)
110116
}
111117
return nil
112118
}
@@ -139,7 +145,7 @@ func (p *CacheJobService) ResumeCacheJob(c echo.Context, resumeCacheJobReq *quer
139145
return err
140146
}
141147
// 将状态重置为进行中
142-
p.schedulerDao.ExecUpdateCacheJobStatus(int(resumeCacheJobReq.Id), consts.StatusCacheJobIng,
148+
p.schedulerDao.ExecUpdateCacheJobStatus(int(resumeCacheJobReq.Id), consts.RunningStatusJobIng,
143149
resumeCacheJobReq.InstanceId, resumeCacheJobReq.Org, resumeCacheJobReq.Repo, "", 0)
144150
task = &task2.PreheatCacheTask{
145151
CacheTask: cacheTask,
@@ -150,7 +156,7 @@ func (p *CacheJobService) ResumeCacheJob(c echo.Context, resumeCacheJobReq *quer
150156
UsedStorage: uint64(resumeCacheJobReq.UsedStorage),
151157
}
152158
if err := p.cachePool.SubmitForTimeout(ctx, task); err != nil {
153-
p.schedulerDao.ExecUpdateCacheJobStatus(int(resumeCacheJobReq.Id), consts.StatusCacheJobWait, resumeCacheJobReq.InstanceId, "", "", consts.TaskMoreErrMsg, 0)
159+
p.schedulerDao.ExecUpdateCacheJobStatus(int(resumeCacheJobReq.Id), consts.RunningStatusJobWait, resumeCacheJobReq.InstanceId, "", "", consts.TaskMoreErrMsg, 0)
154160
return err
155161
}
156162
}

internal/service/task/mount_cache_task.go

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -75,10 +75,10 @@ func (m *MountCacheTask) DoTask() {
7575
errMsg = strings.Join(lines, "\n")
7676
}
7777
}
78-
m.SchedulerDao.ExecUpdateRepositoryMountStatus(m.TaskNo, consts.StatusCacheJobBreak, errMsg)
78+
m.SchedulerDao.ExecUpdateRepositoryMountStatus(m.TaskNo, m.RunningStatus, errMsg)
7979
} else {
8080
zap.S().Infof("command success.%d", m.TaskNo)
81-
m.SchedulerDao.ExecUpdateRepositoryMountStatus(m.TaskNo, consts.StatusCacheJobComplete, "")
81+
m.SchedulerDao.ExecUpdateRepositoryMountStatus(m.TaskNo, consts.RunningStatusJobComplete, "")
8282
}
8383
}
8484

internal/service/task/preheat_cache_task.go

Lines changed: 8 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -19,11 +19,12 @@ import (
1919
)
2020

2121
type CacheTask struct {
22-
TaskNo int
23-
Ctx context.Context
24-
CancelFunc context.CancelFunc
25-
Job *query.CreateCacheJobReq
26-
SchedulerDao *dao.SchedulerDao
22+
TaskNo int
23+
Ctx context.Context
24+
CancelFunc context.CancelFunc
25+
Job *query.CreateCacheJobReq
26+
SchedulerDao *dao.SchedulerDao
27+
RunningStatus int32
2728
}
2829

2930
func (c *CacheTask) GetTaskNo() int {
@@ -53,11 +54,11 @@ func (p *PreheatCacheTask) DoTask() {
5354
go p.realTimeSpeed(ctx)
5455
err := p.preheatProcess(orgRepo)
5556
if err != nil {
56-
p.SchedulerDao.ExecUpdateCacheJobStatus(p.TaskNo, consts.StatusCacheJobBreak, p.Job.InstanceId, p.Job.Org, p.Job.Repo, err.Error(), p.StockProcess)
57+
p.SchedulerDao.ExecUpdateCacheJobStatus(p.TaskNo, p.RunningStatus, p.Job.InstanceId, p.Job.Org, p.Job.Repo, err.Error(), p.StockProcess)
5758
return
5859
}
5960
p.StockProcess = 100
60-
p.SchedulerDao.ExecUpdateCacheJobStatus(p.TaskNo, consts.StatusCacheJobComplete, p.Job.InstanceId, p.Job.Org, p.Job.Repo, "", p.StockProcess)
61+
p.SchedulerDao.ExecUpdateCacheJobStatus(p.TaskNo, consts.RunningStatusJobComplete, p.Job.InstanceId, p.Job.Org, p.Job.Repo, "", p.StockProcess)
6162
}
6263

6364
func (p *PreheatCacheTask) preheatProcess(orgRepo string) error {

pkg/consts/const.go

Lines changed: 6 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -80,11 +80,12 @@ const (
8080
CacheTypePreheat = 1
8181
CacheTypeMount = 2
8282

83-
StatusCacheJobIng = 1
84-
StatusCacheJobBreak = 2
85-
StatusCacheJobComplete = 3
86-
StatusCacheJobStopping = 4
87-
StatusCacheJobWait = 5
83+
RunningStatusJobIng = 1
84+
RunningStatusJobBreak = 2
85+
RunningStatusJobComplete = 3
86+
RunningStatusJobStopping = 4
87+
RunningStatusJobStop = 5
88+
RunningStatusJobWait = 6
8889

8990
OperationProcess int = 1
9091
OperationPreheat int = 2

0 commit comments

Comments
 (0)