feat(記憶與情緒快取): 記憶與情緒狀態加入 Redis 快取層,連不到時自動退回原行為
This commit is contained in:
Vendored
+26
@@ -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;
|
||||||
|
}
|
||||||
@@ -5,6 +5,9 @@ import { RuleBasedEmotionTagger } from "./emotion-tagger.js";
|
|||||||
import { nextEmotionState, type TransitionContext } from "./emotion-state-machine.js";
|
import { nextEmotionState, type TransitionContext } from "./emotion-state-machine.js";
|
||||||
import { toResponseStyle, type ResponseStyle } from "./response-style.js";
|
import { toResponseStyle, type ResponseStyle } from "./response-style.js";
|
||||||
import { getArchetypeParams } from "../personality/archetype-params.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 群組)覆寫。
|
// D-2 時間衰減:預設半衰期 2 小時,實際速度由角色參數(G 群組)覆寫。
|
||||||
const DEFAULT_HALF_LIFE_MS = 2 * 60 * 60 * 1000;
|
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) };
|
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<void> {
|
||||||
|
try {
|
||||||
|
await getRedisClient().set(
|
||||||
|
REDIS_KEY_PREFIX + characterId,
|
||||||
|
JSON.stringify({ ...dims, updatedAt: updatedAt.toISOString() }),
|
||||||
|
"EX",
|
||||||
|
60 * 60 * 24,
|
||||||
|
);
|
||||||
|
} catch {
|
||||||
|
// Redis 不可用時,Prisma 仍是持久真實來源,快取寫入失敗不影響正確性。
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
private async decayedDimensions(
|
private async decayedDimensions(
|
||||||
characterId: string,
|
characterId: string,
|
||||||
now: Date,
|
now: Date,
|
||||||
halfLifeMs = DEFAULT_HALF_LIFE_MS,
|
halfLifeMs = DEFAULT_HALF_LIFE_MS,
|
||||||
): Promise<{ dims: Dimensions; updatedAt: Date }> {
|
): 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({
|
const row = await this.prisma.client.emotionState.upsert({
|
||||||
where: { characterId },
|
where: { characterId },
|
||||||
update: {},
|
update: {},
|
||||||
@@ -154,6 +191,7 @@ export class EmotionService {
|
|||||||
where: { characterId },
|
where: { characterId },
|
||||||
data: { ...dims, updatedAt: now },
|
data: { ...dims, updatedAt: now },
|
||||||
});
|
});
|
||||||
|
await this.writeCache(characterId, dims, now);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -17,7 +17,7 @@ export class MemoryConsolidationService {
|
|||||||
) {}
|
) {}
|
||||||
|
|
||||||
async consolidate(characterId: string, sessionId: string): Promise<void> {
|
async consolidate(characterId: string, sessionId: string): Promise<void> {
|
||||||
const entries = this.workingMemory.getContext(sessionId);
|
const entries = await this.workingMemory.getContext(sessionId);
|
||||||
|
|
||||||
const contentCounts = new Map<string, number>();
|
const contentCounts = new Map<string, number>();
|
||||||
for (const entry of entries) {
|
for (const entry of entries) {
|
||||||
@@ -53,6 +53,6 @@ export class MemoryConsolidationService {
|
|||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
this.workingMemory.clear(sessionId);
|
await this.workingMemory.clear(sessionId);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -27,20 +27,20 @@ export class MemoryController {
|
|||||||
) {}
|
) {}
|
||||||
|
|
||||||
@Post(":characterId/sessions/:sessionId/messages")
|
@Post(":characterId/sessions/:sessionId/messages")
|
||||||
appendMessage(@Param("sessionId") sessionId: string, @Body() body: AppendMessageBody) {
|
async appendMessage(@Param("sessionId") sessionId: string, @Body() body: AppendMessageBody) {
|
||||||
this.workingMemory.append(sessionId, {
|
await this.workingMemory.append(sessionId, {
|
||||||
role: body.role,
|
role: body.role,
|
||||||
content: body.content,
|
content: body.content,
|
||||||
timestamp: new Date(),
|
timestamp: new Date(),
|
||||||
emotionTag: body.emotionTag,
|
emotionTag: body.emotionTag,
|
||||||
emotionIntensity: body.emotionIntensity,
|
emotionIntensity: body.emotionIntensity,
|
||||||
});
|
});
|
||||||
return { context: this.workingMemory.getContext(sessionId) };
|
return { context: await this.workingMemory.getContext(sessionId) };
|
||||||
}
|
}
|
||||||
|
|
||||||
@Get(":characterId/sessions/:sessionId/messages")
|
@Get(":characterId/sessions/:sessionId/messages")
|
||||||
getContext(@Param("sessionId") sessionId: string) {
|
async getContext(@Param("sessionId") sessionId: string) {
|
||||||
return { context: this.workingMemory.getContext(sessionId) };
|
return { context: await this.workingMemory.getContext(sessionId) };
|
||||||
}
|
}
|
||||||
|
|
||||||
@Post(":characterId/sessions/:sessionId/consolidate")
|
@Post(":characterId/sessions/:sessionId/consolidate")
|
||||||
|
|||||||
@@ -1,6 +1,7 @@
|
|||||||
import { Injectable } from "@nestjs/common";
|
import { Injectable } from "@nestjs/common";
|
||||||
import type { EmotionTag } from "@kokorone/shared";
|
import type { EmotionTag } from "@kokorone/shared";
|
||||||
import { HIGH_EMOTION_THRESHOLD } from "./constants.js";
|
import { HIGH_EMOTION_THRESHOLD } from "./constants.js";
|
||||||
|
import { getRedisClient } from "../cache/redis-client.js";
|
||||||
|
|
||||||
export interface WorkingMemoryEntry {
|
export interface WorkingMemoryEntry {
|
||||||
role: "user" | "character";
|
role: "user" | "character";
|
||||||
@@ -10,7 +11,14 @@ export interface WorkingMemoryEntry {
|
|||||||
emotionIntensity?: number; // 0~1,未標記情緒的訊息(如快速通道問候)可省略
|
emotionIntensity?: number; // 0~1,未標記情緒的訊息(如快速通道問候)可省略
|
||||||
}
|
}
|
||||||
|
|
||||||
|
interface SerializedEntry extends Omit<WorkingMemoryEntry, "timestamp"> {
|
||||||
|
timestamp: string;
|
||||||
|
}
|
||||||
|
|
||||||
const DEFAULT_TOKEN_LIMIT = 200;
|
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 近似。
|
// 粗略估算:Mock 階段不需要真實 tokenizer,先以字元數/2 近似。
|
||||||
function estimateTokens(text: string): number {
|
function estimateTokens(text: string): number {
|
||||||
@@ -21,32 +29,11 @@ function isHighEmotion(entry: WorkingMemoryEntry): boolean {
|
|||||||
return (entry.emotionIntensity ?? 0) >= HIGH_EMOTION_THRESHOLD;
|
return (entry.emotionIntensity ?? 0) >= HIGH_EMOTION_THRESHOLD;
|
||||||
}
|
}
|
||||||
|
|
||||||
// C-1 工作記憶:session 對話上下文緩衝,只存在於行程記憶體中(如同海馬迴暫存),
|
function evictOverflow(entries: WorkingMemoryEntry[], tokenLimit: number): WorkingMemoryEntry[] {
|
||||||
// 不落地到任何長期記憶表——長期寫入只能由 C-4 睡眠固化觸發。
|
|
||||||
@Injectable()
|
|
||||||
export class WorkingMemoryService {
|
|
||||||
private readonly sessions = new Map<string, WorkingMemoryEntry[]>();
|
|
||||||
private readonly tokenLimit = DEFAULT_TOKEN_LIMIT;
|
|
||||||
|
|
||||||
append(sessionId: string, entry: WorkingMemoryEntry): void {
|
|
||||||
const entries = this.sessions.get(sessionId) ?? [];
|
|
||||||
entries.push(entry);
|
|
||||||
this.sessions.set(sessionId, this.evictOverflow(entries));
|
|
||||||
}
|
|
||||||
|
|
||||||
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 result = entries;
|
||||||
let totalTokens = result.reduce((sum, e) => sum + estimateTokens(e.content), 0);
|
let totalTokens = result.reduce((sum, e) => sum + estimateTokens(e.content), 0);
|
||||||
|
|
||||||
while (totalTokens > this.tokenLimit && result.length > 0) {
|
while (totalTokens > tokenLimit && result.length > 0) {
|
||||||
// 最舊的低情緒段落先被裁掉;若全部都是高情緒,最後才犧牲最舊的一則以確保不超出上限。
|
// 最舊的低情緒段落先被裁掉;若全部都是高情緒,最後才犧牲最舊的一則以確保不超出上限。
|
||||||
let evictIndex = result.findIndex((e) => !isHighEmotion(e));
|
let evictIndex = result.findIndex((e) => !isHighEmotion(e));
|
||||||
if (evictIndex === -1) {
|
if (evictIndex === -1) {
|
||||||
@@ -57,5 +44,59 @@ export class WorkingMemoryService {
|
|||||||
}
|
}
|
||||||
|
|
||||||
return result;
|
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 fallback = new Map<string, WorkingMemoryEntry[]>();
|
||||||
|
private readonly tokenLimit = DEFAULT_TOKEN_LIMIT;
|
||||||
|
|
||||||
|
async append(sessionId: string, entry: WorkingMemoryEntry): Promise<void> {
|
||||||
|
const entries = await this.getContext(sessionId);
|
||||||
|
entries.push(entry);
|
||||||
|
const evicted = evictOverflow(entries, this.tokenLimit);
|
||||||
|
|
||||||
|
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);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
async getContext(sessionId: string): Promise<WorkingMemoryEntry[]> {
|
||||||
|
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<void> {
|
||||||
|
this.fallback.delete(sessionId);
|
||||||
|
const redis = getRedisClient();
|
||||||
|
try {
|
||||||
|
await redis.del(REDIS_KEY_PREFIX + sessionId);
|
||||||
|
} catch {
|
||||||
|
// Redis 不可用時 fallback 已經在上面清掉,沒有需要額外處理的狀態。
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user