| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493 |
- package task_queue
- import (
- "errors"
- "fmt"
- "github.com/allanpk716/ChineseSubFinder/internal/pkg/global_value"
- "github.com/allanpk716/ChineseSubFinder/internal/pkg/my_util"
- "github.com/allanpk716/ChineseSubFinder/internal/pkg/settings"
- "github.com/allanpk716/ChineseSubFinder/internal/types/common"
- "github.com/allanpk716/ChineseSubFinder/internal/types/task_queue"
- taskQueue2 "github.com/allanpk716/ChineseSubFinder/internal/types/task_queue"
- "github.com/dgraph-io/badger/v3"
- "github.com/emirpasic/gods/maps/treemap"
- "github.com/emirpasic/gods/sets/treeset"
- "github.com/sirupsen/logrus"
- "os"
- "path/filepath"
- "sync"
- "time"
- )
- type TaskQueue struct {
- queueName string // 队列的名称
- settings *settings.Settings // 设置
- log *logrus.Logger // 日志
- taskPriorityMapList []*treemap.Map // 这里有 0-10 个优先级划分的存储 List,每Add一个数据的时候需要切换到这个 List 中去 save
- taskKeyMap *treemap.Map // 以每个任务的唯一 JobID 来存储每个 Job,这样可以快速查询
- taskGroupBySeries *treemap.Map // 以每个任务的 SeriesRootPath 来存储每个任务,然后内层是一个 treeset,后续可以遍历删除即可
- queueLock sync.Mutex // 公用这个锁
- }
- func NewTaskQueue(queueName string, settings *settings.Settings, log *logrus.Logger) *TaskQueue {
- tq := &TaskQueue{queueName: queueName, settings: settings, log: log,
- taskPriorityMapList: make([]*treemap.Map, 0),
- taskKeyMap: treemap.NewWithStringComparator(),
- taskGroupBySeries: treemap.NewWithStringComparator(),
- }
- for i := 0; i <= taskPriorityCount; i++ {
- tq.taskPriorityMapList = append(tq.taskPriorityMapList, treemap.NewWithStringComparator())
- }
- tq.read()
- tq.afterRead()
- return tq
- }
- func (t *TaskQueue) QueueName() string {
- return t.queueName
- }
- func (t *TaskQueue) Clear() error {
- defer t.queueLock.Unlock()
- t.queueLock.Lock()
- err := GetDb().Update(
- func(tx *badger.Txn) error {
- var err error
- for i := 0; i <= taskPriorityCount; i++ {
- key := []byte(MergeBucketAndKeyName(BucketNamePrefixVideoSubDownloadQueue,
- fmt.Sprintf("%s_%d", t.queueName, i)))
- // 因为已经查询了一次,确保一定存在,所以直接更新+1,TTL 多加 5s 确保今天过去,暂时去除 TTL uint32(restOfDaySecond.Seconds())+5
- if err = tx.Delete(key); err != nil {
- return err
- }
- }
- return nil
- })
- if err != nil {
- return err
- }
- for i := 0; i <= taskPriorityCount; i++ {
- t.taskPriorityMapList[i].Clear()
- }
- t.taskKeyMap.Clear()
- t.taskGroupBySeries.Clear()
- return nil
- }
- // Size 队列的长度,对外暴露,有锁
- func (t *TaskQueue) Size() int {
- defer t.queueLock.Unlock()
- t.queueLock.Lock()
- return t.taskKeyMap.Size()
- }
- // checkPriority 检测优先级,会校验范围
- func (t *TaskQueue) checkPriority(oneJob taskQueue2.OneJob) taskQueue2.OneJob {
- if oneJob.TaskPriority > taskPriorityCount {
- oneJob.TaskPriority = taskPriorityCount
- }
- if oneJob.TaskPriority < 0 {
- oneJob.TaskPriority = 0
- }
- return oneJob
- }
- // degrade 降一级,会校验范围
- func (t *TaskQueue) degrade(oneJob taskQueue2.OneJob) taskQueue2.OneJob {
- oneJob.TaskPriority -= 1
- return t.checkPriority(oneJob)
- }
- // Add 放入元素,放入的时候会根据 TaskPriority 进行归类,存在的不会新增和更新
- func (t *TaskQueue) Add(oneJob task_queue.OneJob) (bool, error) {
- defer t.queueLock.Unlock()
- t.queueLock.Lock()
- if t.isExist(oneJob.Id) == true {
- return false, nil
- }
- // 检查权限范围
- oneJob = t.checkPriority(oneJob)
- // 插入到统一的 KeyMap
- t.taskKeyMap.Put(oneJob.Id, oneJob.TaskPriority)
- // 分配到具体的优先级 map 中
- t.taskPriorityMapList[oneJob.TaskPriority].Put(oneJob.Id, oneJob)
- // 如果是连续剧,则需要存储到 taskGroupBySeries 中
- jobIDSet, found := t.taskGroupBySeries.Get(oneJob.SeriesRootDirPath)
- if found == false {
- // 不存在
- nowJobIDSet := treeset.NewWithStringComparator()
- nowJobIDSet.Add(oneJob.Id)
- t.taskGroupBySeries.Put(oneJob.SeriesRootDirPath, nowJobIDSet)
- } else {
- // 存在
- nowJobIDSet := jobIDSet.(*treeset.Set)
- nowJobIDSet.Add(oneJob.Id)
- t.taskGroupBySeries.Put(oneJob.SeriesRootDirPath, nowJobIDSet)
- }
- err := t.save(oneJob.TaskPriority)
- if err != nil {
- return false, err
- }
- return true, nil
- }
- // update 更新素,不存在则会失败,内部用,没有锁
- func (t *TaskQueue) update(oneJob task_queue.OneJob) (bool, error) {
- if t.isExist(oneJob.Id) == false {
- return false, nil
- }
- // 自动更新时间
- oneJob.UpdateTime = time.Now()
- // 这里需要判断是否有优先级的 Update,如果有就需要把之前缓存的表给更新
- // 然后再插入到新的表中
- taskPriorityIndex, _ := t.taskKeyMap.Get(oneJob.Id)
- // 检查权限范围
- oneJob = t.checkPriority(oneJob)
- if oneJob.TaskPriority != taskPriorityIndex {
- // 优先级修改
- // 先删除原有的优先级
- t.taskPriorityMapList[taskPriorityIndex.(int)].Remove(oneJob.Id)
- err := t.save(taskPriorityIndex.(int))
- if err != nil {
- return false, err
- }
- }
- // 插入到统一的 KeyMap
- t.taskKeyMap.Put(oneJob.Id, oneJob.TaskPriority)
- // 分配到具体的优先级 map 中
- t.taskPriorityMapList[oneJob.TaskPriority].Put(oneJob.Id, oneJob)
- err := t.save(oneJob.TaskPriority)
- if err != nil {
- return false, err
- }
- return true, nil
- }
- // Update 更新素,不存在则会失败
- func (t *TaskQueue) Update(oneJob task_queue.OneJob) (bool, error) {
- defer t.queueLock.Unlock()
- t.queueLock.Lock()
- return t.update(oneJob)
- }
- // AutoDetectUpdateJobStatus 根据任务的生命周期图,进行自动判断更新,见《任务的生命周期》流程图
- func (t *TaskQueue) AutoDetectUpdateJobStatus(oneJob task_queue.OneJob, inErr error) {
- defer t.queueLock.Unlock()
- t.queueLock.Lock()
- // 检查权限范围
- oneJob = t.checkPriority(oneJob)
- if inErr == nil {
- // 没有错误就是完成
- oneJob.TaskPriority = DefaultTaskPriorityLevel
- oneJob.JobStatus = taskQueue2.Done
- oneJob.DownloadTimes += 1
- } else {
- // 超过了时间限制,默认是 90 天, A.Before(B) : A < B == true
- if oneJob.AddedTime.AddDate(0, 0, t.settings.AdvancedSettings.TaskQueue.ExpirationTime).Before(time.Now()) == true {
- // 超过 90 天了
- oneJob.JobStatus = taskQueue2.Failed
- } else {
- // 还在 90 天内
- // 是否是首次,那么就看它的 Level 是否是在 5,然后 retry == 0
- if oneJob.TaskPriority == DefaultTaskPriorityLevel && oneJob.RetryTimes == 0 {
- // 需要重置到 L6
- oneJob.RetryTimes = 0
- oneJob.TaskPriority = FirstRetryTaskPriorityLevel
- } else {
- if oneJob.RetryTimes > t.settings.AdvancedSettings.TaskQueue.MaxRetryTimes {
- // 超过重试次数会进行一次降级,然后重置这个次数
- oneJob.RetryTimes = 0
- oneJob = t.degrade(oneJob)
- }
- }
- // 强制为 waiting
- oneJob.JobStatus = taskQueue2.Waiting
- }
- // 传入的错误需要放进来
- oneJob.ErrorInfo = inErr.Error()
- oneJob.DownloadTimes += 1
- }
- // 这里不要用错了,要用无锁的,不然会阻塞
- bok, err := t.update(oneJob)
- if err != nil {
- t.log.Errorln("AutoDetectUpdateJobStatus", oneJob.VideoFPath, err)
- return
- }
- if bok == false {
- t.log.Warningln("AutoDetectUpdateJobStatus ==", oneJob.VideoFPath, "Job.ID", oneJob.Id, "Not Found")
- return
- }
- }
- func (t *TaskQueue) GetJobsByStatus(status task_queue.JobStatus) (bool, []task_queue.OneJob, error) {
- defer t.queueLock.Unlock()
- t.queueLock.Lock()
- outOneJobs := make([]task_queue.OneJob, 0)
- // 如果队列里面没有东西,则返回 false
- if t.isEmpty() == true {
- return false, nil, nil
- }
- for TaskPriority := 0; TaskPriority <= taskPriorityCount; TaskPriority++ {
- t.taskPriorityMapList[TaskPriority].Each(func(key interface{}, value interface{}) {
- tOneJob := task_queue.OneJob{}
- tOneJob = value.(task_queue.OneJob)
- if tOneJob.JobStatus == status {
- // 找到加入列表
- outOneJobs = append(outOneJobs, tOneJob)
- }
- })
- }
- return true, outOneJobs, nil
- }
- // GetJobsByPriorityAndStatus 根据任务优先级和状态获取任务列表
- func (t *TaskQueue) GetJobsByPriorityAndStatus(taskPriority int, status task_queue.JobStatus) (bool, []task_queue.OneJob, error) {
- defer t.queueLock.Unlock()
- t.queueLock.Lock()
- outOneJobs := make([]task_queue.OneJob, 0)
- // 如果队列里面没有东西,则返回 false
- if t.isEmpty() == true {
- return false, nil, nil
- }
- t.taskPriorityMapList[taskPriority].Each(func(key interface{}, value interface{}) {
- tOneJob := task_queue.OneJob{}
- tOneJob = value.(task_queue.OneJob)
- if tOneJob.JobStatus == status {
- // 找到加入列表
- outOneJobs = append(outOneJobs, tOneJob)
- }
- })
- return true, outOneJobs, nil
- }
- func (t *TaskQueue) del(jobId string) (bool, error) {
- if t.isExist(jobId) == false {
- return false, nil
- }
- taskPriority, bok := t.taskKeyMap.Get(jobId)
- if bok == false {
- return false, nil
- }
- // 删除连续剧的 tree.Map 里面的 tree.Set 的元素
- needDelJobObj, bok := t.taskPriorityMapList[taskPriority.(int)].Get(jobId)
- if bok == false {
- return false, nil
- }
- needDelJob := needDelJobObj.(task_queue.OneJob)
- jobSetsObj, bok := t.taskGroupBySeries.Get(needDelJob.SeriesRootDirPath)
- if bok == false {
- return false, nil
- }
- jobSets := jobSetsObj.(*treeset.Set)
- jobSets.Remove(jobId)
- // 删除任务
- t.taskKeyMap.Remove(jobId)
- t.taskPriorityMapList[taskPriority.(int)].Remove(jobId)
- err := t.save(taskPriority.(int))
- if err != nil {
- return false, err
- }
- // 删除任务的时候也需要删除对应的日志
- pathRoot := filepath.Join(global_value.ConfigRootDirFPath(), "Logs")
- fileFPath := filepath.Join(pathRoot, common.OnceLogPrefix+jobId+".log")
- if my_util.IsFile(fileFPath) == true {
- err = os.Remove(fileFPath)
- if err != nil {
- t.log.Errorln("del job", jobId, "logfile,error:", err)
- }
- }
- return true, nil
- }
- // Del 删除一个元素
- func (t *TaskQueue) Del(jobId string) (bool, error) {
- defer t.queueLock.Unlock()
- t.queueLock.Lock()
- return t.del(jobId)
- }
- func (t *TaskQueue) read() {
- err := GetDb().View(
- func(tx *badger.Txn) error {
- var err error
- for i := 0; i <= taskPriorityCount; i++ {
- key := []byte(MergeBucketAndKeyName(BucketNamePrefixVideoSubDownloadQueue,
- fmt.Sprintf("%s_%d", t.queueName, i)))
- var item *badger.Item
- item, err = tx.Get(key)
- if err != nil {
- if err == badger.ErrKeyNotFound {
- continue
- } else {
- return err
- }
- }
- valCopy, err := item.ValueCopy(nil)
- if err != nil {
- return err
- }
- err = t.taskPriorityMapList[i].FromJSON(valCopy)
- if err != nil {
- return err
- }
- }
- return nil
- })
- if err != nil {
- t.log.Panicln(err)
- }
- // 需要把几个优先级的map中的key汇总
- for i := 0; i < taskPriorityCount; i++ {
- // JobID - OneJob
- t.taskPriorityMapList[i].Each(func(key interface{}, value interface{}) {
- // JobID -- taskPriority
- t.taskKeyMap.Put(key, i)
- // SeriesRootDirPath -- tree.Set(JobID)
- oneJob := value.(task_queue.OneJob)
- jobIDSet, found := t.taskGroupBySeries.Get(oneJob.SeriesRootDirPath)
- if found == false {
- // 不存在
- nowJobIDSet := treeset.NewWithStringComparator()
- nowJobIDSet.Add(oneJob.Id)
- t.taskGroupBySeries.Put(oneJob.SeriesRootDirPath, nowJobIDSet)
- } else {
- // 存在
- nowJobIDSet := jobIDSet.(*treeset.Set)
- nowJobIDSet.Add(oneJob.Id)
- t.taskGroupBySeries.Put(oneJob.SeriesRootDirPath, nowJobIDSet)
- }
- })
- }
- }
- func (t *TaskQueue) afterRead() {
- // 将 downloading 的任务重置为 failed
- for TaskPriority := 0; TaskPriority <= taskPriorityCount; TaskPriority++ {
- t.taskPriorityMapList[TaskPriority].Each(func(key interface{}, value interface{}) {
- nowOneJob := value.(task_queue.OneJob)
- if nowOneJob.JobStatus == task_queue.Downloading {
- nowOneJob.JobStatus = task_queue.Failed
- bok, err := t.update(nowOneJob)
- if err != nil {
- t.log.Errorln("afterRead.update failed", err)
- return
- }
- if bok == false {
- t.log.Errorln("afterRead.update failed")
- return
- }
- }
- })
- }
- }
- // save 需要把改变的数据保持到 K/V 数据库中,这个没有锁,所以需要在 Sync 中使用,不对外开放
- func (t *TaskQueue) save(taskPriority int) error {
- err := GetDb().Update(
- func(tx *badger.Txn) error {
- var err error
- key := []byte(MergeBucketAndKeyName(
- BucketNamePrefixVideoSubDownloadQueue,
- fmt.Sprintf("%s_%d", t.queueName, taskPriority)))
- b, err := t.taskPriorityMapList[taskPriority].ToJSON()
- if err != nil {
- return err
- }
- e := badger.NewEntry(key, b)
- err = tx.SetEntry(e)
- if err != nil {
- return err
- }
- return nil
- })
- if err != nil {
- return err
- }
- return nil
- }
- // isExist 是否已经存在,对内,无锁
- func (t *TaskQueue) isExist(jobID string) bool {
- _, bok := t.taskKeyMap.Get(jobID)
- return bok
- }
- // IsExist 是否已经存在,对外,有锁
- func (t *TaskQueue) IsExist(jobID string) bool {
- defer t.queueLock.Unlock()
- t.queueLock.Lock()
- _, bok := t.taskKeyMap.Get(jobID)
- return bok
- }
- // isEmpty 对内,无锁
- func (t *TaskQueue) isEmpty() bool {
- return t.taskKeyMap.Empty()
- }
- const (
- taskPriorityCount = 10
- HighTaskPriorityLevel = 3
- DefaultTaskPriorityLevel = 5
- FirstRetryTaskPriorityLevel = 6
- )
- var (
- ErrNotSubFound = errors.New("Not Sub Found")
- )
|