From 6ce7b2e8007817dc830f3bf60786c19988ceec06 Mon Sep 17 00:00:00 2001 From: Jeffery Date: Fri, 14 Aug 2026 12:44:41 +0800 Subject: [PATCH] =?UTF-8?q?feat(=E8=A8=98=E6=86=B6=E8=88=87=E6=83=85?= =?UTF-8?q?=E7=B7=92=E5=BF=AB=E5=8F=96):=20=E8=A8=98=E6=86=B6=E8=88=87?= =?UTF-8?q?=E6=83=85=E7=B7=92=E7=8B=80=E6=85=8B=E5=8A=A0=E5=85=A5=20Redis?= =?UTF-8?q?=20=E5=BF=AB=E5=8F=96=E5=B1=A4,=E9=80=A3=E4=B8=8D=E5=88=B0?= =?UTF-8?q?=E6=99=82=E8=87=AA=E5=8B=95=E9=80=80=E5=9B=9E=E5=8E=9F=E8=A1=8C?= =?UTF-8?q?=E7=82=BA?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- apps/api/src/cache/redis-client.ts | 26 +++++ apps/api/src/emotion/emotion.service.ts | 38 ++++++++ .../memory/memory-consolidation.service.ts | 4 +- apps/api/src/memory/memory.controller.ts | 10 +- apps/api/src/memory/working-memory.service.ts | 97 +++++++++++++------ 5 files changed, 140 insertions(+), 35 deletions(-) create mode 100644 apps/api/src/cache/redis-client.ts diff --git a/apps/api/src/cache/redis-client.ts b/apps/api/src/cache/redis-client.ts new file mode 100644 index 0000000..9ed64bf --- /dev/null +++ b/apps/api/src/cache/redis-client.ts @@ -0,0 +1,26 @@ +import { Redis } from "ioredis"; +import { log } from "@kokorone/shared"; + +let client: Redis | null = null; +let loggedUnavailable = false; + +// R-2:session 快取/情緒狀態即時讀寫改走 Redis。連線失敗時呼叫端一律要能安全退回原本的行為 +// (工作記憶退回行程記憶體 Map、情緒狀態退回直接讀寫 Prisma),Redis 只是加速層,不是新的必要依賴, +// 這樣即使部署環境還沒有 Redis(例如這個沒有 Docker/sudo 的沙盒),既有行為也不會被打斷。 +export function getRedisClient(): Redis { + if (!client) { + client = new Redis(process.env.REDIS_URL ?? "redis://127.0.0.1:6379", { + lazyConnect: true, + maxRetriesPerRequest: 1, + retryStrategy: () => null, + connectTimeout: 300, + }); + client.on("error", () => { + if (!loggedUnavailable) { + loggedUnavailable = true; + log("啟動", "WRN", `Redis 無法連線(REDIS_URL=${process.env.REDIS_URL ?? "redis://127.0.0.1:6379"}),相關快取將退回原本行為`); + } + }); + } + return client; +} diff --git a/apps/api/src/emotion/emotion.service.ts b/apps/api/src/emotion/emotion.service.ts index 777788d..c19d58e 100644 --- a/apps/api/src/emotion/emotion.service.ts +++ b/apps/api/src/emotion/emotion.service.ts @@ -5,6 +5,9 @@ import { RuleBasedEmotionTagger } from "./emotion-tagger.js"; import { nextEmotionState, type TransitionContext } from "./emotion-state-machine.js"; import { toResponseStyle, type ResponseStyle } from "./response-style.js"; import { getArchetypeParams } from "../personality/archetype-params.js"; +import { getRedisClient } from "../cache/redis-client.js"; + +const REDIS_KEY_PREFIX = "kokorone:emotion-state:"; // D-2 時間衰減:預設半衰期 2 小時,實際速度由角色參數(G 群組)覆寫。 const DEFAULT_HALF_LIFE_MS = 2 * 60 * 60 * 1000; @@ -133,11 +136,45 @@ export class EmotionService { return { state: toDomain(characterId, updated, now), style: toResponseStyle(next) }; } + // R-2:情緒狀態機的即時讀寫改走 Redis 當熱路徑快取,Prisma/SQLite(未來 Postgres) + // 仍是唯一的持久真實來源——快取只省掉命中時的資料庫往返,Redis 不可用或未命中都直接退回原本 + // 讀 Prisma 的行為,因此「重啟 api 後情緒狀態不遺失」這件事完全不受快取層是否存在影響。 + private async readCached(characterId: string): Promise<{ dims: Dimensions; updatedAt: Date } | null> { + try { + const raw = await getRedisClient().get(REDIS_KEY_PREFIX + characterId); + if (!raw) return null; + const parsed = JSON.parse(raw) as Dimensions & { updatedAt: string }; + const { updatedAt, ...dims } = parsed; + return { dims, updatedAt: new Date(updatedAt) }; + } catch { + return null; + } + } + + private async writeCache(characterId: string, dims: Dimensions, updatedAt: Date): Promise { + try { + await getRedisClient().set( + REDIS_KEY_PREFIX + characterId, + JSON.stringify({ ...dims, updatedAt: updatedAt.toISOString() }), + "EX", + 60 * 60 * 24, + ); + } catch { + // Redis 不可用時,Prisma 仍是持久真實來源,快取寫入失敗不影響正確性。 + } + } + private async decayedDimensions( characterId: string, now: Date, halfLifeMs = DEFAULT_HALF_LIFE_MS, ): Promise<{ dims: Dimensions; updatedAt: Date }> { + const cached = await this.readCached(characterId); + if (cached) { + const elapsedMs = now.getTime() - cached.updatedAt.getTime(); + return { dims: decayDimensions(cached.dims, elapsedMs, halfLifeMs), updatedAt: cached.updatedAt }; + } + const row = await this.prisma.client.emotionState.upsert({ where: { characterId }, update: {}, @@ -154,6 +191,7 @@ export class EmotionService { where: { characterId }, data: { ...dims, updatedAt: now }, }); + await this.writeCache(characterId, dims, now); } } diff --git a/apps/api/src/memory/memory-consolidation.service.ts b/apps/api/src/memory/memory-consolidation.service.ts index 5e48e82..4004299 100644 --- a/apps/api/src/memory/memory-consolidation.service.ts +++ b/apps/api/src/memory/memory-consolidation.service.ts @@ -17,7 +17,7 @@ export class MemoryConsolidationService { ) {} async consolidate(characterId: string, sessionId: string): Promise { - const entries = this.workingMemory.getContext(sessionId); + const entries = await this.workingMemory.getContext(sessionId); const contentCounts = new Map(); for (const entry of entries) { @@ -53,6 +53,6 @@ export class MemoryConsolidationService { }); } - this.workingMemory.clear(sessionId); + await this.workingMemory.clear(sessionId); } } diff --git a/apps/api/src/memory/memory.controller.ts b/apps/api/src/memory/memory.controller.ts index eea56d4..c85b54a 100644 --- a/apps/api/src/memory/memory.controller.ts +++ b/apps/api/src/memory/memory.controller.ts @@ -27,20 +27,20 @@ export class MemoryController { ) {} @Post(":characterId/sessions/:sessionId/messages") - appendMessage(@Param("sessionId") sessionId: string, @Body() body: AppendMessageBody) { - this.workingMemory.append(sessionId, { + async appendMessage(@Param("sessionId") sessionId: string, @Body() body: AppendMessageBody) { + await this.workingMemory.append(sessionId, { role: body.role, content: body.content, timestamp: new Date(), emotionTag: body.emotionTag, emotionIntensity: body.emotionIntensity, }); - return { context: this.workingMemory.getContext(sessionId) }; + return { context: await this.workingMemory.getContext(sessionId) }; } @Get(":characterId/sessions/:sessionId/messages") - getContext(@Param("sessionId") sessionId: string) { - return { context: this.workingMemory.getContext(sessionId) }; + async getContext(@Param("sessionId") sessionId: string) { + return { context: await this.workingMemory.getContext(sessionId) }; } @Post(":characterId/sessions/:sessionId/consolidate") diff --git a/apps/api/src/memory/working-memory.service.ts b/apps/api/src/memory/working-memory.service.ts index a3bdaa2..6b79a20 100644 --- a/apps/api/src/memory/working-memory.service.ts +++ b/apps/api/src/memory/working-memory.service.ts @@ -1,6 +1,7 @@ import { Injectable } from "@nestjs/common"; import type { EmotionTag } from "@kokorone/shared"; import { HIGH_EMOTION_THRESHOLD } from "./constants.js"; +import { getRedisClient } from "../cache/redis-client.js"; export interface WorkingMemoryEntry { role: "user" | "character"; @@ -10,7 +11,14 @@ export interface WorkingMemoryEntry { emotionIntensity?: number; // 0~1,未標記情緒的訊息(如快速通道問候)可省略 } +interface SerializedEntry extends Omit { + timestamp: string; +} + const DEFAULT_TOKEN_LIMIT = 200; +const REDIS_KEY_PREFIX = "kokorone:working-memory:"; +// session 閒置這麼久沒有新訊息就任由 Redis 自然過期,避免無限累積殭屍 session。 +const SESSION_TTL_SECONDS = 60 * 60 * 6; // 粗略估算:Mock 階段不需要真實 tokenizer,先以字元數/2 近似。 function estimateTokens(text: string): number { @@ -21,41 +29,74 @@ function isHighEmotion(entry: WorkingMemoryEntry): boolean { return (entry.emotionIntensity ?? 0) >= HIGH_EMOTION_THRESHOLD; } -// C-1 工作記憶:session 對話上下文緩衝,只存在於行程記憶體中(如同海馬迴暫存), -// 不落地到任何長期記憶表——長期寫入只能由 C-4 睡眠固化觸發。 +function evictOverflow(entries: WorkingMemoryEntry[], tokenLimit: number): WorkingMemoryEntry[] { + let result = entries; + let totalTokens = result.reduce((sum, e) => sum + estimateTokens(e.content), 0); + + while (totalTokens > tokenLimit && result.length > 0) { + // 最舊的低情緒段落先被裁掉;若全部都是高情緒,最後才犧牲最舊的一則以確保不超出上限。 + let evictIndex = result.findIndex((e) => !isHighEmotion(e)); + if (evictIndex === -1) { + evictIndex = 0; + } + const [evicted] = result.splice(evictIndex, 1); + totalTokens -= estimateTokens(evicted.content); + } + + return result; +} + +function serialize(entries: WorkingMemoryEntry[]): string { + const payload: SerializedEntry[] = entries.map((e) => ({ ...e, timestamp: e.timestamp.toISOString() })); + return JSON.stringify(payload); +} + +function deserialize(raw: string): WorkingMemoryEntry[] { + const payload = JSON.parse(raw) as SerializedEntry[]; + return payload.map((e) => ({ ...e, timestamp: new Date(e.timestamp) })); +} + +// C-1 工作記憶:session 對話上下文緩衝,如同海馬迴暫存,長期寫入只能由 C-4 睡眠固化觸發。 +// R-2:改走 Redis(session 快取),Redis 不可用時退回行程記憶體 Map,行為(token 上限、高情緒優先保留)不變。 @Injectable() export class WorkingMemoryService { - private readonly sessions = new Map(); + private readonly fallback = new Map(); private readonly tokenLimit = DEFAULT_TOKEN_LIMIT; - append(sessionId: string, entry: WorkingMemoryEntry): void { - const entries = this.sessions.get(sessionId) ?? []; + async append(sessionId: string, entry: WorkingMemoryEntry): Promise { + const entries = await this.getContext(sessionId); entries.push(entry); - this.sessions.set(sessionId, this.evictOverflow(entries)); - } + const evicted = evictOverflow(entries, this.tokenLimit); - getContext(sessionId: string): WorkingMemoryEntry[] { - return [...(this.sessions.get(sessionId) ?? [])]; - } - - clear(sessionId: string): void { - this.sessions.delete(sessionId); - } - - private evictOverflow(entries: WorkingMemoryEntry[]): WorkingMemoryEntry[] { - let result = entries; - let totalTokens = result.reduce((sum, e) => sum + estimateTokens(e.content), 0); - - while (totalTokens > this.tokenLimit && result.length > 0) { - // 最舊的低情緒段落先被裁掉;若全部都是高情緒,最後才犧牲最舊的一則以確保不超出上限。 - let evictIndex = result.findIndex((e) => !isHighEmotion(e)); - if (evictIndex === -1) { - evictIndex = 0; - } - const [evicted] = result.splice(evictIndex, 1); - totalTokens -= estimateTokens(evicted.content); + const redis = getRedisClient(); + try { + await redis.set(REDIS_KEY_PREFIX + sessionId, serialize(evicted), "EX", SESSION_TTL_SECONDS); + this.fallback.delete(sessionId); + return; + } catch { + this.fallback.set(sessionId, evicted); } + } - return result; + async getContext(sessionId: string): Promise { + const redis = getRedisClient(); + try { + const raw = await redis.get(REDIS_KEY_PREFIX + sessionId); + if (raw !== null) return deserialize(raw); + if (this.fallback.has(sessionId)) return [...this.fallback.get(sessionId)!]; + return []; + } catch { + return [...(this.fallback.get(sessionId) ?? [])]; + } + } + + async clear(sessionId: string): Promise { + this.fallback.delete(sessionId); + const redis = getRedisClient(); + try { + await redis.del(REDIS_KEY_PREFIX + sessionId); + } catch { + // Redis 不可用時 fallback 已經在上面清掉,沒有需要額外處理的狀態。 + } } }