task_queue.go 12 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439
  1. package task_queue
  2. import (
  3. "encoding/json"
  4. "errors"
  5. "github.com/allanpk716/ChineseSubFinder/internal/pkg/cache_center"
  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/emby"
  11. "github.com/allanpk716/ChineseSubFinder/internal/types/task_queue"
  12. taskQueue2 "github.com/allanpk716/ChineseSubFinder/internal/types/task_queue"
  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. center *cache_center.CacheCenter // 缓存中心
  26. taskPriorityMapList []*treemap.Map // 这里有 0-10 个优先级划分的存储 List,每Add一个数据的时候需要切换到这个 List 中去 save
  27. taskKeyMap *treemap.Map // 以每个任务的唯一 JobID 来存储每个 Job 的 优先级在哪里,这样可以快速查询
  28. taskGroupBySeries *treemap.Map // 以每个任务的 SeriesRootPath 来存储每个任务,然后内层是一个 treeset,后续可以遍历删除即可
  29. queueLock sync.Mutex // 公用这个锁
  30. }
  31. func NewTaskQueue(center *cache_center.CacheCenter) *TaskQueue {
  32. tq := &TaskQueue{queueName: center.GetName(),
  33. settings: center.Settings,
  34. log: center.Log,
  35. center: center,
  36. taskPriorityMapList: make([]*treemap.Map, 0),
  37. taskKeyMap: treemap.NewWithStringComparator(),
  38. taskGroupBySeries: treemap.NewWithStringComparator(),
  39. }
  40. for i := 0; i <= taskPriorityCount; i++ {
  41. tq.taskPriorityMapList = append(tq.taskPriorityMapList, treemap.NewWithStringComparator())
  42. }
  43. tq.read()
  44. tq.afterRead()
  45. return tq
  46. }
  47. func (t *TaskQueue) Close() {
  48. t.center.Close()
  49. }
  50. func (t *TaskQueue) QueueName() string {
  51. return t.queueName
  52. }
  53. func (t *TaskQueue) Clear() error {
  54. defer t.queueLock.Unlock()
  55. t.queueLock.Lock()
  56. err := t.center.TaskQueueClear()
  57. if err != nil {
  58. return err
  59. }
  60. for i := 0; i <= taskPriorityCount; i++ {
  61. t.taskPriorityMapList[i].Clear()
  62. }
  63. t.taskKeyMap.Clear()
  64. t.taskGroupBySeries.Clear()
  65. return nil
  66. }
  67. // Size 队列的长度,对外暴露,有锁
  68. func (t *TaskQueue) Size() int {
  69. defer t.queueLock.Unlock()
  70. t.queueLock.Lock()
  71. return t.taskKeyMap.Size()
  72. }
  73. // checkPriority 检测优先级,会校验范围
  74. func (t *TaskQueue) checkPriority(oneJob taskQueue2.OneJob) taskQueue2.OneJob {
  75. if oneJob.TaskPriority > taskPriorityCount {
  76. oneJob.TaskPriority = taskPriorityCount
  77. }
  78. if oneJob.TaskPriority < 0 {
  79. oneJob.TaskPriority = 0
  80. }
  81. return oneJob
  82. }
  83. // degrade 降一级,会校验范围
  84. func (t *TaskQueue) degrade(oneJob taskQueue2.OneJob) taskQueue2.OneJob {
  85. oneJob.TaskPriority -= 1
  86. return t.checkPriority(oneJob)
  87. }
  88. // Add 放入元素,放入的时候会根据 TaskPriority 进行归类,存在的不会新增和更新
  89. func (t *TaskQueue) Add(oneJob task_queue.OneJob) (bool, error) {
  90. defer t.queueLock.Unlock()
  91. t.queueLock.Lock()
  92. if t.isExist(oneJob.Id) == true {
  93. return false, nil
  94. }
  95. // 检查权限范围
  96. oneJob = t.checkPriority(oneJob)
  97. // 插入到统一的 KeyMap
  98. t.taskKeyMap.Put(oneJob.Id, oneJob.TaskPriority)
  99. // 分配到具体的优先级 map 中
  100. t.taskPriorityMapList[oneJob.TaskPriority].Put(oneJob.Id, oneJob)
  101. // 如果是连续剧,则需要存储到 taskGroupBySeries 中
  102. jobIDSet, found := t.taskGroupBySeries.Get(oneJob.SeriesRootDirPath)
  103. if found == false {
  104. // 不存在
  105. nowJobIDSet := treeset.NewWithStringComparator()
  106. nowJobIDSet.Add(oneJob.Id)
  107. t.taskGroupBySeries.Put(oneJob.SeriesRootDirPath, nowJobIDSet)
  108. } else {
  109. // 存在
  110. nowJobIDSet := jobIDSet.(*treeset.Set)
  111. nowJobIDSet.Add(oneJob.Id)
  112. t.taskGroupBySeries.Put(oneJob.SeriesRootDirPath, nowJobIDSet)
  113. }
  114. err := t.save(oneJob.TaskPriority)
  115. if err != nil {
  116. return false, err
  117. }
  118. return true, nil
  119. }
  120. // update 更新素,不存在则会失败,内部用,没有锁
  121. func (t *TaskQueue) update(oneJob task_queue.OneJob) (bool, error) {
  122. if t.isExist(oneJob.Id) == false {
  123. return false, nil
  124. }
  125. // 自动更新时间
  126. oneJob.UpdateTime = (emby.Time)(time.Now())
  127. // 这里需要判断是否有优先级的 Update,如果有就需要把之前缓存的表给更新
  128. // 然后再插入到新的表中
  129. taskPriorityIndex, _ := t.taskKeyMap.Get(oneJob.Id)
  130. // 检查权限范围
  131. oneJob = t.checkPriority(oneJob)
  132. if oneJob.TaskPriority != taskPriorityIndex {
  133. // 优先级修改
  134. // 先删除原有的优先级
  135. t.taskPriorityMapList[taskPriorityIndex.(int)].Remove(oneJob.Id)
  136. err := t.save(taskPriorityIndex.(int))
  137. if err != nil {
  138. return false, err
  139. }
  140. }
  141. // 插入到统一的 KeyMap
  142. t.taskKeyMap.Put(oneJob.Id, oneJob.TaskPriority)
  143. // 分配到具体的优先级 map 中
  144. t.taskPriorityMapList[oneJob.TaskPriority].Put(oneJob.Id, oneJob)
  145. err := t.save(oneJob.TaskPriority)
  146. if err != nil {
  147. return false, err
  148. }
  149. return true, nil
  150. }
  151. // Update 更新素,不存在则会失败
  152. func (t *TaskQueue) Update(oneJob task_queue.OneJob) (bool, error) {
  153. defer t.queueLock.Unlock()
  154. t.queueLock.Lock()
  155. return t.update(oneJob)
  156. }
  157. // AutoDetectUpdateJobStatus 根据任务的生命周期图,进行自动判断更新,见《任务的生命周期》流程图
  158. func (t *TaskQueue) AutoDetectUpdateJobStatus(oneJob task_queue.OneJob, inErr error) {
  159. defer t.queueLock.Unlock()
  160. t.queueLock.Lock()
  161. // 检查权限范围
  162. oneJob = t.checkPriority(oneJob)
  163. if inErr == nil {
  164. // 如果任务的优先级是 0,那么这个任务就认为是一次性任务,下载完毕不管如何都会设置为 ignore
  165. if oneJob.TaskPriority == 0 {
  166. oneJob.JobStatus = task_queue.Ignore
  167. }
  168. // 没有错误就是完成
  169. oneJob.TaskPriority = DefaultTaskPriorityLevel
  170. oneJob.JobStatus = taskQueue2.Done
  171. oneJob.DownloadTimes += 1
  172. } else {
  173. // 超过了时间限制,默认是 90 天, A.Before(B) : A < B == true
  174. if (time.Time)(oneJob.AddedTime).AddDate(0, 0, t.settings.AdvancedSettings.TaskQueue.ExpirationTime).Before(time.Now()) == true {
  175. // 超过 90 天了
  176. oneJob.JobStatus = taskQueue2.Failed
  177. } else {
  178. // 还在 90 天内
  179. // 是否是首次,那么就看它的 Level 是否是在 5,然后 retry == 0
  180. if oneJob.TaskPriority == DefaultTaskPriorityLevel && oneJob.RetryTimes == 0 {
  181. // 需要重置到 L6
  182. oneJob.RetryTimes = 0
  183. oneJob.TaskPriority = FirstRetryTaskPriorityLevel
  184. } else {
  185. if oneJob.RetryTimes > t.settings.AdvancedSettings.TaskQueue.MaxRetryTimes {
  186. // 超过重试次数会进行一次降级,然后重置这个次数
  187. oneJob.RetryTimes = 0
  188. oneJob = t.degrade(oneJob)
  189. }
  190. }
  191. // 强制为 waiting
  192. oneJob.JobStatus = taskQueue2.Waiting
  193. }
  194. // 如果任务的优先级是 0,那么这个任务就认为是一次性任务,下载完毕不管如何都会设置为 ignore
  195. if oneJob.TaskPriority == 0 {
  196. oneJob.JobStatus = task_queue.Ignore
  197. }
  198. // 传入的错误需要放进来
  199. oneJob.ErrorInfo = inErr.Error()
  200. oneJob.DownloadTimes += 1
  201. }
  202. // 只要是进入完成标记流程的任务,如果优先级还是很高,那么就需要重置到默认优先级上
  203. if oneJob.TaskPriority < DefaultTaskPriorityLevel {
  204. oneJob.TaskPriority = DefaultTaskPriorityLevel
  205. }
  206. // 这里不要用错了,要用无锁的,不然会阻塞
  207. bok, err := t.update(oneJob)
  208. if err != nil {
  209. t.log.Errorln("AutoDetectUpdateJobStatus", oneJob.VideoFPath, err)
  210. return
  211. }
  212. if bok == false {
  213. t.log.Warningln("AutoDetectUpdateJobStatus ==", oneJob.VideoFPath, "Job.ID", oneJob.Id, "Not Found")
  214. return
  215. }
  216. }
  217. func (t *TaskQueue) del(jobId string) (bool, error) {
  218. if t.isExist(jobId) == false {
  219. return false, nil
  220. }
  221. taskPriority, bok := t.taskKeyMap.Get(jobId)
  222. if bok == false {
  223. return false, nil
  224. }
  225. // 删除连续剧的 tree.Map 里面的 tree.Set 的元素
  226. needDelJobObj, bok := t.taskPriorityMapList[taskPriority.(int)].Get(jobId)
  227. if bok == false {
  228. return false, nil
  229. }
  230. needDelJob := needDelJobObj.(task_queue.OneJob)
  231. jobSetsObj, bok := t.taskGroupBySeries.Get(needDelJob.SeriesRootDirPath)
  232. if bok == false {
  233. return false, nil
  234. }
  235. jobSets := jobSetsObj.(*treeset.Set)
  236. jobSets.Remove(jobId)
  237. // 删除任务
  238. t.taskKeyMap.Remove(jobId)
  239. t.taskPriorityMapList[taskPriority.(int)].Remove(jobId)
  240. err := t.save(taskPriority.(int))
  241. if err != nil {
  242. return false, err
  243. }
  244. // 删除任务的时候也需要删除对应的日志
  245. pathRoot := filepath.Join(global_value.ConfigRootDirFPath(), "Logs")
  246. fileFPath := filepath.Join(pathRoot, common.OnceLogPrefix+jobId+".log")
  247. if my_util.IsFile(fileFPath) == true {
  248. err = os.Remove(fileFPath)
  249. if err != nil {
  250. t.log.Errorln("del job", jobId, "logfile,error:", err)
  251. }
  252. }
  253. return true, nil
  254. }
  255. // Del 删除一个元素
  256. func (t *TaskQueue) Del(jobId string) (bool, error) {
  257. defer t.queueLock.Unlock()
  258. t.queueLock.Lock()
  259. return t.del(jobId)
  260. }
  261. func (t *TaskQueue) read() {
  262. taskQueueRead, err := t.center.TaskQueueRead()
  263. if err != nil {
  264. t.log.Errorln("read task queue TaskQueueRead error:", err)
  265. return
  266. }
  267. for i := 0; i <= taskPriorityCount; i++ {
  268. value, bok := taskQueueRead[i]
  269. if bok == false {
  270. continue
  271. }
  272. err = t.taskPriorityMapList[i].FromJSON(value)
  273. if err != nil {
  274. t.log.Errorln("read task queue FromJSON error:", err)
  275. }
  276. // 上面的操作仅仅是把 OneJob 的 JSON 弄了出来,还需要转换为 OneJob 的结构体
  277. // JobID - OneJob
  278. t.taskPriorityMapList[i].Each(func(key interface{}, value interface{}) {
  279. jsonString, err := json.Marshal(value)
  280. if err != nil {
  281. t.log.Panicln(err)
  282. }
  283. nowOneJob := task_queue.OneJob{}
  284. err = json.Unmarshal(jsonString, &nowOneJob)
  285. if err != nil {
  286. t.log.Panicln(err)
  287. }
  288. t.taskPriorityMapList[i].Put(key, nowOneJob)
  289. })
  290. // 需要把几个优先级的map中的key汇总
  291. // JobID - OneJob
  292. t.taskPriorityMapList[i].Each(func(key interface{}, value interface{}) {
  293. // JobID -- taskPriority
  294. t.taskKeyMap.Put(key, i)
  295. // SeriesRootDirPath -- tree.Set(JobID)
  296. oneJob := value.(task_queue.OneJob)
  297. jobIDSet, found := t.taskGroupBySeries.Get(oneJob.SeriesRootDirPath)
  298. if found == false {
  299. // 不存在
  300. nowJobIDSet := treeset.NewWithStringComparator()
  301. nowJobIDSet.Add(oneJob.Id)
  302. t.taskGroupBySeries.Put(oneJob.SeriesRootDirPath, nowJobIDSet)
  303. } else {
  304. // 存在
  305. nowJobIDSet := jobIDSet.(*treeset.Set)
  306. nowJobIDSet.Add(oneJob.Id)
  307. t.taskGroupBySeries.Put(oneJob.SeriesRootDirPath, nowJobIDSet)
  308. }
  309. })
  310. }
  311. }
  312. func (t *TaskQueue) afterRead() {
  313. // 将 downloading 的任务重置为 waiting
  314. for TaskPriority := 0; TaskPriority <= taskPriorityCount; TaskPriority++ {
  315. t.taskPriorityMapList[TaskPriority].Each(func(key interface{}, value interface{}) {
  316. nowOneJob := value.(task_queue.OneJob)
  317. if nowOneJob.JobStatus == task_queue.Downloading {
  318. nowOneJob.JobStatus = task_queue.Waiting
  319. nowOneJob.DownloadTimes += 1
  320. bok, err := t.update(nowOneJob)
  321. if err != nil {
  322. t.log.Errorln("afterRead.update failed", err)
  323. return
  324. }
  325. if bok == false {
  326. t.log.Errorln("afterRead.update failed")
  327. return
  328. }
  329. }
  330. })
  331. }
  332. }
  333. // save 需要把改变的数据保持到 K/V 数据库中,这个没有锁,所以需要在 Sync 中使用,不对外开放
  334. func (t *TaskQueue) save(taskPriority int) error {
  335. b, err := t.taskPriorityMapList[taskPriority].ToJSON()
  336. if err != nil {
  337. return err
  338. }
  339. err = t.center.TaskQueueSave(taskPriority, b)
  340. if err != nil {
  341. return err
  342. }
  343. return nil
  344. }
  345. // isExist 是否已经存在,对内,无锁
  346. func (t *TaskQueue) isExist(jobID string) bool {
  347. _, bok := t.taskKeyMap.Get(jobID)
  348. return bok
  349. }
  350. // IsExist 是否已经存在,对外,有锁
  351. func (t *TaskQueue) IsExist(jobID string) bool {
  352. defer t.queueLock.Unlock()
  353. t.queueLock.Lock()
  354. _, bok := t.taskKeyMap.Get(jobID)
  355. return bok
  356. }
  357. // isEmpty 对内,无锁
  358. func (t *TaskQueue) isEmpty() bool {
  359. return t.taskKeyMap.Empty()
  360. }
  361. const (
  362. taskPriorityCount = 10
  363. HighTaskPriorityLevel = 3
  364. DefaultTaskPriorityLevel = 5
  365. FirstRetryTaskPriorityLevel = 6
  366. LowTaskPriorityLevel = 7
  367. )
  368. var (
  369. ErrNoSubFound = errors.New("No Sub Found")
  370. )