1#![forbid(unsafe_code)]
9
10use async_trait::async_trait;
11
12use chrono::{DateTime, NaiveDate, TimeZone, Utc};
13use serde_json::{Value, json};
14use sha2::{Digest, Sha256};
15use std::fmt::Write as _;
16use std::sync::Arc;
17use wm_core::{Context, EffectRow, Galaxy, Gana, Resource, Tool, ToolStats};
18use wm_memory::{Memory, MemoryStore};
19
20fn parse_time_bound(v: &Value, end_of_day: bool) -> Option<DateTime<Utc>> {
24 if let Some(secs) = v
25 .as_i64()
26 .or_else(|| v.as_u64().and_then(|u| i64::try_from(u).ok()))
27 {
28 return Utc.timestamp_opt(secs, 0).single();
29 }
30 let s = v.as_str()?;
31 if let Ok(dt) = DateTime::parse_from_rfc3339(s) {
32 return Some(dt.with_timezone(&Utc));
33 }
34 let day = NaiveDate::parse_from_str(s, "%Y-%m-%d").ok()?;
35 let naive = if end_of_day {
36 day.and_hms_opt(23, 59, 59)?
37 } else {
38 day.and_hms_opt(0, 0, 0)?
39 };
40 Some(naive.and_utc())
41}
42
43fn filter_by_time(
46 turns: Vec<(Memory, Value)>,
47 args: &Value,
48) -> wm_core::Result<Vec<(Memory, Value)>> {
49 let since = match args.get("since") {
50 Some(v) if !v.is_null() => Some(parse_time_bound(v, false).ok_or_else(|| {
51 wm_core::CoreError::InvalidArgs(
52 "invalid 'since' — use epoch seconds, RFC 3339, or YYYY-MM-DD".into(),
53 )
54 })?),
55 _ => None,
56 };
57 let until = match args.get("until") {
58 Some(v) if !v.is_null() => Some(parse_time_bound(v, true).ok_or_else(|| {
59 wm_core::CoreError::InvalidArgs(
60 "invalid 'until' — use epoch seconds, RFC 3339, or YYYY-MM-DD".into(),
61 )
62 })?),
63 _ => None,
64 };
65 Ok(turns
66 .into_iter()
67 .filter(|(m, _)| {
68 since.is_none_or(|t| m.metadata.created_at >= t)
69 && until.is_none_or(|t| m.metadata.created_at <= t)
70 })
71 .collect())
72}
73
74fn turn_json(mem: &Memory) -> Option<Value> {
75 let v: Value = serde_json::from_str(&mem.content).ok()?;
76 if v.get("type").and_then(Value::as_str) == Some("session_turn") {
77 Some(v)
78 } else {
79 None
80 }
81}
82
83fn load_turns(
89 store: &MemoryStore,
90 session_id: Option<&str>,
91 limit: usize,
92 include_superseded: bool,
93) -> wm_core::Result<Vec<(Memory, Value)>> {
94 let memories = store.scan_all(Galaxy::Sessions)?;
95 let mut turns: Vec<(Memory, Value)> = memories
96 .iter()
97 .filter(|m| {
98 include_superseded
99 || !m
100 .metadata
101 .tags
102 .iter()
103 .any(|t| t.starts_with("superseded-by:"))
104 })
105 .filter_map(|m| turn_json(m).map(|v| (m.clone(), v)))
106 .filter(|(_, v)| {
107 session_id.is_none_or(|sid| v.get("session_id").and_then(Value::as_str) == Some(sid))
108 })
109 .collect();
110 turns.sort_by_key(|(_, v)| {
111 (
112 v.get("sequence").and_then(Value::as_u64).unwrap_or(0),
113 v.get("timestamp").and_then(Value::as_i64).unwrap_or(0),
114 )
115 });
116 turns.truncate(limit);
117 Ok(turns)
118}
119
120fn latest_checkpoint_handoff(
125 store: &MemoryStore,
126 session_id: &str,
127) -> wm_core::Result<Option<(String, DateTime<Utc>, Value)>> {
128 Ok(store
129 .scan_all(Galaxy::Sessions)?
130 .iter()
131 .filter(|m| {
132 m.metadata.tags.contains(&"checkpoint".to_string()) && m.content.contains(session_id)
133 })
134 .filter_map(|m| {
135 let parsed: Value = serde_json::from_str(&m.content).ok()?;
136 parsed
137 .get("handoff")
138 .filter(|h| !h.is_null())
139 .cloned()
140 .map(|h| (m.metadata.id.to_string(), m.metadata.created_at, h))
141 })
142 .max_by_key(|(_, created_at, _)| *created_at))
143}
144
145fn format_turn(v: &Value, full: bool) -> Value {
146 let role = v.get("role").and_then(Value::as_str).unwrap_or("?");
147 let content = v.get("content").and_then(Value::as_str).unwrap_or("");
148 if full {
149 json!({
150 "session_id": v.get("session_id"),
151 "sequence": v.get("sequence"),
152 "role": role,
153 "turn_type": v.get("turn_type"),
154 "importance": v.get("importance"),
155 "content": content,
156 })
157 } else {
158 json!({
159 "sequence": v.get("sequence"),
160 "role": role,
161 "turn_type": v.get("turn_type"),
162 "preview": content.chars().take(120).collect::<String>(),
163 })
164 }
165}
166
167const LOSSLESS_MAX_PAGE_SIZE: usize = 64;
168const LOSSLESS_DEFAULT_PAGE_SIZE: usize = 16;
169const LOSSLESS_MIN_WIRE_BYTES: usize = 1024;
170const LOSSLESS_MAX_WIRE_BYTES: usize = 49_152;
171
172fn hex_encode(bytes: &[u8]) -> String {
173 use std::fmt::Write as _;
174 bytes
175 .iter()
176 .fold(String::with_capacity(bytes.len() * 2), |mut out, b| {
177 let _ = write!(out, "{b:02x}");
178 out
179 })
180}
181
182fn hex_decode(value: &str) -> Option<Vec<u8>> {
183 if value.len() > 4096
184 || value.len() % 2 != 0
185 || !value
186 .bytes()
187 .all(|b| b.is_ascii_digit() || (b'a'..=b'f').contains(&b))
188 {
189 return None;
190 }
191 (0..value.len())
192 .step_by(2)
193 .map(|i| u8::from_str_radix(&value[i..i + 2], 16).ok())
194 .collect()
195}
196
197fn base64_encode(bytes: &[u8]) -> String {
198 const TABLE: &[u8; 64] = b"ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz0123456789+/";
199 let mut out = String::with_capacity(bytes.len().div_ceil(3) * 4);
200 for chunk in bytes.chunks(3) {
201 let n = u32::from(chunk[0]) << 16
202 | u32::from(*chunk.get(1).unwrap_or(&0)) << 8
203 | u32::from(*chunk.get(2).unwrap_or(&0));
204 out.push(char::from(TABLE[((n >> 18) & 63) as usize]));
205 out.push(char::from(TABLE[((n >> 12) & 63) as usize]));
206 out.push(if chunk.len() > 1 {
207 char::from(TABLE[((n >> 6) & 63) as usize])
208 } else {
209 '='
210 });
211 out.push(if chunk.len() > 2 {
212 char::from(TABLE[(n & 63) as usize])
213 } else {
214 '='
215 });
216 }
217 out
218}
219
220#[cfg(test)]
221fn base64_decode(value: &str) -> Option<Vec<u8>> {
222 const fn digit(byte: u8) -> Option<u8> {
223 match byte {
224 b'A'..=b'Z' => Some(byte - b'A'),
225 b'a'..=b'z' => Some(byte - b'a' + 26),
226 b'0'..=b'9' => Some(byte - b'0' + 52),
227 b'+' => Some(62),
228 b'/' => Some(63),
229 _ => None,
230 }
231 }
232 if value.len() % 4 != 0 {
233 return None;
234 }
235 let mut out = Vec::new();
236 for chunk in value.as_bytes().chunks_exact(4) {
237 let a = digit(chunk[0])?;
238 let b = digit(chunk[1])?;
239 let c = if chunk[2] == b'=' {
240 0
241 } else {
242 digit(chunk[2])?
243 };
244 let d = if chunk[3] == b'=' {
245 0
246 } else {
247 digit(chunk[3])?
248 };
249 out.push((a << 2) | (b >> 4));
250 if chunk[2] != b'=' {
251 out.push((b << 4) | (c >> 2));
252 }
253 if chunk[3] != b'=' {
254 out.push((c << 6) | d);
255 }
256 }
257 Some(out)
258}
259
260#[derive(Clone)]
261struct LosslessTurn {
262 memory: Memory,
263 turn: Value,
264 content: String,
265 content_hash: String,
266}
267
268fn lossless_error(kind: &str) -> wm_core::CoreError {
269 wm_core::CoreError::InvalidArgs(format!("lossless_{kind}"))
270}
271
272fn lossless_cursor(
273 session_id: &str,
274 include_superseded: bool,
275 page_size: usize,
276 max_wire: usize,
277 view: &str,
278 index: usize,
279 offset: usize,
280) -> String {
281 let value = json!({"v":1,"session_id":session_id,"include_superseded":include_superseded,"page_size":page_size,"max_wire_bytes":max_wire,"view":view,"index":index,"offset":offset});
284 hex_encode(value.to_string().as_bytes())
285}
286
287fn parse_lossless_cursor(
288 cursor: &str,
289 session_id: &str,
290 include_superseded: bool,
291 page_size: usize,
292 max_wire: usize,
293) -> wm_core::Result<(String, usize, usize)> {
294 let bytes = hex_decode(cursor).ok_or_else(|| lossless_error("invalid_cursor"))?;
295 let text = String::from_utf8(bytes).map_err(|_| lossless_error("invalid_cursor"))?;
296 let value: Value = serde_json::from_str(&text).map_err(|_| lossless_error("invalid_cursor"))?;
297 let canonical = json!({"v":value.get("v"),"session_id":value.get("session_id"),"include_superseded":value.get("include_superseded"),"page_size":value.get("page_size"),"max_wire_bytes":value.get("max_wire_bytes"),"view":value.get("view"),"index":value.get("index"),"offset":value.get("offset")});
298 let canonical_text =
299 serde_json::to_string(&canonical).map_err(|_| lossless_error("invalid_cursor"))?;
300 if canonical_text != text
301 || value.get("v").and_then(Value::as_u64) != Some(1)
302 || value.get("session_id").and_then(Value::as_str) != Some(session_id)
303 || value.get("include_superseded").and_then(Value::as_bool) != Some(include_superseded)
304 || value.get("page_size").and_then(Value::as_u64) != Some(page_size as u64)
305 || value.get("max_wire_bytes").and_then(Value::as_u64) != Some(max_wire as u64)
306 {
307 return Err(lossless_error("invalid_cursor"));
308 }
309 let view = value
310 .get("view")
311 .and_then(Value::as_str)
312 .filter(|v| {
313 v.len() == 64
314 && v.bytes()
315 .all(|b| b.is_ascii_digit() || (b'a'..=b'f').contains(&b))
316 })
317 .ok_or_else(|| lossless_error("invalid_cursor"))?;
318 let index = value
319 .get("index")
320 .and_then(Value::as_u64)
321 .and_then(|v| usize::try_from(v).ok())
322 .ok_or_else(|| lossless_error("invalid_cursor"))?;
323 let offset = value
324 .get("offset")
325 .and_then(Value::as_u64)
326 .and_then(|v| usize::try_from(v).ok())
327 .ok_or_else(|| lossless_error("invalid_cursor"))?;
328 Ok((view.to_string(), index, offset))
329}
330
331pub const TURN_TYPES: &[&str] = &[
336 "message",
337 "decision",
338 "breakthrough",
339 "question",
340 "answer",
341 "code_change",
342 "error",
343 "summary",
344 "context",
345];
346
347pub const TRACK_MAX_LEN: usize = 64;
349
350pub(crate) fn validate_track(track: &str) -> wm_core::Result<()> {
357 let valid = !track.is_empty()
358 && track.len() <= TRACK_MAX_LEN
359 && track
360 .chars()
361 .next()
362 .is_some_and(|c| c.is_ascii_lowercase() || c.is_ascii_digit())
363 && track.chars().all(|c| {
364 c.is_ascii_lowercase() || c.is_ascii_digit() || matches!(c, '-' | '_' | '.' | '/')
365 });
366 if valid {
367 Ok(())
368 } else {
369 Err(wm_core::CoreError::InvalidArgs(format!(
370 "invalid track {track:?} — must start with a-z0-9, contain only a-z0-9, '-', '_', '.', '/', \
371 and be at most {TRACK_MAX_LEN} chars"
372 )))
373 }
374}
375
376pub struct SessionRecordTool {
378 store: Arc<MemoryStore>,
379 stats: ToolStats,
380 effects: EffectRow,
381 search: Option<Arc<wm_memory::SearchEngine>>,
382}
383
384impl SessionRecordTool {
385 #[must_use]
386 pub fn new(store: Arc<MemoryStore>) -> Self {
387 Self {
388 store,
389 stats: ToolStats::default(),
390 effects: EffectRow {
391 writes: vec![Resource::Galaxy("sessions".into())],
392 ..Default::default()
393 },
394 search: None,
395 }
396 }
397
398 #[must_use]
402 pub fn with_search(mut self, search: Option<Arc<wm_memory::SearchEngine>>) -> Self {
403 self.search = search;
404 self
405 }
406}
407
408#[async_trait]
409impl Tool for SessionRecordTool {
410 fn name(&self) -> &str {
411 "session.record"
412 }
413 fn gana(&self) -> Gana {
414 Gana::StraddlingLegs
415 }
416 fn effects(&self) -> &EffectRow {
417 &self.effects
418 }
419 fn input_schema(&self) -> Value {
420 super::common::schema(
421 &json!({
422 "content": super::common::str_prop("Turn content"),
423 "role": super::common::str_prop("user | ai (default user)"),
424 "turn_type": json!({
425 "type": "string",
426 "enum": TURN_TYPES,
427 "description": "Turn type (default message)",
428 }),
429 "importance": super::common::bounded_num_prop("0-1 importance (default 0.5)", 0.0, 1.0),
430 "session_id": super::common::str_prop("Target session (default: most recent session)"),
431 "supersedes": super::common::str_prop("Memory id of an earlier turn this record corrects/replaces (amend-with-supersede)"),
432 "track": super::common::str_prop("Optional track slug (lowercase; a-z0-9 start, then a-z0-9-_. /) — tags this turn into that track's implementation log (session.track_log)"),
433 }),
434 &["content"],
435 )
436 }
437 fn description(&self) -> &str {
438 "Record a conversation turn as persistent session memory. Args: content (required), role (user|ai, default user), turn_type (default message), importance (0-1, default 0.5), session_id (optional — defaults to the most recent session), supersedes (optional turn memory-id — marks the old turn superseded so replay/continuity/digest use the new record), track (optional slug — tags the turn into that track's log, read back with session.track_log)."
439 }
440 async fn call(&self, _ctx: &mut Context, args: Value) -> wm_core::Result<Value> {
441 let role = args.get("role").and_then(Value::as_str).unwrap_or("user");
442 if !matches!(role, "user" | "ai") {
443 return Err(wm_core::CoreError::InvalidArgs(
444 "role must be 'user' or 'ai'".into(),
445 ));
446 }
447 let content = args
450 .get("content")
451 .and_then(Value::as_str)
452 .filter(|s| !s.trim().is_empty())
453 .ok_or_else(|| {
454 wm_core::CoreError::InvalidArgs("content is required and must not be blank".into())
455 })?;
456 let turn_type = args
460 .get("turn_type")
461 .and_then(Value::as_str)
462 .unwrap_or("message");
463 if !TURN_TYPES.contains(&turn_type) {
464 return Err(wm_core::CoreError::InvalidArgs(format!(
465 "turn_type must be one of: {}",
466 TURN_TYPES.join(", ")
467 )));
468 }
469 let importance = wm_dispatch::write_gate::parse_importance_value(args.get("importance"))
474 .map_err(wm_core::CoreError::InvalidArgs)?
475 .unwrap_or(0.5);
476 let session_id = args.get("session_id").and_then(Value::as_str);
477 let track = args.get("track").and_then(Value::as_str);
481 if let Some(track) = track {
482 validate_track(track)?;
483 }
484
485 let session_id: String = if let Some(sid) = session_id {
491 sid.to_string()
492 } else {
493 self.store
494 .scan_all(Galaxy::Sessions)?
495 .iter()
496 .filter(|m| m.metadata.tags.contains(&"start".to_string()))
497 .max_by_key(|m| m.metadata.created_at)
498 .map(|m| m.metadata.id.to_string())
499 .ok_or_else(|| {
500 wm_core::CoreError::Tool("no session found — run session.start first".into())
501 })?
502 };
503
504 let supersedes = match args.get("supersedes").and_then(Value::as_str) {
507 Some(old_id_str) => {
508 let old_id = uuid::Uuid::parse_str(old_id_str).map_err(|e| {
509 wm_core::CoreError::InvalidArgs(format!("invalid 'supersedes' id: {e}"))
510 })?;
511 if self.store.get(Galaxy::Sessions, old_id)?.is_none() {
512 return Err(wm_core::CoreError::NotFound(format!(
513 "superseded turn {old_id} not found"
514 )));
515 }
516 Some(old_id)
517 }
518 None => None,
519 };
520
521 let timestamp = wm_core::time::now_unix_millis();
526 let (sequence, mem) = self.store.put_session_turn(&session_id, |sequence| {
527 let mut turn = json!({
528 "type": "session_turn",
529 "session_id": session_id,
530 "sequence": sequence,
531 "role": role,
532 "turn_type": turn_type,
533 "importance": importance,
534 "content": content,
535 "timestamp": timestamp,
536 });
537 if let Some(track) = track {
538 turn["track"] = json!(track);
539 }
540 let mut mem = Memory::new(Galaxy::Sessions, turn.to_string());
541 mem.metadata.tags = vec![
542 "session".into(),
543 "turn".into(),
544 role.into(),
545 turn_type.into(),
546 format!("session:{session_id}"),
547 ];
548 if let Some(track) = track {
549 mem.metadata.tags.push(format!("track:{track}"));
550 }
551 let (source, trust): (&str, f32) = if role == "user" {
558 ("user", 1.0)
559 } else {
560 ("agent", 0.7)
561 };
562 mem.metadata.source = source.to_string();
563 mem.metadata.source_trust = trust;
564 mem.metadata.importance = importance as f32;
565 if let Some(old_id) = supersedes {
568 mem.metadata.tags.push(format!("supersedes:{old_id}"));
569 }
570 mem
571 })?;
572
573 if let Some(old_id) = supersedes {
580 if let Some(mut old) = self.store.get(Galaxy::Sessions, old_id)? {
581 old.metadata
582 .tags
583 .push(format!("superseded-by:{}", mem.metadata.id));
584 self.store.put(Galaxy::Sessions, &old)?;
585 super::common::index_memory(self.search.as_deref(), &old);
586 }
587 }
588
589 super::common::index_memory(self.search.as_deref(), &mem);
590 Ok(json!({
591 "status": "success",
592 "session_id": session_id,
593 "sequence": sequence,
594 "memory_id": mem.metadata.id.to_string(),
595 }))
596 }
597 fn stats(&self) -> &ToolStats {
598 &self.stats
599 }
600}
601
602pub struct SessionReplayTool {
604 store: Arc<MemoryStore>,
605 stats: ToolStats,
606 effects: EffectRow,
607}
608
609impl SessionReplayTool {
610 #[must_use]
611 pub fn new(store: Arc<MemoryStore>) -> Self {
612 Self {
613 store,
614 stats: ToolStats::default(),
615 effects: EffectRow::read_only(vec![Resource::Galaxy("sessions".into())]),
616 }
617 }
618
619 fn lossless(&self, args: &Value) -> wm_core::Result<Value> {
620 const ALLOWED: &[&str] = &[
621 "mode",
622 "session_id",
623 "include_superseded",
624 "page_size",
625 "max_wire_bytes",
626 "cursor",
627 ];
628 let object = args
629 .as_object()
630 .ok_or_else(|| lossless_error("invalid_args"))?;
631 if object.keys().any(|key| !ALLOWED.contains(&key.as_str())) {
632 return Err(lossless_error("unsupported_selection_args"));
633 }
634 let session_id = args
635 .get("session_id")
636 .and_then(Value::as_str)
637 .filter(|s| !s.is_empty())
638 .ok_or_else(|| lossless_error("session_id_required"))?;
639 uuid::Uuid::parse_str(session_id).map_err(|_| lossless_error("invalid_session_id"))?;
640 let include_superseded = match args.get("include_superseded") {
641 None => false,
642 Some(v) => v.as_bool().ok_or_else(|| lossless_error("invalid_args"))?,
643 };
644 let size_arg = |name: &str, default: usize| -> wm_core::Result<usize> {
645 match args.get(name) {
646 None => Ok(default),
647 Some(v) => v
648 .as_u64()
649 .and_then(|v| usize::try_from(v).ok())
650 .ok_or_else(|| lossless_error("invalid_args")),
651 }
652 };
653 let page_size = size_arg("page_size", LOSSLESS_DEFAULT_PAGE_SIZE)?;
654 let max_wire = size_arg("max_wire_bytes", LOSSLESS_MAX_WIRE_BYTES)?;
655 if !(1..=LOSSLESS_MAX_PAGE_SIZE).contains(&page_size)
656 || !(LOSSLESS_MIN_WIRE_BYTES..=LOSSLESS_MAX_WIRE_BYTES).contains(&max_wire)
657 {
658 return Err(lossless_error("invalid_args"));
659 }
660
661 let placement = match args.get("cursor") {
663 None => None,
664 Some(v) => Some(parse_lossless_cursor(
665 v.as_str().ok_or_else(|| lossless_error("invalid_cursor"))?,
666 session_id,
667 include_superseded,
668 page_size,
669 max_wire,
670 )?),
671 };
672 let memories = self.store.scan_all_strict(Galaxy::Sessions)?;
673 let start = memories.iter().find(|m| {
674 m.metadata.id.to_string() == session_id
675 && m.metadata.tags.contains(&"start".to_string())
676 && !m.metadata.is_private
677 && !m.metadata.model_exclude
678 });
679 if start.is_none() {
680 return Err(wm_core::CoreError::NotFound("session not found".into()));
681 }
682 let mut turns = Vec::new();
683 for memory in memories {
684 if memory.metadata.is_private
685 || memory.metadata.model_exclude
686 || (!include_superseded
687 && memory
688 .metadata
689 .tags
690 .iter()
691 .any(|t| t.starts_with("superseded-by:")))
692 {
693 continue;
694 }
695 let tagged = memory
696 .metadata
697 .tags
698 .contains(&format!("session:{session_id}"));
699 let turn = match serde_json::from_str::<Value>(&memory.content) {
700 Ok(v) => v,
701 Err(_) if tagged => return Err(lossless_error("malformed_selected_turn")),
702 Err(_) => continue,
703 };
704 if turn.get("type").and_then(Value::as_str) != Some("session_turn") {
705 if tagged && memory.metadata.tags.iter().any(|tag| tag == "turn") {
706 return Err(lossless_error("malformed_selected_turn"));
707 }
708 continue;
709 }
710 if !tagged && turn.get("session_id").and_then(Value::as_str) != Some(session_id) {
711 continue;
712 }
713 let content = turn
714 .get("content")
715 .and_then(Value::as_str)
716 .ok_or_else(|| lossless_error("malformed_selected_turn"))?
717 .to_string();
718 if turn.get("session_id").and_then(Value::as_str) != Some(session_id)
721 || turn.get("sequence").and_then(Value::as_u64).is_none()
722 || turn.get("timestamp").and_then(Value::as_i64).is_none()
723 {
724 return Err(lossless_error("malformed_selected_turn"));
725 }
726 turns.push(LosslessTurn {
727 content_hash: hex_encode(&Sha256::digest(content.as_bytes())),
728 memory,
729 turn,
730 content,
731 });
732 }
733 turns.sort_by_key(|turn| {
734 (
735 turn.turn["sequence"].as_u64().unwrap(),
736 turn.turn["timestamp"].as_i64().unwrap(),
737 turn.memory.metadata.id,
738 )
739 });
740 let visible_ids: std::collections::HashSet<String> = turns
741 .iter()
742 .map(|t| t.memory.metadata.id.to_string())
743 .collect();
744 let metadata: Vec<Value> = turns.iter().map(|t| {
745 let relationships: Vec<&str> = t.memory.metadata.tags.iter().filter_map(|tag| {
746 let (_, id) = tag.split_once(':')?;
747 ((tag.starts_with("superseded-by:") || tag.starts_with("supersedes:")) && visible_ids.contains(id)).then_some(tag.as_str())
748 }).collect();
749 json!({"role":t.turn.get("role"),"turn_type":t.turn.get("turn_type"),"importance":t.turn.get("importance"),"created_at":t.memory.metadata.created_at,"source":t.memory.metadata.source,"source_trust":t.memory.metadata.source_trust,"agent_id":t.memory.metadata.agent_id,"supersession":relationships})
750 }).collect();
751 let manifest: Vec<Value> = turns.iter().zip(&metadata).map(|(t,m)| json!({"id":t.memory.metadata.id,"sequence":t.turn["sequence"],"timestamp":t.turn["timestamp"],"hash":t.content_hash,"metadata":m})).collect();
752 let view = hex_encode(&Sha256::digest(json!({"v":1,"galaxy":"sessions","session_id":session_id,"include_superseded":include_superseded,"page_size":page_size,"max_wire_bytes":max_wire,"records":manifest}).to_string().as_bytes()));
753 let (mut index, mut offset) = match placement {
754 Some((token_view, index, offset)) => {
755 if token_view != view {
756 return Err(lossless_error("stale_view"));
757 }
758 (index, offset)
759 }
760 None => (0, 0),
761 };
762 if index > turns.len() || (index == turns.len() && offset != 0) {
763 return Err(lossless_error("invalid_placement"));
764 }
765 if index < turns.len() && offset != 0 && offset >= turns[index].content.len() {
766 return Err(lossless_error("invalid_placement"));
767 }
768 let mut records = Vec::new();
769 while index < turns.len() && records.len() < page_size {
770 let turn = &turns[index];
771 let bytes = turn.content.as_bytes();
772 let whole = json!({"record_id":turn.memory.metadata.id.to_string(),"sequence":turn.turn["sequence"],"timestamp":turn.turn["timestamp"],"metadata":metadata[index],"content_encoding":"utf-8","content":turn.content,"content_sha256":turn.content_hash,"complete":true});
773 let candidate = json!({"status":"success","mode":"lossless","session_id":session_id,"galaxy":"sessions","view_fingerprint":view,"records":records.iter().cloned().chain(std::iter::once(whole.clone())).collect::<Vec<_>>(),"has_more":index+1<turns.len(),"next_cursor":if index+1==turns.len() {Value::Null} else {json!(lossless_cursor(session_id,include_superseded,page_size,max_wire,&view,index+1,0))},"complete":index+1==turns.len()});
774 if offset == 0 && serde_json::to_vec(&candidate).unwrap().len() <= max_wire {
775 records.push(whole);
776 index += 1;
777 offset = 0;
778 continue;
779 }
780 if !records.is_empty() {
781 break;
782 }
783 let mut take = bytes.len().saturating_sub(offset);
784 while take > 0 {
785 let end = offset + take;
786 let chunk = json!({"record_id":turn.memory.metadata.id.to_string(),"sequence":turn.turn["sequence"],"timestamp":turn.turn["timestamp"],"metadata":metadata[index],"content_sha256":turn.content_hash,"chunk":{"encoding":"base64","byte_offset":offset,"total_bytes":bytes.len(),"data_b64":base64_encode(&bytes[offset..end]),"complete":end==bytes.len()}});
787 let next = if end == bytes.len() {
788 lossless_cursor(
789 session_id,
790 include_superseded,
791 page_size,
792 max_wire,
793 &view,
794 index + 1,
795 0,
796 )
797 } else {
798 lossless_cursor(
799 session_id,
800 include_superseded,
801 page_size,
802 max_wire,
803 &view,
804 index,
805 end,
806 )
807 };
808 let final_chunk = end == bytes.len() && index + 1 == turns.len();
809 let candidate = json!({"status":"success","mode":"lossless","session_id":session_id,"galaxy":"sessions","view_fingerprint":view,"records":[chunk.clone()],"has_more":!final_chunk,"next_cursor":if final_chunk {Value::Null} else {json!(next)},"complete":final_chunk});
810 if serde_json::to_vec(&candidate).unwrap().len() <= max_wire {
811 records.push(chunk);
812 if end == bytes.len() {
813 index += 1;
814 offset = 0;
815 } else {
816 offset = end;
817 }
818 break;
819 }
820 take /= 2;
821 }
822 if records.is_empty() {
823 return Err(lossless_error("wire_ceiling_too_small"));
824 }
825 break;
826 }
827 let complete = index == turns.len() && offset == 0;
828 let response = json!({"status":"success","mode":"lossless","session_id":session_id,"galaxy":"sessions","view_fingerprint":view,"records":records,"has_more":!complete,"next_cursor":if complete { Value::Null } else { json!(lossless_cursor(session_id,include_superseded,page_size,max_wire,&view,index,offset)) },"complete":complete});
829 if serde_json::to_vec(&response).unwrap().len() > max_wire {
830 return Err(lossless_error("wire_ceiling_too_small"));
831 }
832 Ok(response)
833 }
834}
835
836#[async_trait]
837impl Tool for SessionReplayTool {
838 fn name(&self) -> &str {
839 "session.replay"
840 }
841 fn gana(&self) -> Gana {
842 Gana::StraddlingLegs
843 }
844 fn effects(&self) -> &EffectRow {
845 &self.effects
846 }
847 fn input_schema(&self) -> Value {
848 super::common::schema(
849 &json!({
850 "mode": super::common::str_prop("full | selective | progressive | lossless (default full)"),
851 "session_id": super::common::str_prop("Target session (required explicit UUID for lossless; otherwise default most recent)"),
852 "n": super::common::int_prop("Maximum turns (default 50)"),
853 "since": super::common::str_prop("Time-range floor: epoch seconds, RFC 3339, or YYYY-MM-DD"),
854 "until": super::common::str_prop("Time-range ceiling: epoch seconds, RFC 3339, or YYYY-MM-DD"),
855 "include_superseded": {
856 "type": "boolean",
857 "description": "Also return turns replaced via supersedes (default false)."
858 },
859 "turn_types": super::common::str_array_prop("Selective mode: turn types to keep"),
860 "min_importance": super::common::num_prop("Selective mode floor (default 0.7)"),
861 "token_budget": super::common::int_prop("Progressive mode token budget (default 2000)"),
862 "page_size": super::common::int_prop("Lossless mode records per page (1-64, default 16)"),
863 "max_wire_bytes": super::common::int_prop("Lossless mode serialized JSON ceiling (1024-49152)"),
864 "cursor": super::common::str_prop("Lossless mode opaque placement cursor"),
865 }),
866 &[],
867 )
868 }
869 fn description(&self) -> &str {
870 "Replay session turns. Legacy full/selective/progressive use optional session_id, n and since/until. Lossless requires an explicit session UUID and permits only mode, session_id, include_superseded, page_size (1-64, default16), max_wire_bytes (1024-49152, default49152), cursor. Follow returned cursors with unchanged settings for byte-exact content; base64 chunks are raw bytes, concatenate before UTF-8 decoding. Changed visible views refuse stale_view. Cursors are unsigned non-authority placement hints. Tool-result JSON budget is not the outer transport or model budget."
871 }
872 async fn call(&self, _ctx: &mut Context, args: Value) -> wm_core::Result<Value> {
873 let mode = args.get("mode").and_then(Value::as_str).unwrap_or("full");
874 if mode == "lossless" {
875 return self.lossless(&args);
876 }
877 let requested_session_id = args
878 .get("session_id")
879 .and_then(Value::as_str)
880 .filter(|sid| !sid.is_empty());
881 let session_id = match requested_session_id {
886 Some(sid) => Some(sid.to_string()),
887 None => self
888 .store
889 .scan_all(Galaxy::Sessions)?
890 .iter()
891 .filter(|m| m.metadata.tags.contains(&"start".to_string()))
892 .max_by_key(|m| m.metadata.created_at)
893 .map(|m| m.metadata.id.to_string()),
894 };
895 let n = args.get("n").and_then(Value::as_u64).unwrap_or(50) as usize;
896 let include_superseded = args
897 .get("include_superseded")
898 .and_then(Value::as_bool)
899 .unwrap_or(false);
900 let loaded_turns = match session_id.as_deref() {
904 Some(sid) => load_turns(&self.store, Some(sid), 10_000, include_superseded)?,
905 None => Vec::new(),
906 };
907 let turns = filter_by_time(loaded_turns, &args)?;
908
909 if requested_session_id.is_some() && turns.is_empty() {
912 return Err(wm_core::CoreError::InvalidArgs(format!(
913 "no session found with id {requested_session_id:?}"
914 )));
915 }
916
917 let selected: Vec<(Memory, Value)> = match mode {
918 "selective" => {
919 let min_importance = args
920 .get("min_importance")
921 .and_then(Value::as_f64)
922 .unwrap_or(0.7);
923 let turn_types: Vec<String> = args
924 .get("turn_types")
925 .and_then(Value::as_array)
926 .map_or_else(
927 || vec!["decision".into(), "breakthrough".into(), "answer".into()],
928 |a| {
929 a.iter()
930 .filter_map(Value::as_str)
931 .map(str::to_string)
932 .collect()
933 },
934 );
935 turns
936 .into_iter()
937 .filter(|(_, v)| {
938 v.get("importance").and_then(Value::as_f64).unwrap_or(0.0) >= min_importance
939 && v.get("turn_type")
940 .and_then(Value::as_str)
941 .is_some_and(|t| turn_types.contains(&t.to_string()))
942 })
943 .collect()
944 }
945 "progressive" => {
946 let budget = args
947 .get("token_budget")
948 .and_then(Value::as_u64)
949 .unwrap_or(2000) as usize;
950 let mut used = 0usize;
951 let mut out = Vec::new();
952 for (m, v) in turns.into_iter().rev() {
953 let approx = v
954 .get("content")
955 .and_then(Value::as_str)
956 .map_or(0, |c| c.len() / 4);
957 if used + approx > budget {
958 break;
959 }
960 used += approx;
961 out.push((m, v));
962 }
963 out.reverse();
964 out
965 }
966 _ => turns
967 .into_iter()
968 .rev()
969 .take(n)
970 .collect::<Vec<_>>()
971 .into_iter()
972 .rev()
973 .collect(),
974 };
975
976 let full = mode != "progressive";
977 let formatted: Vec<Value> = selected.iter().map(|(_, v)| format_turn(v, full)).collect();
978 Ok(json!({
979 "status": "success",
980 "mode": mode,
981 "count": formatted.len(),
982 "session_id": session_id,
983 "turns": formatted,
984 }))
985 }
986 fn stats(&self) -> &ToolStats {
987 &self.stats
988 }
989}
990
991pub struct SessionContinuityTool {
993 store: Arc<MemoryStore>,
994 stats: ToolStats,
995 effects: EffectRow,
996}
997
998impl SessionContinuityTool {
999 #[must_use]
1000 pub fn new(store: Arc<MemoryStore>) -> Self {
1001 Self {
1002 store,
1003 stats: ToolStats::default(),
1004 effects: EffectRow::read_only(vec![Resource::Galaxy("sessions".into())]),
1005 }
1006 }
1007}
1008
1009#[async_trait]
1010impl Tool for SessionContinuityTool {
1011 fn name(&self) -> &str {
1012 "session.continuity"
1013 }
1014 fn gana(&self) -> Gana {
1015 Gana::StraddlingLegs
1016 }
1017 fn effects(&self) -> &EffectRow {
1018 &self.effects
1019 }
1020 fn input_schema(&self) -> Value {
1021 super::common::schema(
1022 &json!({
1023 "current_session_id": super::common::str_prop("Session to exclude (optional)"),
1024 "n": super::common::int_prop("Number of prior turns (default 10)"),
1025 "since": super::common::str_prop("Time-range floor: epoch seconds, RFC 3339, or YYYY-MM-DD"),
1026 "until": super::common::str_prop("Time-range ceiling: epoch seconds, RFC 3339, or YYYY-MM-DD"),
1027 }),
1028 &[],
1029 )
1030 }
1031 fn description(&self) -> &str {
1032 "Get cross-session continuity — the last N turns of the most recent prior session ('where we left off') plus that session's latest checkpoint handoff (next_queue, open_flags, git state, tests_green, lease_id) when one exists. Args: current_session_id (optional, excluded), n (default 10), since/until (epoch seconds | RFC 3339 | YYYY-MM-DD)."
1033 }
1034 async fn call(&self, _ctx: &mut Context, args: Value) -> wm_core::Result<Value> {
1035 let current = args
1036 .get("current_session_id")
1037 .or_else(|| args.get("session_id"))
1038 .and_then(Value::as_str);
1039 let n = args.get("n").and_then(Value::as_u64).unwrap_or(10) as usize;
1040
1041 let memories = self.store.scan_all(Galaxy::Sessions)?;
1052 let mut starts: Vec<_> = memories
1053 .iter()
1054 .filter(|m| {
1055 m.metadata.tags.contains(&"start".to_string())
1056 && current.is_none_or(|c| m.metadata.id.to_string() != c)
1057 })
1058 .collect();
1059 starts.sort_by_key(|m| std::cmp::Reverse(m.metadata.created_at));
1060 let previous = starts
1061 .iter()
1062 .find(|m| {
1063 let sid = m.metadata.id.to_string();
1064 load_turns(&self.store, Some(&sid), 1, false).is_ok_and(|turns| !turns.is_empty())
1065 })
1066 .copied()
1067 .or_else(|| starts.first().copied());
1068
1069 let Some(prev) = previous else {
1070 let project = std::env::var("WM_PROJECT").ok().filter(|s| !s.is_empty());
1077 let store = self.store.path().display().to_string();
1078 let scope = project.map_or_else(
1079 || format!("store {store}"),
1080 |p| format!("store {store}, project '{p}'"),
1081 );
1082 return Ok(json!({
1083 "status": "success",
1084 "previous_session": null,
1085 "turns": [],
1086 "count": 0,
1087 "message": "no previous session found",
1088 "hint": format!(
1089 "memory is project-scoped and this server's scope ({scope}) has no sessions. If you expected continuity for your project, your client may be wired to a different store: check the mcp block in your opencode config and compare with GET /status on the fleet (store_path, project). Per-project layout: docs/MULTI_PROJECT_MEMORY.md."
1090 ),
1091 }));
1092 };
1093
1094 let prev_id = prev.metadata.id.to_string();
1095 let mut turns = filter_by_time(
1096 load_turns(&self.store, Some(&prev_id), 10_000, false)?,
1097 &args,
1098 )?;
1099 let total = turns.len();
1100 let tail: Vec<Value> = turns
1101 .split_off(total.saturating_sub(n))
1102 .iter()
1103 .map(|(_, v)| format_turn(v, true))
1104 .collect();
1105 let (checkpoint, checkpoint_id) = match latest_checkpoint_handoff(&self.store, &prev_id)? {
1112 Some((id, _, handoff)) => (handoff, Value::String(id)),
1113 None => (Value::Null, Value::Null),
1114 };
1115 Ok(json!({
1116 "status": "success",
1117 "previous_session": prev_id,
1118 "count": tail.len(),
1119 "turns": tail,
1120 "checkpoint": checkpoint,
1121 "checkpoint_id": checkpoint_id,
1122 }))
1123 }
1124 fn stats(&self) -> &ToolStats {
1125 &self.stats
1126 }
1127}
1128
1129pub struct SessionDigestTool {
1136 store: Arc<MemoryStore>,
1137 stats: ToolStats,
1138 effects: EffectRow,
1139}
1140
1141impl SessionDigestTool {
1142 #[must_use]
1143 pub fn new(store: Arc<MemoryStore>) -> Self {
1144 Self {
1145 store,
1146 stats: ToolStats::default(),
1147 effects: EffectRow::read_only(vec![Resource::Galaxy("sessions".into())]),
1148 }
1149 }
1150}
1151
1152const DIGEST_SECTION_ORDER: &[&str] = &["decision", "breakthrough", "error", "summary"];
1154
1155#[async_trait]
1156impl Tool for SessionDigestTool {
1157 fn name(&self) -> &str {
1158 "session.digest"
1159 }
1160 fn gana(&self) -> Gana {
1161 Gana::StraddlingLegs
1162 }
1163 fn effects(&self) -> &EffectRow {
1164 &self.effects
1165 }
1166 fn input_schema(&self) -> Value {
1167 super::common::schema(
1168 &json!({
1169 "session_id": super::common::str_prop("Session to digest (default: most recent)"),
1170 "min_importance": super::common::num_prop("Importance floor (default 0.5)"),
1171 "include_checkpoint": {
1172 "type": "boolean",
1173 "description": "Append the latest checkpoint's git/handoff state (default true)."
1174 },
1175 "since": super::common::str_prop("Time-range floor: epoch seconds, RFC 3339, or YYYY-MM-DD"),
1176 "until": super::common::str_prop("Time-range ceiling: epoch seconds, RFC 3339, or YYYY-MM-DD"),
1177 }),
1178 &[],
1179 )
1180 }
1181 fn description(&self) -> &str {
1182 "Compile a session into a markdown handoff digest — turns grouped by type and importance-ordered, latest checkpoint state appended. Args: session_id (optional), min_importance (default 0.5), include_checkpoint (default true), since/until."
1183 }
1184 async fn call(&self, _ctx: &mut Context, args: Value) -> wm_core::Result<Value> {
1185 let session_id = match args.get("session_id").and_then(Value::as_str) {
1186 Some(sid) if !sid.is_empty() => sid.to_string(),
1187 _ => self
1188 .store
1189 .scan_all(Galaxy::Sessions)?
1190 .iter()
1191 .filter(|m| m.metadata.tags.contains(&"start".to_string()))
1192 .max_by_key(|m| m.metadata.created_at)
1193 .map(|m| m.metadata.id.to_string())
1194 .ok_or_else(|| {
1195 wm_core::CoreError::Tool("no session found — run session.start first".into())
1196 })?,
1197 };
1198 let min_importance = args
1199 .get("min_importance")
1200 .and_then(Value::as_f64)
1201 .unwrap_or(0.5);
1202 let include_checkpoint = args
1203 .get("include_checkpoint")
1204 .and_then(Value::as_bool)
1205 .unwrap_or(true);
1206
1207 let mut turns: Vec<_> = filter_by_time(
1208 load_turns(&self.store, Some(&session_id), 10_000, false)?,
1209 &args,
1210 )?
1211 .into_iter()
1212 .filter(|(_, v)| {
1213 v.get("importance").and_then(Value::as_f64).unwrap_or(0.0) >= min_importance
1214 })
1215 .collect();
1216 turns.sort_by(|a, b| {
1217 b.1.get("importance")
1218 .and_then(Value::as_f64)
1219 .unwrap_or(0.0)
1220 .total_cmp(&a.1.get("importance").and_then(Value::as_f64).unwrap_or(0.0))
1221 });
1222
1223 let mut groups: Vec<(String, Vec<&Value>)> = Vec::new();
1225 for (_, v) in &turns {
1226 let t = v
1227 .get("turn_type")
1228 .and_then(Value::as_str)
1229 .unwrap_or("message")
1230 .to_string();
1231 match groups.iter_mut().find(|(name, _)| *name == t) {
1232 Some((_, list)) => list.push(v),
1233 None => groups.push((t, vec![v])),
1234 }
1235 }
1236 groups.sort_by_key(|(name, _)| {
1237 (
1238 DIGEST_SECTION_ORDER
1239 .iter()
1240 .position(|k| k == name)
1241 .unwrap_or(DIGEST_SECTION_ORDER.len()),
1242 name.clone(),
1243 )
1244 });
1245
1246 let mut digest = format!("# Session handoff — {session_id}\n");
1247 let mut included = 0usize;
1248 for (turn_type, items) in &groups {
1249 writeln!(
1250 digest,
1251 "\n## {} ({})",
1252 capitalize(&pluralize(turn_type)),
1253 items.len()
1254 )
1255 .expect("write to String cannot fail");
1256 for v in items {
1257 let importance = v.get("importance").and_then(Value::as_f64).unwrap_or(0.0);
1258 let content = v.get("content").and_then(Value::as_str).unwrap_or("");
1259 writeln!(digest, "- ({importance:.2}) {content}")
1260 .expect("write to String cannot fail");
1261 included += 1;
1262 }
1263 }
1264
1265 let mut checkpoint_state = Value::Null;
1267 if include_checkpoint {
1268 if let Some((_, _, cp)) = latest_checkpoint_handoff(&self.store, &session_id)? {
1269 digest.push_str("\n## Checkpoint state\n");
1270 if let Some(git) = cp.get("git") {
1271 writeln!(
1272 digest,
1273 "- commit `{}` on `{}` ({} dirty files)",
1274 git.get("commit").and_then(Value::as_str).unwrap_or("?"),
1275 git.get("branch").and_then(Value::as_str).unwrap_or("?"),
1276 git.get("dirty_count").and_then(Value::as_i64).unwrap_or(0)
1277 )
1278 .expect("write to String cannot fail");
1279 }
1280 if let Some(q) = cp.get("next_queue").and_then(Value::as_array) {
1281 if !q.is_empty() {
1282 writeln!(
1283 digest,
1284 "- next queue: {}",
1285 q.iter()
1286 .filter_map(Value::as_str)
1287 .collect::<Vec<_>>()
1288 .join(" → ")
1289 )
1290 .expect("write to String cannot fail");
1291 }
1292 }
1293 if let Some(f) = cp.get("open_flags").and_then(Value::as_array) {
1294 if !f.is_empty() {
1295 writeln!(
1296 digest,
1297 "- open flags: {}",
1298 f.iter()
1299 .filter_map(Value::as_str)
1300 .collect::<Vec<_>>()
1301 .join("; ")
1302 )
1303 .expect("write to String cannot fail");
1304 }
1305 }
1306 if let Some(tg) = cp.get("tests_green") {
1307 writeln!(digest, "- tests green: {tg}").expect("write to String cannot fail");
1308 }
1309 checkpoint_state = cp;
1310 }
1311 }
1312
1313 Ok(json!({
1314 "status": "success",
1315 "session_id": session_id,
1316 "digest": digest,
1317 "turns_included": included,
1318 "turns_total_scanned": turns.len(),
1319 "sections": groups.iter().map(|(t, items)| json!({"type": t, "count": items.len()})).collect::<Vec<_>>(),
1320 "checkpoint": checkpoint_state,
1321 }))
1322 }
1323 fn stats(&self) -> &ToolStats {
1324 &self.stats
1325 }
1326}
1327
1328fn capitalize(s: &str) -> String {
1330 let mut chars = s.chars();
1331 match chars.next() {
1332 Some(first) => first.to_uppercase().collect::<String>() + chars.as_str(),
1333 None => String::new(),
1334 }
1335}
1336
1337fn pluralize(s: &str) -> String {
1340 if let Some(stem) = s.strip_suffix('y') {
1341 format!("{stem}ies")
1342 } else {
1343 format!("{s}s")
1344 }
1345}
1346
1347pub struct SessionHandoffTool {
1349 store: Arc<MemoryStore>,
1350 stats: ToolStats,
1351 effects: EffectRow,
1352 search: Option<Arc<wm_memory::SearchEngine>>,
1353}
1354
1355impl SessionHandoffTool {
1356 #[must_use]
1357 pub fn new(store: Arc<MemoryStore>) -> Self {
1358 Self {
1359 store,
1360 stats: ToolStats::default(),
1361 effects: EffectRow {
1362 writes: vec![Resource::Galaxy("sessions".into())],
1363 ..Default::default()
1364 },
1365 search: None,
1366 }
1367 }
1368
1369 #[must_use]
1371 pub fn with_search(mut self, search: Option<Arc<wm_memory::SearchEngine>>) -> Self {
1372 self.search = search;
1373 self
1374 }
1375}
1376
1377#[async_trait]
1378impl Tool for SessionHandoffTool {
1379 fn name(&self) -> &str {
1380 "session.handoff"
1381 }
1382 fn gana(&self) -> Gana {
1383 Gana::StraddlingLegs
1384 }
1385 fn effects(&self) -> &EffectRow {
1386 &self.effects
1387 }
1388 fn input_schema(&self) -> Value {
1389 super::common::schema(
1390 &json!({
1391 "action": super::common::str_prop("transfer | accept | list"),
1392 "session_id": super::common::str_prop("transfer: session to hand off"),
1393 "message": super::common::str_prop("transfer: handoff note"),
1394 "handoff_id": super::common::str_prop("accept: handoff to accept"),
1395 }),
1396 &["action"],
1397 )
1398 }
1399 fn description(&self) -> &str {
1400 "Transfer or resume a session across devices (actions: transfer, accept, list). transfer: session_id (required) + message; accept: handoff_id; list: pending handoffs."
1401 }
1402 async fn call(&self, _ctx: &mut Context, args: Value) -> wm_core::Result<Value> {
1403 let action = args.get("action").and_then(Value::as_str).unwrap_or("list");
1404 match action {
1405 "transfer" => {
1406 let session_id =
1407 args.get("session_id")
1408 .and_then(Value::as_str)
1409 .ok_or_else(|| {
1410 wm_core::CoreError::InvalidArgs(
1411 "session_id required for transfer".into(),
1412 )
1413 })?;
1414 let message = args.get("message").and_then(Value::as_str).unwrap_or("");
1415 let turns = load_turns(&self.store, Some(session_id), 10_000, false)?;
1416 if turns.is_empty() {
1417 return Err(wm_core::CoreError::Tool(format!(
1418 "session {session_id} has no recorded turns"
1419 )));
1420 }
1421 let summary: Vec<Value> =
1422 turns.iter().map(|(_, v)| format_turn(v, false)).collect();
1423 let handoff_id = format!("handoff-{}", uuid::Uuid::new_v4());
1424 let mut mem = Memory::new(
1425 Galaxy::Sessions,
1426 json!({
1427 "type": "session_handoff",
1428 "handoff_id": handoff_id,
1429 "session_id": session_id,
1430 "message": message,
1431 "status": "pending",
1432 "turn_count": summary.len(),
1433 "summary": summary,
1434 "created_at": wm_core::time::now_unix_millis(),
1435 })
1436 .to_string(),
1437 );
1438 mem.metadata.tags = vec![
1439 "session".into(),
1440 "handoff".into(),
1441 format!("session:{session_id}"),
1442 ];
1443 mem.metadata.importance = 0.8;
1444 self.store.put(Galaxy::Sessions, &mem)?;
1445 super::common::index_memory(self.search.as_deref(), &mem);
1446 Ok(json!({
1447 "status": "success",
1448 "action": "transfer",
1449 "handoff_id": handoff_id,
1450 "session_id": session_id,
1451 "turn_count": summary.len(),
1452 }))
1453 }
1454 "accept" => {
1455 let handoff_id =
1456 args.get("handoff_id")
1457 .and_then(Value::as_str)
1458 .ok_or_else(|| {
1459 wm_core::CoreError::InvalidArgs("handoff_id required for accept".into())
1460 })?;
1461 let memories = self.store.scan_all(Galaxy::Sessions)?;
1462 let found = memories.iter().find(|m| {
1463 m.metadata.tags.contains(&"handoff".to_string())
1464 && m.content.contains(handoff_id)
1465 });
1466 let Some(mem) = found else {
1467 return Err(wm_core::CoreError::Tool(format!(
1468 "handoff {handoff_id} not found"
1469 )));
1470 };
1471 let mut updated = mem.clone();
1472 if let Ok(mut v) = serde_json::from_str::<Value>(&updated.content) {
1473 v["status"] = json!("accepted");
1474 updated.content = v.to_string();
1475 }
1476 self.store.put(Galaxy::Sessions, &updated)?;
1477 super::common::index_memory(self.search.as_deref(), &updated);
1478 Ok(json!({
1479 "status": "success",
1480 "action": "accept",
1481 "handoff_id": handoff_id,
1482 }))
1483 }
1484 "list" => {
1485 let memories = self.store.scan_all(Galaxy::Sessions)?;
1486 let handoffs: Vec<Value> = memories
1487 .iter()
1488 .filter(|m| m.metadata.tags.contains(&"handoff".to_string()))
1489 .filter_map(|m| serde_json::from_str::<Value>(&m.content).ok())
1490 .filter(|v| v.get("status").and_then(Value::as_str) == Some("pending"))
1491 .map(|v| {
1492 json!({
1493 "handoff_id": v.get("handoff_id"),
1494 "session_id": v.get("session_id"),
1495 "message": v.get("message"),
1496 "turn_count": v.get("turn_count"),
1497 })
1498 })
1499 .collect();
1500 Ok(json!({
1501 "status": "success",
1502 "action": "list",
1503 "pending_count": handoffs.len(),
1504 "handoffs": handoffs,
1505 }))
1506 }
1507 other => Err(wm_core::CoreError::InvalidArgs(format!(
1508 "unknown session.handoff action: {other}"
1509 ))),
1510 }
1511 }
1512 fn stats(&self) -> &ToolStats {
1513 &self.stats
1514 }
1515}
1516
1517#[must_use]
1519pub struct SessionExportTool {
1525 store: Arc<MemoryStore>,
1526 stats: ToolStats,
1527 effects: EffectRow,
1528}
1529
1530impl SessionExportTool {
1531 pub fn new(store: Arc<MemoryStore>) -> Self {
1532 Self {
1533 store,
1534 stats: ToolStats::default(),
1535 effects: EffectRow::read_only(vec![Resource::Galaxy("sessions".into())]),
1536 }
1537 }
1538}
1539
1540#[async_trait]
1541impl Tool for SessionExportTool {
1542 fn name(&self) -> &str {
1543 "session.export"
1544 }
1545 fn gana(&self) -> Gana {
1546 Gana::StraddlingLegs
1547 }
1548 fn effects(&self) -> &EffectRow {
1549 &self.effects
1550 }
1551 fn input_schema(&self) -> Value {
1552 super::common::schema(
1553 &json!({
1554 "session_id": super::common::str_prop("Session to export (default: most recent)"),
1555 "path": super::common::str_prop("Write JSONL to this file instead of returning inline"),
1556 }),
1557 &[],
1558 )
1559 }
1560 fn description(&self) -> &str {
1561 "Export a session as JSONL (start marker + turns + checkpoints, preserving ids/timestamps/tags). Args: session_id (optional), path (optional — writes to file; otherwise returns jsonl inline)."
1562 }
1563 async fn call(&self, _ctx: &mut Context, args: Value) -> wm_core::Result<Value> {
1564 let session_id = match args.get("session_id").and_then(Value::as_str) {
1565 Some(sid) if !sid.is_empty() => sid.to_string(),
1566 _ => self
1567 .store
1568 .scan_all(Galaxy::Sessions)?
1569 .iter()
1570 .filter(|m| m.metadata.tags.contains(&"start".to_string()))
1571 .max_by_key(|m| m.metadata.created_at)
1572 .map(|m| m.metadata.id.to_string())
1573 .ok_or_else(|| {
1574 wm_core::CoreError::Tool("no session found — run session.start first".into())
1575 })?,
1576 };
1577
1578 let mut members: Vec<Memory> = self
1582 .store
1583 .scan_all(Galaxy::Sessions)?
1584 .into_iter()
1585 .filter(|m| m.metadata.id.to_string() == session_id || m.content.contains(&session_id))
1586 .collect();
1587 members.sort_by_key(|m| m.metadata.created_at);
1588
1589 let mut jsonl = String::new();
1590 let header = wm_memory::envelope::EnvelopeHeader::new("session_export", members.len());
1596 jsonl.push_str(&header.header_line());
1597 jsonl.push('\n');
1598 for m in &members {
1599 let line = serde_json::to_string(m)
1600 .map_err(|e| wm_core::CoreError::Tool(format!("export serialize: {e}")))?;
1601 jsonl.push_str(&line);
1602 jsonl.push('\n');
1603 }
1604
1605 let path_arg = args
1606 .get("path")
1607 .and_then(Value::as_str)
1608 .filter(|s| !s.is_empty());
1609 if let Some(dest) = path_arg {
1610 std::fs::write(dest, &jsonl)
1611 .map_err(|e| wm_core::CoreError::Tool(format!("export write {dest}: {e}")))?;
1612 Ok(json!({
1613 "status": "success",
1614 "session_id": session_id,
1615 "records": members.len(),
1616 "path": dest,
1617 }))
1618 } else {
1619 Ok(json!({
1620 "status": "success",
1621 "session_id": session_id,
1622 "records": members.len(),
1623 "jsonl": jsonl,
1624 }))
1625 }
1626 }
1627 fn stats(&self) -> &ToolStats {
1628 &self.stats
1629 }
1630}
1631
1632pub struct SessionImportTool {
1643 store: Arc<MemoryStore>,
1644 search: Option<Arc<wm_memory::SearchEngine>>,
1645 stats: ToolStats,
1646 effects: EffectRow,
1647}
1648
1649impl SessionImportTool {
1650 pub fn new(store: Arc<MemoryStore>, search: Option<Arc<wm_memory::SearchEngine>>) -> Self {
1651 Self {
1652 store,
1653 search,
1654 stats: ToolStats::default(),
1655 effects: EffectRow {
1656 writes: vec![Resource::Galaxy("sessions".into())],
1657 ..Default::default()
1658 },
1659 }
1660 }
1661}
1662
1663#[async_trait]
1664impl Tool for SessionImportTool {
1665 fn name(&self) -> &str {
1666 "session.import"
1667 }
1668 fn gana(&self) -> Gana {
1669 Gana::StraddlingLegs
1670 }
1671 fn effects(&self) -> &EffectRow {
1672 &self.effects
1673 }
1674 fn input_schema(&self) -> Value {
1675 super::common::schema(
1676 &json!({
1677 "path": super::common::str_prop("Read JSONL from this file"),
1678 "jsonl": super::common::str_prop("Or pass the export payload inline"),
1679 }),
1680 &[],
1681 )
1682 }
1683 fn description(&self) -> &str {
1684 "Import sessions from session.export JSONL (path or inline jsonl) — envelope-v2 header validated when present, bare v1 accepted; preserves ids, timestamps, and tags; indexes into the search index; overwrites on id collision."
1685 }
1686 async fn call(&self, _ctx: &mut Context, args: Value) -> wm_core::Result<Value> {
1687 let payload = match args
1688 .get("path")
1689 .and_then(Value::as_str)
1690 .filter(|s| !s.is_empty())
1691 {
1692 Some(path) => std::fs::read_to_string(path)
1693 .map_err(|e| wm_core::CoreError::Tool(format!("import read {path}: {e}")))?,
1694 None => args
1695 .get("jsonl")
1696 .and_then(Value::as_str)
1697 .filter(|s| !s.is_empty())
1698 .ok_or_else(|| {
1699 wm_core::CoreError::InvalidArgs(
1700 "provide either 'path' or inline 'jsonl'".into(),
1701 )
1702 })?
1703 .to_string(),
1704 };
1705
1706 let mut envelope: Option<wm_memory::envelope::EnvelopeHeader> = None;
1711 let mut record_lines: Vec<&str> = Vec::new();
1712 let mut header_consumed = false;
1713 for line in payload.lines() {
1714 if line.trim().is_empty() {
1715 continue;
1716 }
1717 if !header_consumed {
1718 header_consumed = true;
1719 match wm_memory::envelope::read_header_line(line) {
1720 wm_memory::envelope::HeaderRead::Header(h) => {
1721 envelope = Some(h);
1722 continue;
1723 }
1724 wm_memory::envelope::HeaderRead::Refused(msg) => {
1725 return Err(wm_core::CoreError::Tool(msg));
1726 }
1727 wm_memory::envelope::HeaderRead::NotAHeader => {}
1728 }
1729 }
1730 record_lines.push(line);
1731 }
1732
1733 let readonly_engine = self.search.as_ref().is_some_and(|s| s.is_readonly());
1738 if readonly_engine {
1739 tracing::warn!(
1740 "session.import running against a read-only search engine — records land \
1741 in LMDB unindexed; they become searchable at the next writable startup \
1742 (heal_index_drift)"
1743 );
1744 }
1745 let mut writer_slot = match (&self.search, readonly_engine) {
1746 (Some(s), false) => s.writer().ok(),
1747 _ => None,
1748 };
1749
1750 let mut imported = 0usize;
1751 let mut indexed = 0usize;
1752 let mut skipped = 0usize;
1753 let mut session_ids: Vec<String> = Vec::new();
1754 for (lineno, line) in record_lines.iter().enumerate() {
1755 let mem: Memory = match serde_json::from_str(line) {
1756 Ok(m) => m,
1757 Err(e) => {
1758 skipped += 1;
1759 tracing::warn!(line = lineno + 1, error = %e, "skipping unparseable export line");
1760 continue;
1761 }
1762 };
1763 if let Ok(parsed) = serde_json::from_str::<Value>(&mem.content) {
1764 if let Some(sid) = parsed.get("session_id").and_then(Value::as_str) {
1765 if !session_ids.iter().any(|s| s == sid) {
1766 session_ids.push(sid.to_string());
1767 }
1768 }
1769 }
1770 if let (Some(search), Some(writer)) = (&self.search, writer_slot.as_mut()) {
1774 let id_str = mem.metadata.id.to_string();
1775 let _ = search.delete_document(writer, &id_str);
1776 match search.add_document(
1777 writer,
1778 &id_str,
1779 mem.metadata.galaxy.db_name(),
1780 &mem.content,
1781 &mem.metadata.tags,
1782 mem.metadata.created_at.timestamp(),
1783 ) {
1784 Ok(()) => indexed += 1,
1785 Err(e) => {
1786 tracing::warn!(id = %id_str, error = %e, "import index add failed (LMDB record kept)");
1787 }
1788 }
1789 }
1790 self.store.put(Galaxy::Sessions, &mem)?;
1791 imported += 1;
1792 }
1793
1794 if let Some(search) = &self.search {
1795 if let Some(mut writer) = writer_slot {
1796 search
1797 .commit(&mut writer)
1798 .map_err(|e| wm_core::CoreError::Tool(format!("import index commit: {e}")))?;
1799 }
1800 }
1801
1802 let mut warnings: Vec<String> = Vec::new();
1803 if let Some(h) = &envelope {
1804 if h.count != imported {
1805 let msg = format!(
1806 "envelope declares count {} but {} records imported",
1807 h.count, imported
1808 );
1809 tracing::warn!("{msg}");
1810 warnings.push(msg);
1811 }
1812 }
1813
1814 let envelope_info = envelope.as_ref().map(|h| {
1815 json!({
1816 "format_version": h.format_version,
1817 "kind": h.kind,
1818 "generator": h.generator,
1819 "created_at": h.created_at,
1820 "declared_count": h.count,
1821 })
1822 });
1823
1824 Ok(json!({
1825 "status": "success",
1826 "imported": imported,
1827 "skipped": skipped,
1828 "session_ids": session_ids,
1829 "indexed": indexed,
1830 "envelope": envelope_info,
1831 "warnings": warnings,
1832 }))
1833 }
1834 fn stats(&self) -> &ToolStats {
1835 &self.stats
1836 }
1837}
1838
1839fn track_of(mem: &Memory) -> Option<&str> {
1841 mem.metadata
1842 .tags
1843 .iter()
1844 .find_map(|t| t.strip_prefix("track:"))
1845}
1846
1847pub struct SessionTrackLogTool {
1858 store: Arc<MemoryStore>,
1859 stats: ToolStats,
1860 effects: EffectRow,
1861}
1862
1863impl SessionTrackLogTool {
1864 #[must_use]
1865 pub fn new(store: Arc<MemoryStore>) -> Self {
1866 Self {
1867 store,
1868 stats: ToolStats::default(),
1869 effects: EffectRow::read_only(vec![Resource::Galaxy("sessions".into())]),
1870 }
1871 }
1872}
1873
1874#[async_trait]
1875impl Tool for SessionTrackLogTool {
1876 fn name(&self) -> &str {
1877 "session.track_log"
1878 }
1879 fn gana(&self) -> Gana {
1880 Gana::StraddlingLegs
1881 }
1882 fn effects(&self) -> &EffectRow {
1883 &self.effects
1884 }
1885 fn input_schema(&self) -> Value {
1886 super::common::schema(
1887 &json!({
1888 "track": super::common::str_prop("Track slug to read (recorded via session.record/session.checkpoint 'track')"),
1889 "tracks": super::common::str_array_prop("Multiple track slugs — merged chronologically (the related-ticket review view)"),
1890 "since": super::common::str_prop("Time-range floor: epoch seconds, RFC 3339, or YYYY-MM-DD"),
1891 "until": super::common::str_prop("Time-range ceiling: epoch seconds, RFC 3339, or YYYY-MM-DD"),
1892 "limit": super::common::positive_int_prop("Max entries returned per track, most recent kept (default 100)"),
1893 "include_superseded": {
1894 "type": "boolean",
1895 "description": "Include superseded turns (default false — the current story only)."
1896 },
1897 }),
1898 &[],
1899 )
1900 }
1901 fn description(&self) -> &str {
1902 "Read the per-track implementation log — turns and checkpoints tagged with a track slug. Args: track (single) or tracks (array, merged chronologically), since/until, limit (default 100, most recent kept per track), include_superseded. With neither track nor tracks: an overview of every track (entry count, last activity, latest preview) for cross-track alignment review."
1903 }
1904 async fn call(&self, _ctx: &mut Context, args: Value) -> wm_core::Result<Value> {
1905 let mut requested: Vec<String> = Vec::new();
1906 if let Some(track) = args.get("track").and_then(Value::as_str) {
1907 validate_track(track)?;
1908 requested.push(track.to_string());
1909 }
1910 if let Some(tracks) = args.get("tracks").and_then(Value::as_array) {
1911 for t in tracks {
1912 let t = t.as_str().ok_or_else(|| {
1913 wm_core::CoreError::InvalidArgs("'tracks' entries must be strings".into())
1914 })?;
1915 validate_track(t)?;
1916 if !requested.iter().any(|r| r == t) {
1917 requested.push(t.to_string());
1918 }
1919 }
1920 }
1921 let limit = args
1922 .get("limit")
1923 .and_then(Value::as_u64)
1924 .unwrap_or(100)
1925 .clamp(1, 1000) as usize;
1926 let include_superseded = args
1927 .get("include_superseded")
1928 .and_then(Value::as_bool)
1929 .unwrap_or(false);
1930
1931 let since = match args.get("since") {
1934 Some(v) if !v.is_null() => Some(parse_time_bound(v, false).ok_or_else(|| {
1935 wm_core::CoreError::InvalidArgs(
1936 "invalid 'since' — use epoch seconds, RFC 3339, or YYYY-MM-DD".into(),
1937 )
1938 })?),
1939 _ => None,
1940 };
1941 let until = match args.get("until") {
1942 Some(v) if !v.is_null() => Some(parse_time_bound(v, true).ok_or_else(|| {
1943 wm_core::CoreError::InvalidArgs(
1944 "invalid 'until' — use epoch seconds, RFC 3339, or YYYY-MM-DD".into(),
1945 )
1946 })?),
1947 _ => None,
1948 };
1949
1950 type Entries = std::collections::BTreeMap<String, Vec<(DateTime<Utc>, Value)>>;
1951 let mut by_track: Entries = std::collections::BTreeMap::new();
1952 let mut latest_checkpoint: std::collections::BTreeMap<String, (DateTime<Utc>, Value)> =
1953 std::collections::BTreeMap::new();
1954
1955 for mem in &self.store.scan_all(Galaxy::Sessions)? {
1956 let Some(track) = track_of(mem) else { continue };
1957 if !requested.is_empty() && !requested.iter().any(|r| r == track) {
1958 continue;
1959 }
1960 if !include_superseded
1961 && mem
1962 .metadata
1963 .tags
1964 .iter()
1965 .any(|t| t.starts_with("superseded-by:"))
1966 {
1967 continue;
1968 }
1969 let created_at = mem.metadata.created_at;
1970 if since.is_some_and(|t| created_at < t) || until.is_some_and(|t| created_at > t) {
1971 continue;
1972 }
1973 let entry = if let Some(turn) = turn_json(mem) {
1974 Some(json!({
1975 "kind": "turn",
1976 "track": track,
1977 "session_id": turn.get("session_id"),
1978 "sequence": turn.get("sequence"),
1979 "role": turn.get("role"),
1980 "turn_type": turn.get("turn_type"),
1981 "importance": turn.get("importance"),
1982 "created_at": created_at.to_rfc3339(),
1983 "content": turn.get("content"),
1984 }))
1985 } else if mem.metadata.tags.contains(&"checkpoint".to_string()) {
1986 serde_json::from_str::<Value>(&mem.content)
1987 .ok()
1988 .map(|parsed| {
1989 json!({
1990 "kind": "checkpoint",
1991 "track": track,
1992 "session_id": parsed.get("session_id"),
1993 "label": parsed.get("label"),
1994 "created_at": created_at.to_rfc3339(),
1995 "handoff": parsed.get("handoff"),
1996 })
1997 })
1998 } else {
1999 None
2000 };
2001 let Some(entry) = entry else { continue };
2002 if entry["kind"] == "checkpoint" {
2003 latest_checkpoint
2004 .entry(track.to_string())
2005 .and_modify(|(ts, existing)| {
2006 if created_at > *ts {
2007 *ts = created_at;
2008 *existing = entry.clone();
2009 }
2010 })
2011 .or_insert_with(|| (created_at, entry.clone()));
2012 }
2013 by_track
2014 .entry(track.to_string())
2015 .or_default()
2016 .push((created_at, entry));
2017 }
2018
2019 let now = Utc::now();
2020 let overview = requested.is_empty();
2021 let track_names: Vec<String> = if overview {
2022 by_track.keys().cloned().collect()
2023 } else {
2024 requested
2025 };
2026 let mut tracks: Vec<Value> = Vec::with_capacity(track_names.len());
2027 for track in track_names {
2028 let mut list = by_track.remove(&track).unwrap_or_default();
2029 list.sort_by_key(|(ts, _)| *ts);
2030 let entry_count = list.len();
2031 let last_activity = list.last().map(|(ts, _)| *ts);
2032 let latest_entry = list.last().map(|(_, e)| {
2033 json!({
2034 "kind": e.get("kind"),
2035 "turn_type": e.get("turn_type"),
2036 "label": e.get("label"),
2037 "preview": e
2038 .get("content")
2039 .and_then(Value::as_str)
2040 .map(|s| s.chars().take(160).collect::<String>()),
2041 })
2042 });
2043 if list.len() > limit {
2044 list.drain(0..list.len() - limit);
2045 }
2046 let mut out = json!({
2047 "track": track,
2048 "known": entry_count > 0,
2049 "entry_count": entry_count,
2050 "returned": list.len(),
2051 "last_activity_at": last_activity.map(|t| t.to_rfc3339()),
2052 "age_seconds": last_activity.map(|t| (now - t).num_seconds().max(0)),
2053 "latest_checkpoint": latest_checkpoint.get(&track).map(|(ts, cp)| {
2054 json!({
2055 "created_at": ts.to_rfc3339(),
2056 "entry": cp,
2057 })
2058 }),
2059 });
2060 if overview {
2061 out["latest_entry"] = latest_entry.unwrap_or(Value::Null);
2062 } else {
2063 out["entries"] = Value::Array(list.into_iter().map(|(_, e)| e).collect());
2064 }
2065 tracks.push(out);
2066 }
2067
2068 Ok(json!({
2069 "status": "success",
2070 "mode": if overview { "overview" } else { "log" },
2071 "track_count": tracks.len(),
2072 "tracks": tracks,
2073 "disclosure": "directional alignment is the reader's judgment — this view supplies the log, the checkpoint handoff, and staleness facts (age_seconds)",
2074 }))
2075 }
2076 fn stats(&self) -> &ToolStats {
2077 &self.stats
2078 }
2079}
2080
2081pub fn register_session_ops(
2082 registry: &wm_dispatch::ToolRegistry,
2083 store: &Arc<MemoryStore>,
2084 search: Option<Arc<wm_memory::SearchEngine>>,
2085) -> wm_dispatch::ToolRegistry {
2086 registry
2087 .register(Arc::new(
2088 SessionRecordTool::new(store.clone()).with_search(search.clone()),
2089 ))
2090 .register(Arc::new(SessionReplayTool::new(store.clone())))
2091 .register(Arc::new(SessionContinuityTool::new(store.clone())))
2092 .register(Arc::new(SessionTrackLogTool::new(store.clone())))
2093 .register(Arc::new(
2094 SessionHandoffTool::new(store.clone()).with_search(search.clone()),
2095 ))
2096 .register(Arc::new(SessionExportTool::new(store.clone())))
2097 .register(Arc::new(SessionImportTool::new(store.clone(), search)))
2098}
2099
2100#[cfg(test)]
2101mod tests {
2102 use super::*;
2103 use crate::expansion::session::SessionCheckpointNodiscoveryTool;
2104
2105 fn test_store() -> Arc<MemoryStore> {
2106 let dir = tempfile::tempdir().unwrap();
2107 let path = dir.path().join("lmdb");
2108 std::fs::create_dir_all(&path).unwrap();
2109 Arc::new(MemoryStore::open_default(path).unwrap())
2110 }
2111
2112 fn start_session(store: &MemoryStore) -> String {
2113 let mut mem = Memory::new(
2114 Galaxy::Sessions,
2115 json!({"type": "session_start"}).to_string(),
2116 );
2117 mem.metadata.tags = vec!["session".into(), "start".into()];
2118 store.put(Galaxy::Sessions, &mem).unwrap();
2119 mem.metadata.id.to_string()
2120 }
2121
2122 fn start_session_aged(store: &MemoryStore, age_secs: i64) -> String {
2125 let mut mem = Memory::new(
2126 Galaxy::Sessions,
2127 json!({"type": "session_start"}).to_string(),
2128 );
2129 mem.metadata.tags = vec!["session".into(), "start".into()];
2130 mem.metadata.created_at = chrono::Utc::now() - chrono::Duration::seconds(age_secs);
2131 store.put(Galaxy::Sessions, &mem).unwrap();
2132 mem.metadata.id.to_string()
2133 }
2134
2135 fn record_aged_turn(store: &MemoryStore, sid: &str, age_days: u32, content: &str) {
2137 let mut mem = Memory::new(
2138 Galaxy::Sessions,
2139 json!({
2140 "type": "session_turn",
2141 "session_id": sid,
2142 "role": "ai",
2143 "turn_type": "decision",
2144 "importance": 0.9,
2145 "content": content,
2146 })
2147 .to_string(),
2148 );
2149 mem.metadata.tags = vec!["session".into(), "turn".into(), format!("session:{sid}")];
2150 mem.metadata.created_at = Utc::now() - chrono::Duration::days(i64::from(age_days));
2151 store.put(Galaxy::Sessions, &mem).unwrap();
2152 }
2153
2154 #[tokio::test]
2155 async fn replay_time_filters_since_and_until() {
2156 let store = test_store();
2157 let sid = start_session(&store);
2158 record_aged_turn(&store, &sid, 3, "three days ago");
2159 record_aged_turn(&store, &sid, 2, "two days ago");
2160 record_aged_turn(&store, &sid, 1, "yesterday");
2161 record_aged_turn(&store, &sid, 0, "today");
2162
2163 let replay = SessionReplayTool::new(store);
2164 let mut ctx = Context::default();
2165
2166 let two_days_ago = (Utc::now() - chrono::Duration::days(2)).format("%Y-%m-%d");
2168 let v = replay
2169 .call(
2170 &mut ctx,
2171 json!({"session_id": sid, "since": two_days_ago.to_string()}),
2172 )
2173 .await
2174 .unwrap();
2175 assert_eq!(v["count"], 3, "since=date keeps day-of + later: {v}");
2176
2177 let until_epoch = (Utc::now() - chrono::Duration::hours(23)).timestamp();
2179 let v = replay
2180 .call(&mut ctx, json!({"session_id": sid, "until": until_epoch}))
2181 .await
2182 .unwrap();
2183 assert_eq!(
2184 v["count"], 3,
2185 "until=epoch(23h ago) keeps the three older turns: {v}"
2186 );
2187
2188 let since = (Utc::now() - chrono::Duration::hours(60)).to_rfc3339();
2190 let until = (Utc::now() - chrono::Duration::hours(12)).to_rfc3339();
2191 let v = replay
2192 .call(
2193 &mut ctx,
2194 json!({"session_id": sid, "since": since, "until": until}),
2195 )
2196 .await
2197 .unwrap();
2198 assert_eq!(v["count"], 2, "window keeps two/two-days-ago turns: {v}");
2199 for turn in v["turns"].as_array().unwrap() {
2200 assert_ne!(
2201 turn["content"], "today",
2202 "time filters must exclude out-of-window turns"
2203 );
2204 }
2205
2206 assert!(
2208 replay
2209 .call(&mut ctx, json!({"session_id": sid, "since": "not-a-date"}))
2210 .await
2211 .is_err()
2212 );
2213 }
2214
2215 #[tokio::test]
2216 async fn replay_omitted_session_id_uses_latest_session_start() {
2217 let store = test_store();
2218 let older = start_session_aged(&store, 120);
2222 record_aged_turn(&store, &older, 0, "older-session-only");
2223 let latest = start_session_aged(&store, 60);
2224 record_aged_turn(&store, &latest, 0, "latest-session-only");
2225
2226 let replay = SessionReplayTool::new(store);
2227 let mut ctx = Context::default();
2228
2229 let omitted = replay
2230 .call(&mut ctx, json!({"mode": "full"}))
2231 .await
2232 .unwrap();
2233 assert_eq!(omitted["session_id"], latest);
2234 assert_eq!(
2235 omitted["count"], 1,
2236 "omitted id must not combine sessions: {omitted}"
2237 );
2238 assert_eq!(omitted["turns"][0]["content"], "latest-session-only");
2239
2240 let explicit_older = replay
2241 .call(&mut ctx, json!({"session_id": older, "mode": "full"}))
2242 .await
2243 .unwrap();
2244 assert_eq!(explicit_older["session_id"], older);
2245 assert_eq!(explicit_older["count"], 1);
2246 assert_eq!(explicit_older["turns"][0]["content"], "older-session-only");
2247 }
2248
2249 #[tokio::test]
2250 async fn replay_without_any_session_remains_truthfully_empty() {
2251 let replay = SessionReplayTool::new(test_store());
2252 let mut ctx = Context::default();
2253
2254 let value = replay.call(&mut ctx, json!({})).await.unwrap();
2255 assert_eq!(value["status"], "success");
2256 assert_eq!(value["session_id"], Value::Null);
2257 assert_eq!(value["count"], 0);
2258 assert_eq!(value["turns"], json!([]));
2259 }
2260
2261 #[tokio::test]
2262 async fn replay_omitted_id_does_not_combine_orphan_turns_without_a_start() {
2263 let store = test_store();
2264 record_aged_turn(&store, "orphan-a", 0, "orphan-a-only");
2268 record_aged_turn(&store, "orphan-b", 0, "orphan-b-only");
2269
2270 let replay = SessionReplayTool::new(store);
2271 let mut ctx = Context::default();
2272
2273 let omitted = replay.call(&mut ctx, json!({})).await.unwrap();
2274 assert_eq!(omitted["session_id"], Value::Null);
2275 assert_eq!(
2276 omitted["count"], 0,
2277 "omitted id must not combine orphans: {omitted}"
2278 );
2279 assert_eq!(omitted["turns"], json!([]));
2280
2281 let explicit = replay
2282 .call(&mut ctx, json!({"session_id": "orphan-a"}))
2283 .await
2284 .unwrap();
2285 assert_eq!(explicit["session_id"], "orphan-a");
2286 assert_eq!(explicit["count"], 1);
2287 assert_eq!(explicit["turns"][0]["content"], "orphan-a-only");
2288 }
2289
2290 #[tokio::test]
2291 async fn continuity_respects_since_filter() {
2292 let store = test_store();
2293 let sid1 = start_session(&store);
2294 record_aged_turn(&store, &sid1, 5, "ancient decision");
2295 record_aged_turn(&store, &sid1, 0, "fresh decision");
2296 let sid2 = start_session(&store);
2297
2298 let continuity = SessionContinuityTool::new(store);
2299 let mut ctx = Context::default();
2300 let cutoff = (Utc::now() - chrono::Duration::days(1))
2301 .format("%Y-%m-%d")
2302 .to_string();
2303 let v = continuity
2304 .call(
2305 &mut ctx,
2306 json!({"current_session_id": sid2, "since": cutoff, "n": 10}),
2307 )
2308 .await
2309 .unwrap();
2310 assert_eq!(v["count"], 1, "only the fresh turn is in range: {v}");
2311 assert_eq!(v["turns"][0]["content"], "fresh decision");
2312
2313 let all = continuity
2315 .call(&mut ctx, json!({"current_session_id": sid2, "n": 10}))
2316 .await
2317 .unwrap();
2318 assert_eq!(all["count"], 2);
2319 }
2320
2321 #[tokio::test]
2322 async fn continuity_skips_empty_newest_session() {
2323 let store = test_store();
2327 let sid1 = start_session(&store);
2328 record_aged_turn(
2329 &store,
2330 &sid1,
2331 0,
2332 "Decision: use SQLite for the report cache",
2333 );
2334 let sid2 = start_session(&store); let continuity = SessionContinuityTool::new(store);
2337 let mut ctx = Context::default();
2338 let v = continuity.call(&mut ctx, json!({"n": 5})).await.unwrap();
2339 assert_ne!(
2340 v["previous_session"], sid2,
2341 "empty newest session must be skipped: {v}"
2342 );
2343 assert_eq!(v["previous_session"], sid1);
2344 assert_eq!(v["count"], 1);
2345 assert!(
2346 v["turns"][0]["content"]
2347 .as_str()
2348 .unwrap()
2349 .contains("SQLite")
2350 );
2351 }
2352
2353 #[tokio::test]
2354 async fn digest_groups_by_type_and_respects_importance_floor() {
2355 let store = test_store();
2356 let sid = start_session(&store);
2357 for (turn_type, importance, content) in [
2359 ("summary", 0.6, "wrapped up"),
2360 ("decision", 0.9, "picked architecture CO over alternatives"),
2361 ("error", 0.95, "startForce root cause found"),
2362 ("breakthrough", 0.85, "watch resolution insight"),
2363 ("message", 0.3, "low-value chatter"),
2364 ] {
2365 let mut mem = Memory::new(
2366 Galaxy::Sessions,
2367 json!({
2368 "type": "session_turn",
2369 "session_id": sid,
2370 "role": "ai",
2371 "turn_type": turn_type,
2372 "importance": importance,
2373 "content": content,
2374 })
2375 .to_string(),
2376 );
2377 mem.metadata.tags = vec!["session".into(), "turn".into()];
2378 store.put(Galaxy::Sessions, &mem).unwrap();
2379 }
2380
2381 let tool = SessionDigestTool::new(store);
2382 let mut ctx = Context::default();
2383 let v = tool
2384 .call(&mut ctx, json!({"session_id": sid, "min_importance": 0.5}))
2385 .await
2386 .unwrap();
2387
2388 assert_eq!(v["status"], "success");
2389 let digest = v["digest"].as_str().unwrap();
2390 for expected in [
2392 "startForce root cause found",
2393 "picked architecture CO over alternatives",
2394 "watch resolution insight",
2395 ] {
2396 assert!(
2397 digest.contains(expected),
2398 "digest must contain '{expected}': {digest}"
2399 );
2400 }
2401 assert!(!digest.contains("low-value chatter"), "got: {digest}");
2403 let d = digest.find("## Decisions").unwrap();
2405 let b = digest.find("## Breakthroughs").unwrap();
2406 let e = digest.find("## Errors").unwrap();
2407 let s = digest.find("## Summaries").unwrap();
2408 assert!(
2409 d < b && b < e && e < s,
2410 "sections must follow canonical order: {digest}"
2411 );
2412 assert_eq!(v["turns_included"], 4);
2413
2414 let strict = tool
2416 .call(&mut ctx, json!({"session_id": sid, "min_importance": 0.9}))
2417 .await
2418 .unwrap();
2419 let strict_digest = strict["digest"].as_str().unwrap();
2420 assert!(strict_digest.contains("startForce"));
2421 assert!(!strict_digest.contains("watch resolution insight"));
2422 }
2423
2424 #[tokio::test]
2425 async fn digest_appends_checkpoint_state() {
2426 let store = test_store();
2427 let sid = start_session(&store);
2428
2429 let mut cp = Memory::new(
2431 Galaxy::Sessions,
2432 json!({
2433 "type": "checkpoint",
2434 "session_id": sid,
2435 "label": "wrap",
2436 "data": {},
2437 "handoff": {
2438 "git": {
2439 "commit": "abc1234",
2440 "branch": "main",
2441 "dirty_count": 2
2442 },
2443 "tests_green": true,
2444 "next_queue": ["first task", "second task"],
2445 "open_flags": ["flaky probe"]
2446 }
2447 })
2448 .to_string(),
2449 );
2450 cp.metadata.tags = vec!["session".into(), "checkpoint".into()];
2451 store.put(Galaxy::Sessions, &cp).unwrap();
2452
2453 let tool = SessionDigestTool::new(store);
2454 let mut ctx = Context::default();
2455 let v = tool
2456 .call(&mut ctx, json!({"session_id": sid}))
2457 .await
2458 .unwrap();
2459
2460 let digest = v["digest"].as_str().unwrap();
2461 assert!(digest.contains("## Checkpoint state"), "got: {digest}");
2462 assert!(digest.contains("abc1234"));
2463 assert!(digest.contains("first task → second task"));
2464 assert!(digest.contains("flaky probe"));
2465 assert_eq!(v["checkpoint"]["git"]["branch"], "main");
2466 }
2467
2468 #[tokio::test]
2469 async fn supersedes_hides_old_turn_until_requested() {
2470 let store = test_store();
2474 let sid = start_session(&store);
2475 let record = SessionRecordTool::new(store.clone());
2476 let mut ctx = Context::default();
2477
2478 let first = record
2479 .call(
2480 &mut ctx,
2481 json!({"role": "ai", "turn_type": "decision", "importance": 0.9,
2482 "content": "perf: 240ms", "session_id": sid}),
2483 )
2484 .await
2485 .unwrap();
2486 let old_id = first["memory_id"].as_str().unwrap().to_string();
2487
2488 let second = record
2489 .call(
2490 &mut ctx,
2491 json!({"role": "ai", "turn_type": "decision", "importance": 0.9,
2492 "content": "perf revised: 180ms after warm cache",
2493 "session_id": sid, "supersedes": old_id}),
2494 )
2495 .await
2496 .unwrap();
2497 assert_eq!(second["status"], "success");
2498 let _new_id = second["memory_id"].as_str().unwrap().to_string();
2499
2500 let replay = SessionReplayTool::new(store.clone());
2502 let v = replay
2503 .call(&mut ctx, json!({"session_id": sid}))
2504 .await
2505 .unwrap();
2506 assert_eq!(
2507 v["count"], 1,
2508 "superseded turn must be hidden by default: {v}"
2509 );
2510 assert_eq!(
2511 v["turns"][0]["content"],
2512 "perf revised: 180ms after warm cache"
2513 );
2514
2515 let with_history = replay
2517 .call(
2518 &mut ctx,
2519 json!({"session_id": sid, "include_superseded": true}),
2520 )
2521 .await
2522 .unwrap();
2523 assert_eq!(with_history["count"], 2, "got: {with_history}");
2524
2525 let sid2 = start_session(&store);
2527 let continuity = SessionContinuityTool::new(store.clone());
2528 let c = continuity
2529 .call(&mut ctx, json!({"current_session_id": sid2}))
2530 .await
2531 .unwrap();
2532 assert_eq!(c["count"], 1, "continuity must skip superseded turns: {c}");
2533 assert_eq!(
2534 c["turns"][0]["content"],
2535 "perf revised: 180ms after warm cache"
2536 );
2537
2538 let digest = SessionDigestTool::new(store);
2539 let d = digest
2540 .call(&mut ctx, json!({"session_id": sid, "min_importance": 0.5}))
2541 .await
2542 .unwrap();
2543 let digest_text = d["digest"].as_str().unwrap();
2544 assert!(digest_text.contains("180ms"), "got: {digest_text}");
2545 assert!(
2546 !digest_text.contains("240ms"),
2547 "superseded claim must not leak: {digest_text}"
2548 );
2549 }
2550
2551 #[tokio::test]
2552 async fn export_import_roundtrip_preserves_history() {
2553 let store_a = test_store();
2556 let sid = start_session(&store_a);
2557 let record = SessionRecordTool::new(store_a.clone());
2558 let mut ctx = Context::default();
2559 record_aged_turn(&store_a, &sid, 2, "day-one decision");
2560 record_aged_turn(&store_a, &sid, 1, "day-two decision");
2561 let first = record
2562 .call(
2563 &mut ctx,
2564 json!({"role": "ai", "turn_type": "decision", "importance": 0.9,
2565 "content": "original claim", "session_id": sid}),
2566 )
2567 .await
2568 .unwrap();
2569 record
2570 .call(
2571 &mut ctx,
2572 json!({"role": "ai", "turn_type": "decision", "importance": 0.9,
2573 "content": "corrected claim",
2574 "session_id": sid,
2575 "supersedes": first["memory_id"].as_str().unwrap()}),
2576 )
2577 .await
2578 .unwrap();
2579
2580 let export = SessionExportTool::new(store_a.clone());
2582 {
2584 let mut marker: Memory = store_a
2585 .scan_all(Galaxy::Sessions)
2586 .unwrap()
2587 .into_iter()
2588 .find(|m| m.metadata.tags.contains(&"start".to_string()))
2589 .unwrap();
2590 marker.metadata.title = Some("The Big Decision".to_string());
2591 marker.metadata.topic = Some("v8-slices".to_string());
2592 store_a.put(Galaxy::Sessions, &marker).unwrap();
2593 }
2594 let exported = export
2595 .call(&mut ctx, json!({"session_id": sid}))
2596 .await
2597 .unwrap();
2598 assert_eq!(exported["status"], "success");
2599 let jsonl = exported["jsonl"].as_str().unwrap();
2600 assert_eq!(
2602 exported["records"], 5,
2603 "start + 2 aged + 2 claims: {exported}"
2604 );
2605 assert_eq!(jsonl.lines().count(), 6);
2607 let header_line = jsonl.lines().next().unwrap();
2608 match wm_memory::envelope::read_header_line(header_line) {
2609 wm_memory::envelope::HeaderRead::Header(h) => {
2610 assert_eq!(h.kind, "session_export");
2611 assert_eq!(h.count, 5);
2612 }
2613 other => panic!("first line must be the envelope header, got {other:?}"),
2614 }
2615
2616 let store_b = test_store();
2618 let import = SessionImportTool::new(store_b.clone(), None);
2619 let imported = import
2620 .call(&mut ctx, json!({"jsonl": jsonl}))
2621 .await
2622 .unwrap();
2623 assert_eq!(imported["imported"], 5, "got: {imported}");
2624 assert_eq!(imported["skipped"], 0);
2625 assert_eq!(imported["session_ids"], json!([sid]));
2626 assert_eq!(imported["envelope"]["format_version"], 2);
2628 assert_eq!(imported["envelope"]["declared_count"], 5);
2629 assert_eq!(imported["warnings"], json!([]));
2630 let marker_b = store_b
2632 .scan_all(Galaxy::Sessions)
2633 .unwrap()
2634 .into_iter()
2635 .find(|m| m.metadata.tags.contains(&"start".to_string()))
2636 .unwrap();
2637 assert_eq!(marker_b.metadata.title.as_deref(), Some("The Big Decision"));
2638 assert_eq!(marker_b.metadata.topic.as_deref(), Some("v8-slices"));
2639
2640 let replay_b = SessionReplayTool::new(store_b.clone());
2642 let v = replay_b
2643 .call(&mut ctx, json!({"session_id": sid}))
2644 .await
2645 .unwrap();
2646 assert_eq!(v["count"], 3, "two aged turns + correction: {v}");
2647 let contents: Vec<&str> = v["turns"]
2648 .as_array()
2649 .unwrap()
2650 .iter()
2651 .filter_map(|t| t["content"].as_str())
2652 .collect();
2653 assert!(contents.contains(&"day-one decision"));
2654 assert!(contents.contains(&"corrected claim"));
2655 assert!(!contents.contains(&"original claim"));
2656
2657 let full = replay_b
2659 .call(
2660 &mut ctx,
2661 json!({"session_id": sid, "include_superseded": true}),
2662 )
2663 .await
2664 .unwrap();
2665 assert_eq!(full["count"], 4);
2666
2667 let new_sid = start_session(&store_b);
2669 let continuity = SessionContinuityTool::new(store_b);
2670 let c = continuity
2671 .call(
2672 &mut ctx,
2673 json!({"current_session_id": new_sid, "since":
2674 (Utc::now() - chrono::Duration::days(3)).format("%Y-%m-%d").to_string()}),
2675 )
2676 .await
2677 .unwrap();
2678 assert_eq!(
2679 c["previous_session"], sid,
2680 "import must preserve created_at so recency resolution works"
2681 );
2682 assert_eq!(c["count"], 3);
2683 }
2684
2685 #[tokio::test]
2686 async fn import_rejects_missing_payload() {
2687 let store = test_store();
2688 let tool = SessionImportTool::new(store, None);
2689 let mut ctx = Context::default();
2690 assert!(tool.call(&mut ctx, json!({})).await.is_err());
2691 }
2692
2693 #[tokio::test]
2694 async fn import_refuses_newer_envelope_format() {
2695 let store = test_store();
2696 let mut ctx = Context::default();
2697 let header = wm_memory::envelope::EnvelopeHeader {
2698 format_version: wm_memory::envelope::ENVELOPE_FORMAT_VERSION + 1,
2699 kind: "session_export".into(),
2700 created_at: chrono::Utc::now().to_rfc3339(),
2701 count: 1,
2702 generator: "wm 99.0.0".into(),
2703 };
2704 let record =
2705 serde_json::to_string(&Memory::new(Galaxy::Sessions, "future".into())).unwrap();
2706 let payload = format!("{}\n{record}\n", header.header_line());
2707 let tool = SessionImportTool::new(store, None);
2708 let result = tool.call(&mut ctx, json!({"jsonl": payload})).await;
2709 let err = format!("{:?}", result.unwrap_err());
2710 assert!(err.contains("newer than this build supports"), "{err}");
2711 }
2712
2713 #[tokio::test]
2717 async fn session_record_indexes_at_write_time() {
2718 let dir = tempfile::tempdir().unwrap();
2719 let lmdb = dir.path().join("lmdb");
2720 std::fs::create_dir_all(&lmdb).unwrap();
2721 let store = Arc::new(MemoryStore::open_default(&lmdb).unwrap());
2722 let tantivy = dir.path().join("tantivy");
2723 std::fs::create_dir_all(&tantivy).unwrap();
2724 let search = Arc::new(wm_memory::SearchEngine::open(&tantivy).unwrap());
2725
2726 let sid = start_session(&store);
2727 let mut ctx = Context::default();
2728 SessionRecordTool::new(store.clone())
2729 .with_search(Some(search.clone()))
2730 .call(
2731 &mut ctx,
2732 json!({"role": "ai", "turn_type": "decision", "importance": 0.7,
2733 "content": "amber lighthouse protocol engaged", "session_id": sid}),
2734 )
2735 .await
2736 .unwrap();
2737
2738 let docs = search.count_docs_in_galaxy("sessions").unwrap();
2739 assert!(
2740 docs >= 1,
2741 "session.record must index its write immediately (docs={docs})"
2742 );
2743 }
2744
2745 #[tokio::test]
2749 async fn import_indexes_tantivy_no_drift_even_on_reimport() {
2750 let dir = tempfile::tempdir().unwrap();
2751 let lmdb = dir.path().join("lmdb");
2752 std::fs::create_dir_all(&lmdb).unwrap();
2753 let store = Arc::new(MemoryStore::open_default(&lmdb).unwrap());
2754 let tantivy = dir.path().join("tantivy");
2755 std::fs::create_dir_all(&tantivy).unwrap();
2756 let search = Arc::new(wm_memory::SearchEngine::open(&tantivy).unwrap());
2757
2758 let store_a = test_store();
2760 let sid = start_session(&store_a);
2761 let mut ctx = Context::default();
2762 SessionRecordTool::new(store_a.clone())
2763 .call(
2764 &mut ctx,
2765 json!({"role": "ai", "turn_type": "decision", "importance": 0.9,
2766 "content": "kumquat governance ratchet engaged", "session_id": sid}),
2767 )
2768 .await
2769 .unwrap();
2770 let exported = SessionExportTool::new(store_a.clone())
2771 .call(&mut ctx, json!({"session_id": sid}))
2772 .await
2773 .unwrap();
2774 let jsonl = exported["jsonl"].as_str().unwrap().to_string();
2775
2776 let import = SessionImportTool::new(store.clone(), Some(search.clone()));
2778 for round in 1..=2 {
2779 let r = import
2780 .call(&mut ctx, json!({"jsonl": jsonl}))
2781 .await
2782 .unwrap();
2783 assert_eq!(r["status"], "success", "round {round}: {r}");
2784 assert_eq!(r["skipped"], 0);
2785 assert_eq!(
2786 r["indexed"], r["imported"],
2787 "round {round}: every record indexed: {r}"
2788 );
2789 }
2790
2791 let report = wm_memory::reindex::check_consistency(&store, &search);
2794 let drifted: Vec<_> = report
2795 .galaxies
2796 .iter()
2797 .filter(|g| g.drift)
2798 .map(|g| g.galaxy.clone())
2799 .collect();
2800 assert!(
2801 drifted.is_empty(),
2802 "import must leave zero index drift, drifted: {drifted:?}"
2803 );
2804
2805 let needle_id = store
2807 .scan_all(Galaxy::Sessions)
2808 .unwrap()
2809 .iter()
2810 .find(|m| m.content.contains("kumquat"))
2811 .unwrap()
2812 .metadata
2813 .id
2814 .to_string();
2815 let hits = search.search("kumquat governance ratchet", 10).unwrap();
2816 assert!(
2817 hits.iter().any(|h| h.memory_id == needle_id),
2818 "imported record must be searchable via the index: {hits:?}"
2819 );
2820 }
2821
2822 #[tokio::test]
2823 async fn record_then_replay_full() {
2824 let store = test_store();
2825 let sid = start_session(&store);
2826 let record = SessionRecordTool::new(store.clone());
2827 let mut ctx = Context::default();
2828 for (i, role) in [("user", "hello"), ("ai", "hi there")].iter().enumerate() {
2829 let r = record
2830 .call(
2831 &mut ctx,
2832 json!({"role": role.0, "content": role.1, "session_id": sid}),
2833 )
2834 .await
2835 .unwrap();
2836 assert_eq!(r["sequence"], i as u64 + 1);
2837 }
2838
2839 let replay = SessionReplayTool::new(store.clone());
2840 let v = replay
2841 .call(&mut ctx, json!({"mode": "full", "session_id": sid}))
2842 .await
2843 .unwrap();
2844 assert_eq!(v["count"], 2);
2845 assert_eq!(v["turns"][0]["content"], "hello");
2846 assert_eq!(v["turns"][1]["role"], "ai");
2847 }
2848
2849 #[tokio::test(flavor = "multi_thread", worker_threads = 8)]
2853 async fn concurrent_record_writers_get_unique_contiguous_sequences() {
2854 const N: usize = 100;
2855 let store = test_store();
2856 let sid = start_session(&store);
2857 let tool = Arc::new(SessionRecordTool::new(store));
2858
2859 let mut handles = Vec::with_capacity(N);
2860 for i in 0..N {
2861 let tool = Arc::clone(&tool);
2862 let sid = sid.clone();
2863 handles.push(tokio::spawn(async move {
2864 let mut ctx = Context::default();
2865 let v = tool
2866 .call(
2867 &mut ctx,
2868 json!({
2869 "session_id": sid,
2870 "role": "ai",
2871 "turn_type": "message",
2872 "content": format!("concurrent turn {i}"),
2873 }),
2874 )
2875 .await
2876 .unwrap();
2877 v["sequence"].as_u64().expect("sequence in response")
2878 }));
2879 }
2880
2881 let mut sequences = Vec::with_capacity(N);
2882 for handle in handles {
2883 sequences.push(handle.await.unwrap());
2884 }
2885 sequences.sort_unstable();
2886 assert_eq!(
2887 sequences,
2888 (1..=N as u64).collect::<Vec<_>>(),
2889 "{N} concurrent session.record writers must get 1..={N}"
2890 );
2891 }
2892
2893 #[tokio::test]
2897 async fn supersede_marks_old_turn_and_hides_it_by_default() {
2898 let store = test_store();
2899 let sid = start_session(&store);
2900 let record = SessionRecordTool::new(store.clone());
2901 let mut ctx = Context::default();
2902
2903 let first = record
2904 .call(
2905 &mut ctx,
2906 json!({"session_id": sid, "content": "we chose A"}),
2907 )
2908 .await
2909 .unwrap();
2910 let first_id = first["memory_id"].as_str().unwrap().to_string();
2911 let second = record
2912 .call(
2913 &mut ctx,
2914 json!({"session_id": sid, "content": "we chose B", "supersedes": first_id}),
2915 )
2916 .await
2917 .unwrap();
2918 assert_eq!(second["sequence"], 2);
2919
2920 let old = store
2922 .get(
2923 wm_core::Galaxy::Sessions,
2924 uuid::Uuid::parse_str(&first_id).unwrap(),
2925 )
2926 .unwrap()
2927 .unwrap();
2928 assert!(
2929 old.metadata
2930 .tags
2931 .iter()
2932 .any(|t| t.starts_with("superseded-by:")),
2933 "old turn must carry superseded-by: {:?}",
2934 old.metadata.tags
2935 );
2936
2937 let replay = SessionReplayTool::new(store.clone());
2938 let visible = replay
2939 .call(&mut ctx, json!({"mode": "full", "session_id": sid}))
2940 .await
2941 .unwrap();
2942 assert_eq!(visible["count"], 1, "default view shows the current story");
2943 assert_eq!(visible["turns"][0]["content"], "we chose B");
2944
2945 let full = replay
2946 .call(
2947 &mut ctx,
2948 json!({"mode": "full", "session_id": sid, "include_superseded": true}),
2949 )
2950 .await
2951 .unwrap();
2952 assert_eq!(full["count"], 2, "history stays queryable");
2953 }
2954
2955 #[tokio::test]
2956 async fn record_requires_content_and_valid_role() {
2957 let store = test_store();
2958 let sid = start_session(&store);
2959 let tool = SessionRecordTool::new(store);
2960 let mut ctx = Context::default();
2961 assert!(
2962 tool.call(&mut ctx, json!({"role": "system", "content": "x"}))
2963 .await
2964 .is_err()
2965 );
2966 for bad in [json!(""), json!(" "), json!("\n\t")] {
2969 let err = tool
2970 .call(&mut ctx, json!({"content": bad, "session_id": sid}))
2971 .await
2972 .unwrap_err();
2973 assert!(
2974 err.to_string().contains("content"),
2975 "content={bad:?}: {err}"
2976 );
2977 }
2978 let err = tool
2979 .call(
2980 &mut ctx,
2981 json!({"content": "x", "session_id": sid, "turn_type": "observation"}),
2982 )
2983 .await
2984 .unwrap_err();
2985 assert!(err.to_string().contains("turn_type"), "{err}");
2986 let v = tool
2988 .call(
2989 &mut ctx,
2990 json!({"content": "x", "session_id": sid, "turn_type": "context"}),
2991 )
2992 .await
2993 .unwrap();
2994 assert_eq!(v["status"], "success", "{v}");
2995 }
2996
2997 #[tokio::test]
3001 async fn record_rejects_out_of_range_importance() {
3002 let store = test_store();
3003 let sid = start_session(&store);
3004 let record = SessionRecordTool::new(store);
3005 let mut ctx = Context::default();
3006 for bad in [json!(1.5), json!(999), json!(-0.25), json!("2.0")] {
3007 let err = record
3008 .call(
3009 &mut ctx,
3010 json!({"content": "x", "session_id": sid, "importance": bad}),
3011 )
3012 .await
3013 .unwrap_err();
3014 assert!(
3015 err.to_string().contains("importance"),
3016 "importance={bad} must be rejected, got: {err}"
3017 );
3018 }
3019 for good in [json!(0.0), json!(1.0), json!(0.75), json!("0.4")] {
3022 let v = record
3023 .call(
3024 &mut ctx,
3025 json!({"content": "x", "session_id": sid, "importance": good}),
3026 )
3027 .await
3028 .unwrap();
3029 assert_eq!(v["status"], "success", "{v}");
3030 }
3031 }
3032
3033 #[tokio::test]
3037 async fn record_stamps_provenance_from_role() {
3038 let store = test_store();
3039 let sid = start_session(&store);
3040 let record = SessionRecordTool::new(store.clone());
3041 let mut ctx = Context::default();
3042 let ai = record
3043 .call(
3044 &mut ctx,
3045 json!({"role": "ai", "content": "agent turn", "session_id": sid}),
3046 )
3047 .await
3048 .unwrap();
3049 let user = record
3050 .call(
3051 &mut ctx,
3052 json!({"role": "user", "content": "human turn", "session_id": sid}),
3053 )
3054 .await
3055 .unwrap();
3056
3057 let ai_mem = store
3058 .get(
3059 Galaxy::Sessions,
3060 uuid::Uuid::parse_str(ai["memory_id"].as_str().unwrap()).unwrap(),
3061 )
3062 .expect("ai turn stored")
3063 .expect("ai turn present");
3064 assert_eq!(ai_mem.metadata.source, "agent");
3065 assert!((ai_mem.metadata.source_trust - 0.7).abs() < 1e-5);
3066
3067 let user_mem = store
3068 .get(
3069 Galaxy::Sessions,
3070 uuid::Uuid::parse_str(user["memory_id"].as_str().unwrap()).unwrap(),
3071 )
3072 .expect("user turn stored")
3073 .expect("user turn present");
3074 assert_eq!(user_mem.metadata.source, "user");
3075 assert!((user_mem.metadata.source_trust - 1.0).abs() < f32::EPSILON);
3076 }
3077
3078 #[tokio::test]
3079 async fn record_rejects_invalid_track_slug() {
3080 let store = test_store();
3081 let sid = start_session(&store);
3082 let record = SessionRecordTool::new(store);
3083 let mut ctx = Context::default();
3084 let mut bad = vec![
3085 String::new(),
3086 "Track-A".to_string(),
3087 "has space".to_string(),
3088 "-leading".to_string(),
3089 ];
3090 bad.push("x".repeat(TRACK_MAX_LEN + 1));
3091 for bad in &bad {
3092 let err = record
3093 .call(
3094 &mut ctx,
3095 json!({"content": "x", "session_id": sid, "track": bad}),
3096 )
3097 .await
3098 .unwrap_err()
3099 .to_string();
3100 assert!(err.contains("invalid track"), "slug {bad:?}: {err}");
3101 }
3102 let ok = record
3104 .call(
3105 &mut ctx,
3106 json!({"content": "x", "session_id": sid, "track": "wmv9/harness-2.1_a"}),
3107 )
3108 .await
3109 .unwrap();
3110 assert_eq!(ok["status"], "success");
3111 }
3112
3113 #[tokio::test]
3114 async fn track_log_reads_recorded_turns_per_track() {
3115 let store = test_store();
3116 let sid = start_session(&store);
3117 let record = SessionRecordTool::new(store.clone());
3118 let mut ctx = Context::default();
3119 for (track, content) in [
3120 ("harness-2", "harness: schema check landed"),
3121 ("safety-a", "safety: mask gate landed"),
3122 ("harness-2", "harness: manifest regen done"),
3123 ] {
3124 record
3125 .call(
3126 &mut ctx,
3127 json!({
3128 "role": "ai",
3129 "content": content,
3130 "session_id": sid,
3131 "track": track,
3132 "turn_type": "summary",
3133 }),
3134 )
3135 .await
3136 .unwrap();
3137 }
3138
3139 let log = SessionTrackLogTool::new(store);
3140 let v = log
3141 .call(&mut ctx, json!({"track": "harness-2"}))
3142 .await
3143 .unwrap();
3144 assert_eq!(v["mode"], "log");
3145 assert_eq!(v["tracks"][0]["track"], "harness-2");
3146 assert_eq!(v["tracks"][0]["entry_count"], 2);
3147 assert_eq!(v["tracks"][0]["returned"], 2);
3148 assert_eq!(
3149 v["tracks"][0]["entries"][0]["content"],
3150 "harness: schema check landed"
3151 );
3152 assert_eq!(
3153 v["tracks"][0]["entries"][1]["content"],
3154 "harness: manifest regen done"
3155 );
3156 assert!(v["tracks"][0]["age_seconds"].as_i64().unwrap() >= 0);
3157
3158 let all = log.call(&mut ctx, json!({})).await.unwrap();
3160 assert_eq!(all["mode"], "overview");
3161 assert_eq!(all["track_count"], 2);
3162 let names: Vec<&str> = all["tracks"]
3163 .as_array()
3164 .unwrap()
3165 .iter()
3166 .map(|t| t["track"].as_str().unwrap())
3167 .collect();
3168 assert_eq!(names, vec!["harness-2", "safety-a"]);
3169 assert!(all["tracks"][0]["entries"].is_null());
3170 assert_eq!(all["tracks"][0]["latest_entry"]["turn_type"], "summary");
3171 assert!(
3172 all["tracks"][0]["latest_entry"]["preview"]
3173 .as_str()
3174 .unwrap()
3175 .contains("manifest regen")
3176 );
3177
3178 let unknown = log
3180 .call(&mut ctx, json!({"track": "never-seen"}))
3181 .await
3182 .unwrap();
3183 assert_eq!(unknown["tracks"][0]["known"], false);
3184 assert_eq!(unknown["tracks"][0]["entry_count"], 0);
3185 }
3186
3187 #[tokio::test]
3188 async fn track_log_merges_related_tracks_and_filters_time() {
3189 let store = test_store();
3190 let sid = start_session(&store);
3191 let record = SessionRecordTool::new(store.clone());
3192 let mut ctx = Context::default();
3193 for (track, content) in [
3194 ("harness-2", "old harness note"),
3195 ("safety-a", "safety note"),
3196 ("harness-2", "new harness note"),
3197 ] {
3198 record
3199 .call(
3200 &mut ctx,
3201 json!({"role": "ai", "content": content, "session_id": sid, "track": track}),
3202 )
3203 .await
3204 .unwrap();
3205 }
3206 for mut mem in store.scan_all(Galaxy::Sessions).unwrap() {
3208 if mem.content.contains("old harness note") {
3209 mem.metadata.created_at = Utc::now() - chrono::Duration::days(3);
3210 store.put(Galaxy::Sessions, &mem).unwrap();
3211 }
3212 }
3213
3214 let log = SessionTrackLogTool::new(store);
3215 let cutoff = (Utc::now() - chrono::Duration::days(1))
3216 .format("%Y-%m-%d")
3217 .to_string();
3218 let v = log
3219 .call(
3220 &mut ctx,
3221 json!({"tracks": ["harness-2", "safety-a"], "since": cutoff}),
3222 )
3223 .await
3224 .unwrap();
3225 assert_eq!(v["track_count"], 2);
3226 let harness = &v["tracks"][0];
3227 assert_eq!(harness["track"], "harness-2");
3228 assert_eq!(harness["entry_count"], 1, "aged turn filtered: {v}");
3229 assert_eq!(harness["entries"][0]["content"], "new harness note");
3230 let safety = &v["tracks"][1];
3231 assert_eq!(safety["track"], "safety-a");
3232 assert_eq!(safety["entry_count"], 1);
3233
3234 let err = log
3236 .call(&mut ctx, json!({"tracks": ["harness-2", "Bad Slug"]}))
3237 .await
3238 .unwrap_err()
3239 .to_string();
3240 assert!(err.contains("invalid track"), "{err}");
3241 }
3242
3243 #[tokio::test]
3244 async fn track_log_hides_superseded_turns_by_default() {
3245 let store = test_store();
3246 let sid = start_session(&store);
3247 let record = SessionRecordTool::new(store.clone());
3248 let mut ctx = Context::default();
3249 let first = record
3250 .call(
3251 &mut ctx,
3252 json!({"role": "ai", "content": "v1 estimate", "session_id": sid,
3253 "track": "bench", "turn_type": "summary"}),
3254 )
3255 .await
3256 .unwrap();
3257 record
3258 .call(
3259 &mut ctx,
3260 json!({"role": "ai", "content": "v2 estimate", "session_id": sid,
3261 "track": "bench", "turn_type": "summary",
3262 "supersedes": first["memory_id"]}),
3263 )
3264 .await
3265 .unwrap();
3266
3267 let log = SessionTrackLogTool::new(store);
3268 let v = log.call(&mut ctx, json!({"track": "bench"})).await.unwrap();
3269 assert_eq!(v["tracks"][0]["entry_count"], 1);
3270 assert_eq!(v["tracks"][0]["entries"][0]["content"], "v2 estimate");
3271 let all = log
3272 .call(
3273 &mut ctx,
3274 json!({"track": "bench", "include_superseded": true}),
3275 )
3276 .await
3277 .unwrap();
3278 assert_eq!(all["tracks"][0]["entry_count"], 2);
3279 }
3280
3281 #[tokio::test]
3282 async fn checkpoint_track_lands_in_track_log() {
3283 let store = test_store();
3284 let sid = start_session(&store);
3285 let checkpoint = SessionCheckpointNodiscoveryTool::new(store.clone());
3286 let mut ctx = Context::default();
3287 let v = checkpoint
3288 .call(
3289 &mut ctx,
3290 json!({
3291 "session_id": sid,
3292 "track": "harness-2",
3293 "label": "slice2-done",
3294 "next_queue": ["regen manifest", "push after ceremony"],
3295 "open_flags": ["none"],
3296 }),
3297 )
3298 .await
3299 .unwrap();
3300 assert_eq!(v["status"], "success");
3301
3302 let log = SessionTrackLogTool::new(store);
3303 let out = log
3304 .call(&mut ctx, json!({"track": "harness-2"}))
3305 .await
3306 .unwrap();
3307 assert_eq!(out["tracks"][0]["entry_count"], 1);
3308 assert_eq!(out["tracks"][0]["entries"][0]["kind"], "checkpoint");
3309 assert_eq!(
3310 out["tracks"][0]["latest_checkpoint"]["entry"]["label"],
3311 "slice2-done"
3312 );
3313 assert_eq!(
3314 out["tracks"][0]["latest_checkpoint"]["entry"]["handoff"]["next_queue"][0],
3315 "regen manifest"
3316 );
3317
3318 let err = checkpoint
3320 .call(&mut ctx, json!({"session_id": sid, "track": "Bad"}))
3321 .await
3322 .unwrap_err()
3323 .to_string();
3324 assert!(err.contains("invalid track"), "{err}");
3325 }
3326
3327 #[tokio::test]
3328 async fn continuity_returns_previous_session_tail() {
3329 let store = test_store();
3330 let sid1 = start_session(&store);
3331 let record = SessionRecordTool::new(store.clone());
3332 let mut ctx = Context::default();
3333 for i in 0..5 {
3334 record
3335 .call(
3336 &mut ctx,
3337 json!({"role": "user", "content": format!("turn {i}"), "session_id": sid1}),
3338 )
3339 .await
3340 .unwrap();
3341 }
3342 let sid2 = start_session(&store);
3343 let continuity = SessionContinuityTool::new(store);
3344 let v = continuity
3345 .call(&mut ctx, json!({"current_session_id": sid2, "n": 2}))
3346 .await
3347 .unwrap();
3348 assert_eq!(v["previous_session"], sid1);
3349 assert_eq!(v["count"], 2);
3350 assert_eq!(v["turns"][1]["content"], "turn 4");
3351 assert!(
3352 v.get("hint").is_none(),
3353 "non-empty continuity must not carry the scoping hint"
3354 );
3355 assert!(
3356 v["checkpoint"].is_null() && v["checkpoint_id"].is_null(),
3357 "no checkpoint recorded — the fields must be present but null: {v}"
3358 );
3359 }
3360
3361 #[tokio::test]
3362 async fn continuity_surfaces_latest_checkpoint_handoff() {
3363 let store = test_store();
3369 let sid1 = start_session(&store);
3370 let record = SessionRecordTool::new(store.clone());
3371 let mut ctx = Context::default();
3372 record
3373 .call(
3374 &mut ctx,
3375 json!({"role": "ai", "turn_type": "summary", "importance": 0.9,
3376 "content": "seam work done", "session_id": sid1}),
3377 )
3378 .await
3379 .unwrap();
3380
3381 let seed_checkpoint = |handoff: Value, age_secs: i64| {
3382 let mut cp = Memory::new(
3383 Galaxy::Sessions,
3384 json!({
3385 "type": "checkpoint",
3386 "session_id": sid1,
3387 "label": "wrap",
3388 "data": {},
3389 "handoff": handoff,
3390 })
3391 .to_string(),
3392 );
3393 cp.metadata.tags = vec!["session".into(), "checkpoint".into()];
3394 cp.metadata.created_at = Utc::now() - chrono::Duration::seconds(age_secs);
3395 store.put(Galaxy::Sessions, &cp).unwrap();
3396 cp.metadata.id.to_string()
3397 };
3398 seed_checkpoint(
3399 json!({"next_queue": ["stale task"], "open_flags": ["stale flag"]}),
3400 120,
3401 );
3402 let newest_id = seed_checkpoint(
3403 json!({
3404 "git": {"commit": "abc1234", "branch": "main", "dirty_count": 0},
3405 "tests_green": true,
3406 "next_queue": ["fuzz malformed headers", "verify seal"],
3407 "open_flags": ["header length cap unresolved"],
3408 }),
3409 0,
3410 );
3411
3412 let sid2 = start_session(&store);
3413 let continuity = SessionContinuityTool::new(store);
3414 let v = continuity
3415 .call(&mut ctx, json!({"current_session_id": sid2}))
3416 .await
3417 .unwrap();
3418 assert_eq!(v["previous_session"], sid1);
3419 assert_eq!(
3420 v["checkpoint"]["next_queue"][0], "fuzz malformed headers",
3421 "the latest checkpoint must win over an older one: {v}"
3422 );
3423 assert_eq!(
3424 v["checkpoint"]["open_flags"][0], "header length cap unresolved",
3425 "open flags must survive the handoff through continuity: {v}"
3426 );
3427 assert_eq!(v["checkpoint"]["git"]["commit"], "abc1234");
3428 assert_eq!(v["checkpoint"]["tests_green"], true);
3429 assert_eq!(v["checkpoint_id"].as_str().unwrap(), newest_id);
3430 }
3431
3432 #[tokio::test]
3433 async fn continuity_empty_store_discloses_project_scoping() {
3434 let store = test_store();
3439 let continuity = SessionContinuityTool::new(store);
3440 let mut ctx = Context::default();
3441 let v = continuity.call(&mut ctx, json!({})).await.unwrap();
3442
3443 assert_eq!(v["status"], "success");
3444 assert_eq!(v["count"], 0);
3445 let hint = v["hint"].as_str().expect("hint present on empty store");
3446 assert!(hint.contains("project-scoped"), "got: {hint}");
3447 assert!(hint.contains("opencode config"), "got: {hint}");
3448 assert!(hint.contains("GET /status"), "got: {hint}");
3449 assert!(hint.contains("store "), "got: {hint}");
3451 }
3452
3453 #[tokio::test]
3454 async fn record_defaults_to_newest_start_by_time_not_key_order() {
3455 let store = test_store();
3460 for age in [50_000, 40_000, 30_000, 20_000, 10_000] {
3461 start_session_aged(&store, age);
3462 }
3463 let newest = start_session_aged(&store, 0);
3464
3465 let record = SessionRecordTool::new(store);
3466 let mut ctx = Context::default();
3467 let r = record
3468 .call(&mut ctx, json!({"role": "ai", "content": "latest turn"}))
3469 .await
3470 .unwrap();
3471 assert_eq!(
3472 r["session_id"], newest,
3473 "record without explicit session_id must target the newest start by created_at"
3474 );
3475 }
3476
3477 #[tokio::test]
3478 async fn continuity_picks_newest_prior_by_time_not_key_order() {
3479 let store = test_store();
3483 for age in [40_000, 30_000, 20_000] {
3484 start_session_aged(&store, age);
3485 }
3486 let newest_prior = start_session_aged(&store, 10);
3487 let current = start_session_aged(&store, 0);
3488
3489 let continuity = SessionContinuityTool::new(store);
3490 let mut ctx = Context::default();
3491 let v = continuity
3492 .call(&mut ctx, json!({"current_session_id": current, "n": 1}))
3493 .await
3494 .unwrap();
3495 assert_eq!(
3496 v["previous_session"], newest_prior,
3497 "continuity must select the newest prior session by created_at"
3498 );
3499 }
3500
3501 #[tokio::test]
3502 async fn handoff_transfer_accept_list() {
3503 let store = test_store();
3504 let sid = start_session(&store);
3505 let record = SessionRecordTool::new(store.clone());
3506 let mut ctx = Context::default();
3507 record
3508 .call(
3509 &mut ctx,
3510 json!({"role": "ai", "content": "context", "session_id": sid}),
3511 )
3512 .await
3513 .unwrap();
3514
3515 let handoff = SessionHandoffTool::new(store.clone());
3516 let t = handoff
3517 .call(
3518 &mut ctx,
3519 json!({"action": "transfer", "session_id": sid, "message": "take over"}),
3520 )
3521 .await
3522 .unwrap();
3523 assert_eq!(t["status"], "success");
3524 let hid = t["handoff_id"].as_str().unwrap().to_string();
3525
3526 let list = handoff
3527 .call(&mut ctx, json!({"action": "list"}))
3528 .await
3529 .unwrap();
3530 assert_eq!(list["pending_count"], 1);
3531
3532 let a = handoff
3533 .call(&mut ctx, json!({"action": "accept", "handoff_id": hid}))
3534 .await
3535 .unwrap();
3536 assert_eq!(a["status"], "success");
3537
3538 let list2 = handoff
3539 .call(&mut ctx, json!({"action": "list"}))
3540 .await
3541 .unwrap();
3542 assert_eq!(list2["pending_count"], 0);
3543 }
3544
3545 #[tokio::test]
3546 async fn lossless_replay_binds_explicit_session_chunks_and_detects_stale_view() {
3547 let store = test_store();
3548 let older = start_session_aged(&store, 60);
3549 let newer = start_session_aged(&store, 0);
3550 let content = format!("prefix {} DISTINCT-FACT-AFTER-120", "é".repeat(9000));
3551 let mut turn = Memory::new(Galaxy::Sessions, json!({"type":"session_turn","session_id":older,"sequence":1,"timestamp":1_i64,"content":content}).to_string());
3552 turn.metadata.tags = vec!["session".into(), "turn".into(), format!("session:{older}")];
3553 store.put(Galaxy::Sessions, &turn).unwrap();
3554 let mut other = Memory::new(Galaxy::Sessions, json!({"type":"session_turn","session_id":newer,"sequence":1,"timestamp":1_i64,"content":"newer-only"}).to_string());
3555 other.metadata.tags = vec!["session".into(), "turn".into(), format!("session:{newer}")];
3556 store.put(Galaxy::Sessions, &other).unwrap();
3557 let replay = SessionReplayTool::new(store.clone());
3558 let mut ctx = Context::default();
3559 let mut response = replay
3560 .call(
3561 &mut ctx,
3562 json!({"mode":"lossless","session_id":older,"page_size":1,"max_wire_bytes":2048}),
3563 )
3564 .await
3565 .unwrap();
3566 assert_eq!(response["session_id"], older);
3567 assert!(response["records"][0].get("chunk").is_some());
3568 let mut bytes = Vec::new();
3569 loop {
3570 assert!(serde_json::to_vec(&response).unwrap().len() <= 2048);
3571 for record in response["records"].as_array().unwrap() {
3572 if let Some(chunk) = record.get("chunk") {
3573 assert_eq!(chunk["byte_offset"].as_u64().unwrap() as usize, bytes.len());
3574 bytes.extend(base64_decode(chunk["data_b64"].as_str().unwrap()).unwrap());
3575 } else {
3576 bytes.extend(record["content"].as_str().unwrap().as_bytes());
3577 }
3578 }
3579 if response["complete"] == true {
3580 assert!(response["next_cursor"].is_null());
3581 break;
3582 }
3583 let cursor = response["next_cursor"].as_str().unwrap();
3584 response = replay.call(&mut ctx, json!({"mode":"lossless","session_id":older,"page_size":1,"max_wire_bytes":2048,"cursor":cursor})).await.unwrap();
3585 }
3586 assert_eq!(String::from_utf8(bytes).unwrap(), content);
3587 let first = replay
3588 .call(
3589 &mut ctx,
3590 json!({"mode":"lossless","session_id":older,"page_size":1,"max_wire_bytes":2048}),
3591 )
3592 .await
3593 .unwrap();
3594 let stale_cursor = first["next_cursor"].as_str().unwrap().to_string();
3595 let mut appended = Memory::new(Galaxy::Sessions, json!({"type":"session_turn","session_id":older,"sequence":2,"timestamp":2_i64,"content":"later"}).to_string());
3596 appended.metadata.tags = vec!["session".into(), "turn".into(), format!("session:{older}")];
3597 store.put(Galaxy::Sessions, &appended).unwrap();
3598 assert!(replay.call(&mut ctx, json!({"mode":"lossless","session_id":older,"page_size":1,"max_wire_bytes":2048,"cursor":stale_cursor})).await.unwrap_err().to_string().contains("stale_view"));
3599 }
3600
3601 fn lossless_fixture_turn(sid: &str, sequence: u64, text: &str) -> Memory {
3602 let mut m = Memory::new(Galaxy::Sessions, json!({"type":"session_turn","session_id":sid,"sequence":sequence,"timestamp":1_i64,"role":"ai","turn_type":"decision","importance":0.8,"content":text}).to_string());
3603 m.metadata.tags = vec!["turn".into(), format!("session:{sid}")];
3604 m
3605 }
3606
3607 #[tokio::test]
3608 async fn lossless_strict_args_and_cursor_prevalidation_before_corrupt_scan() {
3609 let store = test_store();
3610 let sid = start_session(&store);
3611 store
3612 .put_raw(
3613 Galaxy::Sessions,
3614 uuid::Uuid::new_v4().as_bytes(),
3615 b"not-a-memory",
3616 )
3617 .unwrap();
3618 let replay = SessionReplayTool::new(store);
3619 for (key, value) in [
3620 ("include_superseded", json!("false")),
3621 ("page_size", json!(null)),
3622 ("page_size", json!(-1)),
3623 ("max_wire_bytes", json!("2048")),
3624 ("cursor", json!(false)),
3625 ("cursor", json!("a€")),
3626 ("cursor", json!("🔑")),
3627 ] {
3628 let mut args = json!({"mode":"lossless","session_id":sid});
3629 args[key] = value;
3630 let error = replay
3631 .call(&mut Context::default(), args)
3632 .await
3633 .unwrap_err()
3634 .to_string();
3635 assert!(
3636 error.contains("invalid_args") || error.contains("invalid_cursor"),
3637 "{key}: {error}"
3638 );
3639 assert!(!error.contains("incomplete scan"));
3640 }
3641 let error = replay
3642 .call(
3643 &mut Context::default(),
3644 json!({"mode":"lossless","session_id":sid}),
3645 )
3646 .await
3647 .unwrap_err()
3648 .to_string();
3649 assert!(error.contains("refusing incomplete scan"));
3650 }
3651
3652 #[tokio::test]
3653 async fn lossless_metadata_staleness_settings_seek_and_partial_resume() {
3654 let store = test_store();
3655 let sid = start_session(&store);
3656 let mut turn = lossless_fixture_turn(&sid, 1, &"€\n".repeat(2000));
3657 store.put(Galaxy::Sessions, &turn).unwrap();
3658 let replay = SessionReplayTool::new(store.clone());
3659 let args = json!({"mode":"lossless","session_id":sid,"page_size":1,"max_wire_bytes":2048});
3660 let first = replay
3661 .call(&mut Context::default(), args.clone())
3662 .await
3663 .unwrap();
3664 let cursor = first["next_cursor"].as_str().unwrap();
3665 for (key, value) in [("page_size", json!(2)), ("max_wire_bytes", json!(4096))] {
3666 let mut next = args.clone();
3667 next["cursor"] = json!(cursor);
3668 next[key] = value;
3669 assert!(
3670 replay
3671 .call(&mut Context::default(), next)
3672 .await
3673 .unwrap_err()
3674 .to_string()
3675 .contains("invalid_cursor")
3676 );
3677 }
3678 let (view, _, _) = parse_lossless_cursor(cursor, &sid, false, 1, 2048).unwrap();
3679 let mut seek = args.clone();
3681 seek["cursor"] = json!(lossless_cursor(
3682 &sid,
3683 false,
3684 1,
3685 2048,
3686 &view,
3687 0,
3688 turn.content.len() - 2
3689 ));
3690 let inner = serde_json::from_str::<Value>(&turn.content).unwrap()["content"]
3692 .as_str()
3693 .unwrap()
3694 .to_string();
3695 seek["cursor"] = json!(lossless_cursor(
3696 &sid,
3697 false,
3698 1,
3699 2048,
3700 &view,
3701 0,
3702 inner.len() - 2
3703 ));
3704 let tail = replay.call(&mut Context::default(), seek).await.unwrap();
3705 assert!(tail["records"][0].get("content").is_none());
3706 assert_eq!(
3707 base64_decode(tail["records"][0]["chunk"]["data_b64"].as_str().unwrap()).unwrap(),
3708 inner.as_bytes()[inner.len() - 2..]
3709 );
3710 assert!(tail["next_cursor"].is_null());
3711 let mut changed = serde_json::from_str::<Value>(&turn.content).unwrap();
3712 changed["role"] = json!("human");
3713 turn.content = changed.to_string();
3714 store.put(Galaxy::Sessions, &turn).unwrap();
3715 let mut next = args;
3716 next["cursor"] = json!(cursor);
3717 assert!(
3718 replay
3719 .call(&mut Context::default(), next)
3720 .await
3721 .unwrap_err()
3722 .to_string()
3723 .contains("stale_view")
3724 );
3725 }
3726
3727 #[tokio::test]
3728 async fn lossless_visibility_empty_nonturn_ties_and_supersession() {
3729 let store = test_store();
3730 let sid = start_session(&store);
3731 let replay = SessionReplayTool::new(store.clone());
3732 let mut args = json!({"mode":"lossless","session_id":sid,"page_size":64});
3733 let empty = replay
3734 .call(&mut Context::default(), args.clone())
3735 .await
3736 .unwrap();
3737 assert_eq!(empty["records"], json!([]));
3738 assert!(empty["next_cursor"].is_null());
3739 assert_eq!(empty["complete"], true);
3740 let mut a = lossless_fixture_turn(&sid, 1, "");
3741 let mut b = lossless_fixture_turn(&sid, 1, "second");
3742 b.metadata.tags.clear();
3744 let mut hidden = lossless_fixture_turn(&sid, 2, "PRIVATE-SENTINEL");
3745 hidden.metadata.is_private = true;
3746 let mut excluded = lossless_fixture_turn(&sid, 3, "EXCLUDED-SENTINEL");
3747 excluded.metadata.model_exclude = true;
3748 a.metadata
3749 .tags
3750 .push(format!("supersedes:{}", hidden.metadata.id));
3751 let mut handoff = Memory::new(
3752 Galaxy::Sessions,
3753 json!({"type":"session_handoff"}).to_string(),
3754 );
3755 handoff.metadata.tags = vec![format!("session:{sid}")];
3756 store
3757 .put_batch(
3758 Galaxy::Sessions,
3759 &[a.clone(), b.clone(), hidden.clone(), excluded, handoff],
3760 )
3761 .unwrap();
3762 let result = replay
3763 .call(&mut Context::default(), args.clone())
3764 .await
3765 .unwrap();
3766 let wire = result.to_string();
3767 assert!(!wire.contains("SENTINEL"));
3768 assert!(!wire.contains(&hidden.metadata.id.to_string()));
3769 let mut ids = vec![a.metadata.id.to_string(), b.metadata.id.to_string()];
3770 ids.sort();
3771 assert_eq!(
3772 result["records"]
3773 .as_array()
3774 .unwrap()
3775 .iter()
3776 .map(|r| r["record_id"].as_str().unwrap().to_string())
3777 .collect::<Vec<_>>(),
3778 ids
3779 );
3780 a.metadata
3781 .tags
3782 .push(format!("superseded-by:{}", b.metadata.id));
3783 store.put(Galaxy::Sessions, &a).unwrap();
3784 assert_eq!(
3785 replay
3786 .call(&mut Context::default(), args.clone())
3787 .await
3788 .unwrap()["records"]
3789 .as_array()
3790 .unwrap()
3791 .len(),
3792 1
3793 );
3794 args["include_superseded"] = json!(true);
3795 assert_eq!(
3796 replay.call(&mut Context::default(), args).await.unwrap()["records"]
3797 .as_array()
3798 .unwrap()
3799 .len(),
3800 2
3801 );
3802 let mut start = store
3803 .get(Galaxy::Sessions, uuid::Uuid::parse_str(&sid).unwrap())
3804 .unwrap()
3805 .unwrap();
3806 start.metadata.model_exclude = true;
3807 store.put(Galaxy::Sessions, &start).unwrap();
3808 assert!(
3809 replay
3810 .call(
3811 &mut Context::default(),
3812 json!({"mode":"lossless","session_id":sid})
3813 )
3814 .await
3815 .unwrap_err()
3816 .to_string()
3817 .contains("not found")
3818 );
3819 }
3820
3821 #[tokio::test]
3822 async fn lossless_malformed_selected_turn_fails_and_schema_advertises_mode() {
3823 let store = test_store();
3824 let sid = start_session(&store);
3825 let mut bad = lossless_fixture_turn(&sid, 1, "good");
3826 bad.content = "{broken".into();
3827 store.put(Galaxy::Sessions, &bad).unwrap();
3828 let replay = SessionReplayTool::new(store);
3829 assert!(
3830 replay
3831 .call(
3832 &mut Context::default(),
3833 json!({"mode":"lossless","session_id":sid})
3834 )
3835 .await
3836 .unwrap_err()
3837 .to_string()
3838 .contains("malformed_selected_turn")
3839 );
3840 assert!(replay.input_schema().to_string().contains("lossless"));
3841 assert!(replay.description().contains("explicit session UUID"));
3842 }
3843
3844 #[test]
3845 fn lossless_arbitrary_cursor_text_is_panic_free() {
3846 use proptest::prelude::*;
3847 proptest!(|(text in any::<String>())| { let _ = hex_decode(&text); });
3848 assert!(hex_decode("AB").is_none());
3849 }
3850
3851 #[tokio::test]
3852 async fn lossless_selection_reaches_record_after_ten_thousand_sources() {
3853 let store = test_store();
3854 let sid = start_session(&store);
3855 let mut sources = Vec::new();
3856 for i in 1..=10_001_u128 {
3857 let mut m = Memory::new(Galaxy::Sessions, "unrelated".into());
3858 m.metadata.id = uuid::Uuid::from_u128(i);
3859 sources.push(m);
3860 }
3861 let mut target = lossless_fixture_turn(&sid, 1, "AFTER-TEN-THOUSAND");
3862 target.metadata.id = uuid::Uuid::from_u128(u128::MAX);
3863 sources.push(target);
3864 store.put_batch(Galaxy::Sessions, &sources).unwrap();
3865 let replay = SessionReplayTool::new(store);
3866 let result = replay
3867 .call(
3868 &mut Context::default(),
3869 json!({"mode":"lossless","session_id":sid}),
3870 )
3871 .await
3872 .unwrap();
3873 assert_eq!(result["records"][0]["content"], "AFTER-TEN-THOUSAND");
3874 }
3875
3876 #[tokio::test]
3877 async fn lossless_tag_payload_disagreement_refuses_and_missing_start_is_not_found() {
3878 let store = test_store();
3879 let sid = start_session(&store);
3880 let mut m = lossless_fixture_turn(&uuid::Uuid::new_v4().to_string(), 1, "contradiction");
3881 m.metadata.tags = vec![format!("session:{sid}")];
3882 store.put(Galaxy::Sessions, &m).unwrap();
3883 let replay = SessionReplayTool::new(store);
3884 assert!(
3885 replay
3886 .call(
3887 &mut Context::default(),
3888 json!({"mode":"lossless","session_id":sid})
3889 )
3890 .await
3891 .unwrap_err()
3892 .to_string()
3893 .contains("malformed_selected_turn")
3894 );
3895 assert!(
3896 replay
3897 .call(
3898 &mut Context::default(),
3899 json!({"mode":"lossless","session_id":uuid::Uuid::new_v4().to_string()})
3900 )
3901 .await
3902 .unwrap_err()
3903 .to_string()
3904 .contains("not found")
3905 );
3906 }
3907}