From b72cb8510b33a927fdfdf63d304988bd67d3a439 Mon Sep 17 00:00:00 2001 From: SXP-Simon Date: Tue, 10 Feb 2026 14:54:36 +0800 Subject: [PATCH] =?UTF-8?q?feat(=E5=A2=9E=E9=87=8F=E5=88=86=E6=9E=90):=20?= =?UTF-8?q?=E6=B7=BB=E5=8A=A0=20IncrementalStore=20=E6=8C=81=E4=B9=85?= =?UTF-8?q?=E5=8C=96=E4=BB=93=E5=82=A8=E5=92=8C=20IncrementalMergeService?= =?UTF-8?q?=20=E5=90=88=E5=B9=B6=E6=9C=8D=E5=8A=A1?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- src/domain/services/__init__.py | 3 + .../services/incremental_merge_service.py | 214 ++++++++++++++++++ src/infrastructure/persistence/__init__.py | 7 +- .../persistence/incremental_store.py | 186 +++++++++++++++ 4 files changed, 408 insertions(+), 2 deletions(-) create mode 100644 src/domain/services/incremental_merge_service.py create mode 100644 src/infrastructure/persistence/incremental_store.py diff --git a/src/domain/services/__init__.py b/src/domain/services/__init__.py index 5d880ec..8379c7e 100644 --- a/src/domain/services/__init__.py +++ b/src/domain/services/__init__.py @@ -11,6 +11,7 @@ """ from .golden_quote_analyzer import GoldenQuoteAnalyzerAdapter, IGoldenQuoteAnalyzer +from .incremental_merge_service import IncrementalMergeService from .report_generator import ReportGenerator from .statistics_calculator import StatisticsCalculator from .topic_analyzer import ITopicAnalyzer, TopicAnalyzerAdapter @@ -20,6 +21,8 @@ __all__ = [ # 统计与报告服务 "StatisticsCalculator", "ReportGenerator", + # 增量合并服务 + "IncrementalMergeService", # 话题分析服务 "ITopicAnalyzer", "TopicAnalyzerAdapter", diff --git a/src/domain/services/incremental_merge_service.py b/src/domain/services/incremental_merge_service.py new file mode 100644 index 0000000..3cae380 --- /dev/null +++ b/src/domain/services/incremental_merge_service.py @@ -0,0 +1,214 @@ +""" +增量合并领域服务 + +负责将 IncrementalState 累积数据转换为现有实体类型, +以便复用现有的报告生成器和分发器。 + +核心职责: +- IncrementalState → GroupStatistics(含 ActivityVisualization、EmojiStatistics) +- IncrementalState → list[SummaryTopic] +- IncrementalState → list[GoldenQuote] +""" + +from ...domain.entities.incremental_state import IncrementalState +from ...domain.models.data_models import ( + ActivityVisualization, + EmojiStatistics, + GoldenQuote, + GroupStatistics, + SummaryTopic, + TokenUsage, +) +from ...utils.logger import logger + + +class IncrementalMergeService: + """ + 增量合并服务 + + 将一天内累积的增量分析状态转换为现有报告系统所需的数据结构, + 确保增量模式下生成的最终报告与传统单次分析报告格式完全一致。 + """ + + def build_final_statistics(self, state: IncrementalState) -> GroupStatistics: + """ + 从增量状态构建最终的群组统计数据。 + + 将 IncrementalState 中的累积数据映射到 GroupStatistics, + 包含完整的 24 小时活跃度分布、表情统计和 token 消耗。 + + Args: + state: 当天的增量分析状态 + + Returns: + GroupStatistics: 与传统分析格式一致的统计数据 + """ + # 构建 24 小时活跃度分布 + hourly_activity = {} + for hour in range(24): + hour_key = str(hour) + hourly_activity[hour] = state.hourly_message_counts.get(hour_key, 0) + + # 获取高峰时段 + peak_hours = state.get_peak_hours(3) + + # 构建用户活跃排名 + user_ranking = state.get_user_activity_ranking(10) + + # 构建活跃度可视化数据 + activity_visualization = ActivityVisualization( + hourly_activity=hourly_activity, + daily_activity={state.date_str: state.total_message_count}, + user_activity_ranking=user_ranking, + peak_hours=peak_hours, + activity_heatmap_data={}, + ) + + # 构建表情统计 + emoji_statistics = self._build_emoji_statistics(state) + + # 构建 token 消耗 + token_usage = TokenUsage( + prompt_tokens=state.total_token_usage.get("prompt_tokens", 0), + completion_tokens=state.total_token_usage.get("completion_tokens", 0), + total_tokens=state.total_token_usage.get("total_tokens", 0), + ) + + # 获取最活跃时段描述 + most_active_period = state.get_most_active_period() + + statistics = GroupStatistics( + message_count=state.total_message_count, + total_characters=state.total_character_count, + participant_count=len(state.all_participant_ids), + most_active_period=most_active_period, + golden_quotes=[], # 金句通过 build_quotes_for_report 单独构建 + emoji_count=emoji_statistics.total_emoji_count, + emoji_statistics=emoji_statistics, + activity_visualization=activity_visualization, + token_usage=token_usage, + ) + + logger.debug( + f"从增量状态构建统计: " + f"消息数={state.total_message_count}, " + f"参与人数={len(state.all_participant_ids)}, " + f"话题数={len(state.topics)}, " + f"金句数={len(state.golden_quotes)}" + ) + + return statistics + + def build_topics_for_report(self, state: IncrementalState) -> list[SummaryTopic]: + """ + 从增量状态构建报告用的话题列表。 + + 将 IncrementalState 中累积的话题字典转换为 SummaryTopic 实例列表。 + + Args: + state: 当天的增量分析状态 + + Returns: + list[SummaryTopic]: 话题列表,格式与传统分析结果一致 + """ + topics = [] + for topic_dict in state.topics: + topic = SummaryTopic( + topic=topic_dict.get("topic", "未知话题"), + contributors=topic_dict.get("contributors", []), + detail=topic_dict.get("detail", ""), + ) + topics.append(topic) + + logger.debug(f"从增量状态构建了 {len(topics)} 个话题") + return topics + + def build_quotes_for_report(self, state: IncrementalState) -> list[GoldenQuote]: + """ + 从增量状态构建报告用的金句列表。 + + 将 IncrementalState 中累积的金句字典转换为 GoldenQuote 实例列表。 + + Args: + state: 当天的增量分析状态 + + Returns: + list[GoldenQuote]: 金句列表,格式与传统分析结果一致 + """ + quotes = [] + for quote_dict in state.golden_quotes: + quote = GoldenQuote( + content=quote_dict.get("content", ""), + sender=quote_dict.get("sender", ""), + reason=quote_dict.get("reason", ""), + user_id=str(quote_dict.get("user_id", "")), + ) + quotes.append(quote) + + logger.debug(f"从增量状态构建了 {len(quotes)} 条金句") + return quotes + + def build_analysis_result( + self, + state: IncrementalState, + user_titles: list | None = None, + ) -> dict: + """ + 从增量状态构建完整的 analysis_result 字典。 + + 该字典格式与 AnalysisApplicationService.execute_daily_analysis() + 返回的 analysis_result 完全一致,可直接传入 ReportDispatcher。 + + Args: + state: 当天的增量分析状态 + user_titles: 用户称号列表(由最终报告时 LLM 分析生成) + + Returns: + dict: 包含 statistics、topics、user_titles、user_analysis 的结果字典 + """ + statistics = self.build_final_statistics(state) + topics = self.build_topics_for_report(state) + golden_quotes = self.build_quotes_for_report(state) + + # 将金句回填到 statistics 中(与传统流程一致) + statistics.golden_quotes = golden_quotes + + analysis_result = { + "statistics": statistics, + "topics": topics, + "user_titles": user_titles or [], + "user_analysis": state.user_activities, + } + + logger.info( + f"从增量状态构建完整分析结果: " + f"群={state.group_id}, 日期={state.date_str}, " + f"消息={state.total_message_count}, " + f"话题={len(topics)}, " + f"金句={len(golden_quotes)}, " + f"批次={state.total_analysis_count}" + ) + + return analysis_result + + def _build_emoji_statistics(self, state: IncrementalState) -> EmojiStatistics: + """ + 从增量状态构建表情统计。 + + 将 IncrementalState 中的 emoji_counts 字典映射到 EmojiStatistics 字段。 + + Args: + state: 增量分析状态 + + Returns: + EmojiStatistics: 表情统计实例 + """ + emoji_counts = state.emoji_counts + return EmojiStatistics( + face_count=emoji_counts.get("face_count", 0), + mface_count=emoji_counts.get("mface_count", 0), + bface_count=emoji_counts.get("bface_count", 0), + sface_count=emoji_counts.get("sface_count", 0), + other_emoji_count=emoji_counts.get("other_emoji_count", 0), + face_details=emoji_counts.get("face_details", {}), + ) diff --git a/src/infrastructure/persistence/__init__.py b/src/infrastructure/persistence/__init__.py index 1e7e9d3..d429bd8 100644 --- a/src/infrastructure/persistence/__init__.py +++ b/src/infrastructure/persistence/__init__.py @@ -1,7 +1,10 @@ """ -Persistence Module - Data storage implementations +持久化模块 - 数据存储实现 + +包含历史记录仓储和增量分析状态仓储。 """ from .history_repository import HistoryRepository +from .incremental_store import IncrementalStore -__all__ = ["HistoryRepository"] +__all__ = ["HistoryRepository", "IncrementalStore"] diff --git a/src/infrastructure/persistence/incremental_store.py b/src/infrastructure/persistence/incremental_store.py new file mode 100644 index 0000000..e255db4 --- /dev/null +++ b/src/infrastructure/persistence/incremental_store.py @@ -0,0 +1,186 @@ +""" +增量分析状态持久化存储 - 基础设施持久化层 + +负责增量分析状态的存储和读取。 +使用 AstrBot 的 put_kv_data/get_kv_data 实现, +每个群聊每天对应一个独立的状态键。 + +键格式: incremental_state_{group_id}_{date_str} +""" + +import datetime +from typing import Any + +from ...domain.entities.incremental_state import IncrementalState +from ...utils.logger import logger + + +class IncrementalStore: + """ + 增量分析状态持久化仓储 + + 该类封装了增量分析状态在 KV 存储中的读写操作。 + 每个群组每天的增量状态独立存储,支持创建、读取、更新和删除。 + + 使用方式与 HistoryManager 一致,依赖 star_instance 提供的 + put_kv_data / get_kv_data 异步接口。 + """ + + # KV 存储键前缀 + KEY_PREFIX = "incremental_state" + + def __init__(self, star_instance: Any): + """ + 初始化增量状态仓储。 + + Args: + star_instance: Star 插件实例,用于访问底层 KV 存储引擎 + """ + self.plugin = star_instance + + def _build_key(self, group_id: str, date_str: str | None = None) -> str: + """ + 构建 KV 存储键。 + + Args: + group_id: 群组 ID + date_str: 日期字符串 (YYYY-MM-DD),缺省为当天 + + Returns: + str: 格式为 "incremental_state_{group_id}_{date_str}" 的键 + """ + if not date_str: + date_str = datetime.datetime.now().strftime("%Y-%m-%d") + return f"{self.KEY_PREFIX}_{group_id}_{date_str}" + + async def get_state( + self, group_id: str, date_str: str | None = None + ) -> IncrementalState | None: + """ + 读取指定群组在指定日期的增量分析状态。 + + Args: + group_id: 群组 ID + date_str: 日期字符串 (YYYY-MM-DD),缺省为当天 + + Returns: + IncrementalState | None: 状态实例,不存在则返回 None + """ + if not date_str: + date_str = datetime.datetime.now().strftime("%Y-%m-%d") + + key = self._build_key(group_id, date_str) + + try: + data = await self.plugin.get_kv_data(key, None) + if data is None: + return None + + state = IncrementalState.from_dict(data) + logger.debug(f"已读取群 {group_id} 在 {date_str} 的增量状态 (Key: {key})") + return state + except Exception as e: + logger.error(f"读取增量状态失败 (Key: {key}): {e}", exc_info=True) + return None + + async def save_state(self, state: IncrementalState) -> bool: + """ + 持久化增量分析状态。 + + 将状态序列化为字典后写入 KV 存储。 + 如果已存在同键数据则覆盖更新。 + + Args: + state: 要保存的增量分析状态实例 + + Returns: + bool: 保存是否成功 + """ + key = self._build_key(state.group_id, state.date_str) + + try: + data = state.to_dict() + await self.plugin.put_kv_data(key, data) + logger.debug( + f"已保存群 {state.group_id} 在 {state.date_str} 的增量状态 " + f"(Key: {key}, 批次数: {state.total_analysis_count})" + ) + return True + except Exception as e: + logger.error(f"保存增量状态失败 (Key: {key}): {e}", exc_info=True) + return False + + async def get_or_create_state( + self, group_id: str, date_str: str | None = None + ) -> IncrementalState: + """ + 获取或创建增量分析状态。 + + 如果指定群组在指定日期已有状态则返回现有状态, + 否则创建一个新的空白状态实例(不自动持久化)。 + + Args: + group_id: 群组 ID + date_str: 日期字符串 (YYYY-MM-DD),缺省为当天 + + Returns: + IncrementalState: 现有或新创建的状态实例 + """ + if not date_str: + date_str = datetime.datetime.now().strftime("%Y-%m-%d") + + existing = await self.get_state(group_id, date_str) + if existing is not None: + return existing + + # 创建新的空白状态 + new_state = IncrementalState( + group_id=group_id, + date_str=date_str, + ) + logger.info(f"为群 {group_id} 创建了 {date_str} 的新增量状态") + return new_state + + async def delete_state( + self, group_id: str, date_str: str | None = None + ) -> bool: + """ + 删除指定群组在指定日期的增量分析状态。 + + 通过将键值设为 None 来实现删除效果。 + + Args: + group_id: 群组 ID + date_str: 日期字符串 (YYYY-MM-DD),缺省为当天 + + Returns: + bool: 删除是否成功 + """ + if not date_str: + date_str = datetime.datetime.now().strftime("%Y-%m-%d") + + key = self._build_key(group_id, date_str) + + try: + await self.plugin.put_kv_data(key, None) + logger.info(f"已删除群 {group_id} 在 {date_str} 的增量状态 (Key: {key})") + return True + except Exception as e: + logger.error(f"删除增量状态失败 (Key: {key}): {e}", exc_info=True) + return False + + async def has_state( + self, group_id: str, date_str: str | None = None + ) -> bool: + """ + 判断指定群组在指定日期是否存在增量分析状态。 + + Args: + group_id: 群组 ID + date_str: 日期字符串 (YYYY-MM-DD),缺省为当天 + + Returns: + bool: 是否存在状态 + """ + state = await self.get_state(group_id, date_str) + return state is not None