Files
life-echo/api/app/core/redis.py
Kevin 786ebf8ae6 refactor(api,expo): 多智能体与会话收敛、回忆录兼容层移除、后端测试集大幅删减
- 对齐「多智能体收敛」与「回忆录 stories-first / markdown-first」方向:收紧运行时契约、
  删除过渡兼容路径与双轨逻辑,并同步更新客户端与文档。

- Chat:以 ChatOrchestrator 为实时编排入口;删除独立 conversation_agent,精简 prompts。
- Memoir:删除 memory_agent;MemoirOrchestrator、classification / story_route 与 prompts 收敛到
  prepare_batches + run_story_pipeline_for_category_batch 主链路。
- 将 agents 侧 processor 迁入 feature 层为 background_runner,并移除 features 下重复/过时
  processor 封装。

- 新增 history_store,强化「conversation_messages 为 DB 真源、Redis 为缓存」模型。
- 调整 models、repo、service、session_history;精简 WS message_types,重构 pipeline 与 router。

- 移除章节占位、整章再生等旧路径;章节列表与封面逻辑要求 story 关联;收紧 cover 资格与
  enqueue。
- helpers、repo、service、router、reading_segment_materialize、story_pipeline_sync、pdf_service
  等按 canonical markdown / cover_asset_id 收缩;删除 memoir_images/provider 等冗余。
- tasks:memoir_tasks、chapter_cover_tasks 等大幅瘦身;story_image_tasks 等与当前图片任务对齐。

- core:config、logging、redis、task_tracker 小幅调整。
- auth / user / payment / quota:路由或服务侧删减过时接口或逻辑(如 payment router 行数减少)。

- pyproject.toml、development.sh、.env.example / .env.production、README 等同步说明或变量。

- Alembic 0001_initial_schema 微调(与当前 schema 叙事一致的小改动)。

- 回忆录:types / mappers / api、章节页与 memoir 页与后端契约对齐;markdown-renderer 调整。
- 语音:删除 voice/player,voice-segment-store 相应精简。

- api/tests:删除 conftest 及绝大部分既有测试文件(websocket_baseline、conversation、memoir
  图片、PDF、SMS 等),属有意收缩/待按 backend-test-system 重建的信号。
- docs:新增多智能体收敛与移除兼容层计划摘要;更新 story-first 设计、backend-test-system、
  multi-agent-refactor-plan、实施总结等。

BREAKING CHANGE: 后端对外契约、回忆录章节字段与若干路由/任务行为已变更;大量 API 测试被移除,
  CI 若依赖这些用例需按新策略补测或调整流水线。
2026-03-22 18:10:28 +08:00

