fix chain

This commit is contained in:
jxxghp
2025-06-26 13:16:10 +08:00
parent c440ce3045
commit 2700e639f1
6 changed files with 631 additions and 623 deletions
+6 -7
View File
@@ -202,16 +202,15 @@ class SearchChain(ChainBase):
# 过滤完成 # 过滤完成
progress.update(value=50, text=f'过滤完成,剩余 {len(torrents)} 个资源', key=ProgressKey.Search) progress.update(value=50, text=f'过滤完成,剩余 {len(torrents)} 个资源', key=ProgressKey.Search)
# 开始匹配
_match_torrents = []
# 总数 # 总数
_total = len(torrents) _total = len(torrents)
# 已处理数 # 已处理数
_count = 0 _count = 0
# 开始匹配
_match_torrents = []
torrenthelper = TorrentHelper() torrenthelper = TorrentHelper()
try:
if mediainfo:
# 英文标题应该在别名/原标题中,不需要再匹配 # 英文标题应该在别名/原标题中,不需要再匹配
logger.info(f"开始匹配结果 标题:{mediainfo.title},原标题:{mediainfo.original_title},别名:{mediainfo.names}") logger.info(f"开始匹配结果 标题:{mediainfo.title},原标题:{mediainfo.original_title},别名:{mediainfo.names}")
progress.update(value=51, text=f'开始匹配,总 {_total} 个资源 ...', key=ProgressKey.Search) progress.update(value=51, text=f'开始匹配,总 {_total} 个资源 ...', key=ProgressKey.Search)
@@ -256,16 +255,16 @@ class SearchChain(ChainBase):
progress.update(value=97, progress.update(value=97,
text=f'匹配完成,共匹配到 {len(_match_torrents)} 个资源', text=f'匹配完成,共匹配到 {len(_match_torrents)} 个资源',
key=ProgressKey.Search) key=ProgressKey.Search)
else:
_match_torrents = [(t, MetaInfo(title=t.title, subtitle=t.description)) for t in torrents]
# 去掉mediainfo中多余的数据 # 去掉mediainfo中多余的数据
mediainfo.clear() mediainfo.clear()
# 组装上下文 # 组装上下文
contexts = [Context(torrent_info=t[0], contexts = [Context(torrent_info=t[0],
media_info=mediainfo, media_info=mediainfo,
meta_info=t[1]) for t in _match_torrents] meta_info=t[1]) for t in _match_torrents]
finally:
torrents.clear()
_match_torrents.clear()
# 排序 # 排序
progress.update(value=99, progress.update(value=99,
+1 -2
View File
@@ -92,10 +92,9 @@ class SiteChain(ChainBase):
""" """
刷新所有站点的用户数据 刷新所有站点的用户数据
""" """
sites = SitesHelper().get_indexers()
any_site_updated = False any_site_updated = False
result = {} result = {}
for site in sites: for site in SitesHelper().get_indexers():
if global_vars.is_system_stopped: if global_vars.is_system_stopped:
return None return None
if site.get("is_active"): if site.get("is_active"):
+14 -15
View File
@@ -286,6 +286,8 @@ class SubscribeChain(ChainBase):
subscribes = [subscribe] if subscribe else [] subscribes = [subscribe] if subscribe else []
else: else:
subscribes = subscribeoper.list(self.get_states_for_search(state)) subscribes = subscribeoper.list(self.get_states_for_search(state))
try:
# 遍历订阅 # 遍历订阅
for subscribe in subscribes: for subscribe in subscribes:
if global_vars.is_system_stopped: if global_vars.is_system_stopped:
@@ -426,7 +428,10 @@ class SubscribeChain(ChainBase):
self.messagehelper.put('所有订阅搜索完成!', title="订阅搜索", role="system") self.messagehelper.put('所有订阅搜索完成!', title="订阅搜索", role="system")
else: else:
self.messagehelper.put('没有找到订阅!', title="订阅搜索", role="system") self.messagehelper.put('没有找到订阅!', title="订阅搜索", role="system")
logger.debug(f"search Lock released at {datetime.now()}") logger.debug(f"search Lock released at {datetime.now()}")
finally:
subscribes.clear()
def update_subscribe_priority(self, subscribe: Subscribe, meta: MetaBase, def update_subscribe_priority(self, subscribe: Subscribe, meta: MetaBase,
mediainfo: MediaInfo, downloads: Optional[List[Context]]): mediainfo: MediaInfo, downloads: Optional[List[Context]]):
@@ -527,13 +532,9 @@ class SubscribeChain(ChainBase):
获取订阅中涉及的所有站点清单(节约资源) 获取订阅中涉及的所有站点清单(节约资源)
:return: 返回[]代表所有站点命中,返回None代表没有订阅 :return: 返回[]代表所有站点命中,返回None代表没有订阅
""" """
# 查询所有订阅
subscribes = SubscribeOper().list(self.get_states_for_search('R'))
if not subscribes:
return None
ret_sites = [] ret_sites = []
# 刷新订阅选中的Rss站点 # 刷新订阅选中的Rss站点
for subscribe in subscribes: for subscribe in SubscribeOper().list(self.get_states_for_search('R')):
# 刷新选中的站点 # 刷新选中的站点
ret_sites.extend(self.get_sub_sites(subscribe)) ret_sites.extend(self.get_sub_sites(subscribe))
# 去重 # 去重
@@ -551,9 +552,8 @@ class SubscribeChain(ChainBase):
return return
with self._rlock: with self._rlock:
try:
logger.debug(f"match lock acquired at {datetime.now()}") logger.debug(f"match lock acquired at {datetime.now()}")
# 所有订阅
subscribes = SubscribeOper().list(self.get_states_for_search('R'))
# 预识别所有未识别的种子 # 预识别所有未识别的种子
processed_torrents: Dict[str, List[Context]] = {} processed_torrents: Dict[str, List[Context]] = {}
@@ -576,7 +576,8 @@ class SubscribeChain(ChainBase):
# 添加已预处理 # 添加已预处理
processed_torrents[domain].append(context) processed_torrents[domain].append(context)
# 遍历订阅 # 所有订阅
subscribes = SubscribeOper().list(self.get_states_for_search('R'))
for subscribe in subscribes: for subscribe in subscribes:
if global_vars.is_system_stopped: if global_vars.is_system_stopped:
break break
@@ -803,6 +804,9 @@ class SubscribeChain(ChainBase):
downloads=downloads, lefts=lefts) downloads=downloads, lefts=lefts)
logger.debug(f"match Lock released at {datetime.now()}") logger.debug(f"match Lock released at {datetime.now()}")
finally:
subscribes.clear()
processed_torrents.clear()
def check(self): def check(self):
""" """
@@ -810,12 +814,8 @@ class SubscribeChain(ChainBase):
""" """
# 查询所有订阅 # 查询所有订阅
subscribeoper = SubscribeOper() subscribeoper = SubscribeOper()
subscribes = subscribeoper.list()
if not subscribes:
# 没有订阅不运行
return
# 遍历订阅 # 遍历订阅
for subscribe in subscribes: for subscribe in subscribeoper.list():
if global_vars.is_system_stopped: if global_vars.is_system_stopped:
break break
logger.info(f'开始更新订阅元数据:{subscribe.name} ...') logger.info(f'开始更新订阅元数据:{subscribe.name} ...')
@@ -871,11 +871,10 @@ class SubscribeChain(ChainBase):
follow_users: List[str] = SystemConfigOper().get(SystemConfigKey.FollowSubscribers) follow_users: List[str] = SystemConfigOper().get(SystemConfigKey.FollowSubscribers)
if not follow_users: if not follow_users:
return return
share_subs = SubscribeHelper().get_shares()
logger.info(f'开始刷新follow用户分享订阅 ...') logger.info(f'开始刷新follow用户分享订阅 ...')
success_count = 0 success_count = 0
subscribeoper = SubscribeOper() subscribeoper = SubscribeOper()
for share_sub in share_subs: for share_sub in SubscribeHelper().get_shares():
if global_vars.is_system_stopped: if global_vars.is_system_stopped:
break break
uid = share_sub.get("share_uid") uid = share_sub.get("share_uid")
+8 -5
View File
@@ -98,6 +98,7 @@ class TorrentsChain(ChainBase):
if not site.get("rss"): if not site.get("rss"):
logger.error(f'站点 {domain} 未配置RSS地址!') logger.error(f'站点 {domain} 未配置RSS地址!')
return [] return []
# 解析RSS
rss_items = RssHelper().parse(site.get("rss"), True if site.get("proxy") else False, rss_items = RssHelper().parse(site.get("rss"), True if site.get("proxy") else False,
timeout=int(site.get("timeout") or 30)) timeout=int(site.get("timeout") or 30))
if rss_items is None: if rss_items is None:
@@ -109,6 +110,7 @@ class TorrentsChain(ChainBase):
return [] return []
# 组装种子 # 组装种子
ret_torrents: List[TorrentInfo] = [] ret_torrents: List[TorrentInfo] = []
try:
for item in rss_items: for item in rss_items:
if not item.get("title"): if not item.get("title"):
continue continue
@@ -127,7 +129,8 @@ class TorrentsChain(ChainBase):
pubdate=item["pubdate"].strftime("%Y-%m-%d %H:%M:%S") if item.get("pubdate") else None, pubdate=item["pubdate"].strftime("%Y-%m-%d %H:%M:%S") if item.get("pubdate") else None,
) )
ret_torrents.append(torrentinfo) ret_torrents.append(torrentinfo)
finally:
rss_items.clear()
return ret_torrents return ret_torrents
def refresh(self, stype: Optional[str] = None, sites: List[int] = None) -> Dict[str, List[Context]]: def refresh(self, stype: Optional[str] = None, sites: List[int] = None) -> Dict[str, List[Context]]:
@@ -152,13 +155,10 @@ class TorrentsChain(ChainBase):
torrents_cache[_domain] = [_torrent for _torrent in _torrents torrents_cache[_domain] = [_torrent for _torrent in _torrents
if not TorrentHelper().is_invalid(_torrent.torrent_info.enclosure)] if not TorrentHelper().is_invalid(_torrent.torrent_info.enclosure)]
# 所有站点索引
indexers = SitesHelper().get_indexers()
# 需要刷新的站点domain # 需要刷新的站点domain
domains = [] domains = []
# 遍历站点缓存资源 # 遍历站点缓存资源
for indexer in indexers: for indexer in SitesHelper().get_indexers():
if global_vars.is_system_stopped: if global_vars.is_system_stopped:
break break
# 未开启的站点不刷新 # 未开启的站点不刷新
@@ -187,6 +187,7 @@ class TorrentsChain(ChainBase):
else: else:
logger.info(f'{indexer.get("name")} 没有新种子') logger.info(f'{indexer.get("name")} 没有新种子')
continue continue
try:
for torrent in torrents: for torrent in torrents:
if global_vars.is_system_stopped: if global_vars.is_system_stopped:
break break
@@ -217,6 +218,8 @@ class TorrentsChain(ChainBase):
# 如果超过了限制条数则移除掉前面的 # 如果超过了限制条数则移除掉前面的
if len(torrents_cache[domain]) > settings.CONF["torrents"]: if len(torrents_cache[domain]) > settings.CONF["torrents"]:
torrents_cache[domain] = torrents_cache[domain][-settings.CONF["torrents"]:] torrents_cache[domain] = torrents_cache[domain][-settings.CONF["torrents"]:]
finally:
torrents.clear()
else: else:
logger.info(f'{indexer.get("name")} 没有获取到种子') logger.info(f'{indexer.get("name")} 没有获取到种子')
+10 -1
View File
@@ -788,6 +788,7 @@ class TransferChain(ChainBase, metaclass=Singleton):
for dir_info in download_dirs): for dir_info in download_dirs):
return True return True
logger.info("开始整理下载器中已经完成下载的文件 ...") logger.info("开始整理下载器中已经完成下载的文件 ...")
# 从下载器获取种子列表 # 从下载器获取种子列表
torrents: Optional[List[TransferTorrent]] = self.list_torrents(status=TorrentStatus.TRANSFER) torrents: Optional[List[TransferTorrent]] = self.list_torrents(status=TorrentStatus.TRANSFER)
if not torrents: if not torrents:
@@ -796,6 +797,7 @@ class TransferChain(ChainBase, metaclass=Singleton):
logger.info(f"获取到 {len(torrents)} 个已完成的下载任务") logger.info(f"获取到 {len(torrents)} 个已完成的下载任务")
try:
for torrent in torrents: for torrent in torrents:
if global_vars.is_system_stopped: if global_vars.is_system_stopped:
break break
@@ -860,6 +862,8 @@ class TransferChain(ChainBase, metaclass=Singleton):
if not state: if not state:
logger.warn(f"整理下载器任务失败:{torrent.hash} - {errmsg}") logger.warn(f"整理下载器任务失败:{torrent.hash} - {errmsg}")
self.transfer_completed(hashs=torrent.hash, downloader=torrent.downloader) self.transfer_completed(hashs=torrent.hash, downloader=torrent.downloader)
finally:
torrents.clear()
# 结束 # 结束
logger.info("所有下载器中下载完成的文件已整理完成") logger.info("所有下载器中下载完成的文件已整理完成")
@@ -1032,6 +1036,7 @@ class TransferChain(ChainBase, metaclass=Singleton):
# 整理所有文件 # 整理所有文件
transfer_tasks: List[TransferTask] = [] transfer_tasks: List[TransferTask] = []
try:
for file_item, bluray_dir in file_items: for file_item, bluray_dir in file_items:
if global_vars.is_system_stopped: if global_vars.is_system_stopped:
break break
@@ -1137,6 +1142,8 @@ class TransferChain(ChainBase, metaclass=Singleton):
# 加入列表 # 加入列表
self.__put_to_jobview(transfer_task) self.__put_to_jobview(transfer_task)
transfer_tasks.append(transfer_task) transfer_tasks.append(transfer_task)
finally:
file_items.clear()
# 实时整理 # 实时整理
if transfer_tasks: if transfer_tasks:
@@ -1155,7 +1162,7 @@ class TransferChain(ChainBase, metaclass=Singleton):
progress.update(value=0, progress.update(value=0,
text=__process_msg, text=__process_msg,
key=ProgressKey.FileTransfer) key=ProgressKey.FileTransfer)
try:
for transfer_task in transfer_tasks: for transfer_task in transfer_tasks:
if global_vars.is_system_stopped: if global_vars.is_system_stopped:
break break
@@ -1178,6 +1185,8 @@ class TransferChain(ChainBase, metaclass=Singleton):
fail_num += 1 fail_num += 1
else: else:
processed_num += 1 processed_num += 1
finally:
transfer_tasks.clear()
# 整理结束 # 整理结束
__end_msg = f"整理队列处理完成,共整理 {total_num} 个文件,失败 {fail_num}" __end_msg = f"整理队列处理完成,共整理 {total_num} 个文件,失败 {fail_num}"
+1 -2
View File
@@ -221,8 +221,7 @@ class Base:
@classmethod @classmethod
@db_query @db_query
def list(cls, db: Session) -> List[Self]: def list(cls, db: Session) -> List[Self]:
result = db.query(cls).all() return db.query(cls).all()
return list(result)
def to_dict(self): def to_dict(self):
return {c.name: getattr(self, c.name, None) for c in self.__table__.columns} # noqa return {c.name: getattr(self, c.name, None) for c in self.__table__.columns} # noqa