task_queue.go 13 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495
  1. package task_queue
  2. import (
  3. "errors"
  4. "fmt"
  5. "github.com/allanpk716/ChineseSubFinder/internal/pkg/badger_err_check"
  6. "github.com/allanpk716/ChineseSubFinder/internal/pkg/global_value"
  7. "github.com/allanpk716/ChineseSubFinder/internal/pkg/my_util"
  8. "github.com/allanpk716/ChineseSubFinder/internal/pkg/settings"
  9. "github.com/allanpk716/ChineseSubFinder/internal/types/common"
  10. "github.com/allanpk716/ChineseSubFinder/internal/types/task_queue"
  11. taskQueue2 "github.com/allanpk716/ChineseSubFinder/internal/types/task_queue"
  12. "github.com/dgraph-io/badger/v3"
  13. "github.com/emirpasic/gods/maps/treemap"
  14. "github.com/emirpasic/gods/sets/treeset"
  15. "github.com/sirupsen/logrus"
  16. "os"
  17. "path/filepath"
  18. "sync"
  19. "time"
  20. )
  21. type TaskQueue struct {
  22. queueName string // 队列的名称
  23. settings *settings.Settings // 设置
  24. log *logrus.Logger // 日志
  25. taskPriorityMapList []*treemap.Map // 这里有 0-10 个优先级划分的存储 List,每Add一个数据的时候需要切换到这个 List 中去 save
  26. taskKeyMap *treemap.Map // 以每个任务的唯一 JobID 来存储每个 Job,这样可以快速查询
  27. taskGroupBySeries *treemap.Map // 以每个任务的 SeriesRootPath 来存储每个任务,然后内层是一个 treeset,后续可以遍历删除即可
  28. queueLock sync.Mutex // 公用这个锁
  29. }
  30. func NewTaskQueue(queueName string, settings *settings.Settings, log *logrus.Logger) *TaskQueue {
  31. tq := &TaskQueue{queueName: queueName, settings: settings, log: log,
  32. taskPriorityMapList: make([]*treemap.Map, 0),
  33. taskKeyMap: treemap.NewWithStringComparator(),
  34. taskGroupBySeries: treemap.NewWithStringComparator(),
  35. }
  36. for i := 0; i <= taskPriorityCount; i++ {
  37. tq.taskPriorityMapList = append(tq.taskPriorityMapList, treemap.NewWithStringComparator())
  38. }
  39. tq.read()
  40. tq.afterRead()
  41. return tq
  42. }
  43. func (t *TaskQueue) QueueName() string {
  44. return t.queueName
  45. }
  46. func (t *TaskQueue) Clear() error {
  47. defer t.queueLock.Unlock()
  48. t.queueLock.Lock()
  49. err := GetDb().Update(
  50. func(tx *badger.Txn) error {
  51. var err error
  52. for i := 0; i <= taskPriorityCount; i++ {
  53. key := []byte(MergeBucketAndKeyName(BucketNamePrefixVideoSubDownloadQueue,
  54. fmt.Sprintf("%s_%d", t.queueName, i)))
  55. // 因为已经查询了一次,确保一定存在,所以直接更新+1,TTL 多加 5s 确保今天过去,暂时去除 TTL uint32(restOfDaySecond.Seconds())+5
  56. if err = tx.Delete(key); err != nil {
  57. return err
  58. }
  59. }
  60. return nil
  61. })
  62. if err != nil {
  63. return err
  64. }
  65. for i := 0; i <= taskPriorityCount; i++ {
  66. t.taskPriorityMapList[i].Clear()
  67. }
  68. t.taskKeyMap.Clear()
  69. t.taskGroupBySeries.Clear()
  70. return nil
  71. }
  72. // Size 队列的长度,对外暴露,有锁
  73. func (t *TaskQueue) Size() int {
  74. defer t.queueLock.Unlock()
  75. t.queueLock.Lock()
  76. return t.taskKeyMap.Size()
  77. }
  78. // checkPriority 检测优先级,会校验范围
  79. func (t *TaskQueue) checkPriority(oneJob taskQueue2.OneJob) taskQueue2.OneJob {
  80. if oneJob.TaskPriority > taskPriorityCount {
  81. oneJob.TaskPriority = taskPriorityCount
  82. }
  83. if oneJob.TaskPriority < 0 {
  84. oneJob.TaskPriority = 0
  85. }
  86. return oneJob
  87. }
  88. // degrade 降一级,会校验范围
  89. func (t *TaskQueue) degrade(oneJob taskQueue2.OneJob) taskQueue2.OneJob {
  90. oneJob.TaskPriority -= 1
  91. return t.checkPriority(oneJob)
  92. }
  93. // Add 放入元素,放入的时候会根据 TaskPriority 进行归类,存在的不会新增和更新
  94. func (t *TaskQueue) Add(oneJob task_queue.OneJob) (bool, error) {
  95. defer t.queueLock.Unlock()
  96. t.queueLock.Lock()
  97. if t.isExist(oneJob.Id) == true {
  98. return false, nil
  99. }
  100. // 检查权限范围
  101. oneJob = t.checkPriority(oneJob)
  102. // 插入到统一的 KeyMap
  103. t.taskKeyMap.Put(oneJob.Id, oneJob.TaskPriority)
  104. // 分配到具体的优先级 map 中
  105. t.taskPriorityMapList[oneJob.TaskPriority].Put(oneJob.Id, oneJob)
  106. // 如果是连续剧,则需要存储到 taskGroupBySeries 中
  107. jobIDSet, found := t.taskGroupBySeries.Get(oneJob.SeriesRootDirPath)
  108. if found == false {
  109. // 不存在
  110. nowJobIDSet := treeset.NewWithStringComparator()
  111. nowJobIDSet.Add(oneJob.Id)
  112. t.taskGroupBySeries.Put(oneJob.SeriesRootDirPath, nowJobIDSet)
  113. } else {
  114. // 存在
  115. nowJobIDSet := jobIDSet.(*treeset.Set)
  116. nowJobIDSet.Add(oneJob.Id)
  117. t.taskGroupBySeries.Put(oneJob.SeriesRootDirPath, nowJobIDSet)
  118. }
  119. err := t.save(oneJob.TaskPriority)
  120. if err != nil {
  121. return false, err
  122. }
  123. return true, nil
  124. }
  125. // update 更新素,不存在则会失败,内部用,没有锁
  126. func (t *TaskQueue) update(oneJob task_queue.OneJob) (bool, error) {
  127. if t.isExist(oneJob.Id) == false {
  128. return false, nil
  129. }
  130. // 自动更新时间
  131. oneJob.UpdateTime = time.Now()
  132. // 这里需要判断是否有优先级的 Update,如果有就需要把之前缓存的表给更新
  133. // 然后再插入到新的表中
  134. taskPriorityIndex, _ := t.taskKeyMap.Get(oneJob.Id)
  135. // 检查权限范围
  136. oneJob = t.checkPriority(oneJob)
  137. if oneJob.TaskPriority != taskPriorityIndex {
  138. // 优先级修改
  139. // 先删除原有的优先级
  140. t.taskPriorityMapList[taskPriorityIndex.(int)].Remove(oneJob.Id)
  141. err := t.save(taskPriorityIndex.(int))
  142. if err != nil {
  143. return false, err
  144. }
  145. }
  146. // 插入到统一的 KeyMap
  147. t.taskKeyMap.Put(oneJob.Id, oneJob.TaskPriority)
  148. // 分配到具体的优先级 map 中
  149. t.taskPriorityMapList[oneJob.TaskPriority].Put(oneJob.Id, oneJob)
  150. err := t.save(oneJob.TaskPriority)
  151. if err != nil {
  152. return false, err
  153. }
  154. return true, nil
  155. }
  156. // Update 更新素,不存在则会失败
  157. func (t *TaskQueue) Update(oneJob task_queue.OneJob) (bool, error) {
  158. defer t.queueLock.Unlock()
  159. t.queueLock.Lock()
  160. return t.update(oneJob)
  161. }
  162. // AutoDetectUpdateJobStatus 根据任务的生命周期图,进行自动判断更新,见《任务的生命周期》流程图
  163. func (t *TaskQueue) AutoDetectUpdateJobStatus(oneJob task_queue.OneJob, inErr error) {
  164. defer t.queueLock.Unlock()
  165. t.queueLock.Lock()
  166. // 检查权限范围
  167. oneJob = t.checkPriority(oneJob)
  168. if inErr == nil {
  169. // 没有错误就是完成
  170. oneJob.TaskPriority = DefaultTaskPriorityLevel
  171. oneJob.JobStatus = taskQueue2.Done
  172. oneJob.DownloadTimes += 1
  173. } else {
  174. // 超过了时间限制,默认是 90 天, A.Before(B) : A < B == true
  175. if oneJob.AddedTime.AddDate(0, 0, t.settings.AdvancedSettings.TaskQueue.ExpirationTime).Before(time.Now()) == true {
  176. // 超过 90 天了
  177. oneJob.JobStatus = taskQueue2.Failed
  178. } else {
  179. // 还在 90 天内
  180. // 是否是首次,那么就看它的 Level 是否是在 5,然后 retry == 0
  181. if oneJob.TaskPriority == DefaultTaskPriorityLevel && oneJob.RetryTimes == 0 {
  182. // 需要重置到 L6
  183. oneJob.RetryTimes = 0
  184. oneJob.TaskPriority = FirstRetryTaskPriorityLevel
  185. } else {
  186. if oneJob.RetryTimes > t.settings.AdvancedSettings.TaskQueue.MaxRetryTimes {
  187. // 超过重试次数会进行一次降级,然后重置这个次数
  188. oneJob.RetryTimes = 0
  189. oneJob = t.degrade(oneJob)
  190. }
  191. }
  192. // 强制为 waiting
  193. oneJob.JobStatus = taskQueue2.Waiting
  194. }
  195. // 传入的错误需要放进来
  196. oneJob.ErrorInfo = inErr.Error()
  197. oneJob.DownloadTimes += 1
  198. }
  199. // 这里不要用错了,要用无锁的,不然会阻塞
  200. bok, err := t.update(oneJob)
  201. if err != nil {
  202. t.log.Errorln("AutoDetectUpdateJobStatus", oneJob.VideoFPath, err)
  203. return
  204. }
  205. if bok == false {
  206. t.log.Warningln("AutoDetectUpdateJobStatus ==", oneJob.VideoFPath, "Job.ID", oneJob.Id, "Not Found")
  207. return
  208. }
  209. }
  210. func (t *TaskQueue) GetJobsByStatus(status task_queue.JobStatus) (bool, []task_queue.OneJob, error) {
  211. defer t.queueLock.Unlock()
  212. t.queueLock.Lock()
  213. outOneJobs := make([]task_queue.OneJob, 0)
  214. // 如果队列里面没有东西,则返回 false
  215. if t.isEmpty() == true {
  216. return false, nil, nil
  217. }
  218. for TaskPriority := 0; TaskPriority <= taskPriorityCount; TaskPriority++ {
  219. t.taskPriorityMapList[TaskPriority].Each(func(key interface{}, value interface{}) {
  220. tOneJob := task_queue.OneJob{}
  221. tOneJob = value.(task_queue.OneJob)
  222. if tOneJob.JobStatus == status {
  223. // 找到加入列表
  224. outOneJobs = append(outOneJobs, tOneJob)
  225. }
  226. })
  227. }
  228. return true, outOneJobs, nil
  229. }
  230. // GetJobsByPriorityAndStatus 根据任务优先级和状态获取任务列表
  231. func (t *TaskQueue) GetJobsByPriorityAndStatus(taskPriority int, status task_queue.JobStatus) (bool, []task_queue.OneJob, error) {
  232. defer t.queueLock.Unlock()
  233. t.queueLock.Lock()
  234. outOneJobs := make([]task_queue.OneJob, 0)
  235. // 如果队列里面没有东西,则返回 false
  236. if t.isEmpty() == true {
  237. return false, nil, nil
  238. }
  239. t.taskPriorityMapList[taskPriority].Each(func(key interface{}, value interface{}) {
  240. tOneJob := task_queue.OneJob{}
  241. tOneJob = value.(task_queue.OneJob)
  242. if tOneJob.JobStatus == status {
  243. // 找到加入列表
  244. outOneJobs = append(outOneJobs, tOneJob)
  245. }
  246. })
  247. return true, outOneJobs, nil
  248. }
  249. func (t *TaskQueue) del(jobId string) (bool, error) {
  250. if t.isExist(jobId) == false {
  251. return false, nil
  252. }
  253. taskPriority, bok := t.taskKeyMap.Get(jobId)
  254. if bok == false {
  255. return false, nil
  256. }
  257. // 删除连续剧的 tree.Map 里面的 tree.Set 的元素
  258. needDelJobObj, bok := t.taskPriorityMapList[taskPriority.(int)].Get(jobId)
  259. if bok == false {
  260. return false, nil
  261. }
  262. needDelJob := needDelJobObj.(task_queue.OneJob)
  263. jobSetsObj, bok := t.taskGroupBySeries.Get(needDelJob.SeriesRootDirPath)
  264. if bok == false {
  265. return false, nil
  266. }
  267. jobSets := jobSetsObj.(*treeset.Set)
  268. jobSets.Remove(jobId)
  269. // 删除任务
  270. t.taskKeyMap.Remove(jobId)
  271. t.taskPriorityMapList[taskPriority.(int)].Remove(jobId)
  272. err := t.save(taskPriority.(int))
  273. if err != nil {
  274. return false, err
  275. }
  276. // 删除任务的时候也需要删除对应的日志
  277. pathRoot := filepath.Join(global_value.ConfigRootDirFPath(), "Logs")
  278. fileFPath := filepath.Join(pathRoot, common.OnceLogPrefix+jobId+".log")
  279. if my_util.IsFile(fileFPath) == true {
  280. err = os.Remove(fileFPath)
  281. if err != nil {
  282. t.log.Errorln("del job", jobId, "logfile,error:", err)
  283. }
  284. }
  285. return true, nil
  286. }
  287. // Del 删除一个元素
  288. func (t *TaskQueue) Del(jobId string) (bool, error) {
  289. defer t.queueLock.Unlock()
  290. t.queueLock.Lock()
  291. return t.del(jobId)
  292. }
  293. func (t *TaskQueue) read() {
  294. err := GetDb().View(
  295. func(tx *badger.Txn) error {
  296. var err error
  297. for i := 0; i <= taskPriorityCount; i++ {
  298. key := []byte(MergeBucketAndKeyName(BucketNamePrefixVideoSubDownloadQueue,
  299. fmt.Sprintf("%s_%d", t.queueName, i)))
  300. var item *badger.Item
  301. item, err = tx.Get(key)
  302. if err != nil {
  303. if badger_err_check.IsErrOk(err) == true {
  304. return nil
  305. }
  306. return err
  307. }
  308. valCopy, err := item.ValueCopy(nil)
  309. if err != nil {
  310. return err
  311. }
  312. err = t.taskPriorityMapList[i].FromJSON(valCopy)
  313. if err != nil {
  314. return err
  315. }
  316. }
  317. return nil
  318. })
  319. if err != nil {
  320. t.log.Panicln(err)
  321. }
  322. // 需要把几个优先级的map中的key汇总
  323. for i := 0; i < taskPriorityCount; i++ {
  324. // JobID - OneJob
  325. t.taskPriorityMapList[i].Each(func(key interface{}, value interface{}) {
  326. // JobID -- taskPriority
  327. t.taskKeyMap.Put(key, i)
  328. // SeriesRootDirPath -- tree.Set(JobID)
  329. oneJob := value.(task_queue.OneJob)
  330. jobIDSet, found := t.taskGroupBySeries.Get(oneJob.SeriesRootDirPath)
  331. if found == false {
  332. // 不存在
  333. nowJobIDSet := treeset.NewWithStringComparator()
  334. nowJobIDSet.Add(oneJob.Id)
  335. t.taskGroupBySeries.Put(oneJob.SeriesRootDirPath, nowJobIDSet)
  336. } else {
  337. // 存在
  338. nowJobIDSet := jobIDSet.(*treeset.Set)
  339. nowJobIDSet.Add(oneJob.Id)
  340. t.taskGroupBySeries.Put(oneJob.SeriesRootDirPath, nowJobIDSet)
  341. }
  342. })
  343. }
  344. }
  345. func (t *TaskQueue) afterRead() {
  346. // 将 downloading 的任务重置为 failed
  347. for TaskPriority := 0; TaskPriority <= taskPriorityCount; TaskPriority++ {
  348. t.taskPriorityMapList[TaskPriority].Each(func(key interface{}, value interface{}) {
  349. nowOneJob := value.(task_queue.OneJob)
  350. if nowOneJob.JobStatus == task_queue.Downloading {
  351. nowOneJob.JobStatus = task_queue.Failed
  352. bok, err := t.update(nowOneJob)
  353. if err != nil {
  354. t.log.Errorln("afterRead.update failed", err)
  355. return
  356. }
  357. if bok == false {
  358. t.log.Errorln("afterRead.update failed")
  359. return
  360. }
  361. }
  362. })
  363. }
  364. }
  365. // save 需要把改变的数据保持到 K/V 数据库中,这个没有锁,所以需要在 Sync 中使用,不对外开放
  366. func (t *TaskQueue) save(taskPriority int) error {
  367. err := GetDb().Update(
  368. func(tx *badger.Txn) error {
  369. var err error
  370. key := []byte(MergeBucketAndKeyName(BucketNamePrefixVideoSubDownloadQueue,
  371. fmt.Sprintf("%s_%d", t.queueName, taskPriority)))
  372. if err != nil {
  373. return err
  374. }
  375. b, err := t.taskPriorityMapList[taskPriority].ToJSON()
  376. if err != nil {
  377. return err
  378. }
  379. e := badger.NewEntry(key, b)
  380. err = tx.SetEntry(e)
  381. if err != nil {
  382. return err
  383. }
  384. return nil
  385. })
  386. if err != nil {
  387. return err
  388. }
  389. return nil
  390. }
  391. // isExist 是否已经存在,对内,无锁
  392. func (t *TaskQueue) isExist(jobID string) bool {
  393. _, bok := t.taskKeyMap.Get(jobID)
  394. return bok
  395. }
  396. // IsExist 是否已经存在,对外,有锁
  397. func (t *TaskQueue) IsExist(jobID string) bool {
  398. defer t.queueLock.Unlock()
  399. t.queueLock.Lock()
  400. _, bok := t.taskKeyMap.Get(jobID)
  401. return bok
  402. }
  403. // isEmpty 对内,无锁
  404. func (t *TaskQueue) isEmpty() bool {
  405. return t.taskKeyMap.Empty()
  406. }
  407. const (
  408. taskPriorityCount = 10
  409. HighTaskPriorityLevel = 3
  410. DefaultTaskPriorityLevel = 5
  411. FirstRetryTaskPriorityLevel = 6
  412. )
  413. var (
  414. ErrNotSubFound = errors.New("Not Sub Found")
  415. )