Skip to content

Commit d8604b5

Browse files
committed
feat(session): 使用 active/historical 双列表优化会话历史存储
本次调整将会话历史从原先依赖 Event model_flags 区分模型可见性,改为使用 Session.events 与 Session.historical_events 分别存储 active 模型窗口和历史数据,减少内部遍历与过滤开销,并统一 Redis、SQL、InMemory 后端的持久化语义。 主要变更: - Session.events 现在只保留模型可见的 active event window。 - 新增 Session.historical_events,用于按需保存被 max_events、TTL 或 summary 移出 active window 的历史事件。 - 新增 SessionServiceConfig.store_historical_events,控制是否保存历史事件。 - 保留 Event.is_model_visible() / set_model_visible() 公共接口用于兼容,但 SDK 内部不再依赖该 flag 进行模型历史过滤。 - 调整 Session.add_event() / apply_event_filtering() 内部实现,新增内部 helper 返回本次被裁剪的事件,避免影响公共返回值兼容。 - Summary 压缩后使用 [summary_event, recent_events...] 作为新的 active window,旧 summary 会参与下一轮 re-summary,避免历史摘要丢失。 - Summarizer checker 只跳过首位 summary anchor,避免每次反向扫描 events,降低大窗口下的触发检查成本。 - RedisSessionService 写回裁剪后的 active events 与 historical_events,同时保留 app/user/session/temp state 的分层存储语义。 - SqlSessionService 将 active events 存在 events 表,将 historical_events 存在 sessions 表 JSON 字段;append 时仅在本次发生裁剪时删除对应旧 event rows,避免每次全量同步。 - InMemory、Redis、SQL 的 list_sessions() 均不返回 events 和 historical_events,避免列表接口返回大量历史数据。 - MemPalace、Mem0、SQL/Redis/InMemory memory service 适配新的 active events 语义,不再使用 is_model_visible() 过滤。 - 更新相关单测,覆盖 Redis/SQL active-historical 持久化、state 分层、summary re-compress、list_sessions 不返回 historical_events 等场景。 兼容性说明: - Event.is_model_visible() / set_model_visible() 仍保留,避免业务调用时报错。 - SDK 内部模型可见性以 Session.events / Session.historical_events 的归属为准。 - Redis/SQL 后端仅在未显式传入 session_config 时默认开启 store_historical_events;显式传入 SessionServiceConfig() 时仍按默认 False,不保存历史事件。 - SQL stale reload 场景保留现有 best-effort 行为,并通过注释说明并发写入下可能存在窗口不一致,需要更重的版本控制/锁机制才能彻底解决。
1 parent ea2e4bb commit d8604b5

34 files changed

Lines changed: 640 additions & 283 deletions

docs/mkdocs/en/memory.md

Lines changed: 0 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1248,7 +1248,6 @@ memory_service = MempalaceMemoryService(
12481248
),
12491249
wing="my_app_user",
12501250
room="conversations",
1251-
store_only_model_visible=True,
12521251
)
12531252
```
12541253

docs/mkdocs/zh/memory.md

Lines changed: 0 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -619,7 +619,6 @@ memory_service = MempalaceMemoryService(
619619
memory_service_config=memory_service_config,
620620
wing="my_app_user", # 可选:记忆命名空间;不传则默认由 save_key 推导
621621
room="conversations", # 可选:记忆类别;默认 conversations
622-
store_only_model_visible=True,
623622
)
624623
```
625624

examples/memory_service_with_mempalace/README.md

