1use std::path::PathBuf;
7use std::sync::Arc;
8
9use crate::error::RuntimeError;
10use crate::event::EventEnvelope;
11use crate::index::AnchorIndex;
12use crate::message::Message;
13use crate::projection::message_window::replay_all_messages_with_seq;
14use crate::session::SessionOpenError;
15
16pub enum SearchScope {
17 Session,
18 Project,
19}
20
21pub struct HistoryQuery {
22 pub session_id: String,
23 pub offset: usize,
24 pub limit: usize,
25 pub role_filter: Option<Vec<String>>,
26}
27
28#[derive(Debug)]
29pub struct HistoryPage {
30 pub total: u64,
31 pub offset: usize,
32 pub limit: usize,
33 pub items: Vec<Message>,
34}
35
36#[derive(Debug)]
37pub struct SearchHit {
38 pub session_id: String,
39 pub seq: u64,
40 pub ts: String,
41 pub kind: String,
42 pub snippet: String,
43}
44
45#[derive(Debug)]
46pub struct SearchResult {
47 pub total: u64,
48 pub hits: Vec<SearchHit>,
49}
50
51pub trait HistoryStore: Send + Sync {
52 fn count(&self, session_id: &str, role_filter: Option<&[&str]>) -> Result<u64, RuntimeError>;
53 fn read(&self, query: HistoryQuery) -> Result<HistoryPage, RuntimeError>;
54 fn search(
55 &self,
56 query: &str,
57 scope: SearchScope,
58 limit: usize,
59 ) -> Result<SearchResult, RuntimeError>;
60 fn recent(&self, n: usize) -> Result<(u64, u64, Vec<Message>), RuntimeError>;
61}
62
63pub struct HistoryStoreImpl {
64 project_index: Option<Arc<AnchorIndex>>,
65 session: Option<Arc<crate::session::Session>>,
66 current_session_id: Option<String>,
67 sessions_root: Option<PathBuf>,
68}
69
70impl HistoryStoreImpl {
71 pub fn new(
72 project_index: Option<Arc<AnchorIndex>>,
73 session: Option<Arc<crate::session::Session>>,
74 current_session_id: Option<String>,
75 sessions_root: Option<PathBuf>,
76 ) -> Self {
77 Self {
78 project_index,
79 session,
80 current_session_id,
81 sessions_root,
82 }
83 }
84
85 fn is_current_session(&self, session_id: &str) -> bool {
86 self.current_session_id
87 .as_deref()
88 .is_some_and(|sid| sid == session_id)
89 }
90
91 fn sqlite_available(&self) -> bool {
92 self.project_index.is_some()
93 }
94
95 fn session_dir(&self, session_id: &str) -> Result<PathBuf, RuntimeError> {
96 self.sessions_root
97 .as_ref()
98 .map(|root| root.join(session_id))
99 .ok_or_else(|| {
100 RuntimeError::ToolFailed(format!(
101 "history: no sessions_root to resolve session `{session_id}`"
102 ))
103 })
104 }
105
106 fn replay_from_jsonl(
107 &self,
108 session_id: &str,
109 role_filter: Option<&[&str]>,
110 ) -> Result<Vec<Message>, RuntimeError> {
111 let dir = self.session_dir(session_id)?;
112 let path = dir.join("events.jsonl");
113 let msgs = replay_all_messages_with_seq(&path).map_err(|e| {
114 RuntimeError::ToolFailed(format!("history: replay {}: {e}", path.display()))
115 })?;
116 let msgs: Vec<Message> = msgs.into_iter().map(|(_, m)| m).collect();
117 Ok(filter_messages_by_role(msgs, role_filter))
118 }
119}
120
121fn role_to_kind(role: &str) -> Option<&'static str> {
122 match role {
123 "user" => Some("user_msg"),
124 "assistant" => Some("assistant_msg"),
125 "tool" => Some("tool_result_msg"),
126 "system" => Some("system_msg"),
127 _ => None,
128 }
129}
130
131const MESSAGE_KINDS: &[&str] = &["user_msg", "assistant_msg", "tool_result_msg", "system_msg"];
132
133fn roles_to_kinds(roles: Option<&[&str]>) -> Vec<&'static str> {
134 match roles {
135 Some(rs) if !rs.is_empty() => rs
136 .iter()
137 .filter_map(|r| role_to_kind(r))
138 .collect::<Vec<_>>(),
139 _ => MESSAGE_KINDS.to_vec(),
140 }
141}
142
143fn filter_messages_by_role(msgs: Vec<Message>, roles: Option<&[&str]>) -> Vec<Message> {
144 match roles {
145 Some(rs) if !rs.is_empty() => msgs
146 .into_iter()
147 .filter(|m| rs.iter().any(|r| *r == m.role.as_str()))
148 .collect(),
149 _ => msgs,
150 }
151}
152
153fn extract_message_from_payload(payload: &str) -> Option<Message> {
154 let env: EventEnvelope = serde_json::from_str(payload).ok()?;
155 match env.event {
156 crate::event::Event::UserMsg { message, .. }
157 | crate::event::Event::AssistantMsg { message, .. }
158 | crate::event::Event::ToolResultMsg { message, .. }
159 | crate::event::Event::SystemMsg { message, .. } => Some(message),
160 _ => None,
161 }
162}
163
164fn rows_to_messages(rows: Vec<crate::index::ProjectEventRow>) -> Vec<Message> {
165 rows.into_iter()
166 .filter_map(|r| extract_message_from_payload(&r.payload))
167 .collect()
168}
169
170impl HistoryStore for HistoryStoreImpl {
171 fn count(&self, session_id: &str, role_filter: Option<&[&str]>) -> Result<u64, RuntimeError> {
172 let kinds = roles_to_kinds(role_filter);
173
174 if let Some(idx) = &self.project_index
175 && self.sqlite_available()
176 {
177 return idx
178 .count_events(session_id, Some(&kinds))
179 .map_err(|e| RuntimeError::ToolFailed(format!("history.count: {e}")));
180 }
181
182 if self.is_current_session(session_id)
183 && let Some(session) = &self.session
184 {
185 let msgs = session.messages_full();
186 let filtered = filter_messages_by_role(msgs.to_vec(), role_filter);
187 return Ok(filtered.len() as u64);
188 }
189
190 let msgs = self.replay_from_jsonl(session_id, role_filter)?;
191 Ok(msgs.len() as u64)
192 }
193
194 fn read(&self, query: HistoryQuery) -> Result<HistoryPage, RuntimeError> {
195 let HistoryQuery {
196 session_id,
197 offset,
198 limit,
199 role_filter,
200 } = query;
201 let role_strs: Option<Vec<&str>> = role_filter
202 .as_ref()
203 .map(|rs| rs.iter().map(|s| s.as_str()).collect());
204 let kinds = roles_to_kinds(role_strs.as_deref());
205 let offset0 = offset.saturating_sub(1);
206
207 if let Some(idx) = &self.project_index
208 && self.sqlite_available()
209 {
210 let total = idx
211 .count_events(&session_id, Some(&kinds))
212 .map_err(|e| RuntimeError::ToolFailed(format!("history.read count: {e}")))?;
213 let rows = idx
214 .read_events_paginated(&session_id, offset0, limit, Some(&kinds))
215 .map_err(|e| RuntimeError::ToolFailed(format!("history.read: {e}")))?;
216 let items = rows_to_messages(rows);
217 return Ok(HistoryPage {
218 total,
219 offset,
220 limit,
221 items,
222 });
223 }
224
225 let msgs = if self.is_current_session(&session_id)
226 && let Some(session) = &self.session
227 {
228 session.messages_full().to_vec()
229 } else {
230 self.replay_from_jsonl(&session_id, None)?
231 };
232 let filtered = filter_messages_by_role(msgs, role_strs.as_deref());
233 let total = filtered.len() as u64;
234 let end = (offset0 + limit).min(filtered.len());
235 let items = if offset0 >= filtered.len() {
236 Vec::new()
237 } else {
238 filtered[offset0..end].to_vec()
239 };
240 Ok(HistoryPage {
241 total,
242 offset,
243 limit,
244 items,
245 })
246 }
247
248 fn search(
249 &self,
250 query: &str,
251 scope: SearchScope,
252 limit: usize,
253 ) -> Result<SearchResult, RuntimeError> {
254 if query.trim().is_empty() {
255 return Err(RuntimeError::ToolFailed(
256 "history.search: empty query".into(),
257 ));
258 }
259
260 let session_filter = match &scope {
261 SearchScope::Project => None,
262 SearchScope::Session => self.current_session_id.clone(),
263 };
264
265 if let Some(idx) = &self.project_index {
266 let total = idx
267 .count_search_hits(query, session_filter.as_deref())
268 .map_err(|e| RuntimeError::ToolFailed(format!("history.search count: {e}")))?;
269 let rows = idx
270 .fts_search_project_events(query, session_filter.as_deref(), limit)
271 .map_err(|e| RuntimeError::ToolFailed(format!("history.search: {e}")))?;
272 let hits = rows
273 .into_iter()
274 .map(|r| {
275 let snippet: String = r
276 .payload
277 .chars()
278 .take(200)
279 .collect::<String>()
280 .replace('\n', " ");
281 SearchHit {
282 session_id: r.session_id,
283 seq: r.seq,
284 ts: r.ts,
285 kind: r.kind,
286 snippet,
287 }
288 })
289 .collect();
290 return Ok(SearchResult { total, hits });
291 }
292
293 if matches!(scope, SearchScope::Project) {
294 return Err(RuntimeError::ToolFailed(
295 "history.search: project scope requires project index".into(),
296 ));
297 }
298
299 if let Some(session) = &self.session {
300 let msgs = session.messages_full();
301 let query_lower = query.to_lowercase();
302 let mut hits = Vec::new();
303 let sid = self.current_session_id.clone().unwrap_or_default();
304 for (i, msg) in msgs.iter().enumerate() {
305 let text = msg.text_concat();
306 if text.to_lowercase().contains(&query_lower) {
307 let snippet: String = text.chars().take(200).collect();
308 hits.push(SearchHit {
309 session_id: sid.clone(),
310 seq: i as u64,
311 ts: String::new(),
312 kind: msg.role.as_str().to_string(),
313 snippet,
314 });
315 if hits.len() >= limit {
316 break;
317 }
318 }
319 }
320 let total = hits.len() as u64;
321 return Ok(SearchResult { total, hits });
322 }
323
324 Err(RuntimeError::ToolFailed(
325 "history.search: no project index or session on context".into(),
326 ))
327 }
328
329 fn recent(&self, n: usize) -> Result<(u64, u64, Vec<Message>), RuntimeError> {
330 let msgs = if let Some(session) = &self.session {
331 session.messages_full().to_vec()
332 } else if let Some(session_id) = &self.current_session_id {
333 self.replay_from_jsonl(session_id, None)?
334 } else {
335 return Err(RuntimeError::ToolFailed(
336 "memory.recent_turns: no session available".into(),
337 ));
338 };
339 let (turn_count, items) = recent_turn_messages(&msgs, n);
340 Ok((msgs.len() as u64, turn_count, items))
341 }
342}
343
344pub(crate) fn recent_turn_messages(messages: &[Message], n: usize) -> (u64, Vec<Message>) {
345 if messages.is_empty() || n == 0 {
346 let mut turn_ids = Vec::new();
347 for message in messages {
348 if !turn_ids.contains(&message.turn_id) {
349 turn_ids.push(message.turn_id.clone());
350 }
351 }
352 let total = turn_ids.len() as u64;
353 return (total, Vec::new());
354 }
355 let mut turns: Vec<(crate::event::TurnId, Vec<Message>)> = Vec::new();
356 for message in messages {
357 if let Some((turn_id, items)) = turns.last_mut()
358 && *turn_id == message.turn_id
359 {
360 items.push(message.clone());
361 } else {
362 turns.push((message.turn_id.clone(), vec![message.clone()]));
363 }
364 }
365 let total = turns.len() as u64;
366 let start = turns.len().saturating_sub(n);
367 let items = turns[start..]
368 .iter()
369 .flat_map(|(_, items)| items.iter().cloned())
370 .collect();
371 (total, items)
372}
373
374impl From<SessionOpenError> for RuntimeError {
375 fn from(e: SessionOpenError) -> Self {
376 RuntimeError::ToolFailed(format!("{e}"))
377 }
378}
379
380#[cfg(test)]
381mod tests {
382 use super::*;
383 use crate::index::{AnchorIndex, ProjectEventInsert};
384 use crate::message::{Message, MessageOrigin, MessagePart, MessageRole};
385 use tempfile::TempDir;
386
387 fn user_msg(text: &str) -> Message {
388 Message {
389 role: MessageRole::User,
390 parts: vec![MessagePart::Text {
391 text: text.to_string(),
392 }],
393 turn_id: crate::event::TurnId::now(),
394 origin: MessageOrigin::User,
395 }
396 }
397
398 fn assistant_msg(text: &str) -> Message {
399 Message {
400 role: MessageRole::Assistant,
401 parts: vec![MessagePart::Text {
402 text: text.to_string(),
403 }],
404 turn_id: crate::event::TurnId::now(),
405 origin: MessageOrigin::User,
406 }
407 }
408
409 fn seed(idx: &AnchorIndex, sid: &str, seq: i64, kind: &str, text: &str) {
410 let payload = serde_json::json!({
411 "type": kind,
412 "seq": seq,
413 "turn_id": "019f0000-0000-7000-0000-000000000001",
414 "message": {
415 "role": if kind == "user_msg" { "user" } else { "assistant" },
416 "parts": [{"type": "text", "text": text}],
417 "turn_id": "019f0000-0000-7000-0000-000000000001"
418 },
419 "ts": "2026-07-08T00:00:00Z"
420 });
421 idx.insert_project_event_raw(ProjectEventInsert {
422 session_id: sid,
423 seq,
424 ts: "2026-07-08T00:00:00Z",
425 kind,
426 turn_id: Some("019f0000-0000-7000-0000-000000000001"),
427 flow_run_id: None,
428 text_content: text,
429 payload_json: &payload.to_string(),
430 })
431 .unwrap();
432 }
433
434 #[test]
435 fn recent_groups_messages_by_turn() {
436 let turn_a = crate::event::TurnId::now();
437 let turn_b = crate::event::TurnId::now();
438 let messages = vec![
439 Message {
440 turn_id: turn_a.clone(),
441 ..user_msg("request")
442 },
443 Message {
444 turn_id: turn_a.clone(),
445 ..assistant_msg("answer")
446 },
447 Message {
448 turn_id: turn_a,
449 role: MessageRole::Tool,
450 ..user_msg("tool output")
451 },
452 Message {
453 turn_id: turn_b,
454 ..user_msg("next request")
455 },
456 ];
457 let (turns, recent) = recent_turn_messages(&messages, 1);
458 assert_eq!(turns, 2);
459 assert_eq!(recent.len(), 1);
460 assert_eq!(recent[0].text_concat(), "next request");
461 let (_, recent) = recent_turn_messages(&messages, 2);
462 assert_eq!(recent.len(), 4);
463 }
464
465 #[test]
466 fn count_via_sqlite() {
467 let dir = TempDir::new().unwrap();
468 let idx = Arc::new(AnchorIndex::open_project(dir.path()).unwrap());
469 let store = HistoryStoreImpl::new(Some(idx), None, Some("s1".into()), None);
470 seed(
471 store.project_index.as_ref().unwrap(),
472 "s1",
473 1,
474 "user_msg",
475 "hello",
476 );
477 seed(
478 store.project_index.as_ref().unwrap(),
479 "s1",
480 2,
481 "assistant_msg",
482 "hi",
483 );
484 seed(
485 store.project_index.as_ref().unwrap(),
486 "s1",
487 3,
488 "user_msg",
489 "bye",
490 );
491 assert_eq!(store.count("s1", None).unwrap(), 3);
492 assert_eq!(store.count("s1", Some(&["user"])).unwrap(), 2);
493 assert_eq!(store.count("s1", Some(&["assistant"])).unwrap(), 1);
494 }
495
496 #[test]
497 fn read_via_sqlite_paginated() {
498 let dir = TempDir::new().unwrap();
499 let idx = Arc::new(AnchorIndex::open_project(dir.path()).unwrap());
500 let store = HistoryStoreImpl::new(Some(idx), None, Some("s1".into()), None);
501 for i in 1..=5 {
502 seed(
503 store.project_index.as_ref().unwrap(),
504 "s1",
505 i,
506 "user_msg",
507 &format!("msg {i}"),
508 );
509 }
510 let page = store
511 .read(HistoryQuery {
512 session_id: "s1".into(),
513 offset: 2,
514 limit: 2,
515 role_filter: None,
516 })
517 .unwrap();
518 assert_eq!(page.total, 5);
519 assert_eq!(page.offset, 2);
520 assert_eq!(page.limit, 2);
521 assert_eq!(page.items.len(), 2);
522 assert_eq!(page.items[0].text_concat(), "msg 2");
523 assert_eq!(page.items[1].text_concat(), "msg 3");
524 }
525
526 #[test]
527 fn read_via_sqlite_role_filter() {
528 let dir = TempDir::new().unwrap();
529 let idx = Arc::new(AnchorIndex::open_project(dir.path()).unwrap());
530 let store = HistoryStoreImpl::new(Some(idx), None, Some("s1".into()), None);
531 seed(
532 store.project_index.as_ref().unwrap(),
533 "s1",
534 1,
535 "user_msg",
536 "u1",
537 );
538 seed(
539 store.project_index.as_ref().unwrap(),
540 "s1",
541 2,
542 "assistant_msg",
543 "a1",
544 );
545 seed(
546 store.project_index.as_ref().unwrap(),
547 "s1",
548 3,
549 "user_msg",
550 "u2",
551 );
552 let page = store
553 .read(HistoryQuery {
554 session_id: "s1".into(),
555 offset: 1,
556 limit: 100,
557 role_filter: Some(vec!["user".into()]),
558 })
559 .unwrap();
560 assert_eq!(page.total, 2);
561 assert_eq!(page.items.len(), 2);
562 assert_eq!(page.items[0].text_concat(), "u1");
563 assert_eq!(page.items[1].text_concat(), "u2");
564 }
565
566 #[test]
567 fn search_returns_total_and_hits() {
568 let dir = TempDir::new().unwrap();
569 let idx = Arc::new(AnchorIndex::open_project(dir.path()).unwrap());
570 let store = HistoryStoreImpl::new(Some(idx), None, Some("s1".into()), None);
571 seed(
572 store.project_index.as_ref().unwrap(),
573 "s1",
574 1,
575 "user_msg",
576 "hello world",
577 );
578 seed(
579 store.project_index.as_ref().unwrap(),
580 "s1",
581 2,
582 "user_msg",
583 "hello again",
584 );
585 seed(
586 store.project_index.as_ref().unwrap(),
587 "s1",
588 3,
589 "user_msg",
590 "goodbye",
591 );
592 let result = store.search("hello", SearchScope::Session, 10).unwrap();
593 assert_eq!(result.total, 2);
594 assert_eq!(result.hits.len(), 2);
595 }
596
597 #[test]
598 fn search_cjk() {
599 let dir = TempDir::new().unwrap();
600 let idx = Arc::new(AnchorIndex::open_project(dir.path()).unwrap());
601 let store = HistoryStoreImpl::new(Some(idx), None, Some("s1".into()), None);
602 seed(
603 store.project_index.as_ref().unwrap(),
604 "s1",
605 1,
606 "user_msg",
607 "浮动面板设计",
608 );
609 let result = store.search("浮动", SearchScope::Session, 10).unwrap();
610 assert_eq!(result.total, 1);
611 assert_eq!(result.hits.len(), 1);
612 }
613
614 #[test]
615 fn count_empty_session_via_sqlite() {
616 let dir = TempDir::new().unwrap();
617 let idx = Arc::new(AnchorIndex::open_project(dir.path()).unwrap());
618 let store = HistoryStoreImpl::new(Some(idx), None, Some("s1".into()), None);
619 assert_eq!(store.count("no-such", None).unwrap(), 0);
620 }
621
622 #[test]
623 fn read_empty_session_via_sqlite() {
624 let dir = TempDir::new().unwrap();
625 let idx = Arc::new(AnchorIndex::open_project(dir.path()).unwrap());
626 let store = HistoryStoreImpl::new(Some(idx), None, Some("s1".into()), None);
627 let page = store
628 .read(HistoryQuery {
629 session_id: "no-such".into(),
630 offset: 1,
631 limit: 10,
632 role_filter: None,
633 })
634 .unwrap();
635 assert_eq!(page.total, 0);
636 assert!(page.items.is_empty());
637 }
638
639 #[test]
645 fn fallback_replay_from_jsonl() {
646 let dir = TempDir::new().unwrap();
647 let sessions_root = dir.path().join("sessions");
648 let sid = "test-sid";
649 let session_dir = sessions_root.join(sid);
650 std::fs::create_dir_all(&session_dir).unwrap();
651 let events_path = session_dir.join("events.jsonl");
652 let u1 = user_msg("first");
653 let a1 = assistant_msg("second");
654 let env1 = crate::event::EventEnvelope::new(
655 1,
656 crate::event::Event::UserMsg {
657 turn_id: u1.turn_id.clone(),
658 flow_run_id: None,
659 message: u1,
660 },
661 );
662 let env2 = crate::event::EventEnvelope::new(
663 2,
664 crate::event::Event::AssistantMsg {
665 turn_id: a1.turn_id.clone(),
666 flow_run_id: None,
667 message: a1,
668 },
669 );
670 let line1 = serde_json::to_string(&env1).unwrap();
671 let line2 = serde_json::to_string(&env2).unwrap();
672 std::fs::write(&events_path, format!("{line1}\n{line2}\n")).unwrap();
673
674 let store = HistoryStoreImpl::new(None, None, Some(sid.into()), Some(sessions_root));
675 assert_eq!(store.count(sid, None).unwrap(), 2);
676 assert_eq!(store.count(sid, Some(&["user"])).unwrap(), 1);
677 assert_eq!(store.count(sid, Some(&["assistant"])).unwrap(), 1);
678
679 let page = store
680 .read(HistoryQuery {
681 session_id: sid.into(),
682 offset: 1,
683 limit: 10,
684 role_filter: None,
685 })
686 .unwrap();
687 assert_eq!(page.total, 2);
688 assert_eq!(page.items.len(), 2);
689 assert_eq!(page.items[0].text_concat(), "first");
690 assert_eq!(page.items[1].text_concat(), "second");
691 }
692
693 #[test]
694 fn fallback_replay_paginated() {
695 let dir = TempDir::new().unwrap();
696 let sessions_root = dir.path().join("sessions");
697 let sid = "test-sid";
698 let session_dir = sessions_root.join(sid);
699 std::fs::create_dir_all(&session_dir).unwrap();
700 let events_path = session_dir.join("events.jsonl");
701 let mut lines = Vec::new();
702 for i in 1..=5 {
703 let m = user_msg(&format!("msg {i}"));
704 let env = crate::event::EventEnvelope::new(
705 i,
706 crate::event::Event::UserMsg {
707 turn_id: m.turn_id.clone(),
708 flow_run_id: None,
709 message: m,
710 },
711 );
712 lines.push(serde_json::to_string(&env).unwrap());
713 }
714 std::fs::write(&events_path, lines.join("\n") + "\n").unwrap();
715
716 let store = HistoryStoreImpl::new(None, None, Some(sid.into()), Some(sessions_root));
717 let page = store
718 .read(HistoryQuery {
719 session_id: sid.into(),
720 offset: 2,
721 limit: 2,
722 role_filter: None,
723 })
724 .unwrap();
725 assert_eq!(page.total, 5);
726 assert_eq!(page.items.len(), 2);
727 assert_eq!(page.items[0].text_concat(), "msg 2");
728 assert_eq!(page.items[1].text_concat(), "msg 3");
729 }
730
731 #[test]
732 fn search_empty_query_errors() {
733 let dir = TempDir::new().unwrap();
734 let idx = Arc::new(AnchorIndex::open_project(dir.path()).unwrap());
735 let store = HistoryStoreImpl::new(Some(idx), None, None, None);
736 let err = store.search("", SearchScope::Session, 10).unwrap_err();
737 assert!(matches!(err, RuntimeError::ToolFailed(_)));
738 }
739
740 #[test]
741 fn search_project_scope_without_index_errors() {
742 let store = HistoryStoreImpl::new(None, None, None, None);
743 let err = store.search("hello", SearchScope::Project, 10).unwrap_err();
744 assert!(matches!(err, RuntimeError::ToolFailed(_)));
745 }
746}