task_queue.go 13 KB

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