Lines changed: 0 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -62,15 +62,13 @@ memory_service = MempalaceMemoryService(
6262
memory_service_config=memory_service_config,
6363
wing="trpc-agent",
6464
room="conversations",
65-
store_only_model_visible=True,
6665
)
6766
```
6867

6968
这里的含义是:
7069

7170
- `wing="trpc-agent"`:把示例记忆固定写入 `trpc-agent` 这个 wing。
7271
- `room="conversations"`:把普通对话记忆写入 `conversations` room。
73-
- `store_only_model_visible=True`:只存模型可见的事件。
7472
- `ttl_seconds=20`:超过 20 秒的记忆会被后台 cleanup 删除。
7573
- `cleanup_interval_seconds=20`:每 20 秒执行一次清理。
7674

examples/memory_service_with_mempalace/run_agent.py

Lines changed: 0 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -44,7 +44,6 @@ def create_memory_service():
4444
memory_service_config=memory_service_config,
4545
wing="trpc-agent",
4646
room="conversations",
47-
store_only_model_visible=True,
4847
)
4948
return memory_service
5049

tests/memory/test_mempalace_memory_service.py

Lines changed: 23 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -91,25 +91,43 @@ def fake_store(session, events_to_store, wing, room):
9191
assert "remember this" in events_to_store[0][1]
9292
assert events_to_store[0][2] in svc._stored_drawer_ids
9393

94-
async def test_store_session_skips_invisible_events(self, monkeypatch):
94+
async def test_store_session_ignores_model_visible_flag(self, monkeypatch):
9595
calls = []
9696

9797
def fake_store(session, events_to_store, wing, room):
9898
calls.append(events_to_store)
9999
return {drawer_id for _, _, drawer_id in events_to_store}
100100

101101
visible_event = _make_event("visible")
102-
invisible_event = _make_event("hidden")
103-
invisible_event.set_model_visible(False)
102+
flagged_event = _make_event("hidden")
103+
flagged_event.set_model_visible(False)
104104
svc = MempalaceMemoryService(memory_service_config=_make_config())
105105
monkeypatch.setattr(svc, "_store_events", fake_store)
106106

107-
await svc.store_session(_make_session(events=[visible_event, invisible_event]))
107+
await svc.store_session(_make_session(events=[visible_event, flagged_event]))
108108
await svc.close()
109109

110110
assert len(calls) == 1
111-
assert len(calls[0]) == 1
111+
assert len(calls[0]) == 2
112112
assert "visible" in calls[0][0][1]
113+
assert "hidden" in calls[0][1][1]
114+
115+
async def test_store_only_model_visible_flag_is_compatibility_noop(self, monkeypatch):
116+
calls = []
117+
118+
def fake_store(session, events_to_store, wing, room):
119+
calls.append(events_to_store)
120+
return {drawer_id for _, _, drawer_id in events_to_store}
121+
122+
svc = MempalaceMemoryService(memory_service_config=_make_config(), store_only_model_visible=False)
123+
monkeypatch.setattr(svc, "_store_events", fake_store)
124+
125+
await svc.store_session(_make_session(events=[_make_event("active event")]))
126+
await svc.close()
127+
128+
assert len(calls) == 1
129+
assert len(calls[0]) == 1
130+
assert "active event" in calls[0][0][1]
113131

114132
async def test_store_session_is_incremental(self, monkeypatch):
115133
calls = []

tests/sessions/test_base_session_service.py

Lines changed: 41 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -141,6 +141,18 @@ async def test_append_event_empty_state_delta(self):
141141
await svc.append_event(session, event)
142142
assert len(session.events) == 1
143143

144+
async def test_append_event_stores_filtered_events_when_configured(self):
145+
config = SessionServiceConfig(max_events=2, store_historical_events=True)
146+
svc = ConcreteSessionService(session_config=config)
147+
session = _make_session()
148+
149+
for i in range(6):
150+
event = _make_event(author="user" if i == 2 else "agent", text=f"msg{i}")
151+
await svc.append_event(session, event)
152+
153+
assert [event.get_text() for event in session.events] == ["msg2", "msg5"]
154+
assert [event.get_text() for event in session.historical_events] == ["msg0", "msg1", "msg3", "msg4"]
155+
144156

145157
class TestBaseSessionServiceTrimTempDeltaState:
146158
"""Test _trim_temp_delta_state method."""
@@ -172,10 +184,23 @@ def test_filter_by_num_recent_events(self):
172184
for i in range(10):
173185
author = "user" if i == 7 else "agent"
174186
session.events.append(_make_event(author=author, text=f"msg{i}"))
175-
svc.filter_events(session)
187+
filtered_session = svc.filter_events(session)
188+
assert filtered_session is session
189+
assert [event.get_text() for event in session.events] == ["msg7", "msg8", "msg9"]
190+
191+
def test_filter_by_num_recent_events_with_copy(self):
192+
config = SessionServiceConfig(num_recent_events=3)
193+
svc = ConcreteSessionService(session_config=config)
194+
session = _make_session()
195+
for i in range(10):
196+
author = "user" if i == 7 else "agent"
197+
session.events.append(_make_event(author=author, text=f"msg{i}"))
198+
199+
filtered_session = svc.filter_events(session, need_copy=True)
200+
201+
assert filtered_session is not session
176202
assert len(session.events) == 10
177-
visible_events = [event for event in session.events if event.is_model_visible()]
178-
assert [event.get_text() for event in visible_events] == ["msg7", "msg8", "msg9"]
203+
assert [event.get_text() for event in filtered_session.events] == ["msg7", "msg8", "msg9"]
179204

180205
def test_filter_by_event_ttl(self):
181206
config = SessionServiceConfig(event_ttl_seconds=5.0)
@@ -190,18 +215,18 @@ def test_filter_by_event_ttl(self):
190215
new_event.timestamp = time.time()
191216
session.events.append(new_event)
192217

193-
svc.filter_events(session)
194-
assert len(session.events) == 2
195-
visible_events = [event for event in session.events if event.is_model_visible()]
196-
assert len(visible_events) == 1
197-
assert visible_events[0].get_text() == "new"
218+
filtered_session = svc.filter_events(session)
219+
assert filtered_session is session
220+
assert len(session.events) == 1
221+
assert session.events[0].get_text() == "new"
198222

199223
def test_filter_no_config(self):
200224
svc = ConcreteSessionService()
201225
session = _make_session()
202226
for i in range(5):
203227
session.events.append(_make_event(text=f"msg{i}"))
204-
svc.filter_events(session)
228+
filtered_session = svc.filter_events(session)
229+
assert filtered_session is session
205230
assert len(session.events) == 5
206231

207232
def test_filter_ttl_removes_all_old(self):
@@ -212,9 +237,9 @@ def test_filter_ttl_removes_all_old(self):
212237
e = _make_event(text=f"old{i}")
213238
e.timestamp = time.time() - 100
214239
session.events.append(e)
215-
svc.filter_events(session)
216-
assert len(session.events) == 5
217-
assert all(not event.is_model_visible() for event in session.events)
240+
filtered_session = svc.filter_events(session)
241+
assert filtered_session is session
242+
assert session.events == []
218243

219244
def test_filter_by_num_recent_events_preserves_summary_anchor(self):
220245
config = SessionServiceConfig(num_recent_events=3)
@@ -227,11 +252,11 @@ def test_filter_by_num_recent_events_preserves_summary_anchor(self):
227252
for i in range(5):
228253
session.events.append(_make_event(text=f"agent{i}"))
229254

230-
svc.filter_events(session)
255+
filtered_session = svc.filter_events(session)
231256

232-
visible_events = [event for event in session.events if event.is_model_visible()]
233-
assert len(visible_events) == 1
234-
assert visible_events[0].is_summary_event()
257+
assert filtered_session is session
258+
assert len(session.events) == 1
259+
assert session.events[0].is_summary_event()
235260

236261

237262
class TestBaseSessionServiceSetSummarizerManager:

tests/sessions/test_in_memory_session_service.py

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -235,6 +235,7 @@ async def test_list_sessions_have_no_events(self):
235235
result = await svc.list_sessions(app_name="app", user_id="user")
236236
assert len(result.sessions) == 1
237237
assert result.sessions[0].events == []
238+
assert result.sessions[0].historical_events == []
238239
await svc.close()
239240

240241
async def test_list_nonexistent_app(self):

tests/sessions/test_redis_session_service.py

Lines changed: 42 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -27,8 +27,12 @@
2727
from trpc_agent_sdk.types import Content, EventActions, Part, State
2828

2929

30-
def _make_config(ttl_seconds=0, cleanup_interval=0.0, enable_ttl=False):
31-
config = SessionServiceConfig()
30+
def _make_config(ttl_seconds=0,
31+
cleanup_interval=0.0,
32+
enable_ttl=False,
33+
max_events=0,
34+
store_historical_events=False):
35+
config = SessionServiceConfig(max_events=max_events, store_historical_events=store_historical_events)
3236
if enable_ttl:
3337
config.ttl = SessionServiceConfig.create_ttl_config(
3438
enable=True, ttl_seconds=ttl_seconds, cleanup_interval_seconds=cleanup_interval)
@@ -203,6 +207,7 @@ async def test_list_sessions_have_no_events(self):
203207
result = await svc.list_sessions(app_name="app", user_id="user")
204208
for s in result.sessions:
205209
assert s.events == []
210+
assert s.historical_events == []
206211
await svc.close()
207212

208213

@@ -256,6 +261,28 @@ async def test_append_with_state_delta(self):
256261
assert stored.state[f"{State.USER_PREFIX}user_key"] == "uv"
257262
await svc.close()
258263

264+
async def test_append_does_not_persist_merged_or_temp_state_in_session_json(self):
265+
svc = _create_service()
266+
session = await svc.create_session(app_name="app", user_id="user", session_id="s1")
267+
event = _make_event(state_delta={
268+
"session_key": "sv",
269+
f"{State.APP_PREFIX}app_key": "av",
270+
f"{State.USER_PREFIX}user_key": "uv",
271+
f"{State.TEMP_PREFIX}temp_key": "tv",
272+
})
273+
274+
await svc.append_event(session, event)
275+
276+
stored_json = svc._redis_storage._store["session:app:user:s1"]
277+
raw_session = Session.model_validate_json(stored_json)
278+
assert raw_session.state == {"session_key": "sv"}
279+
stored = await svc.get_session(app_name="app", user_id="user", session_id="s1")
280+
assert stored.state["session_key"] == "sv"
281+
assert stored.state[f"{State.APP_PREFIX}app_key"] == "av"
282+
assert stored.state[f"{State.USER_PREFIX}user_key"] == "uv"
283+
assert f"{State.TEMP_PREFIX}temp_key" not in raw_session.state
284+
await svc.close()
285+
259286
async def test_append_to_nonexistent_session(self):
260287
svc = _create_service()
261288
session = _make_session_obj(id="nonexistent")
@@ -264,6 +291,19 @@ async def test_append_to_nonexistent_session(self):
264291
assert result is event
265292
await svc.close()
266293

294+
async def test_append_persists_filtered_active_and_historical_events(self):
295+
config = _make_config(max_events=2, store_historical_events=True)
296+
svc = _create_service(config=config)
297+
session = await svc.create_session(app_name="app", user_id="user", session_id="s1")
298+
299+
for i in range(4):
300+
await svc.append_event(session, _make_event(author="user" if i == 2 else "agent", text=f"msg{i}"))
301+
302+
stored = await svc.get_session(app_name="app", user_id="user", session_id="s1")
303+
assert [event.get_text() for event in stored.events] == ["msg2", "msg3"]
304+
assert [event.get_text() for event in stored.historical_events] == ["msg0", "msg1"]
305+
await svc.close()
306+
267307

268308
class TestRedisUpdateSession:
269309
async def test_update_existing(self):

0 commit comments

Comments
 (0)