fix(subscribe): add RLock to prevent duplicate subscription downloads

This commit is contained in:
InfinityPacer
2024-12-04 11:07:45 +08:00
parent 6aa684d6a5
commit 92dacdf6a2
+10 -1
View File
@@ -1,6 +1,7 @@
import copy
import random
import time
import threading
from datetime import datetime
from typing import Dict, List, Optional, Union, Tuple
@@ -28,15 +29,17 @@ from app.log import logger
from app.schemas import NotExistMediaInfo, Notification, SubscrbieInfo, SubscribeEpisodeInfo, SubscribeDownloadFileInfo, \
SubscribeLibraryFileInfo
from app.schemas.types import MediaType, SystemConfigKey, MessageChannel, NotificationType, EventType
from app.utils.singleton import Singleton
class SubscribeChain(ChainBase):
class SubscribeChain(ChainBase, metaclass=Singleton):
"""
订阅管理处理链
"""
def __init__(self):
super().__init__()
self._rlock = threading.RLock()
self.downloadchain = DownloadChain()
self.downloadhis = DownloadHistoryOper()
self.searchchain = SearchChain()
@@ -238,6 +241,8 @@ class SubscribeChain(ChainBase):
:param manual: 是否手动搜索
:return: 更新订阅状态为R或删除订阅
"""
with self._rlock:
logger.debug(f"search lock acquired at {datetime.now()}")
if sid:
subscribe = self.subscribeoper.get(sid)
subscribes = [subscribe] if subscribe else []
@@ -412,6 +417,7 @@ class SubscribeChain(ChainBase):
self.message.put('所有订阅搜索完成!', title="订阅搜索", role="system")
else:
self.message.put('没有找到订阅!', title="订阅搜索", role="system")
logger.debug(f"search Lock released at {datetime.now()}")
def update_subscribe_priority(self, subscribe: Subscribe, meta: MetaInfo,
mediainfo: MediaInfo, downloads: List[Context]):
@@ -539,6 +545,8 @@ class SubscribeChain(ChainBase):
# 记录重新识别过的种子
_recognize_cached = []
with self._rlock:
logger.debug(f"match lock acquired at {datetime.now()}")
# 所有订阅
subscribes = self.subscribeoper.list('R')
# 遍历订阅
@@ -783,6 +791,7 @@ class SubscribeChain(ChainBase):
# 判断是否要完成订阅
self.finish_subscribe_or_not(subscribe=subscribe, meta=meta, mediainfo=mediainfo,
downloads=downloads, lefts=lefts)
logger.debug(f"match Lock released at {datetime.now()}")
def check(self):
"""