task_queue.go 13 KB

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