238 lines
8.3 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""
Redis 客户端与会话/缓存能力:供应用生命周期、会话历史、任务追踪等使用。
配置从 app.core.config.settings 读取,禁止业务层散落 os.getenv。
"""
import json
from datetime import datetime, timezone
from typing import Any, Dict, List, Optional
import redis.asyncio as aioredis
from app.core.config import settings
from app.core.logging import get_logger
logger = get_logger(__name__)
class RedisService:
"""Redis 服务:连接管理、对话历史、通用缓存。"""
def __init__(self) -> None:
self.redis_url = settings.redis_url
self._client: Optional[aioredis.Redis] = None
self.session_ttl = settings.redis_session_ttl
async def get_client(self) -> aioredis.Redis:
"""获取 Redis 客户端(延迟初始化)。"""
if self._client is None:
try:
self._client = await aioredis.from_url(
self.redis_url,
encoding="utf-8",
decode_responses=True,
)
await self._client.ping()
logger.info("Redis 连接成功")
logger.debug("Redis 连接 URL: %s", self.redis_url)
except Exception as e:
logger.error("Redis 连接失败: %s", e)
raise
return self._client
async def close(self) -> None:
"""关闭 Redis 连接。"""
if self._client:
await self._client.close()
self._client = None
def _conversation_key(self, conversation_id: str) -> str:
return f"conversation:history:{conversation_id}"
async def get_conversation_history(
self, conversation_id: str
) -> List[Dict[str, Any]]:
try:
client = await self.get_client()
key = self._conversation_key(conversation_id)
data = await client.get(key)
if data:
return json.loads(data)
return []
except Exception as e:
logger.error("获取对话历史失败: %s", e)
return []
async def set_conversation_history(
self, conversation_id: str, history: List[Dict[str, Any]]
) -> bool:
"""整表覆盖会话历史(用于从 DB 回填),应用 session_ttl。"""
try:
client = await self.get_client()
key = self._conversation_key(conversation_id)
await client.setex(
key, self.session_ttl, json.dumps(history, ensure_ascii=False)
)
return True
except Exception as e:
logger.error("写入对话历史失败: %s", e)
return False
async def add_message(
self,
conversation_id: str,
role: str,
content: str,
message_type: str = "text",
voice_session_id: str | None = None,
timestamp: str | int | None = None,
audio_duration_seconds: int | None = None,
) -> bool:
try:
client = await self.get_client()
key = self._conversation_key(conversation_id)
history = await self.get_conversation_history(conversation_id)
item = {
"role": role,
"content": content,
"messageType": message_type,
"timestamp": timestamp or datetime.now(timezone.utc).isoformat(),
}
if voice_session_id:
item["voiceSessionId"] = voice_session_id
if (
audio_duration_seconds is not None
and audio_duration_seconds > 0
and message_type == "audio"
):
item["durationSeconds"] = int(audio_duration_seconds)
history.append(item)
await client.setex(
key, self.session_ttl, json.dumps(history, ensure_ascii=False)
)
return True
except Exception as e:
logger.error("添加消息失败: %s", e)
return False
async def append_tts_audio_url_to_last_ai_message(
self, conversation_id: str, url: str
) -> bool:
"""向最近一条 AI 消息的 ttsAudioUrls 追加 COS 公开 URL。"""
if not url:
return False
try:
client = await self.get_client()
key = self._conversation_key(conversation_id)
history = await self.get_conversation_history(conversation_id)
for i in range(len(history) - 1, -1, -1):
if history[i].get("role") == "ai":
existing = history[i].get("ttsAudioUrls")
urls: List[str] = (
[x for x in existing if isinstance(x, str)]
if isinstance(existing, list)
else []
)
urls.append(url)
history[i]["ttsAudioUrls"] = urls
break
else:
logger.warning(
"append_tts_audio_url: no ai message in history conversation_id=%s",
conversation_id,
)
return False
await client.setex(
key, self.session_ttl, json.dumps(history, ensure_ascii=False)
)
return True
except Exception as e:
logger.error("append_tts_audio_url 失败: %s", e)
return False
async def clear_conversation_history(self, conversation_id: str) -> bool:
try:
client = await self.get_client()
key = self._conversation_key(conversation_id)
await client.delete(key)
return True
except Exception as e:
logger.error("清除对话历史失败: %s", e)
return False
async def delete_keys_matching_pattern(self, pattern: str) -> int:
"""按 SCAN 批量删除 key避免阻塞式 KEYS *。"""
try:
client = await self.get_client()
batch: list[str] = []
deleted = 0
async for key in client.scan_iter(match=pattern):
batch.append(key)
if len(batch) >= 200:
deleted += int(await client.delete(*batch))
batch.clear()
if batch:
deleted += int(await client.delete(*batch))
return deleted
except Exception as e:
logger.error("按 pattern 删除 Redis key 失败: %s", e)
return 0
async def extend_session_ttl(self, conversation_id: str) -> bool:
try:
client = await self.get_client()
key = self._conversation_key(conversation_id)
await client.expire(key, self.session_ttl)
return True
except Exception as e:
logger.error("延长会话TTL失败: %s", e)
return False
async def set_cache(self, key: str, value: Any, ttl: Optional[int] = None) -> bool:
try:
client = await self.get_client()
data = (
json.dumps(value, ensure_ascii=False)
if not isinstance(value, str)
else value
)
if ttl:
await client.setex(key, ttl, data)
else:
await client.set(key, data)
return True
except Exception as e:
logger.error("设置缓存失败: %s", e)
return False
async def get_cache(self, key: str) -> Optional[Any]:
try:
client = await self.get_client()
data = await client.get(key)
if data:
try:
return json.loads(data)
except json.JSONDecodeError:
return data
return None
except Exception as e:
logger.error("获取缓存失败: %s", e)
return None
async def delete_cache(self, key: str) -> bool:
try:
client = await self.get_client()
await client.delete(key)
return True
except Exception as e:
logger.error("删除缓存失败: %s", e)
return False
def is_available(self) -> bool:
return self._client is not None
# 全局单例,供 main 生命周期与各 feature 通过 get_redis_service 或直接引用使用
redis_service = RedisService()