From b45597866b3b7d13e1c813099ffa06c7484e0cf5 Mon Sep 17 00:00:00 2001 From: jxxghp Date: Sat, 22 Aug 2026 21:07:18 +0800 Subject: [PATCH] refactor: split subscription creation stages --- app/chain/subscribe.py | 780 ++++++++++-------- .../backend-architecture-next-stage.md | 5 + .../architecture/complexity-baseline.json | 2 - 3 files changed, 452 insertions(+), 335 deletions(-) diff --git a/app/chain/subscribe.py b/app/chain/subscribe.py index e551e69b7..4a9d2ce20 100644 --- a/app/chain/subscribe.py +++ b/app/chain/subscribe.py @@ -110,6 +110,30 @@ class _SubscribePostCommitContext: message: bool +@dataclass(slots=True) +class _SubscribeCreateContext: + """订阅新增各阶段共享的显式状态,避免同步与异步入口各自维护散落变量。""" + + title: str + year: str + mtype: Optional[MediaType] + episode_group: Optional[str] + season: Optional[int] + channel: Optional[NotificationChannel] + source: Optional[str] + userid: Optional[str] + username: Optional[str] + message: bool + exist_ok: bool + options: Dict[str, Any] + explicit_identity: bool + media_source: Optional[MediaSource] + media_id: Optional[str] + requested_music_type: Optional[str] + metainfo: MetaBase + mediainfo: Optional[MediaInfo] = None + + def _system_config(): """返回配置端口,并兼容旧测试对本地别名的替换。""" if SystemConfigOper is not _DEFAULT_SYSTEM_CONFIG_PROVIDER: @@ -973,6 +997,403 @@ class SubscribeChain(MusicSubscribeMixin, InteractionChainMixin, ChainBase): self.__subscribe_report_payload(context) ) + @staticmethod + def __build_subscribe_create_context( + title: str, + year: str, + mtype: Optional[MediaType], + episode_group: Optional[str], + season: Optional[int], + channel: Optional[NotificationChannel], + source: Optional[str], + userid: Optional[str], + username: Optional[str], + message: bool, + exist_ok: bool, + media_source: Optional[MediaSource], + media_id: Optional[str], + options: Dict[str, Any], + ) -> Tuple[Optional[_SubscribeCreateContext], Optional[str]]: + """规范订阅新增输入,并在任何识别 I/O 前拒绝不完整的显式媒体身份。""" + explicit_identity = media_source is not None or media_id is not None + media_source, media_id = resolve_media_identity( + media_source=media_source, + media_id=media_id, + ) + if explicit_identity and (not media_source or not media_id): + return None, "媒体来源和媒体 ID 必须同时提供" + + metainfo = MetaMusic.parse_query(title) if mtype == MediaType.MUSIC else MetaInfo(title) + if year: + metainfo.year = year + if mtype: + metainfo.type = mtype + if season is not None: + metainfo.type = MediaType.TV + metainfo.begin_season = season + if mtype == MediaType.MUSIC and media_id: + metainfo.media_id = str(media_id) + + return _SubscribeCreateContext( + title=title, + year=year, + mtype=mtype, + episode_group=episode_group, + season=season, + channel=channel, + source=source, + userid=userid, + username=username, + message=bool(message), + exist_ok=bool(exist_ok), + options=options, + explicit_identity=explicit_identity, + media_source=media_source, + media_id=media_id, + requested_music_type=options.get("music_type"), + metainfo=metainfo, + ), None + + @staticmethod + def __normalize_recognized_subscribe_media(context: _SubscribeCreateContext) -> None: + """保留非 TMDB 影视标题清洗与从标题补季号的历史行为。""" + mediainfo = context.mediainfo + if ( + context.mtype == MediaType.MUSIC + or not mediainfo + or mediainfo.media_source == MediaSource.TMDB + ): + return + meta = MetaInfo(mediainfo.title) + mediainfo.title = meta.name + if context.season is None: + context.season = meta.begin_season + + def __recognize_subscribe_media(self, context: _SubscribeCreateContext) -> Optional[str]: + """同步识别订阅目标;显式身份失败时禁止按标题换成另一个媒体。""" + if context.media_source and context.media_id: + context.mediainfo = MediaChain().recognize_media( + meta=context.metainfo, + mtype=context.mtype, + media_source=context.media_source, + media_id=context.media_id, + music_type=context.requested_music_type, + episode_group=context.episode_group, + cache=False, + ) + self.__normalize_recognized_subscribe_media(context) + if not context.mediainfo and not context.explicit_identity: + context.mediainfo = MediaChain().recognize_by_meta( + context.metainfo, + media_source=context.media_source, + episode_group=context.episode_group, + obtain_images=False, + music_type=context.requested_music_type, + ) + if ( + context.mtype == MediaType.MUSIC + and context.mediainfo + and not context.mediainfo.media_source + ): + context.mediainfo = None + return self.__validate_recognized_subscribe_media(context) + + async def __async_recognize_subscribe_media( + self, + context: _SubscribeCreateContext, + ) -> Optional[str]: + """异步识别订阅目标,并与同步入口共享相同的规范化和校验规则。""" + if context.media_source and context.media_id: + context.mediainfo = await MediaChain().async_recognize_media( + meta=context.metainfo, + mtype=context.mtype, + media_source=context.media_source, + media_id=context.media_id, + music_type=context.requested_music_type, + episode_group=context.episode_group, + cache=False, + ) + self.__normalize_recognized_subscribe_media(context) + if not context.mediainfo and not context.explicit_identity: + context.mediainfo = await MediaChain().async_recognize_by_meta( + context.metainfo, + media_source=context.media_source, + episode_group=context.episode_group, + obtain_images=False, + music_type=context.requested_music_type, + ) + if ( + context.mtype == MediaType.MUSIC + and context.mediainfo + and not context.mediainfo.media_source + ): + context.mediainfo = None + return self.__validate_recognized_subscribe_media(context) + + def __validate_recognized_subscribe_media( + self, + context: _SubscribeCreateContext, + ) -> Optional[str]: + """校验识别结果以及音乐订阅实体,返回兼容旧入口的错误文案。""" + if not context.mediainfo: + logger.warning( + f"未识别到媒体信息,标题:{context.title},媒体来源:{context.media_source}," + f"媒体 ID:{context.media_id}" + ) + return "未识别到媒体信息" + if context.mtype != MediaType.MUSIC: + return None + music_error = self._validate_music_subscribe_target( + context.mediainfo, + requested_music_type=context.requested_music_type, + ) + if music_error: + logger.warning(f"音乐订阅目标校验失败:{context.title} - {music_error}") + return music_error + + def __prepare_subscribe_episodes(self, context: _SubscribeCreateContext) -> Optional[str]: + """同步补齐电视剧季集信息,并保持外部集数刷新只能扩展创建目标。""" + mediainfo = context.mediainfo + if mediainfo.type != MediaType.TV: + context.season = None + return None + if context.season is None: + context.season = 1 + if not context.options.get("total_episode"): + if not mediainfo.seasons or context.episode_group: + mediainfo = MediaChain().recognize_media( + mtype=mediainfo.type, + **_media_recognize_kwargs(mediainfo), + episode_group=context.episode_group, + cache=False, + ) + context.mediainfo = mediainfo + error = self.__validate_subscribe_seasons(context) + if error: + return error + current_total_episode = len(mediainfo.seasons.get(context.season) or []) + total_episode = self.__apply_episodes_refresh( + current_total_episode, + season=context.season, + mediainfo=mediainfo, + media_source=resolve_media_identity(media=mediainfo)[0], + media_id=resolve_media_identity(media=mediainfo)[1], + scene="create", + ) + error = self.__store_subscribe_episode_total( + context, + current_total_episode, + total_episode, + ) + if error: + return error + self.__fill_subscribe_lack_episode(context) + return None + + async def __async_prepare_subscribe_episodes( + self, + context: _SubscribeCreateContext, + ) -> Optional[str]: + """异步补齐电视剧季集信息,并复用同步入口的结果校验和字段写入规则。""" + mediainfo = context.mediainfo + if mediainfo.type != MediaType.TV: + context.season = None + return None + if context.season is None: + context.season = 1 + if not context.options.get("total_episode"): + if not mediainfo.seasons or context.episode_group: + mediainfo = await MediaChain().async_recognize_media( + mtype=mediainfo.type, + **_media_recognize_kwargs(mediainfo), + episode_group=context.episode_group, + cache=False, + ) + context.mediainfo = mediainfo + error = self.__validate_subscribe_seasons(context) + if error: + return error + current_total_episode = len(mediainfo.seasons.get(context.season) or []) + total_episode = await self.__async_apply_episodes_refresh( + current_total_episode, + season=context.season, + mediainfo=mediainfo, + media_source=resolve_media_identity(media=mediainfo)[0], + media_id=resolve_media_identity(media=mediainfo)[1], + scene="create", + ) + error = self.__store_subscribe_episode_total( + context, + current_total_episode, + total_episode, + ) + if error: + return error + self.__fill_subscribe_lack_episode(context) + return None + + @staticmethod + def __validate_subscribe_seasons(context: _SubscribeCreateContext) -> Optional[str]: + """校验补充识别结果是否仍包含创建电视剧订阅所需的季集信息。""" + if not context.mediainfo: + logger.error("媒体信息识别失败!") + return "媒体信息识别失败" + if not context.mediainfo.seasons: + logger.error(f"媒体信息中没有季集信息,标题:{context.title}") + return "媒体信息中没有季集信息" + return None + + @staticmethod + def __store_subscribe_episode_total( + context: _SubscribeCreateContext, + current_total_episode: int, + total_episode: int, + ) -> Optional[str]: + """写入创建场景最终集数,阻止外部刷新把可靠的当前集数向下覆盖。""" + if current_total_episode and total_episode < current_total_episode: + total_episode = current_total_episode + if not total_episode: + logger.error(f"未获取到总集数,标题:{context.title}") + return f"未获取到第 {context.season} 季的总集数" + context.options["total_episode"] = total_episode + return None + + @staticmethod + def __fill_subscribe_lack_episode(context: _SubscribeCreateContext) -> None: + """未显式指定缺失集数时沿用总集数,保持旧创建默认值。""" + if not context.options.get("lack_episode"): + context.options["lack_episode"] = context.options.get("total_episode") + + def __finalize_subscribe_create_context(self, context: _SubscribeCreateContext) -> None: + """同步补图并写入规范媒体身份和订阅默认配置。""" + if context.mediainfo.type != MediaType.MUSIC: + self.obtain_images(mediainfo=context.mediainfo) + self.__apply_subscribe_create_defaults(context) + + async def __async_finalize_subscribe_create_context( + self, + context: _SubscribeCreateContext, + ) -> None: + """异步补图并写入规范媒体身份和订阅默认配置。""" + if context.mediainfo.type != MediaType.MUSIC: + await self.async_obtain_images(mediainfo=context.mediainfo) + self.__apply_subscribe_create_defaults(context) + + def __apply_subscribe_create_defaults(self, context: _SubscribeCreateContext) -> None: + """以最终识别结果覆盖身份字段,并补齐当前媒体类型的默认订阅参数。""" + context.media_source, context.media_id = resolve_media_identity(media=context.mediainfo) + context.options.update({ + "media_source": context.media_source, + "media_id": context.media_id, + }) + context.options.update( + self.__get_default_kwargs(context.mediainfo.type, **context.options) + ) + + @staticmethod + def __subscribe_post_commit_context( + context: _SubscribeCreateContext, + ) -> _SubscribePostCommitContext: + """从创建阶段状态冻结提交后副作用需要的最小快照。""" + return _SubscribePostCommitContext( + title=context.title, + year=context.year, + metainfo=context.metainfo, + mediainfo=context.mediainfo, + media_source=context.media_source, + media_id=context.media_id, + season=context.season, + channel=context.channel, + source=context.source, + userid=context.userid, + username=context.username, + message=context.message, + ) + + def __persist_subscribe_create(self, context: _SubscribeCreateContext) -> Tuple[Optional[int], str]: + """同步提交订阅,并在提交成功后按原顺序执行消息、事件和统计。""" + post_commit_context = self.__subscribe_post_commit_context(context) + + def _after_commit(subscribe_id: int) -> None: + """把同步提交后的副作用委托给单一顺序实现。""" + self.__post_subscribe_added(subscribe_id, post_commit_context) + + sid, err_msg = add_subscribe( + mediainfo=context.mediainfo, + season=context.season, + username=context.username, + after_commit=_after_commit, + **context.options, + ) + if not sid: + self.__notify_subscribe_create_failure(context, err_msg) + return None, err_msg + return sid, err_msg + + async def __async_persist_subscribe_create( + self, + context: _SubscribeCreateContext, + ) -> Tuple[Optional[int], str]: + """异步提交订阅,并在提交成功后按原顺序执行消息、事件和统计。""" + post_commit_context = self.__subscribe_post_commit_context(context) + + async def _after_commit(subscribe_id: int) -> None: + """把异步提交后的副作用委托给单一顺序实现。""" + await self.__async_post_subscribe_added(subscribe_id, post_commit_context) + + sid, err_msg = await async_add_subscribe( + mediainfo=context.mediainfo, + season=context.season, + username=context.username, + after_commit=_after_commit, + **context.options, + ) + if not sid: + await self.__async_notify_subscribe_create_failure(context, err_msg) + return None, err_msg + return sid, err_msg + + def __notify_subscribe_create_failure( + self, + context: _SubscribeCreateContext, + err_msg: str, + ) -> None: + """同步记录持久化失败,并按旧规则向原用户反馈。""" + logger.error(f"{context.mediainfo.title_year} {err_msg}") + if context.exist_ok or not context.message: + return + self.post_message(self.__subscribe_create_failure_message(context, err_msg)) + + async def __async_notify_subscribe_create_failure( + self, + context: _SubscribeCreateContext, + err_msg: str, + ) -> None: + """异步记录持久化失败,并按旧规则向原用户反馈。""" + logger.error(f"{context.mediainfo.title_year} {err_msg}") + if context.exist_ok or not context.message: + return + await self.async_post_message(self.__subscribe_create_failure_message(context, err_msg)) + + @staticmethod + def __subscribe_create_failure_message( + context: _SubscribeCreateContext, + err_msg: str, + ) -> _SchemaMessage: + """构造保持旧标题、图片和接收人字段的订阅失败消息。""" + return _SchemaMessage( + channel=context.channel, + source=context.source, + mtype=MessageType.Subscribe, + title=( + f"{context.mediainfo.title_year} {context.metainfo.season} " + "添加订阅失败!" + ), + text=err_msg, + image=context.mediainfo.get_message_image(), + userid=context.userid, + ) + def add(self, title: str, year: str, mtype: MediaType = None, episode_group: Optional[str] = None, @@ -989,173 +1410,21 @@ class SubscribeChain(MusicSubscribeMixin, InteractionChainMixin, ChainBase): """ 识别媒体信息并添加订阅 """ - logger.info(f'开始添加订阅,标题:{title} ...') - - explicit_identity = media_source is not None or media_id is not None - media_source, media_id = resolve_media_identity( - media_source=media_source, - media_id=media_id, + context, error = self.__build_subscribe_create_context( + title, year, mtype, episode_group, season, channel, source, userid, + username, message, exist_ok, media_source, media_id, kwargs, ) - if explicit_identity and (not media_source or not media_id): - return None, "媒体来源和媒体 ID 必须同时提供" - - mediainfo = None - requested_music_type = kwargs.get("music_type") - metainfo = MetaMusic.parse_query(title) if mtype == MediaType.MUSIC else MetaInfo(title) - if year: - metainfo.year = year - if mtype: - metainfo.type = mtype - if season is not None: - metainfo.type = MediaType.TV - metainfo.begin_season = season - # 音乐身份同步落到 meta;显式来源与 ID 直接走统一识别入口,不允许失败后换目标。 - if mtype == MediaType.MUSIC and media_id: - metainfo.media_id = str(media_id) - if media_source and media_id: - mediainfo = MediaChain().recognize_media( - meta=metainfo, - mtype=mtype, - media_source=media_source, - media_id=media_id, - music_type=requested_music_type, - episode_group=episode_group, - cache=False, - ) - - if ( - mtype != MediaType.MUSIC - and mediainfo - and mediainfo.media_source != MediaSource.TMDB - ): - meta = MetaInfo(mediainfo.title) - mediainfo.title = meta.name - if season is None: - season = meta.begin_season - - # 没有稳定音乐身份时才允许按名称识别;影视保留原有同源兜底行为。 - if not mediainfo and not explicit_identity: - mediainfo = MediaChain().recognize_by_meta( - metainfo, - media_source=media_source, - episode_group=episode_group, - obtain_images=False, - music_type=requested_music_type, - ) - # 音乐 recognize_by_meta 未命中远端时返回离线兜底,订阅创建要求真实命中 - if mtype == MediaType.MUSIC and mediainfo and not mediainfo.media_source: - mediainfo = None - - # 识别失败 - if not mediainfo: - logger.warn( - f"未识别到媒体信息,标题:{title},媒体来源:{media_source}," - f"媒体 ID:{media_id}" - ) - return None, "未识别到媒体信息" - - if mtype == MediaType.MUSIC: - music_error = self._validate_music_subscribe_target( - mediainfo, - requested_music_type=requested_music_type, - ) - if music_error: - logger.warning(f"音乐订阅目标校验失败:{title} - {music_error}") - return None, music_error - - # 总集数 - if mediainfo.type == MediaType.TV: - if season is None: - season = 1 - # 总集数 - if not kwargs.get('total_episode'): - if not mediainfo.seasons or episode_group: - # 补充媒体信息 - mediainfo = MediaChain().recognize_media(mtype=mediainfo.type, - **_media_recognize_kwargs(mediainfo), - episode_group=episode_group, - cache=False) - if not mediainfo: - logger.error(f"媒体信息识别失败!") - return None, "媒体信息识别失败" - if not mediainfo.seasons: - logger.error(f"媒体信息中没有季集信息,标题:{title}") - return None, "媒体信息中没有季集信息" - current_total_episode = len(mediainfo.seasons.get(season) or []) - # 创建场景没有旧订阅事实,仅允许外部补正未知或扩展总集数。 - total_episode = self.__apply_episodes_refresh( - current_total_episode, season=season, mediainfo=mediainfo, - media_source=resolve_media_identity(media=mediainfo)[0], - media_id=resolve_media_identity(media=mediainfo)[1], scene="create") - if current_total_episode and total_episode < current_total_episode: - total_episode = current_total_episode - if not total_episode: - logger.error(f'未获取到总集数,标题:{title}') - return None, f"未获取到第 {season} 季的总集数" - kwargs.update({ - 'total_episode': total_episode - }) - # 缺失集 - if not kwargs.get('lack_episode'): - kwargs.update({ - 'lack_episode': kwargs.get('total_episode') - }) - else: - # 避免season为0的问题 - season = None - - # 更新媒体图片 - if mediainfo.type != MediaType.MUSIC: - self.obtain_images(mediainfo=mediainfo) - media_source, media_id = resolve_media_identity(media=mediainfo) - kwargs.update({"media_source": media_source, "media_id": media_id}) - - # 添加订阅 - kwargs.update(self.__get_default_kwargs(mediainfo.type, **kwargs)) - - post_commit_context = _SubscribePostCommitContext( - title=title, - year=year, - metainfo=metainfo, - mediainfo=mediainfo, - media_source=media_source, - media_id=media_id, - season=season, - channel=channel, - source=source, - userid=userid, - username=username, - message=bool(message), - ) - - def _after_commit(subscribe_id: int) -> None: - """把同步提交后的副作用委托给单一顺序实现。""" - self.__post_subscribe_added(subscribe_id, post_commit_context) - - # 操作数据库 - sid, err_msg = add_subscribe( - mediainfo=mediainfo, - season=season, - username=username, - after_commit=_after_commit, - **kwargs, - ) - if not sid: - logger.error(f'{mediainfo.title_year} {err_msg}') - if not exist_ok and message: - # 失败发回原用户 - self.post_message(_SchemaMessage(channel=channel, - source=source, - mtype=MessageType.Subscribe, - title=f"{mediainfo.title_year} {metainfo.season} " - f"添加订阅失败!", - text=f"{err_msg}", - image=mediainfo.get_message_image(), - userid=userid)) - return None, err_msg - # 返回结果 - return sid, err_msg + if error: + return None, error + error = self.__recognize_subscribe_media(context) + if error: + return None, error + error = self.__prepare_subscribe_episodes(context) + if error: + return None, error + self.__finalize_subscribe_create_context(context) + return self.__persist_subscribe_create(context) async def async_add(self, title: str, year: str, mtype: MediaType = None, @@ -1173,176 +1442,21 @@ class SubscribeChain(MusicSubscribeMixin, InteractionChainMixin, ChainBase): """ 异步识别媒体信息并添加订阅 """ - logger.info(f'开始添加订阅,标题:{title} ...') - - explicit_identity = media_source is not None or media_id is not None - media_source, media_id = resolve_media_identity( - media_source=media_source, - media_id=media_id, + context, error = self.__build_subscribe_create_context( + title, year, mtype, episode_group, season, channel, source, userid, + username, message, exist_ok, media_source, media_id, kwargs, ) - if explicit_identity and (not media_source or not media_id): - return None, "媒体来源和媒体 ID 必须同时提供" - - mediainfo = None - requested_music_type = kwargs.get("music_type") - metainfo = MetaMusic.parse_query(title) if mtype == MediaType.MUSIC else MetaInfo(title) - if year: - metainfo.year = year - if mtype: - metainfo.type = mtype - if season is not None: - metainfo.type = MediaType.TV - metainfo.begin_season = season - # 音乐身份同步落到 meta;显式来源与 ID 直接走统一识别入口,不允许失败后换目标。 - if mtype == MediaType.MUSIC and media_id: - metainfo.media_id = str(media_id) - if media_source and media_id: - mediainfo = await MediaChain().async_recognize_media( - meta=metainfo, - mtype=mtype, - media_source=media_source, - media_id=media_id, - music_type=requested_music_type, - episode_group=episode_group, - cache=False, - ) - - if ( - mtype != MediaType.MUSIC - and mediainfo - and mediainfo.media_source != MediaSource.TMDB - ): - meta = MetaInfo(mediainfo.title) - mediainfo.title = meta.name - if season is None: - season = meta.begin_season - - # 没有稳定音乐身份时才允许按名称识别;影视保留原有同源兜底行为。 - if not mediainfo and not explicit_identity: - mediainfo = await MediaChain().async_recognize_by_meta( - metainfo, - media_source=media_source, - episode_group=episode_group, - obtain_images=False, - music_type=requested_music_type, - ) - # 音乐 recognize_by_meta 未命中远端时返回离线兜底,订阅创建要求真实命中 - if mtype == MediaType.MUSIC and mediainfo and not mediainfo.media_source: - mediainfo = None - - # 识别失败 - if not mediainfo: - logger.warn( - f"未识别到媒体信息,标题:{title},媒体来源:{media_source}," - f"媒体 ID:{media_id}" - ) - return None, "未识别到媒体信息" - - if mtype == MediaType.MUSIC: - music_error = self._validate_music_subscribe_target( - mediainfo, - requested_music_type=requested_music_type, - ) - if music_error: - logger.warning(f"音乐订阅目标校验失败:{title} - {music_error}") - return None, music_error - - # 总集数 - if mediainfo.type == MediaType.TV: - if season is None: - season = 1 - # 总集数 - if not kwargs.get('total_episode'): - if not mediainfo.seasons or episode_group: - # 补充媒体信息 - mediainfo = await MediaChain().async_recognize_media(mtype=mediainfo.type, - **_media_recognize_kwargs(mediainfo), - episode_group=episode_group, - cache=False) - if not mediainfo: - logger.error(f"媒体信息识别失败!") - return None, "媒体信息识别失败" - if not mediainfo.seasons: - logger.error(f"媒体信息中没有季集信息,标题:{title}") - return None, "媒体信息中没有季集信息" - current_total_episode = len(mediainfo.seasons.get(season) or []) - # 创建场景没有旧订阅事实,仅允许外部补正未知或扩展总集数。 - total_episode = await self.__async_apply_episodes_refresh( - current_total_episode, season=season, mediainfo=mediainfo, - media_source=resolve_media_identity(media=mediainfo)[0], - media_id=resolve_media_identity(media=mediainfo)[1], scene="create") - if current_total_episode and total_episode < current_total_episode: - total_episode = current_total_episode - if not total_episode: - logger.error(f'未获取到总集数,标题:{title}') - return None, f"未获取到第 {season} 季的总集数" - kwargs.update({ - 'total_episode': total_episode - }) - # 缺失集 - if not kwargs.get('lack_episode'): - kwargs.update({ - 'lack_episode': kwargs.get('total_episode') - }) - else: - # 避免season为0的问题 - season = None - - # 更新媒体图片 - if mediainfo.type != MediaType.MUSIC: - await self.async_obtain_images(mediainfo=mediainfo) - media_source, media_id = resolve_media_identity(media=mediainfo) - kwargs.update({"media_source": media_source, "media_id": media_id}) - - # 列新默认参数 - kwargs.update(self.__get_default_kwargs(mediainfo.type, **kwargs)) - - post_commit_context = _SubscribePostCommitContext( - title=title, - year=year, - metainfo=metainfo, - mediainfo=mediainfo, - media_source=media_source, - media_id=media_id, - season=season, - channel=channel, - source=source, - userid=userid, - username=username, - message=bool(message), - ) - - async def _after_commit(subscribe_id: int) -> None: - """把异步提交后的副作用委托给单一顺序实现。""" - await self.__async_post_subscribe_added( - subscribe_id, - post_commit_context, - ) - - # 操作数据库 - sid, err_msg = await async_add_subscribe( - mediainfo=mediainfo, - season=season, - username=username, - after_commit=_after_commit, - **kwargs, - ) - if not sid: - logger.error(f'{mediainfo.title_year} {err_msg}') - if not exist_ok and message: - # 失败发回原用户 - await self.async_post_message(_SchemaMessage(channel=channel, - source=source, - mtype=MessageType.Subscribe, - title=f"{mediainfo.title_year} {metainfo.season} " - f"添加订阅失败!", - text=f"{err_msg}", - image=mediainfo.get_message_image(), - userid=userid)) - return None, err_msg - # 返回结果 - return sid, err_msg + if error: + return None, error + error = await self.__async_recognize_subscribe_media(context) + if error: + return None, error + error = await self.__async_prepare_subscribe_episodes(context) + if error: + return None, error + await self.__async_finalize_subscribe_create_context(context) + return await self.__async_persist_subscribe_create(context) @staticmethod def _subscription_query() -> SubscriptionQueryService: diff --git a/docs/refactor/backend-architecture-next-stage.md b/docs/refactor/backend-architecture-next-stage.md index 0436947ce..3e58bcf20 100644 --- a/docs/refactor/backend-architecture-next-stage.md +++ b/docs/refactor/backend-architecture-next-stage.md @@ -987,6 +987,11 @@ MFA/Passkey 专项测试与架构门禁通过,密钥类配置仍保留在安 鉴权与 MFA 专项测试通过。 `DownloadChain.download_single` 的下载成功结算已提取为独立阶段,入口从 255 行降至 167 行; 历史、文件明细、durable intent、post-commit 通知和旧测试 fallback 语义保持,下载专项测试通过。 +`SubscribeChain.add/async_add` 的同步/异步重复编排随后收口到显式创建上下文和阶段方法:输入规范化、媒体识别、 +电视剧集数准备、默认字段/图片处理、事务提交和失败反馈分别拥有明确边界;订阅重复检测、owner scope、 +`SubscribeAdded` payload、outbox stage/commit/post-commit 顺序仍由既有 `application/subscription/write.py` 负责。 +两个公开入口均降至 150 行预算内,复杂度基线移除对应债务项;订阅识别、音乐订阅、写入事务和搜索来源专项 +共 280 项测试通过,架构、复杂度与异步阻塞门禁通过。 #### ARCH-272:异步阻塞检测 diff --git a/tests/fixtures/architecture/complexity-baseline.json b/tests/fixtures/architecture/complexity-baseline.json index ba3aad4b8..9592de0aa 100644 --- a/tests/fixtures/architecture/complexity-baseline.json +++ b/tests/fixtures/architecture/complexity-baseline.json @@ -20,8 +20,6 @@ "app/chain/download.py:DownloadChain.batch_download": 572, "app/chain/download.py:DownloadChain.download_single": 167, "app/chain/mediaserver.py:MediaServerChain.sync": 292, - "app/chain/subscribe.py:SubscribeChain.add": 183, - "app/chain/subscribe.py:SubscribeChain.async_add": 186, "app/chain/subscribe.py:SubscribeChain.match": 415, "app/chain/subscribe.py:SubscribeChain.search": 246, "app/chain/transfer.py:TransferChain.do_transfer": 885