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
74pub(crate) fn 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
83pub(crate) fn 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 CONTINUITY_DEFAULT_MAX_CONTENT_BYTES: usize = 8 * 1024;
174const CONTINUITY_DEFAULT_MAX_RESPONSE_BYTES: usize = 48 * 1024;
175const CONTINUITY_MIN_BUDGET_BYTES: usize = 256;
176const CONTINUITY_MAX_BUDGET_BYTES: usize = 256 * 1024;
177const CONTINUITY_WRAPPER_ALLOWANCE: usize = 256;
181
182fn floor_char_boundary_index(s: &str, max: usize) -> usize {
184 let mut end = max.min(s.len());
185 while end > 0 && !s.is_char_boundary(end) {
186 end -= 1;
187 }
188 end
189}
190
191fn format_continuity_turn(mem: &Memory, v: &Value, max_content_bytes: usize) -> Value {
193 let role = v.get("role").and_then(Value::as_str).unwrap_or("?");
194 let content = v.get("content").and_then(Value::as_str).unwrap_or("");
195 let end = floor_char_boundary_index(content, max_content_bytes);
196 json!({
197 "memory_id": mem.metadata.id.to_string(),
198 "session_id": v.get("session_id"),
199 "sequence": v.get("sequence"),
200 "role": role,
201 "turn_type": v.get("turn_type"),
202 "importance": v.get("importance"),
203 "content": &content[..end],
204 "content_bytes": content.len(),
205 "content_truncated": end < content.len(),
206 })
207}
208
209const LOSSLESS_MAX_PAGE_SIZE: usize = 64;
210const LOSSLESS_DEFAULT_PAGE_SIZE: usize = 16;
211const LOSSLESS_MIN_WIRE_BYTES: usize = 1024;
212const LOSSLESS_MAX_WIRE_BYTES: usize = 49_152;
213
214fn hex_encode(bytes: &[u8]) -> String {
215 use std::fmt::Write as _;
216 bytes
217 .iter()
218 .fold(String::with_capacity(bytes.len() * 2), |mut out, b| {
219 let _ = write!(out, "{b:02x}");
220 out
221 })
222}
223
224fn hex_decode(value: &str) -> Option<Vec<u8>> {
225 if value.len() > 4096
226 || value.len() % 2 != 0
227 || !value
228 .bytes()
229 .all(|b| b.is_ascii_digit() || (b'a'..=b'f').contains(&b))
230 {
231 return None;
232 }
233 (0..value.len())
234 .step_by(2)
235 .map(|i| u8::from_str_radix(&value[i..i + 2], 16).ok())
236 .collect()
237}
238
239fn base64_encode(bytes: &[u8]) -> String {
240 const TABLE: &[u8; 64] = b"ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz0123456789+/";
241 let mut out = String::with_capacity(bytes.len().div_ceil(3) * 4);
242 for chunk in bytes.chunks(3) {
243 let n = u32::from(chunk[0]) << 16
244 | u32::from(*chunk.get(1).unwrap_or(&0)) << 8
245 | u32::from(*chunk.get(2).unwrap_or(&0));
246 out.push(char::from(TABLE[((n >> 18) & 63) as usize]));
247 out.push(char::from(TABLE[((n >> 12) & 63) as usize]));
248 out.push(if chunk.len() > 1 {
249 char::from(TABLE[((n >> 6) & 63) as usize])
250 } else {
251 '='
252 });
253 out.push(if chunk.len() > 2 {
254 char::from(TABLE[(n & 63) as usize])
255 } else {
256 '='
257 });
258 }
259 out
260}
261
262#[cfg(test)]
263fn base64_decode(value: &str) -> Option<Vec<u8>> {
264 const fn digit(byte: u8) -> Option<u8> {
265 match byte {
266 b'A'..=b'Z' => Some(byte - b'A'),
267 b'a'..=b'z' => Some(byte - b'a' + 26),
268 b'0'..=b'9' => Some(byte - b'0' + 52),
269 b'+' => Some(62),
270 b'/' => Some(63),
271 _ => None,
272 }
273 }
274 if value.len() % 4 != 0 {
275 return None;
276 }
277 let mut out = Vec::new();
278 for chunk in value.as_bytes().chunks_exact(4) {
279 let a = digit(chunk[0])?;
280 let b = digit(chunk[1])?;
281 let c = if chunk[2] == b'=' {
282 0
283 } else {
284 digit(chunk[2])?
285 };
286 let d = if chunk[3] == b'=' {
287 0
288 } else {
289 digit(chunk[3])?
290 };
291 out.push((a << 2) | (b >> 4));
292 if chunk[2] != b'=' {
293 out.push((b << 4) | (c >> 2));
294 }
295 if chunk[3] != b'=' {
296 out.push((c << 6) | d);
297 }
298 }
299 Some(out)
300}
301
302#[derive(Clone)]
303struct LosslessTurn {
304 memory: Memory,
305 turn: Value,
306 content: String,
307 content_hash: String,
308}
309
310fn lossless_error(kind: &str) -> wm_core::CoreError {
311 wm_core::CoreError::InvalidArgs(format!("lossless_{kind}"))
312}
313
314fn lossless_cursor(
315 session_id: &str,
316 include_superseded: bool,
317 page_size: usize,
318 max_wire: usize,
319 view: &str,
320 index: usize,
321 offset: usize,
322) -> String {
323 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});
326 hex_encode(value.to_string().as_bytes())
327}
328
329fn parse_lossless_cursor(
330 cursor: &str,
331 session_id: &str,
332 include_superseded: bool,
333 page_size: usize,
334 max_wire: usize,
335) -> wm_core::Result<(String, usize, usize)> {
336 let bytes = hex_decode(cursor).ok_or_else(|| lossless_error("invalid_cursor"))?;
337 let text = String::from_utf8(bytes).map_err(|_| lossless_error("invalid_cursor"))?;
338 let value: Value = serde_json::from_str(&text).map_err(|_| lossless_error("invalid_cursor"))?;
339 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")});
340 let canonical_text =
341 serde_json::to_string(&canonical).map_err(|_| lossless_error("invalid_cursor"))?;
342 if canonical_text != text
343 || value.get("v").and_then(Value::as_u64) != Some(1)
344 || value.get("session_id").and_then(Value::as_str) != Some(session_id)
345 || value.get("include_superseded").and_then(Value::as_bool) != Some(include_superseded)
346 || value.get("page_size").and_then(Value::as_u64) != Some(page_size as u64)
347 || value.get("max_wire_bytes").and_then(Value::as_u64) != Some(max_wire as u64)
348 {
349 return Err(lossless_error("invalid_cursor"));
350 }
351 let view = value
352 .get("view")
353 .and_then(Value::as_str)
354 .filter(|v| {
355 v.len() == 64
356 && v.bytes()
357 .all(|b| b.is_ascii_digit() || (b'a'..=b'f').contains(&b))
358 })
359 .ok_or_else(|| lossless_error("invalid_cursor"))?;
360 let index = value
361 .get("index")
362 .and_then(Value::as_u64)
363 .and_then(|v| usize::try_from(v).ok())
364 .ok_or_else(|| lossless_error("invalid_cursor"))?;
365 let offset = value
366 .get("offset")
367 .and_then(Value::as_u64)
368 .and_then(|v| usize::try_from(v).ok())
369 .ok_or_else(|| lossless_error("invalid_cursor"))?;
370 Ok((view.to_string(), index, offset))
371}
372
373pub const TURN_TYPES: &[&str] = &[
378 "message",
379 "decision",
380 "breakthrough",
381 "question",
382 "answer",
383 "code_change",
384 "error",
385 "summary",
386 "context",
387];
388
389pub const TRACK_MAX_LEN: usize = 64;
391
392pub(crate) fn validate_track(track: &str) -> wm_core::Result<()> {
399 let valid = !track.is_empty()
400 && track.len() <= TRACK_MAX_LEN
401 && track
402 .chars()
403 .next()
404 .is_some_and(|c| c.is_ascii_lowercase() || c.is_ascii_digit())
405 && track.chars().all(|c| {
406 c.is_ascii_lowercase() || c.is_ascii_digit() || matches!(c, '-' | '_' | '.' | '/')
407 });
408 if valid {
409 Ok(())
410 } else {
411 Err(wm_core::CoreError::InvalidArgs(format!(
412 "invalid track {track:?} — must start with a-z0-9, contain only a-z0-9, '-', '_', '.', '/', \
413 and be at most {TRACK_MAX_LEN} chars"
414 )))
415 }
416}
417
418pub struct SessionRecordTool {
420 store: Arc<MemoryStore>,
421 stats: ToolStats,
422 effects: EffectRow,
423 search: Option<Arc<wm_memory::SearchEngine>>,
424}
425
426impl SessionRecordTool {
427 #[must_use]
428 pub fn new(store: Arc<MemoryStore>) -> Self {
429 Self {
430 store,
431 stats: ToolStats::default(),
432 effects: EffectRow {
433 writes: vec![Resource::Galaxy("sessions".into())],
434 ..Default::default()
435 },
436 search: None,
437 }
438 }
439
440 #[must_use]
444 pub fn with_search(mut self, search: Option<Arc<wm_memory::SearchEngine>>) -> Self {
445 self.search = search;
446 self
447 }
448}
449
450#[async_trait]
451impl Tool for SessionRecordTool {
452 fn name(&self) -> &str {
453 "session.record"
454 }
455 fn gana(&self) -> Gana {
456 Gana::StraddlingLegs
457 }
458 fn effects(&self) -> &EffectRow {
459 &self.effects
460 }
461 fn input_schema(&self) -> Value {
462 super::common::schema(
463 &json!({
464 "content": super::common::str_prop("Turn content"),
465 "role": super::common::str_prop("user | ai (default user)"),
466 "turn_type": json!({
467 "type": "string",
468 "enum": TURN_TYPES,
469 "description": "Turn type (default message)",
470 }),
471 "importance": super::common::bounded_num_prop("0-1 importance (default 0.5)", 0.0, 1.0),
472 "session_id": super::common::str_prop("Target session (default: most recent session)"),
473 "supersedes": super::common::str_prop("Memory id of an earlier turn this record corrects/replaces (amend-with-supersede)"),
474 "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)"),
475 }),
476 &["content"],
477 )
478 }
479 fn description(&self) -> &str {
480 "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)."
481 }
482 async fn call(&self, _ctx: &mut Context, args: Value) -> wm_core::Result<Value> {
483 let role = args.get("role").and_then(Value::as_str).unwrap_or("user");
484 if !matches!(role, "user" | "ai") {
485 return Err(wm_core::CoreError::InvalidArgs(
486 "role must be 'user' or 'ai'".into(),
487 ));
488 }
489 let content = args
492 .get("content")
493 .and_then(Value::as_str)
494 .filter(|s| !s.trim().is_empty())
495 .ok_or_else(|| {
496 wm_core::CoreError::InvalidArgs("content is required and must not be blank".into())
497 })?;
498 let turn_type = args
502 .get("turn_type")
503 .and_then(Value::as_str)
504 .unwrap_or("message");
505 if !TURN_TYPES.contains(&turn_type) {
506 return Err(wm_core::CoreError::InvalidArgs(format!(
507 "turn_type must be one of: {}",
508 TURN_TYPES.join(", ")
509 )));
510 }
511 let importance = wm_dispatch::write_gate::parse_importance_value(args.get("importance"))
516 .map_err(wm_core::CoreError::InvalidArgs)?
517 .unwrap_or(0.5);
518 let session_id = args.get("session_id").and_then(Value::as_str);
519 let track = args.get("track").and_then(Value::as_str);
523 if let Some(track) = track {
524 validate_track(track)?;
525 }
526
527 let session_id: String = if let Some(sid) = session_id {
533 let parsed = uuid::Uuid::parse_str(sid).map_err(|e| {
539 wm_core::CoreError::InvalidArgs(format!("invalid session_id {sid:?}: {e}"))
540 })?;
541 let is_start = self
542 .store
543 .get(Galaxy::Sessions, parsed)?
544 .is_some_and(|m| m.metadata.tags.iter().any(|t| t == "start"));
545 if !is_start {
546 return Err(wm_core::CoreError::NotFound(format!(
547 "no session found with id {sid} — run session.start first \
548 (session.import is the recovery path for orphan/imported turns)"
549 )));
550 }
551 sid.to_string()
552 } else {
553 self.store
554 .scan_all(Galaxy::Sessions)?
555 .iter()
556 .filter(|m| m.metadata.tags.contains(&"start".to_string()))
557 .max_by_key(|m| m.metadata.created_at)
558 .map(|m| m.metadata.id.to_string())
559 .ok_or_else(|| {
560 wm_core::CoreError::Tool("no session found — run session.start first".into())
561 })?
562 };
563
564 let supersedes = match args.get("supersedes").and_then(Value::as_str) {
567 Some(old_id_str) => {
568 let old_id = uuid::Uuid::parse_str(old_id_str).map_err(|e| {
569 wm_core::CoreError::InvalidArgs(format!("invalid 'supersedes' id: {e}"))
570 })?;
571 if self.store.get(Galaxy::Sessions, old_id)?.is_none() {
572 return Err(wm_core::CoreError::NotFound(format!(
573 "superseded turn {old_id} not found"
574 )));
575 }
576 Some(old_id)
577 }
578 None => None,
579 };
580
581 let timestamp = wm_core::time::now_unix_millis();
586 let (sequence, mem) = self.store.put_session_turn(&session_id, |sequence| {
587 let mut turn = json!({
588 "type": "session_turn",
589 "session_id": session_id,
590 "sequence": sequence,
591 "role": role,
592 "turn_type": turn_type,
593 "importance": importance,
594 "content": content,
595 "timestamp": timestamp,
596 });
597 if let Some(track) = track {
598 turn["track"] = json!(track);
599 }
600 let mut mem = Memory::new(Galaxy::Sessions, turn.to_string());
601 mem.metadata.tags = vec![
602 "session".into(),
603 "turn".into(),
604 role.into(),
605 turn_type.into(),
606 format!("session:{session_id}"),
607 ];
608 if let Some(track) = track {
609 mem.metadata.tags.push(format!("track:{track}"));
610 }
611 let (source, trust): (&str, f32) = if role == "user" {
618 ("user", 1.0)
619 } else {
620 ("agent", 0.7)
621 };
622 mem.metadata.source = source.to_string();
623 mem.metadata.source_trust = trust;
624 mem.metadata.importance = importance as f32;
625 if let Some(old_id) = supersedes {
628 mem.metadata.tags.push(format!("supersedes:{old_id}"));
629 }
630 mem
631 })?;
632
633 if let Some(old_id) = supersedes {
640 if let Some(mut old) = self.store.get(Galaxy::Sessions, old_id)? {
641 old.metadata
642 .tags
643 .push(format!("superseded-by:{}", mem.metadata.id));
644 self.store.put(Galaxy::Sessions, &old)?;
645 super::common::index_memory(&self.store, self.search.as_deref(), &old);
646 }
647 }
648
649 super::common::index_memory(&self.store, self.search.as_deref(), &mem);
650 let mut response = json!({
651 "status": "success",
652 "session_id": session_id,
653 "sequence": sequence,
654 "memory_id": mem.metadata.id.to_string(),
655 });
656 if let Some(pattern) = wm_memory::detect_injection(content) {
660 response["warnings"] = json!([format!(
661 "turn content contains an instruction-shaped pattern ({pattern}) — stored as data; \
662 review it before trusting it as context"
663 )]);
664 }
665 super::common::append_savings_row(
668 &self.store,
669 &json!({
670 "ts_ms": timestamp,
671 "op": "record",
672 "session_id": session_id,
673 "sequence": sequence,
674 "role": role,
675 "turn_type": turn_type,
676 "bytes_stored": content.len(),
677 }),
678 );
679 Ok(response)
680 }
681 fn stats(&self) -> &ToolStats {
682 &self.stats
683 }
684}
685
686pub struct SessionReplayTool {
688 store: Arc<MemoryStore>,
689 stats: ToolStats,
690 effects: EffectRow,
691}
692
693impl SessionReplayTool {
694 #[must_use]
695 pub fn new(store: Arc<MemoryStore>) -> Self {
696 Self {
697 store,
698 stats: ToolStats::default(),
699 effects: EffectRow::read_only(vec![Resource::Galaxy("sessions".into())]),
700 }
701 }
702
703 fn lossless(&self, args: &Value) -> wm_core::Result<Value> {
704 const ALLOWED: &[&str] = &[
705 "mode",
706 "session_id",
707 "include_superseded",
708 "page_size",
709 "max_wire_bytes",
710 "cursor",
711 ];
712 let object = args
713 .as_object()
714 .ok_or_else(|| lossless_error("invalid_args"))?;
715 if object.keys().any(|key| !ALLOWED.contains(&key.as_str())) {
716 return Err(lossless_error("unsupported_selection_args"));
717 }
718 let session_id = args
719 .get("session_id")
720 .and_then(Value::as_str)
721 .filter(|s| !s.is_empty())
722 .ok_or_else(|| lossless_error("session_id_required"))?;
723 uuid::Uuid::parse_str(session_id).map_err(|_| lossless_error("invalid_session_id"))?;
724 let include_superseded = match args.get("include_superseded") {
725 None => false,
726 Some(v) => v.as_bool().ok_or_else(|| lossless_error("invalid_args"))?,
727 };
728 let size_arg = |name: &str, default: usize| -> wm_core::Result<usize> {
729 match args.get(name) {
730 None => Ok(default),
731 Some(v) => v
732 .as_u64()
733 .and_then(|v| usize::try_from(v).ok())
734 .ok_or_else(|| lossless_error("invalid_args")),
735 }
736 };
737 let page_size = size_arg("page_size", LOSSLESS_DEFAULT_PAGE_SIZE)?;
738 let max_wire = size_arg("max_wire_bytes", LOSSLESS_MAX_WIRE_BYTES)?;
739 if !(1..=LOSSLESS_MAX_PAGE_SIZE).contains(&page_size)
740 || !(LOSSLESS_MIN_WIRE_BYTES..=LOSSLESS_MAX_WIRE_BYTES).contains(&max_wire)
741 {
742 return Err(lossless_error("invalid_args"));
743 }
744
745 let placement = match args.get("cursor") {
747 None => None,
748 Some(v) => Some(parse_lossless_cursor(
749 v.as_str().ok_or_else(|| lossless_error("invalid_cursor"))?,
750 session_id,
751 include_superseded,
752 page_size,
753 max_wire,
754 )?),
755 };
756 let memories = self.store.scan_all_strict(Galaxy::Sessions)?;
757 let start = memories.iter().find(|m| {
758 m.metadata.id.to_string() == session_id
759 && m.metadata.tags.contains(&"start".to_string())
760 && !m.metadata.is_private
761 && !m.metadata.model_exclude
762 });
763 if start.is_none() {
764 return Err(wm_core::CoreError::NotFound("session not found".into()));
765 }
766 let mut turns = Vec::new();
767 for memory in memories {
768 if memory.metadata.is_private
769 || memory.metadata.model_exclude
770 || (!include_superseded
771 && memory
772 .metadata
773 .tags
774 .iter()
775 .any(|t| t.starts_with("superseded-by:")))
776 {
777 continue;
778 }
779 let tagged = memory
780 .metadata
781 .tags
782 .contains(&format!("session:{session_id}"));
783 let turn = match serde_json::from_str::<Value>(&memory.content) {
784 Ok(v) => v,
785 Err(_) if tagged => return Err(lossless_error("malformed_selected_turn")),
786 Err(_) => continue,
787 };
788 if turn.get("type").and_then(Value::as_str) != Some("session_turn") {
789 if tagged && memory.metadata.tags.iter().any(|tag| tag == "turn") {
790 return Err(lossless_error("malformed_selected_turn"));
791 }
792 continue;
793 }
794 if !tagged && turn.get("session_id").and_then(Value::as_str) != Some(session_id) {
795 continue;
796 }
797 let content = turn
798 .get("content")
799 .and_then(Value::as_str)
800 .ok_or_else(|| lossless_error("malformed_selected_turn"))?
801 .to_string();
802 if turn.get("session_id").and_then(Value::as_str) != Some(session_id)
805 || turn.get("sequence").and_then(Value::as_u64).is_none()
806 || turn.get("timestamp").and_then(Value::as_i64).is_none()
807 {
808 return Err(lossless_error("malformed_selected_turn"));
809 }
810 turns.push(LosslessTurn {
811 content_hash: hex_encode(&Sha256::digest(content.as_bytes())),
812 memory,
813 turn,
814 content,
815 });
816 }
817 turns.sort_by_key(|turn| {
818 (
819 turn.turn["sequence"].as_u64().unwrap(),
820 turn.turn["timestamp"].as_i64().unwrap(),
821 turn.memory.metadata.id,
822 )
823 });
824 let visible_ids: std::collections::HashSet<String> = turns
825 .iter()
826 .map(|t| t.memory.metadata.id.to_string())
827 .collect();
828 let metadata: Vec<Value> = turns.iter().map(|t| {
829 let relationships: Vec<&str> = t.memory.metadata.tags.iter().filter_map(|tag| {
830 let (_, id) = tag.split_once(':')?;
831 ((tag.starts_with("superseded-by:") || tag.starts_with("supersedes:")) && visible_ids.contains(id)).then_some(tag.as_str())
832 }).collect();
833 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})
834 }).collect();
835 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();
836 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()));
837 let (mut index, mut offset) = match placement {
838 Some((token_view, index, offset)) => {
839 if token_view != view {
840 return Err(lossless_error("stale_view"));
841 }
842 (index, offset)
843 }
844 None => (0, 0),
845 };
846 if index > turns.len() || (index == turns.len() && offset != 0) {
847 return Err(lossless_error("invalid_placement"));
848 }
849 if index < turns.len() && offset != 0 && offset >= turns[index].content.len() {
850 return Err(lossless_error("invalid_placement"));
851 }
852 let mut records = Vec::new();
853 while index < turns.len() && records.len() < page_size {
854 let turn = &turns[index];
855 let bytes = turn.content.as_bytes();
856 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});
857 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()});
858 if offset == 0 && serde_json::to_vec(&candidate).unwrap().len() <= max_wire {
859 records.push(whole);
860 index += 1;
861 offset = 0;
862 continue;
863 }
864 if !records.is_empty() {
865 break;
866 }
867 let mut take = bytes.len().saturating_sub(offset);
868 while take > 0 {
869 let end = offset + take;
870 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()}});
871 let next = if end == bytes.len() {
872 lossless_cursor(
873 session_id,
874 include_superseded,
875 page_size,
876 max_wire,
877 &view,
878 index + 1,
879 0,
880 )
881 } else {
882 lossless_cursor(
883 session_id,
884 include_superseded,
885 page_size,
886 max_wire,
887 &view,
888 index,
889 end,
890 )
891 };
892 let final_chunk = end == bytes.len() && index + 1 == turns.len();
893 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});
894 if serde_json::to_vec(&candidate).unwrap().len() <= max_wire {
895 records.push(chunk);
896 if end == bytes.len() {
897 index += 1;
898 offset = 0;
899 } else {
900 offset = end;
901 }
902 break;
903 }
904 take /= 2;
905 }
906 if records.is_empty() {
907 return Err(lossless_error("wire_ceiling_too_small"));
908 }
909 break;
910 }
911 let complete = index == turns.len() && offset == 0;
912 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});
913 if serde_json::to_vec(&response).unwrap().len() > max_wire {
914 return Err(lossless_error("wire_ceiling_too_small"));
915 }
916 Ok(response)
917 }
918}
919
920#[async_trait]
921impl Tool for SessionReplayTool {
922 fn name(&self) -> &str {
923 "session.replay"
924 }
925 fn gana(&self) -> Gana {
926 Gana::StraddlingLegs
927 }
928 fn effects(&self) -> &EffectRow {
929 &self.effects
930 }
931 fn input_schema(&self) -> Value {
932 super::common::schema(
933 &json!({
934 "mode": super::common::str_prop("full | selective | progressive | lossless (default full)"),
935 "session_id": super::common::str_prop("Target session (required explicit UUID for lossless; otherwise default most recent)"),
936 "n": super::common::int_prop("Maximum turns (default 50)"),
937 "since": super::common::str_prop("Time-range floor: epoch seconds, RFC 3339, or YYYY-MM-DD"),
938 "until": super::common::str_prop("Time-range ceiling: epoch seconds, RFC 3339, or YYYY-MM-DD"),
939 "include_superseded": {
940 "type": "boolean",
941 "description": "Also return turns replaced via supersedes (default false)."
942 },
943 "turn_types": super::common::str_array_prop("Selective mode: turn types to keep"),
944 "min_importance": super::common::num_prop("Selective mode floor (default 0.7)"),
945 "token_budget": super::common::int_prop("Progressive mode token budget (default 2000)"),
946 "page_size": super::common::int_prop("Lossless mode records per page (1-64, default 16)"),
947 "max_wire_bytes": super::common::int_prop("Lossless mode serialized JSON ceiling (1024-49152)"),
948 "cursor": super::common::str_prop("Lossless mode opaque placement cursor"),
949 }),
950 &[],
951 )
952 }
953 fn description(&self) -> &str {
954 "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."
955 }
956 async fn call(&self, _ctx: &mut Context, args: Value) -> wm_core::Result<Value> {
957 let mode = args.get("mode").and_then(Value::as_str).unwrap_or("full");
958 if mode == "lossless" {
959 return self.lossless(&args);
960 }
961 let requested_session_id = args
962 .get("session_id")
963 .and_then(Value::as_str)
964 .filter(|sid| !sid.is_empty());
965 let session_id = match requested_session_id {
970 Some(sid) => Some(sid.to_string()),
971 None => self
972 .store
973 .scan_all(Galaxy::Sessions)?
974 .iter()
975 .filter(|m| m.metadata.tags.contains(&"start".to_string()))
976 .max_by_key(|m| m.metadata.created_at)
977 .map(|m| m.metadata.id.to_string()),
978 };
979 let n = args.get("n").and_then(Value::as_u64).unwrap_or(50) as usize;
980 let include_superseded = args
981 .get("include_superseded")
982 .and_then(Value::as_bool)
983 .unwrap_or(false);
984 let loaded_turns = match session_id.as_deref() {
988 Some(sid) => load_turns(&self.store, Some(sid), 10_000, include_superseded)?,
989 None => Vec::new(),
990 };
991 let turns = filter_by_time(loaded_turns, &args)?;
992
993 if requested_session_id.is_some() && turns.is_empty() {
996 return Err(wm_core::CoreError::InvalidArgs(format!(
997 "no session found with id {requested_session_id:?}"
998 )));
999 }
1000
1001 let selected: Vec<(Memory, Value)> = match mode {
1002 "selective" => {
1003 let min_importance = args
1004 .get("min_importance")
1005 .and_then(Value::as_f64)
1006 .unwrap_or(0.7);
1007 let turn_types: Vec<String> = args
1008 .get("turn_types")
1009 .and_then(Value::as_array)
1010 .map_or_else(
1011 || vec!["decision".into(), "breakthrough".into(), "answer".into()],
1012 |a| {
1013 a.iter()
1014 .filter_map(Value::as_str)
1015 .map(str::to_string)
1016 .collect()
1017 },
1018 );
1019 turns
1020 .into_iter()
1021 .filter(|(_, v)| {
1022 v.get("importance").and_then(Value::as_f64).unwrap_or(0.0) >= min_importance
1023 && v.get("turn_type")
1024 .and_then(Value::as_str)
1025 .is_some_and(|t| turn_types.contains(&t.to_string()))
1026 })
1027 .collect()
1028 }
1029 "progressive" => {
1030 let budget = args
1031 .get("token_budget")
1032 .and_then(Value::as_u64)
1033 .unwrap_or(2000) as usize;
1034 let mut used = 0usize;
1035 let mut out = Vec::new();
1036 for (m, v) in turns.into_iter().rev() {
1037 let approx = v
1038 .get("content")
1039 .and_then(Value::as_str)
1040 .map_or(0, |c| c.len() / 4);
1041 if used + approx > budget {
1042 break;
1043 }
1044 used += approx;
1045 out.push((m, v));
1046 }
1047 out.reverse();
1048 out
1049 }
1050 _ => turns
1051 .into_iter()
1052 .rev()
1053 .take(n)
1054 .collect::<Vec<_>>()
1055 .into_iter()
1056 .rev()
1057 .collect(),
1058 };
1059
1060 let full = mode != "progressive";
1061 let formatted: Vec<Value> = selected.iter().map(|(_, v)| format_turn(v, full)).collect();
1062 Ok(json!({
1063 "status": "success",
1064 "mode": mode,
1065 "count": formatted.len(),
1066 "session_id": session_id,
1067 "turns": formatted,
1068 }))
1069 }
1070 fn stats(&self) -> &ToolStats {
1071 &self.stats
1072 }
1073}
1074
1075pub struct SessionContinuityTool {
1077 store: Arc<MemoryStore>,
1078 stats: ToolStats,
1079 effects: EffectRow,
1080}
1081
1082impl SessionContinuityTool {
1083 #[must_use]
1084 pub fn new(store: Arc<MemoryStore>) -> Self {
1085 Self {
1086 store,
1087 stats: ToolStats::default(),
1088 effects: EffectRow::read_only(vec![Resource::Galaxy("sessions".into())]),
1089 }
1090 }
1091}
1092
1093#[async_trait]
1094impl Tool for SessionContinuityTool {
1095 fn name(&self) -> &str {
1096 "session.continuity"
1097 }
1098 fn gana(&self) -> Gana {
1099 Gana::StraddlingLegs
1100 }
1101 fn effects(&self) -> &EffectRow {
1102 &self.effects
1103 }
1104 fn input_schema(&self) -> Value {
1105 super::common::schema(
1106 &json!({
1107 "current_session_id": super::common::str_prop("Session to exclude (optional)"),
1108 "n": super::common::int_prop("Number of prior turns (default 10)"),
1109 "since": super::common::str_prop("Time-range floor: epoch seconds, RFC 3339, or YYYY-MM-DD"),
1110 "until": super::common::str_prop("Time-range ceiling: epoch seconds, RFC 3339, or YYYY-MM-DD"),
1111 "max_content_bytes": super::common::int_prop("Per-turn content cap in bytes (256-262144, default 8192)"),
1112 "max_response_bytes": super::common::int_prop("Turns budget in bytes (256-262144, default 49152)"),
1113 }),
1114 &[],
1115 )
1116 }
1117 fn description(&self) -> &str {
1118 "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), max_content_bytes (per-turn cap, default 8192), max_response_bytes (turns budget, default 49152). Output is bounded: each turn carries memory_id, content_bytes, and content_truncated; the newest turns win when the budget is exceeded (turns_omitted), and an oversized checkpoint handoff is replaced by a truncation marker. Read exact originals with memory.read id=<memory_id>."
1119 }
1120 async fn call(&self, _ctx: &mut Context, args: Value) -> wm_core::Result<Value> {
1121 let current = args
1122 .get("current_session_id")
1123 .or_else(|| args.get("session_id"))
1124 .and_then(Value::as_str);
1125 let n = args.get("n").and_then(Value::as_u64).unwrap_or(10) as usize;
1126 let size_arg = |name: &str, default: usize| -> wm_core::Result<usize> {
1127 match args.get(name) {
1128 None => Ok(default),
1129 Some(v) => v
1130 .as_u64()
1131 .and_then(|v| usize::try_from(v).ok())
1132 .ok_or_else(|| {
1133 wm_core::CoreError::InvalidArgs(format!("{name} must be an integer"))
1134 }),
1135 }
1136 };
1137 let max_content_bytes =
1138 size_arg("max_content_bytes", CONTINUITY_DEFAULT_MAX_CONTENT_BYTES)?;
1139 let max_response_bytes =
1140 size_arg("max_response_bytes", CONTINUITY_DEFAULT_MAX_RESPONSE_BYTES)?;
1141 if !(CONTINUITY_MIN_BUDGET_BYTES..=CONTINUITY_MAX_BUDGET_BYTES).contains(&max_content_bytes)
1142 || !(CONTINUITY_MIN_BUDGET_BYTES..=CONTINUITY_MAX_BUDGET_BYTES)
1143 .contains(&max_response_bytes)
1144 {
1145 return Err(wm_core::CoreError::InvalidArgs(format!(
1146 "continuity budgets must be in {CONTINUITY_MIN_BUDGET_BYTES}..={CONTINUITY_MAX_BUDGET_BYTES} bytes"
1147 )));
1148 }
1149
1150 let memories = self.store.scan_all(Galaxy::Sessions)?;
1161 let mut starts: Vec<_> = memories
1162 .iter()
1163 .filter(|m| {
1164 m.metadata.tags.contains(&"start".to_string())
1165 && current.is_none_or(|c| m.metadata.id.to_string() != c)
1166 })
1167 .collect();
1168 starts.sort_by_key(|m| std::cmp::Reverse(m.metadata.created_at));
1169 let previous = starts
1170 .iter()
1171 .find(|m| {
1172 let sid = m.metadata.id.to_string();
1173 load_turns(&self.store, Some(&sid), 1, false).is_ok_and(|turns| !turns.is_empty())
1174 })
1175 .copied()
1176 .or_else(|| starts.first().copied());
1177
1178 let Some(prev) = previous else {
1179 let project = std::env::var("WM_PROJECT").ok().filter(|s| !s.is_empty());
1186 let store = self.store.path().display().to_string();
1187 let scope = project.map_or_else(
1188 || format!("store {store}"),
1189 |p| format!("store {store}, project '{p}'"),
1190 );
1191 return Ok(json!({
1192 "status": "success",
1193 "previous_session": null,
1194 "turns": [],
1195 "count": 0,
1196 "message": "no previous session found",
1197 "hint": format!(
1198 "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."
1199 ),
1200 }));
1201 };
1202
1203 let prev_id = prev.metadata.id.to_string();
1204 let mut turns = filter_by_time(
1205 load_turns(&self.store, Some(&prev_id), 10_000, false)?,
1206 &args,
1207 )?;
1208 let total = turns.len();
1209 let selected = turns.split_off(total.saturating_sub(n));
1210 let newest_first: Vec<(Memory, Value)> = selected.into_iter().rev().collect();
1215 let per_turn_cap = max_content_bytes.min(
1216 max_response_bytes
1217 .saturating_sub(CONTINUITY_WRAPPER_ALLOWANCE)
1218 .max(1),
1219 );
1220 let mut formatted_rev: Vec<Value> = Vec::new();
1221 let mut used = 0usize;
1222 let mut turns_omitted = 0usize;
1223 for (idx, (mem, v)) in newest_first.iter().enumerate() {
1224 let turn = format_continuity_turn(mem, v, per_turn_cap);
1225 let size = serde_json::to_string(&turn).map_or(0, |s| s.len());
1226 if idx > 0 && used + size > max_response_bytes {
1227 turns_omitted = newest_first.len() - idx;
1228 break;
1229 }
1230 used += size;
1231 formatted_rev.push(turn);
1232 }
1233 let tail: Vec<Value> = formatted_rev.into_iter().rev().collect();
1234 let content_truncated = tail.iter().any(|t| t["content_truncated"] == true);
1235 let truncated = turns_omitted > 0 || content_truncated;
1236
1237 let (checkpoint, checkpoint_id) = match latest_checkpoint_handoff(&self.store, &prev_id)? {
1246 Some((id, _, handoff)) => {
1247 let bytes = serde_json::to_string(&handoff).map_or(0, |s| s.len());
1248 if bytes > max_response_bytes {
1249 (
1250 json!({
1251 "truncated": true,
1252 "bytes": bytes,
1253 "hint": format!(
1254 "checkpoint handoff exceeds the {max_response_bytes} B budget — read the full record with memory.read id={id}"
1255 ),
1256 }),
1257 Value::String(id),
1258 )
1259 } else {
1260 (handoff, Value::String(id))
1261 }
1262 }
1263 None => (Value::Null, Value::Null),
1264 };
1265
1266 let mut response = json!({
1267 "status": "success",
1268 "previous_session": prev_id,
1269 "count": tail.len(),
1270 "turns": tail,
1271 "checkpoint": checkpoint,
1272 "checkpoint_id": checkpoint_id,
1273 "truncated": truncated,
1274 "turns_omitted": turns_omitted,
1275 "max_content_bytes": per_turn_cap,
1276 "max_response_bytes": max_response_bytes,
1277 });
1278 if truncated {
1279 response["hint"] = json!(format!(
1280 "continuity output is bounded (per-turn cap {per_turn_cap} B, turns budget \
1281 {max_response_bytes} B). Every turn carries memory_id; read the exact \
1282 original with memory.read id=<memory_id>."
1283 ));
1284 }
1285 let bytes_available: usize = newest_first.iter().map(|(m, _)| m.content.len()).sum();
1290 super::common::append_savings_row(
1291 &self.store,
1292 &json!({
1293 "ts_ms": wm_core::time::now_unix_millis(),
1294 "op": "continuity",
1295 "previous_session": prev_id,
1296 "turns_available": total,
1297 "turns_returned": tail.len(),
1298 "turns_omitted": turns_omitted,
1299 "bytes_available": bytes_available,
1300 "bytes_injected": used,
1301 "max_response_bytes": max_response_bytes,
1302 }),
1303 );
1304 Ok(response)
1305 }
1306 fn stats(&self) -> &ToolStats {
1307 &self.stats
1308 }
1309}
1310
1311pub struct SessionDigestTool {
1318 store: Arc<MemoryStore>,
1319 stats: ToolStats,
1320 effects: EffectRow,
1321}
1322
1323impl SessionDigestTool {
1324 #[must_use]
1325 pub fn new(store: Arc<MemoryStore>) -> Self {
1326 Self {
1327 store,
1328 stats: ToolStats::default(),
1329 effects: EffectRow::read_only(vec![Resource::Galaxy("sessions".into())]),
1330 }
1331 }
1332}
1333
1334const DIGEST_SECTION_ORDER: &[&str] = &["decision", "breakthrough", "error", "summary"];
1336
1337#[async_trait]
1338impl Tool for SessionDigestTool {
1339 fn name(&self) -> &str {
1340 "session.digest"
1341 }
1342 fn gana(&self) -> Gana {
1343 Gana::StraddlingLegs
1344 }
1345 fn effects(&self) -> &EffectRow {
1346 &self.effects
1347 }
1348 fn input_schema(&self) -> Value {
1349 super::common::schema(
1350 &json!({
1351 "session_id": super::common::str_prop("Session to digest (default: most recent)"),
1352 "min_importance": super::common::num_prop("Importance floor (default 0.5)"),
1353 "include_checkpoint": {
1354 "type": "boolean",
1355 "description": "Append the latest checkpoint's git/handoff state (default true)."
1356 },
1357 "since": super::common::str_prop("Time-range floor: epoch seconds, RFC 3339, or YYYY-MM-DD"),
1358 "until": super::common::str_prop("Time-range ceiling: epoch seconds, RFC 3339, or YYYY-MM-DD"),
1359 }),
1360 &[],
1361 )
1362 }
1363 fn description(&self) -> &str {
1364 "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."
1365 }
1366 async fn call(&self, _ctx: &mut Context, args: Value) -> wm_core::Result<Value> {
1367 let session_id = match args.get("session_id").and_then(Value::as_str) {
1368 Some(sid) if !sid.is_empty() => sid.to_string(),
1369 _ => self
1370 .store
1371 .scan_all(Galaxy::Sessions)?
1372 .iter()
1373 .filter(|m| m.metadata.tags.contains(&"start".to_string()))
1374 .max_by_key(|m| m.metadata.created_at)
1375 .map(|m| m.metadata.id.to_string())
1376 .ok_or_else(|| {
1377 wm_core::CoreError::Tool("no session found — run session.start first".into())
1378 })?,
1379 };
1380 let min_importance = args
1381 .get("min_importance")
1382 .and_then(Value::as_f64)
1383 .unwrap_or(0.5);
1384 let include_checkpoint = args
1385 .get("include_checkpoint")
1386 .and_then(Value::as_bool)
1387 .unwrap_or(true);
1388
1389 let mut turns: Vec<_> = filter_by_time(
1390 load_turns(&self.store, Some(&session_id), 10_000, false)?,
1391 &args,
1392 )?
1393 .into_iter()
1394 .filter(|(_, v)| {
1395 v.get("importance").and_then(Value::as_f64).unwrap_or(0.0) >= min_importance
1396 })
1397 .collect();
1398 turns.sort_by(|a, b| {
1399 b.1.get("importance")
1400 .and_then(Value::as_f64)
1401 .unwrap_or(0.0)
1402 .total_cmp(&a.1.get("importance").and_then(Value::as_f64).unwrap_or(0.0))
1403 });
1404
1405 let mut groups: Vec<(String, Vec<&Value>)> = Vec::new();
1407 for (_, v) in &turns {
1408 let t = v
1409 .get("turn_type")
1410 .and_then(Value::as_str)
1411 .unwrap_or("message")
1412 .to_string();
1413 match groups.iter_mut().find(|(name, _)| *name == t) {
1414 Some((_, list)) => list.push(v),
1415 None => groups.push((t, vec![v])),
1416 }
1417 }
1418 groups.sort_by_key(|(name, _)| {
1419 (
1420 DIGEST_SECTION_ORDER
1421 .iter()
1422 .position(|k| k == name)
1423 .unwrap_or(DIGEST_SECTION_ORDER.len()),
1424 name.clone(),
1425 )
1426 });
1427
1428 let mut digest = format!("# Session handoff — {session_id}\n");
1429 let mut included = 0usize;
1430 for (turn_type, items) in &groups {
1431 writeln!(
1432 digest,
1433 "\n## {} ({})",
1434 capitalize(&pluralize(turn_type)),
1435 items.len()
1436 )
1437 .expect("write to String cannot fail");
1438 for v in items {
1439 let importance = v.get("importance").and_then(Value::as_f64).unwrap_or(0.0);
1440 let content = v.get("content").and_then(Value::as_str).unwrap_or("");
1441 writeln!(digest, "- ({importance:.2}) {content}")
1442 .expect("write to String cannot fail");
1443 included += 1;
1444 }
1445 }
1446
1447 let mut checkpoint_state = Value::Null;
1449 if include_checkpoint {
1450 if let Some((_, _, cp)) = latest_checkpoint_handoff(&self.store, &session_id)? {
1451 digest.push_str("\n## Checkpoint state\n");
1452 if let Some(git) = cp.get("git") {
1453 writeln!(
1454 digest,
1455 "- commit `{}` on `{}` ({} dirty files)",
1456 git.get("commit").and_then(Value::as_str).unwrap_or("?"),
1457 git.get("branch").and_then(Value::as_str).unwrap_or("?"),
1458 git.get("dirty_count").and_then(Value::as_i64).unwrap_or(0)
1459 )
1460 .expect("write to String cannot fail");
1461 }
1462 if let Some(q) = cp.get("next_queue").and_then(Value::as_array) {
1463 if !q.is_empty() {
1464 writeln!(
1465 digest,
1466 "- next queue: {}",
1467 q.iter()
1468 .filter_map(Value::as_str)
1469 .collect::<Vec<_>>()
1470 .join(" → ")
1471 )
1472 .expect("write to String cannot fail");
1473 }
1474 }
1475 if let Some(f) = cp.get("open_flags").and_then(Value::as_array) {
1476 if !f.is_empty() {
1477 writeln!(
1478 digest,
1479 "- open flags: {}",
1480 f.iter()
1481 .filter_map(Value::as_str)
1482 .collect::<Vec<_>>()
1483 .join("; ")
1484 )
1485 .expect("write to String cannot fail");
1486 }
1487 }
1488 if let Some(tg) = cp.get("tests_green") {
1489 writeln!(digest, "- tests green: {tg}").expect("write to String cannot fail");
1490 }
1491 checkpoint_state = cp;
1492 }
1493 }
1494
1495 Ok(json!({
1496 "status": "success",
1497 "session_id": session_id,
1498 "digest": digest,
1499 "turns_included": included,
1500 "turns_total_scanned": turns.len(),
1501 "sections": groups.iter().map(|(t, items)| json!({"type": t, "count": items.len()})).collect::<Vec<_>>(),
1502 "checkpoint": checkpoint_state,
1503 }))
1504 }
1505 fn stats(&self) -> &ToolStats {
1506 &self.stats
1507 }
1508}
1509
1510fn capitalize(s: &str) -> String {
1512 let mut chars = s.chars();
1513 match chars.next() {
1514 Some(first) => first.to_uppercase().collect::<String>() + chars.as_str(),
1515 None => String::new(),
1516 }
1517}
1518
1519fn pluralize(s: &str) -> String {
1522 if let Some(stem) = s.strip_suffix('y') {
1523 format!("{stem}ies")
1524 } else {
1525 format!("{s}s")
1526 }
1527}
1528
1529pub struct SessionHandoffTool {
1531 store: Arc<MemoryStore>,
1532 stats: ToolStats,
1533 effects: EffectRow,
1534 search: Option<Arc<wm_memory::SearchEngine>>,
1535}
1536
1537impl SessionHandoffTool {
1538 #[must_use]
1539 pub fn new(store: Arc<MemoryStore>) -> Self {
1540 Self {
1541 store,
1542 stats: ToolStats::default(),
1543 effects: EffectRow {
1544 writes: vec![Resource::Galaxy("sessions".into())],
1545 ..Default::default()
1546 },
1547 search: None,
1548 }
1549 }
1550
1551 #[must_use]
1553 pub fn with_search(mut self, search: Option<Arc<wm_memory::SearchEngine>>) -> Self {
1554 self.search = search;
1555 self
1556 }
1557}
1558
1559#[async_trait]
1560impl Tool for SessionHandoffTool {
1561 fn name(&self) -> &str {
1562 "session.handoff"
1563 }
1564 fn gana(&self) -> Gana {
1565 Gana::StraddlingLegs
1566 }
1567 fn effects(&self) -> &EffectRow {
1568 &self.effects
1569 }
1570 fn input_schema(&self) -> Value {
1571 super::common::schema(
1572 &json!({
1573 "action": super::common::str_prop("transfer | accept | list"),
1574 "session_id": super::common::str_prop("transfer: session to hand off"),
1575 "message": super::common::str_prop("transfer: handoff note"),
1576 "handoff_id": super::common::str_prop("accept: handoff to accept"),
1577 }),
1578 &["action"],
1579 )
1580 }
1581 fn description(&self) -> &str {
1582 "Transfer or resume a session across devices (actions: transfer, accept, list). transfer: session_id (required) + message; accept: handoff_id; list: pending handoffs."
1583 }
1584 async fn call(&self, _ctx: &mut Context, args: Value) -> wm_core::Result<Value> {
1585 let action = args.get("action").and_then(Value::as_str).unwrap_or("list");
1586 match action {
1587 "transfer" => {
1588 let session_id =
1589 args.get("session_id")
1590 .and_then(Value::as_str)
1591 .ok_or_else(|| {
1592 wm_core::CoreError::InvalidArgs(
1593 "session_id required for transfer".into(),
1594 )
1595 })?;
1596 let message = args.get("message").and_then(Value::as_str).unwrap_or("");
1597 let turns = load_turns(&self.store, Some(session_id), 10_000, false)?;
1598 if turns.is_empty() {
1599 return Err(wm_core::CoreError::Tool(format!(
1600 "session {session_id} has no recorded turns"
1601 )));
1602 }
1603 let summary: Vec<Value> =
1604 turns.iter().map(|(_, v)| format_turn(v, false)).collect();
1605 let handoff_id = format!("handoff-{}", uuid::Uuid::new_v4());
1606 let mut mem = Memory::new(
1607 Galaxy::Sessions,
1608 json!({
1609 "type": "session_handoff",
1610 "handoff_id": handoff_id,
1611 "session_id": session_id,
1612 "message": message,
1613 "status": "pending",
1614 "turn_count": summary.len(),
1615 "summary": summary,
1616 "created_at": wm_core::time::now_unix_millis(),
1617 })
1618 .to_string(),
1619 );
1620 mem.metadata.tags = vec![
1621 "session".into(),
1622 "handoff".into(),
1623 format!("session:{session_id}"),
1624 ];
1625 mem.metadata.importance = 0.8;
1626 self.store.put(Galaxy::Sessions, &mem)?;
1627 super::common::index_memory(&self.store, self.search.as_deref(), &mem);
1628 Ok(json!({
1629 "status": "success",
1630 "action": "transfer",
1631 "handoff_id": handoff_id,
1632 "session_id": session_id,
1633 "turn_count": summary.len(),
1634 }))
1635 }
1636 "accept" => {
1637 let handoff_id =
1638 args.get("handoff_id")
1639 .and_then(Value::as_str)
1640 .ok_or_else(|| {
1641 wm_core::CoreError::InvalidArgs("handoff_id required for accept".into())
1642 })?;
1643 let memories = self.store.scan_all(Galaxy::Sessions)?;
1644 let found = memories.iter().find(|m| {
1645 m.metadata.tags.contains(&"handoff".to_string())
1646 && m.content.contains(handoff_id)
1647 });
1648 let Some(mem) = found else {
1649 return Err(wm_core::CoreError::Tool(format!(
1650 "handoff {handoff_id} not found"
1651 )));
1652 };
1653 let mut updated = mem.clone();
1654 if let Ok(mut v) = serde_json::from_str::<Value>(&updated.content) {
1655 v["status"] = json!("accepted");
1656 updated.content = v.to_string();
1657 }
1658 self.store.put(Galaxy::Sessions, &updated)?;
1659 super::common::index_memory(&self.store, self.search.as_deref(), &updated);
1660 Ok(json!({
1661 "status": "success",
1662 "action": "accept",
1663 "handoff_id": handoff_id,
1664 }))
1665 }
1666 "list" => {
1667 let memories = self.store.scan_all(Galaxy::Sessions)?;
1668 let handoffs: Vec<Value> = memories
1669 .iter()
1670 .filter(|m| m.metadata.tags.contains(&"handoff".to_string()))
1671 .filter_map(|m| serde_json::from_str::<Value>(&m.content).ok())
1672 .filter(|v| v.get("status").and_then(Value::as_str) == Some("pending"))
1673 .map(|v| {
1674 json!({
1675 "handoff_id": v.get("handoff_id"),
1676 "session_id": v.get("session_id"),
1677 "message": v.get("message"),
1678 "turn_count": v.get("turn_count"),
1679 })
1680 })
1681 .collect();
1682 Ok(json!({
1683 "status": "success",
1684 "action": "list",
1685 "pending_count": handoffs.len(),
1686 "handoffs": handoffs,
1687 }))
1688 }
1689 other => Err(wm_core::CoreError::InvalidArgs(format!(
1690 "unknown session.handoff action: {other}"
1691 ))),
1692 }
1693 }
1694 fn stats(&self) -> &ToolStats {
1695 &self.stats
1696 }
1697}
1698
1699#[must_use]
1701pub struct SessionExportTool {
1707 store: Arc<MemoryStore>,
1708 stats: ToolStats,
1709 effects: EffectRow,
1710}
1711
1712impl SessionExportTool {
1713 pub fn new(store: Arc<MemoryStore>) -> Self {
1714 Self {
1715 store,
1716 stats: ToolStats::default(),
1717 effects: EffectRow::read_only(vec![Resource::Galaxy("sessions".into())]),
1718 }
1719 }
1720}
1721
1722#[async_trait]
1723impl Tool for SessionExportTool {
1724 fn name(&self) -> &str {
1725 "session.export"
1726 }
1727 fn gana(&self) -> Gana {
1728 Gana::StraddlingLegs
1729 }
1730 fn effects(&self) -> &EffectRow {
1731 &self.effects
1732 }
1733 fn input_schema(&self) -> Value {
1734 super::common::schema(
1735 &json!({
1736 "session_id": super::common::str_prop("Session to export (default: most recent)"),
1737 "path": super::common::str_prop("Write JSONL to this file instead of returning inline"),
1738 }),
1739 &[],
1740 )
1741 }
1742 fn description(&self) -> &str {
1743 "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)."
1744 }
1745 async fn call(&self, _ctx: &mut Context, args: Value) -> wm_core::Result<Value> {
1746 let session_id = match args.get("session_id").and_then(Value::as_str) {
1747 Some(sid) if !sid.is_empty() => sid.to_string(),
1748 _ => self
1749 .store
1750 .scan_all(Galaxy::Sessions)?
1751 .iter()
1752 .filter(|m| m.metadata.tags.contains(&"start".to_string()))
1753 .max_by_key(|m| m.metadata.created_at)
1754 .map(|m| m.metadata.id.to_string())
1755 .ok_or_else(|| {
1756 wm_core::CoreError::Tool("no session found — run session.start first".into())
1757 })?,
1758 };
1759
1760 let mut members: Vec<Memory> = self
1764 .store
1765 .scan_all(Galaxy::Sessions)?
1766 .into_iter()
1767 .filter(|m| m.metadata.id.to_string() == session_id || m.content.contains(&session_id))
1768 .collect();
1769 members.sort_by_key(|m| m.metadata.created_at);
1770
1771 let mut jsonl = String::new();
1772 let header = wm_memory::envelope::EnvelopeHeader::new("session_export", members.len());
1778 jsonl.push_str(&header.header_line());
1779 jsonl.push('\n');
1780 for m in &members {
1781 let line = serde_json::to_string(m)
1782 .map_err(|e| wm_core::CoreError::Tool(format!("export serialize: {e}")))?;
1783 jsonl.push_str(&line);
1784 jsonl.push('\n');
1785 }
1786
1787 let path_arg = args
1788 .get("path")
1789 .and_then(Value::as_str)
1790 .filter(|s| !s.is_empty());
1791 if let Some(dest) = path_arg {
1792 std::fs::write(dest, &jsonl)
1793 .map_err(|e| wm_core::CoreError::Tool(format!("export write {dest}: {e}")))?;
1794 Ok(json!({
1795 "status": "success",
1796 "session_id": session_id,
1797 "records": members.len(),
1798 "path": dest,
1799 }))
1800 } else {
1801 Ok(json!({
1802 "status": "success",
1803 "session_id": session_id,
1804 "records": members.len(),
1805 "jsonl": jsonl,
1806 }))
1807 }
1808 }
1809 fn stats(&self) -> &ToolStats {
1810 &self.stats
1811 }
1812}
1813
1814pub struct SessionImportTool {
1825 store: Arc<MemoryStore>,
1826 search: Option<Arc<wm_memory::SearchEngine>>,
1827 stats: ToolStats,
1828 effects: EffectRow,
1829}
1830
1831impl SessionImportTool {
1832 pub fn new(store: Arc<MemoryStore>, search: Option<Arc<wm_memory::SearchEngine>>) -> Self {
1833 Self {
1834 store,
1835 search,
1836 stats: ToolStats::default(),
1837 effects: EffectRow {
1838 writes: vec![Resource::Galaxy("sessions".into())],
1839 ..Default::default()
1840 },
1841 }
1842 }
1843}
1844
1845#[async_trait]
1846impl Tool for SessionImportTool {
1847 fn name(&self) -> &str {
1848 "session.import"
1849 }
1850 fn gana(&self) -> Gana {
1851 Gana::StraddlingLegs
1852 }
1853 fn effects(&self) -> &EffectRow {
1854 &self.effects
1855 }
1856 fn input_schema(&self) -> Value {
1857 super::common::schema(
1858 &json!({
1859 "path": super::common::str_prop("Read JSONL from this file"),
1860 "jsonl": super::common::str_prop("Or pass the export payload inline"),
1861 }),
1862 &[],
1863 )
1864 }
1865 fn description(&self) -> &str {
1866 "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."
1867 }
1868 async fn call(&self, _ctx: &mut Context, args: Value) -> wm_core::Result<Value> {
1869 let payload = match args
1870 .get("path")
1871 .and_then(Value::as_str)
1872 .filter(|s| !s.is_empty())
1873 {
1874 Some(path) => std::fs::read_to_string(path)
1875 .map_err(|e| wm_core::CoreError::Tool(format!("import read {path}: {e}")))?,
1876 None => args
1877 .get("jsonl")
1878 .and_then(Value::as_str)
1879 .filter(|s| !s.is_empty())
1880 .ok_or_else(|| {
1881 wm_core::CoreError::InvalidArgs(
1882 "provide either 'path' or inline 'jsonl'".into(),
1883 )
1884 })?
1885 .to_string(),
1886 };
1887
1888 let mut envelope: Option<wm_memory::envelope::EnvelopeHeader> = None;
1893 let mut record_lines: Vec<&str> = Vec::new();
1894 let mut header_consumed = false;
1895 for line in payload.lines() {
1896 if line.trim().is_empty() {
1897 continue;
1898 }
1899 if !header_consumed {
1900 header_consumed = true;
1901 match wm_memory::envelope::read_header_line(line) {
1902 wm_memory::envelope::HeaderRead::Header(h) => {
1903 envelope = Some(h);
1904 continue;
1905 }
1906 wm_memory::envelope::HeaderRead::Refused(msg) => {
1907 return Err(wm_core::CoreError::Tool(msg));
1908 }
1909 wm_memory::envelope::HeaderRead::NotAHeader => {}
1910 }
1911 }
1912 record_lines.push(line);
1913 }
1914
1915 let readonly_engine = self.search.as_ref().is_some_and(|s| s.is_readonly());
1920 if readonly_engine {
1921 tracing::warn!(
1922 "session.import running against a read-only search engine — records land \
1923 in LMDB unindexed; they become searchable at the next writable startup \
1924 (heal_index_drift)"
1925 );
1926 }
1927 let mut writer_slot = match (&self.search, readonly_engine) {
1928 (Some(s), false) => s.writer().ok(),
1929 _ => None,
1930 };
1931
1932 let mut imported = 0usize;
1933 let mut indexed = 0usize;
1934 let mut skipped = 0usize;
1935 let mut session_ids: Vec<String> = Vec::new();
1936 for (lineno, line) in record_lines.iter().enumerate() {
1937 let mem: Memory = match serde_json::from_str(line) {
1938 Ok(m) => m,
1939 Err(e) => {
1940 skipped += 1;
1941 tracing::warn!(line = lineno + 1, error = %e, "skipping unparseable export line");
1942 continue;
1943 }
1944 };
1945 if let Ok(parsed) = serde_json::from_str::<Value>(&mem.content) {
1946 if let Some(sid) = parsed.get("session_id").and_then(Value::as_str) {
1947 if !session_ids.iter().any(|s| s == sid) {
1948 session_ids.push(sid.to_string());
1949 }
1950 }
1951 }
1952 if let (Some(search), Some(writer)) = (&self.search, writer_slot.as_mut()) {
1956 let id_str = mem.metadata.id.to_string();
1957 let _ = search.delete_document(writer, &id_str);
1958 match search.add_document(
1959 writer,
1960 &id_str,
1961 mem.metadata.galaxy.db_name(),
1962 &mem.content,
1963 &mem.metadata.tags,
1964 mem.metadata.created_at.timestamp(),
1965 ) {
1966 Ok(()) => indexed += 1,
1967 Err(e) => {
1968 tracing::warn!(id = %id_str, error = %e, "import index add failed (LMDB record kept)");
1969 }
1970 }
1971 }
1972 self.store.put(Galaxy::Sessions, &mem)?;
1973 imported += 1;
1974 }
1975
1976 if let Some(search) = &self.search {
1977 if let Some(mut writer) = writer_slot {
1978 search
1979 .commit(&mut writer)
1980 .map_err(|e| wm_core::CoreError::Tool(format!("import index commit: {e}")))?;
1981 }
1982 }
1983
1984 let mut warnings: Vec<String> = Vec::new();
1985 if let Some(h) = &envelope {
1986 if h.count != imported {
1987 let msg = format!(
1988 "envelope declares count {} but {} records imported",
1989 h.count, imported
1990 );
1991 tracing::warn!("{msg}");
1992 warnings.push(msg);
1993 }
1994 }
1995
1996 let envelope_info = envelope.as_ref().map(|h| {
1997 json!({
1998 "format_version": h.format_version,
1999 "kind": h.kind,
2000 "generator": h.generator,
2001 "created_at": h.created_at,
2002 "declared_count": h.count,
2003 })
2004 });
2005
2006 Ok(json!({
2007 "status": "success",
2008 "imported": imported,
2009 "skipped": skipped,
2010 "session_ids": session_ids,
2011 "indexed": indexed,
2012 "envelope": envelope_info,
2013 "warnings": warnings,
2014 }))
2015 }
2016 fn stats(&self) -> &ToolStats {
2017 &self.stats
2018 }
2019}
2020
2021fn track_of(mem: &Memory) -> Option<&str> {
2023 mem.metadata
2024 .tags
2025 .iter()
2026 .find_map(|t| t.strip_prefix("track:"))
2027}
2028
2029pub struct SessionTrackLogTool {
2040 store: Arc<MemoryStore>,
2041 stats: ToolStats,
2042 effects: EffectRow,
2043}
2044
2045impl SessionTrackLogTool {
2046 #[must_use]
2047 pub fn new(store: Arc<MemoryStore>) -> Self {
2048 Self {
2049 store,
2050 stats: ToolStats::default(),
2051 effects: EffectRow::read_only(vec![Resource::Galaxy("sessions".into())]),
2052 }
2053 }
2054}
2055
2056#[async_trait]
2057impl Tool for SessionTrackLogTool {
2058 fn name(&self) -> &str {
2059 "session.track_log"
2060 }
2061 fn gana(&self) -> Gana {
2062 Gana::StraddlingLegs
2063 }
2064 fn effects(&self) -> &EffectRow {
2065 &self.effects
2066 }
2067 fn input_schema(&self) -> Value {
2068 super::common::schema(
2069 &json!({
2070 "track": super::common::str_prop("Track slug to read (recorded via session.record/session.checkpoint 'track')"),
2071 "tracks": super::common::str_array_prop("Multiple track slugs — merged chronologically (the related-ticket review view)"),
2072 "since": super::common::str_prop("Time-range floor: epoch seconds, RFC 3339, or YYYY-MM-DD"),
2073 "until": super::common::str_prop("Time-range ceiling: epoch seconds, RFC 3339, or YYYY-MM-DD"),
2074 "limit": super::common::positive_int_prop("Max entries returned per track, most recent kept (default 100)"),
2075 "include_superseded": {
2076 "type": "boolean",
2077 "description": "Include superseded turns (default false — the current story only)."
2078 },
2079 }),
2080 &[],
2081 )
2082 }
2083 fn description(&self) -> &str {
2084 "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."
2085 }
2086 async fn call(&self, _ctx: &mut Context, args: Value) -> wm_core::Result<Value> {
2087 let mut requested: Vec<String> = Vec::new();
2088 if let Some(track) = args.get("track").and_then(Value::as_str) {
2089 validate_track(track)?;
2090 requested.push(track.to_string());
2091 }
2092 if let Some(tracks) = args.get("tracks").and_then(Value::as_array) {
2093 for t in tracks {
2094 let t = t.as_str().ok_or_else(|| {
2095 wm_core::CoreError::InvalidArgs("'tracks' entries must be strings".into())
2096 })?;
2097 validate_track(t)?;
2098 if !requested.iter().any(|r| r == t) {
2099 requested.push(t.to_string());
2100 }
2101 }
2102 }
2103 let limit = args
2104 .get("limit")
2105 .and_then(Value::as_u64)
2106 .unwrap_or(100)
2107 .clamp(1, 1000) as usize;
2108 let include_superseded = args
2109 .get("include_superseded")
2110 .and_then(Value::as_bool)
2111 .unwrap_or(false);
2112
2113 let since = match args.get("since") {
2116 Some(v) if !v.is_null() => Some(parse_time_bound(v, false).ok_or_else(|| {
2117 wm_core::CoreError::InvalidArgs(
2118 "invalid 'since' — use epoch seconds, RFC 3339, or YYYY-MM-DD".into(),
2119 )
2120 })?),
2121 _ => None,
2122 };
2123 let until = match args.get("until") {
2124 Some(v) if !v.is_null() => Some(parse_time_bound(v, true).ok_or_else(|| {
2125 wm_core::CoreError::InvalidArgs(
2126 "invalid 'until' — use epoch seconds, RFC 3339, or YYYY-MM-DD".into(),
2127 )
2128 })?),
2129 _ => None,
2130 };
2131
2132 type Entries = std::collections::BTreeMap<String, Vec<(DateTime<Utc>, Value)>>;
2133 let mut by_track: Entries = std::collections::BTreeMap::new();
2134 let mut latest_checkpoint: std::collections::BTreeMap<String, (DateTime<Utc>, Value)> =
2135 std::collections::BTreeMap::new();
2136
2137 for mem in &self.store.scan_all(Galaxy::Sessions)? {
2138 let Some(track) = track_of(mem) else { continue };
2139 if !requested.is_empty() && !requested.iter().any(|r| r == track) {
2140 continue;
2141 }
2142 if !include_superseded
2143 && mem
2144 .metadata
2145 .tags
2146 .iter()
2147 .any(|t| t.starts_with("superseded-by:"))
2148 {
2149 continue;
2150 }
2151 let created_at = mem.metadata.created_at;
2152 if since.is_some_and(|t| created_at < t) || until.is_some_and(|t| created_at > t) {
2153 continue;
2154 }
2155 let entry = if let Some(turn) = turn_json(mem) {
2156 Some(json!({
2157 "kind": "turn",
2158 "track": track,
2159 "session_id": turn.get("session_id"),
2160 "sequence": turn.get("sequence"),
2161 "role": turn.get("role"),
2162 "turn_type": turn.get("turn_type"),
2163 "importance": turn.get("importance"),
2164 "created_at": created_at.to_rfc3339(),
2165 "content": turn.get("content"),
2166 }))
2167 } else if mem.metadata.tags.contains(&"checkpoint".to_string()) {
2168 serde_json::from_str::<Value>(&mem.content)
2169 .ok()
2170 .map(|parsed| {
2171 json!({
2172 "kind": "checkpoint",
2173 "track": track,
2174 "session_id": parsed.get("session_id"),
2175 "label": parsed.get("label"),
2176 "created_at": created_at.to_rfc3339(),
2177 "handoff": parsed.get("handoff"),
2178 })
2179 })
2180 } else {
2181 None
2182 };
2183 let Some(entry) = entry else { continue };
2184 if entry["kind"] == "checkpoint" {
2185 latest_checkpoint
2186 .entry(track.to_string())
2187 .and_modify(|(ts, existing)| {
2188 if created_at > *ts {
2189 *ts = created_at;
2190 *existing = entry.clone();
2191 }
2192 })
2193 .or_insert_with(|| (created_at, entry.clone()));
2194 }
2195 by_track
2196 .entry(track.to_string())
2197 .or_default()
2198 .push((created_at, entry));
2199 }
2200
2201 let now = Utc::now();
2202 let overview = requested.is_empty();
2203 let track_names: Vec<String> = if overview {
2204 by_track.keys().cloned().collect()
2205 } else {
2206 requested
2207 };
2208 let mut tracks: Vec<Value> = Vec::with_capacity(track_names.len());
2209 for track in track_names {
2210 let mut list = by_track.remove(&track).unwrap_or_default();
2211 list.sort_by_key(|(ts, _)| *ts);
2212 let entry_count = list.len();
2213 let last_activity = list.last().map(|(ts, _)| *ts);
2214 let latest_entry = list.last().map(|(_, e)| {
2215 json!({
2216 "kind": e.get("kind"),
2217 "turn_type": e.get("turn_type"),
2218 "label": e.get("label"),
2219 "preview": e
2220 .get("content")
2221 .and_then(Value::as_str)
2222 .map(|s| s.chars().take(160).collect::<String>()),
2223 })
2224 });
2225 if list.len() > limit {
2226 list.drain(0..list.len() - limit);
2227 }
2228 let mut out = json!({
2229 "track": track,
2230 "known": entry_count > 0,
2231 "entry_count": entry_count,
2232 "returned": list.len(),
2233 "last_activity_at": last_activity.map(|t| t.to_rfc3339()),
2234 "age_seconds": last_activity.map(|t| (now - t).num_seconds().max(0)),
2235 "latest_checkpoint": latest_checkpoint.get(&track).map(|(ts, cp)| {
2236 json!({
2237 "created_at": ts.to_rfc3339(),
2238 "entry": cp,
2239 })
2240 }),
2241 });
2242 if overview {
2243 out["latest_entry"] = latest_entry.unwrap_or(Value::Null);
2244 } else {
2245 out["entries"] = Value::Array(list.into_iter().map(|(_, e)| e).collect());
2246 }
2247 tracks.push(out);
2248 }
2249
2250 Ok(json!({
2251 "status": "success",
2252 "mode": if overview { "overview" } else { "log" },
2253 "track_count": tracks.len(),
2254 "tracks": tracks,
2255 "disclosure": "directional alignment is the reader's judgment — this view supplies the log, the checkpoint handoff, and staleness facts (age_seconds)",
2256 }))
2257 }
2258 fn stats(&self) -> &ToolStats {
2259 &self.stats
2260 }
2261}
2262
2263pub fn register_session_ops(
2264 registry: &wm_dispatch::ToolRegistry,
2265 store: &Arc<MemoryStore>,
2266 search: Option<Arc<wm_memory::SearchEngine>>,
2267) -> wm_dispatch::ToolRegistry {
2268 registry
2269 .register(Arc::new(
2270 SessionRecordTool::new(store.clone()).with_search(search.clone()),
2271 ))
2272 .register(Arc::new(SessionReplayTool::new(store.clone())))
2273 .register(Arc::new(SessionContinuityTool::new(store.clone())))
2274 .register(Arc::new(SessionTrackLogTool::new(store.clone())))
2275 .register(Arc::new(
2276 SessionHandoffTool::new(store.clone()).with_search(search.clone()),
2277 ))
2278 .register(Arc::new(SessionExportTool::new(store.clone())))
2279 .register(Arc::new(SessionImportTool::new(store.clone(), search)))
2280}
2281
2282#[cfg(test)]
2283mod tests {
2284 use super::*;
2285 use crate::expansion::session::SessionCheckpointNodiscoveryTool;
2286
2287 fn test_store() -> Arc<MemoryStore> {
2288 let dir = tempfile::tempdir().unwrap();
2289 let path = dir.path().join("lmdb");
2290 std::fs::create_dir_all(&path).unwrap();
2291 Arc::new(MemoryStore::open_default(path).unwrap())
2292 }
2293
2294 fn start_session(store: &MemoryStore) -> String {
2295 let mut mem = Memory::new(
2296 Galaxy::Sessions,
2297 json!({"type": "session_start"}).to_string(),
2298 );
2299 mem.metadata.tags = vec!["session".into(), "start".into()];
2300 store.put(Galaxy::Sessions, &mem).unwrap();
2301 mem.metadata.id.to_string()
2302 }
2303
2304 fn start_session_aged(store: &MemoryStore, age_secs: i64) -> String {
2307 let mut mem = Memory::new(
2308 Galaxy::Sessions,
2309 json!({"type": "session_start"}).to_string(),
2310 );
2311 mem.metadata.tags = vec!["session".into(), "start".into()];
2312 mem.metadata.created_at = chrono::Utc::now() - chrono::Duration::seconds(age_secs);
2313 store.put(Galaxy::Sessions, &mem).unwrap();
2314 mem.metadata.id.to_string()
2315 }
2316
2317 fn record_aged_turn(store: &MemoryStore, sid: &str, age_days: u32, content: &str) {
2319 let mut mem = Memory::new(
2320 Galaxy::Sessions,
2321 json!({
2322 "type": "session_turn",
2323 "session_id": sid,
2324 "role": "ai",
2325 "turn_type": "decision",
2326 "importance": 0.9,
2327 "content": content,
2328 })
2329 .to_string(),
2330 );
2331 mem.metadata.tags = vec!["session".into(), "turn".into(), format!("session:{sid}")];
2332 mem.metadata.created_at = Utc::now() - chrono::Duration::days(i64::from(age_days));
2333 store.put(Galaxy::Sessions, &mem).unwrap();
2334 }
2335
2336 #[tokio::test]
2337 async fn replay_time_filters_since_and_until() {
2338 let store = test_store();
2339 let sid = start_session(&store);
2340 record_aged_turn(&store, &sid, 3, "three days ago");
2341 record_aged_turn(&store, &sid, 2, "two days ago");
2342 record_aged_turn(&store, &sid, 1, "yesterday");
2343 record_aged_turn(&store, &sid, 0, "today");
2344
2345 let replay = SessionReplayTool::new(store);
2346 let mut ctx = Context::default();
2347
2348 let two_days_ago = (Utc::now() - chrono::Duration::days(2)).format("%Y-%m-%d");
2350 let v = replay
2351 .call(
2352 &mut ctx,
2353 json!({"session_id": sid, "since": two_days_ago.to_string()}),
2354 )
2355 .await
2356 .unwrap();
2357 assert_eq!(v["count"], 3, "since=date keeps day-of + later: {v}");
2358
2359 let until_epoch = (Utc::now() - chrono::Duration::hours(23)).timestamp();
2361 let v = replay
2362 .call(&mut ctx, json!({"session_id": sid, "until": until_epoch}))
2363 .await
2364 .unwrap();
2365 assert_eq!(
2366 v["count"], 3,
2367 "until=epoch(23h ago) keeps the three older turns: {v}"
2368 );
2369
2370 let since = (Utc::now() - chrono::Duration::hours(60)).to_rfc3339();
2372 let until = (Utc::now() - chrono::Duration::hours(12)).to_rfc3339();
2373 let v = replay
2374 .call(
2375 &mut ctx,
2376 json!({"session_id": sid, "since": since, "until": until}),
2377 )
2378 .await
2379 .unwrap();
2380 assert_eq!(v["count"], 2, "window keeps two/two-days-ago turns: {v}");
2381 for turn in v["turns"].as_array().unwrap() {
2382 assert_ne!(
2383 turn["content"], "today",
2384 "time filters must exclude out-of-window turns"
2385 );
2386 }
2387
2388 assert!(
2390 replay
2391 .call(&mut ctx, json!({"session_id": sid, "since": "not-a-date"}))
2392 .await
2393 .is_err()
2394 );
2395 }
2396
2397 #[tokio::test]
2398 async fn replay_omitted_session_id_uses_latest_session_start() {
2399 let store = test_store();
2400 let older = start_session_aged(&store, 120);
2404 record_aged_turn(&store, &older, 0, "older-session-only");
2405 let latest = start_session_aged(&store, 60);
2406 record_aged_turn(&store, &latest, 0, "latest-session-only");
2407
2408 let replay = SessionReplayTool::new(store);
2409 let mut ctx = Context::default();
2410
2411 let omitted = replay
2412 .call(&mut ctx, json!({"mode": "full"}))
2413 .await
2414 .unwrap();
2415 assert_eq!(omitted["session_id"], latest);
2416 assert_eq!(
2417 omitted["count"], 1,
2418 "omitted id must not combine sessions: {omitted}"
2419 );
2420 assert_eq!(omitted["turns"][0]["content"], "latest-session-only");
2421
2422 let explicit_older = replay
2423 .call(&mut ctx, json!({"session_id": older, "mode": "full"}))
2424 .await
2425 .unwrap();
2426 assert_eq!(explicit_older["session_id"], older);
2427 assert_eq!(explicit_older["count"], 1);
2428 assert_eq!(explicit_older["turns"][0]["content"], "older-session-only");
2429 }
2430
2431 #[tokio::test]
2432 async fn replay_without_any_session_remains_truthfully_empty() {
2433 let replay = SessionReplayTool::new(test_store());
2434 let mut ctx = Context::default();
2435
2436 let value = replay.call(&mut ctx, json!({})).await.unwrap();
2437 assert_eq!(value["status"], "success");
2438 assert_eq!(value["session_id"], Value::Null);
2439 assert_eq!(value["count"], 0);
2440 assert_eq!(value["turns"], json!([]));
2441 }
2442
2443 #[tokio::test]
2444 async fn replay_omitted_id_does_not_combine_orphan_turns_without_a_start() {
2445 let store = test_store();
2446 record_aged_turn(&store, "orphan-a", 0, "orphan-a-only");
2450 record_aged_turn(&store, "orphan-b", 0, "orphan-b-only");
2451
2452 let replay = SessionReplayTool::new(store);
2453 let mut ctx = Context::default();
2454
2455 let omitted = replay.call(&mut ctx, json!({})).await.unwrap();
2456 assert_eq!(omitted["session_id"], Value::Null);
2457 assert_eq!(
2458 omitted["count"], 0,
2459 "omitted id must not combine orphans: {omitted}"
2460 );
2461 assert_eq!(omitted["turns"], json!([]));
2462
2463 let explicit = replay
2464 .call(&mut ctx, json!({"session_id": "orphan-a"}))
2465 .await
2466 .unwrap();
2467 assert_eq!(explicit["session_id"], "orphan-a");
2468 assert_eq!(explicit["count"], 1);
2469 assert_eq!(explicit["turns"][0]["content"], "orphan-a-only");
2470 }
2471
2472 #[tokio::test]
2473 async fn continuity_respects_since_filter() {
2474 let store = test_store();
2475 let sid1 = start_session(&store);
2476 record_aged_turn(&store, &sid1, 5, "ancient decision");
2477 record_aged_turn(&store, &sid1, 0, "fresh decision");
2478 let sid2 = start_session(&store);
2479
2480 let continuity = SessionContinuityTool::new(store);
2481 let mut ctx = Context::default();
2482 let cutoff = (Utc::now() - chrono::Duration::days(1))
2483 .format("%Y-%m-%d")
2484 .to_string();
2485 let v = continuity
2486 .call(
2487 &mut ctx,
2488 json!({"current_session_id": sid2, "since": cutoff, "n": 10}),
2489 )
2490 .await
2491 .unwrap();
2492 assert_eq!(v["count"], 1, "only the fresh turn is in range: {v}");
2493 assert_eq!(v["turns"][0]["content"], "fresh decision");
2494
2495 let all = continuity
2497 .call(&mut ctx, json!({"current_session_id": sid2, "n": 10}))
2498 .await
2499 .unwrap();
2500 assert_eq!(all["count"], 2);
2501 }
2502
2503 #[tokio::test]
2504 async fn continuity_skips_empty_newest_session() {
2505 let store = test_store();
2509 let sid1 = start_session(&store);
2510 record_aged_turn(
2511 &store,
2512 &sid1,
2513 0,
2514 "Decision: use SQLite for the report cache",
2515 );
2516 let sid2 = start_session(&store); let continuity = SessionContinuityTool::new(store);
2519 let mut ctx = Context::default();
2520 let v = continuity.call(&mut ctx, json!({"n": 5})).await.unwrap();
2521 assert_ne!(
2522 v["previous_session"], sid2,
2523 "empty newest session must be skipped: {v}"
2524 );
2525 assert_eq!(v["previous_session"], sid1);
2526 assert_eq!(v["count"], 1);
2527 assert!(
2528 v["turns"][0]["content"]
2529 .as_str()
2530 .unwrap()
2531 .contains("SQLite")
2532 );
2533 }
2534
2535 #[tokio::test]
2536 async fn digest_groups_by_type_and_respects_importance_floor() {
2537 let store = test_store();
2538 let sid = start_session(&store);
2539 for (turn_type, importance, content) in [
2541 ("summary", 0.6, "wrapped up"),
2542 ("decision", 0.9, "picked architecture CO over alternatives"),
2543 ("error", 0.95, "startForce root cause found"),
2544 ("breakthrough", 0.85, "watch resolution insight"),
2545 ("message", 0.3, "low-value chatter"),
2546 ] {
2547 let mut mem = Memory::new(
2548 Galaxy::Sessions,
2549 json!({
2550 "type": "session_turn",
2551 "session_id": sid,
2552 "role": "ai",
2553 "turn_type": turn_type,
2554 "importance": importance,
2555 "content": content,
2556 })
2557 .to_string(),
2558 );
2559 mem.metadata.tags = vec!["session".into(), "turn".into()];
2560 store.put(Galaxy::Sessions, &mem).unwrap();
2561 }
2562
2563 let tool = SessionDigestTool::new(store);
2564 let mut ctx = Context::default();
2565 let v = tool
2566 .call(&mut ctx, json!({"session_id": sid, "min_importance": 0.5}))
2567 .await
2568 .unwrap();
2569
2570 assert_eq!(v["status"], "success");
2571 let digest = v["digest"].as_str().unwrap();
2572 for expected in [
2574 "startForce root cause found",
2575 "picked architecture CO over alternatives",
2576 "watch resolution insight",
2577 ] {
2578 assert!(
2579 digest.contains(expected),
2580 "digest must contain '{expected}': {digest}"
2581 );
2582 }
2583 assert!(!digest.contains("low-value chatter"), "got: {digest}");
2585 let d = digest.find("## Decisions").unwrap();
2587 let b = digest.find("## Breakthroughs").unwrap();
2588 let e = digest.find("## Errors").unwrap();
2589 let s = digest.find("## Summaries").unwrap();
2590 assert!(
2591 d < b && b < e && e < s,
2592 "sections must follow canonical order: {digest}"
2593 );
2594 assert_eq!(v["turns_included"], 4);
2595
2596 let strict = tool
2598 .call(&mut ctx, json!({"session_id": sid, "min_importance": 0.9}))
2599 .await
2600 .unwrap();
2601 let strict_digest = strict["digest"].as_str().unwrap();
2602 assert!(strict_digest.contains("startForce"));
2603 assert!(!strict_digest.contains("watch resolution insight"));
2604 }
2605
2606 #[tokio::test]
2607 async fn digest_appends_checkpoint_state() {
2608 let store = test_store();
2609 let sid = start_session(&store);
2610
2611 let mut cp = Memory::new(
2613 Galaxy::Sessions,
2614 json!({
2615 "type": "checkpoint",
2616 "session_id": sid,
2617 "label": "wrap",
2618 "data": {},
2619 "handoff": {
2620 "git": {
2621 "commit": "abc1234",
2622 "branch": "main",
2623 "dirty_count": 2
2624 },
2625 "tests_green": true,
2626 "next_queue": ["first task", "second task"],
2627 "open_flags": ["flaky probe"]
2628 }
2629 })
2630 .to_string(),
2631 );
2632 cp.metadata.tags = vec!["session".into(), "checkpoint".into()];
2633 store.put(Galaxy::Sessions, &cp).unwrap();
2634
2635 let tool = SessionDigestTool::new(store);
2636 let mut ctx = Context::default();
2637 let v = tool
2638 .call(&mut ctx, json!({"session_id": sid}))
2639 .await
2640 .unwrap();
2641
2642 let digest = v["digest"].as_str().unwrap();
2643 assert!(digest.contains("## Checkpoint state"), "got: {digest}");
2644 assert!(digest.contains("abc1234"));
2645 assert!(digest.contains("first task → second task"));
2646 assert!(digest.contains("flaky probe"));
2647 assert_eq!(v["checkpoint"]["git"]["branch"], "main");
2648 }
2649
2650 #[tokio::test]
2651 async fn supersedes_hides_old_turn_until_requested() {
2652 let store = test_store();
2656 let sid = start_session(&store);
2657 let record = SessionRecordTool::new(store.clone());
2658 let mut ctx = Context::default();
2659
2660 let first = record
2661 .call(
2662 &mut ctx,
2663 json!({"role": "ai", "turn_type": "decision", "importance": 0.9,
2664 "content": "perf: 240ms", "session_id": sid}),
2665 )
2666 .await
2667 .unwrap();
2668 let old_id = first["memory_id"].as_str().unwrap().to_string();
2669
2670 let second = record
2671 .call(
2672 &mut ctx,
2673 json!({"role": "ai", "turn_type": "decision", "importance": 0.9,
2674 "content": "perf revised: 180ms after warm cache",
2675 "session_id": sid, "supersedes": old_id}),
2676 )
2677 .await
2678 .unwrap();
2679 assert_eq!(second["status"], "success");
2680 let _new_id = second["memory_id"].as_str().unwrap().to_string();
2681
2682 let replay = SessionReplayTool::new(store.clone());
2684 let v = replay
2685 .call(&mut ctx, json!({"session_id": sid}))
2686 .await
2687 .unwrap();
2688 assert_eq!(
2689 v["count"], 1,
2690 "superseded turn must be hidden by default: {v}"
2691 );
2692 assert_eq!(
2693 v["turns"][0]["content"],
2694 "perf revised: 180ms after warm cache"
2695 );
2696
2697 let with_history = replay
2699 .call(
2700 &mut ctx,
2701 json!({"session_id": sid, "include_superseded": true}),
2702 )
2703 .await
2704 .unwrap();
2705 assert_eq!(with_history["count"], 2, "got: {with_history}");
2706
2707 let sid2 = start_session(&store);
2709 let continuity = SessionContinuityTool::new(store.clone());
2710 let c = continuity
2711 .call(&mut ctx, json!({"current_session_id": sid2}))
2712 .await
2713 .unwrap();
2714 assert_eq!(c["count"], 1, "continuity must skip superseded turns: {c}");
2715 assert_eq!(
2716 c["turns"][0]["content"],
2717 "perf revised: 180ms after warm cache"
2718 );
2719
2720 let digest = SessionDigestTool::new(store);
2721 let d = digest
2722 .call(&mut ctx, json!({"session_id": sid, "min_importance": 0.5}))
2723 .await
2724 .unwrap();
2725 let digest_text = d["digest"].as_str().unwrap();
2726 assert!(digest_text.contains("180ms"), "got: {digest_text}");
2727 assert!(
2728 !digest_text.contains("240ms"),
2729 "superseded claim must not leak: {digest_text}"
2730 );
2731 }
2732
2733 #[tokio::test]
2738 async fn continuity_bounds_pathological_turn_content() {
2739 let store = test_store();
2740 let sid = start_session(&store);
2741 let huge = "A".repeat(100_000);
2742 let first = SessionRecordTool::new(store.clone())
2743 .call(
2744 &mut Context::default(),
2745 json!({"role": "ai", "content": huge, "session_id": sid}),
2746 )
2747 .await
2748 .unwrap();
2749 let turn_id = first["memory_id"].as_str().unwrap().to_string();
2750 let current = start_session(&store);
2751
2752 let continuity = SessionContinuityTool::new(store.clone());
2753 let v = continuity
2754 .call(
2755 &mut Context::default(),
2756 json!({"current_session_id": current}),
2757 )
2758 .await
2759 .unwrap();
2760
2761 assert_eq!(v["count"], 1);
2762 assert_eq!(v["truncated"], true);
2763 assert_eq!(v["turns"][0]["content_truncated"], true);
2764 assert_eq!(v["turns"][0]["content_bytes"], 100_000);
2765 assert_eq!(v["turns"][0]["memory_id"], turn_id);
2766 let cap = v["max_content_bytes"].as_u64().unwrap() as usize;
2767 assert!(v["turns"][0]["content"].as_str().unwrap().len() <= cap);
2768 let wire = serde_json::to_string(&v).unwrap();
2769 assert!(
2770 wire.len() <= CONTINUITY_DEFAULT_MAX_RESPONSE_BYTES + 2048,
2771 "response must stay near the default budget ({} bytes)",
2772 wire.len()
2773 );
2774 assert!(v["hint"].as_str().unwrap().contains("memory.read"));
2775
2776 let id = uuid::Uuid::parse_str(&turn_id).unwrap();
2778 let mem = store.get(Galaxy::Sessions, id).unwrap().unwrap();
2779 assert_eq!(mem.content.matches('A').count(), 100_000);
2780 }
2781
2782 #[tokio::test]
2783 async fn continuity_budget_keeps_newest_turns_and_reports_omissions() {
2784 let store = test_store();
2785 let sid = start_session(&store);
2786 let record = SessionRecordTool::new(store.clone());
2787 let mut ctx = Context::default();
2788 for i in 0..10 {
2789 record
2790 .call(
2791 &mut ctx,
2792 json!({
2793 "role": "ai",
2794 "content": format!("turn-{i}-{}", "x".repeat(5_000)),
2795 "session_id": sid,
2796 }),
2797 )
2798 .await
2799 .unwrap();
2800 }
2801 let current = start_session(&store);
2802 let continuity = SessionContinuityTool::new(store.clone());
2803 let v = continuity
2804 .call(
2805 &mut ctx,
2806 json!({"current_session_id": current, "max_response_bytes": 12_288}),
2807 )
2808 .await
2809 .unwrap();
2810
2811 let omitted = v["turns_omitted"].as_u64().unwrap();
2812 assert!(omitted > 0, "budget must omit older turns: {v}");
2813 assert_eq!(v["truncated"], true);
2814 let kept = v["turns"].as_array().unwrap();
2815 assert!(!kept.is_empty());
2816 let seqs: Vec<u64> = kept
2817 .iter()
2818 .map(|t| t["sequence"].as_u64().unwrap())
2819 .collect();
2820 assert_eq!(*seqs.last().unwrap(), 10, "newest turn must be kept");
2821 assert!(
2822 seqs.windows(2).all(|w| w[0] < w[1]),
2823 "turns stay chronological: {seqs:?}"
2824 );
2825 for t in kept {
2826 assert!(t["memory_id"].as_str().is_some(), "exact-read id required");
2827 }
2828 let wire = serde_json::to_string(&v).unwrap();
2829 assert!(
2830 wire.len() <= 12_288 + 2048,
2831 "response must stay near the requested budget ({} bytes)",
2832 wire.len()
2833 );
2834 }
2835
2836 #[tokio::test]
2837 async fn continuity_budget_args_are_validated() {
2838 let store = test_store();
2839 let continuity = SessionContinuityTool::new(store);
2840 let mut ctx = Context::default();
2841 let err = continuity
2842 .call(&mut ctx, json!({"max_response_bytes": 1}))
2843 .await
2844 .unwrap_err();
2845 assert!(err.to_string().contains("continuity budgets"), "got: {err}");
2846 let err = continuity
2847 .call(&mut ctx, json!({"max_content_bytes": "big"}))
2848 .await
2849 .unwrap_err();
2850 assert!(err.to_string().contains("must be an integer"), "got: {err}");
2851 }
2852
2853 #[tokio::test]
2856 async fn record_rejects_unknown_session_ids() {
2857 let store = test_store();
2858 let record = SessionRecordTool::new(store.clone());
2859 let mut ctx = Context::default();
2860
2861 let err = record
2862 .call(
2863 &mut ctx,
2864 json!({"content": "orphan turn probe",
2865 "session_id": "00000000-0000-0000-0000-000000000999"}),
2866 )
2867 .await
2868 .unwrap_err();
2869 assert!(
2870 err.to_string().contains("no session found"),
2871 "unknown id must be refused: {err}"
2872 );
2873 assert!(
2874 store.scan_all(Galaxy::Sessions).unwrap().is_empty(),
2875 "refused record must not write anything"
2876 );
2877
2878 let err = record
2879 .call(
2880 &mut ctx,
2881 json!({"content": "x", "session_id": "not-a-uuid"}),
2882 )
2883 .await
2884 .unwrap_err();
2885 assert!(
2886 err.to_string().contains("invalid session_id"),
2887 "malformed id must be a caller error: {err}"
2888 );
2889
2890 let sid = start_session(&store);
2892 let ok = record
2893 .call(&mut ctx, json!({"content": "real turn", "session_id": sid}))
2894 .await
2895 .unwrap();
2896 assert_eq!(ok["status"], "success");
2897 }
2898
2899 #[tokio::test]
2900 async fn export_import_roundtrip_preserves_history() {
2901 let store_a = test_store();
2904 let sid = start_session(&store_a);
2905 let record = SessionRecordTool::new(store_a.clone());
2906 let mut ctx = Context::default();
2907 record_aged_turn(&store_a, &sid, 2, "day-one decision");
2908 record_aged_turn(&store_a, &sid, 1, "day-two decision");
2909 let first = record
2910 .call(
2911 &mut ctx,
2912 json!({"role": "ai", "turn_type": "decision", "importance": 0.9,
2913 "content": "original claim", "session_id": sid}),
2914 )
2915 .await
2916 .unwrap();
2917 record
2918 .call(
2919 &mut ctx,
2920 json!({"role": "ai", "turn_type": "decision", "importance": 0.9,
2921 "content": "corrected claim",
2922 "session_id": sid,
2923 "supersedes": first["memory_id"].as_str().unwrap()}),
2924 )
2925 .await
2926 .unwrap();
2927
2928 let export = SessionExportTool::new(store_a.clone());
2930 {
2932 let mut marker: Memory = store_a
2933 .scan_all(Galaxy::Sessions)
2934 .unwrap()
2935 .into_iter()
2936 .find(|m| m.metadata.tags.contains(&"start".to_string()))
2937 .unwrap();
2938 marker.metadata.title = Some("The Big Decision".to_string());
2939 marker.metadata.topic = Some("v8-slices".to_string());
2940 store_a.put(Galaxy::Sessions, &marker).unwrap();
2941 }
2942 let exported = export
2943 .call(&mut ctx, json!({"session_id": sid}))
2944 .await
2945 .unwrap();
2946 assert_eq!(exported["status"], "success");
2947 let jsonl = exported["jsonl"].as_str().unwrap();
2948 assert_eq!(
2950 exported["records"], 5,
2951 "start + 2 aged + 2 claims: {exported}"
2952 );
2953 assert_eq!(jsonl.lines().count(), 6);
2955 let header_line = jsonl.lines().next().unwrap();
2956 match wm_memory::envelope::read_header_line(header_line) {
2957 wm_memory::envelope::HeaderRead::Header(h) => {
2958 assert_eq!(h.kind, "session_export");
2959 assert_eq!(h.count, 5);
2960 }
2961 other => panic!("first line must be the envelope header, got {other:?}"),
2962 }
2963
2964 let store_b = test_store();
2966 let import = SessionImportTool::new(store_b.clone(), None);
2967 let imported = import
2968 .call(&mut ctx, json!({"jsonl": jsonl}))
2969 .await
2970 .unwrap();
2971 assert_eq!(imported["imported"], 5, "got: {imported}");
2972 assert_eq!(imported["skipped"], 0);
2973 assert_eq!(imported["session_ids"], json!([sid]));
2974 assert_eq!(imported["envelope"]["format_version"], 2);
2976 assert_eq!(imported["envelope"]["declared_count"], 5);
2977 assert_eq!(imported["warnings"], json!([]));
2978 let marker_b = store_b
2980 .scan_all(Galaxy::Sessions)
2981 .unwrap()
2982 .into_iter()
2983 .find(|m| m.metadata.tags.contains(&"start".to_string()))
2984 .unwrap();
2985 assert_eq!(marker_b.metadata.title.as_deref(), Some("The Big Decision"));
2986 assert_eq!(marker_b.metadata.topic.as_deref(), Some("v8-slices"));
2987
2988 let replay_b = SessionReplayTool::new(store_b.clone());
2990 let v = replay_b
2991 .call(&mut ctx, json!({"session_id": sid}))
2992 .await
2993 .unwrap();
2994 assert_eq!(v["count"], 3, "two aged turns + correction: {v}");
2995 let contents: Vec<&str> = v["turns"]
2996 .as_array()
2997 .unwrap()
2998 .iter()
2999 .filter_map(|t| t["content"].as_str())
3000 .collect();
3001 assert!(contents.contains(&"day-one decision"));
3002 assert!(contents.contains(&"corrected claim"));
3003 assert!(!contents.contains(&"original claim"));
3004
3005 let full = replay_b
3007 .call(
3008 &mut ctx,
3009 json!({"session_id": sid, "include_superseded": true}),
3010 )
3011 .await
3012 .unwrap();
3013 assert_eq!(full["count"], 4);
3014
3015 let new_sid = start_session(&store_b);
3017 let continuity = SessionContinuityTool::new(store_b);
3018 let c = continuity
3019 .call(
3020 &mut ctx,
3021 json!({"current_session_id": new_sid, "since":
3022 (Utc::now() - chrono::Duration::days(3)).format("%Y-%m-%d").to_string()}),
3023 )
3024 .await
3025 .unwrap();
3026 assert_eq!(
3027 c["previous_session"], sid,
3028 "import must preserve created_at so recency resolution works"
3029 );
3030 assert_eq!(c["count"], 3);
3031 }
3032
3033 #[tokio::test]
3034 async fn import_rejects_missing_payload() {
3035 let store = test_store();
3036 let tool = SessionImportTool::new(store, None);
3037 let mut ctx = Context::default();
3038 assert!(tool.call(&mut ctx, json!({})).await.is_err());
3039 }
3040
3041 #[tokio::test]
3042 async fn import_refuses_newer_envelope_format() {
3043 let store = test_store();
3044 let mut ctx = Context::default();
3045 let header = wm_memory::envelope::EnvelopeHeader {
3046 format_version: wm_memory::envelope::ENVELOPE_FORMAT_VERSION + 1,
3047 kind: "session_export".into(),
3048 created_at: chrono::Utc::now().to_rfc3339(),
3049 count: 1,
3050 generator: "wm 99.0.0".into(),
3051 };
3052 let record =
3053 serde_json::to_string(&Memory::new(Galaxy::Sessions, "future".into())).unwrap();
3054 let payload = format!("{}\n{record}\n", header.header_line());
3055 let tool = SessionImportTool::new(store, None);
3056 let result = tool.call(&mut ctx, json!({"jsonl": payload})).await;
3057 let err = format!("{:?}", result.unwrap_err());
3058 assert!(err.contains("newer than this build supports"), "{err}");
3059 }
3060
3061 #[tokio::test]
3065 async fn session_record_indexes_at_write_time() {
3066 let dir = tempfile::tempdir().unwrap();
3067 let lmdb = dir.path().join("lmdb");
3068 std::fs::create_dir_all(&lmdb).unwrap();
3069 let store = Arc::new(MemoryStore::open_default(&lmdb).unwrap());
3070 let tantivy = dir.path().join("tantivy");
3071 std::fs::create_dir_all(&tantivy).unwrap();
3072 let search = Arc::new(wm_memory::SearchEngine::open(&tantivy).unwrap());
3073
3074 let sid = start_session(&store);
3075 let mut ctx = Context::default();
3076 SessionRecordTool::new(store.clone())
3077 .with_search(Some(search.clone()))
3078 .call(
3079 &mut ctx,
3080 json!({"role": "ai", "turn_type": "decision", "importance": 0.7,
3081 "content": "amber lighthouse protocol engaged", "session_id": sid}),
3082 )
3083 .await
3084 .unwrap();
3085
3086 let docs = search.count_docs_in_galaxy("sessions").unwrap();
3087 assert!(
3088 docs >= 1,
3089 "session.record must index its write immediately (docs={docs})"
3090 );
3091 }
3092
3093 #[tokio::test]
3097 async fn session_record_flags_instruction_shaped_content() {
3098 let store = test_store();
3099 let sid = start_session(&store);
3100 let mut ctx = Context::default();
3101
3102 let clean = SessionRecordTool::new(store.clone())
3103 .call(
3104 &mut ctx,
3105 json!({"role": "ai", "content": "ordinary status update", "session_id": sid}),
3106 )
3107 .await
3108 .unwrap();
3109 assert!(clean.get("warnings").is_none(), "clean turn: {clean}");
3110
3111 let flagged = SessionRecordTool::new(store.clone())
3112 .call(
3113 &mut ctx,
3114 json!({"role": "ai",
3115 "content": "incident note: review the jailbreak attempt before retrying",
3116 "session_id": sid}),
3117 )
3118 .await
3119 .unwrap();
3120 assert_eq!(flagged["status"], "success", "flag, not refusal: {flagged}");
3121 let warnings = flagged["warnings"].as_array().unwrap();
3122 assert!(
3123 warnings[0].as_str().unwrap().contains("instruction-shaped"),
3124 "got: {warnings:?}"
3125 );
3126 }
3127
3128 #[tokio::test]
3131 async fn record_and_continuity_write_savings_ledger() {
3132 let dir = tempfile::tempdir().unwrap();
3136 let lmdb = dir.path().join("lmdb");
3137 std::fs::create_dir_all(&lmdb).unwrap();
3138 let store = Arc::new(MemoryStore::open_default(&lmdb).unwrap());
3139 let sid = start_session(&store);
3140 let mut ctx = Context::default();
3141 SessionRecordTool::new(store.clone())
3142 .call(
3143 &mut ctx,
3144 json!({"role": "ai", "content": "x".repeat(2000), "session_id": sid}),
3145 )
3146 .await
3147 .unwrap();
3148 SessionContinuityTool::new(store.clone())
3149 .call(
3150 &mut ctx,
3151 json!({"current_session_id": "00000000-0000-4000-8000-000000000000"}),
3152 )
3153 .await
3154 .unwrap();
3155
3156 let path = store.path().join("savings_ledger.jsonl");
3157 let text = std::fs::read_to_string(&path).unwrap();
3158 let rows: Vec<serde_json::Value> = text
3159 .lines()
3160 .map(|l| serde_json::from_str(l).unwrap())
3161 .collect();
3162 assert_eq!(
3163 rows.len(),
3164 2,
3165 "one record row + one continuity row: {rows:?}"
3166 );
3167 assert_eq!(rows[0]["op"], "record");
3168 assert_eq!(rows[0]["bytes_stored"], 2000);
3169 assert_eq!(rows[1]["op"], "continuity");
3170 assert!(
3171 rows[1]["bytes_available"].as_u64().unwrap() >= 2000,
3172 "stored state must be counted: {rows:?}"
3173 );
3174 assert!(
3175 rows[1]["bytes_injected"].as_u64().unwrap() > 0,
3176 "injected envelope must be counted: {rows:?}"
3177 );
3178 }
3179
3180 #[tokio::test]
3184 async fn import_indexes_tantivy_no_drift_even_on_reimport() {
3185 let dir = tempfile::tempdir().unwrap();
3186 let lmdb = dir.path().join("lmdb");
3187 std::fs::create_dir_all(&lmdb).unwrap();
3188 let store = Arc::new(MemoryStore::open_default(&lmdb).unwrap());
3189 let tantivy = dir.path().join("tantivy");
3190 std::fs::create_dir_all(&tantivy).unwrap();
3191 let search = Arc::new(wm_memory::SearchEngine::open(&tantivy).unwrap());
3192
3193 let store_a = test_store();
3195 let sid = start_session(&store_a);
3196 let mut ctx = Context::default();
3197 SessionRecordTool::new(store_a.clone())
3198 .call(
3199 &mut ctx,
3200 json!({"role": "ai", "turn_type": "decision", "importance": 0.9,
3201 "content": "kumquat governance ratchet engaged", "session_id": sid}),
3202 )
3203 .await
3204 .unwrap();
3205 let exported = SessionExportTool::new(store_a.clone())
3206 .call(&mut ctx, json!({"session_id": sid}))
3207 .await
3208 .unwrap();
3209 let jsonl = exported["jsonl"].as_str().unwrap().to_string();
3210
3211 let import = SessionImportTool::new(store.clone(), Some(search.clone()));
3213 for round in 1..=2 {
3214 let r = import
3215 .call(&mut ctx, json!({"jsonl": jsonl}))
3216 .await
3217 .unwrap();
3218 assert_eq!(r["status"], "success", "round {round}: {r}");
3219 assert_eq!(r["skipped"], 0);
3220 assert_eq!(
3221 r["indexed"], r["imported"],
3222 "round {round}: every record indexed: {r}"
3223 );
3224 }
3225
3226 let report = wm_memory::reindex::check_consistency(&store, &search);
3229 let drifted: Vec<_> = report
3230 .galaxies
3231 .iter()
3232 .filter(|g| g.drift)
3233 .map(|g| g.galaxy.clone())
3234 .collect();
3235 assert!(
3236 drifted.is_empty(),
3237 "import must leave zero index drift, drifted: {drifted:?}"
3238 );
3239
3240 let needle_id = store
3242 .scan_all(Galaxy::Sessions)
3243 .unwrap()
3244 .iter()
3245 .find(|m| m.content.contains("kumquat"))
3246 .unwrap()
3247 .metadata
3248 .id
3249 .to_string();
3250 let hits = search.search("kumquat governance ratchet", 10).unwrap();
3251 assert!(
3252 hits.iter().any(|h| h.memory_id == needle_id),
3253 "imported record must be searchable via the index: {hits:?}"
3254 );
3255 }
3256
3257 #[tokio::test]
3258 async fn record_then_replay_full() {
3259 let store = test_store();
3260 let sid = start_session(&store);
3261 let record = SessionRecordTool::new(store.clone());
3262 let mut ctx = Context::default();
3263 for (i, role) in [("user", "hello"), ("ai", "hi there")].iter().enumerate() {
3264 let r = record
3265 .call(
3266 &mut ctx,
3267 json!({"role": role.0, "content": role.1, "session_id": sid}),
3268 )
3269 .await
3270 .unwrap();
3271 assert_eq!(r["sequence"], i as u64 + 1);
3272 }
3273
3274 let replay = SessionReplayTool::new(store.clone());
3275 let v = replay
3276 .call(&mut ctx, json!({"mode": "full", "session_id": sid}))
3277 .await
3278 .unwrap();
3279 assert_eq!(v["count"], 2);
3280 assert_eq!(v["turns"][0]["content"], "hello");
3281 assert_eq!(v["turns"][1]["role"], "ai");
3282 }
3283
3284 #[tokio::test(flavor = "multi_thread", worker_threads = 8)]
3288 async fn concurrent_record_writers_get_unique_contiguous_sequences() {
3289 const N: usize = 100;
3290 let store = test_store();
3291 let sid = start_session(&store);
3292 let tool = Arc::new(SessionRecordTool::new(store));
3293
3294 let mut handles = Vec::with_capacity(N);
3295 for i in 0..N {
3296 let tool = Arc::clone(&tool);
3297 let sid = sid.clone();
3298 handles.push(tokio::spawn(async move {
3299 let mut ctx = Context::default();
3300 let v = tool
3301 .call(
3302 &mut ctx,
3303 json!({
3304 "session_id": sid,
3305 "role": "ai",
3306 "turn_type": "message",
3307 "content": format!("concurrent turn {i}"),
3308 }),
3309 )
3310 .await
3311 .unwrap();
3312 v["sequence"].as_u64().expect("sequence in response")
3313 }));
3314 }
3315
3316 let mut sequences = Vec::with_capacity(N);
3317 for handle in handles {
3318 sequences.push(handle.await.unwrap());
3319 }
3320 sequences.sort_unstable();
3321 assert_eq!(
3322 sequences,
3323 (1..=N as u64).collect::<Vec<_>>(),
3324 "{N} concurrent session.record writers must get 1..={N}"
3325 );
3326 }
3327
3328 #[tokio::test]
3332 async fn supersede_marks_old_turn_and_hides_it_by_default() {
3333 let store = test_store();
3334 let sid = start_session(&store);
3335 let record = SessionRecordTool::new(store.clone());
3336 let mut ctx = Context::default();
3337
3338 let first = record
3339 .call(
3340 &mut ctx,
3341 json!({"session_id": sid, "content": "we chose A"}),
3342 )
3343 .await
3344 .unwrap();
3345 let first_id = first["memory_id"].as_str().unwrap().to_string();
3346 let second = record
3347 .call(
3348 &mut ctx,
3349 json!({"session_id": sid, "content": "we chose B", "supersedes": first_id}),
3350 )
3351 .await
3352 .unwrap();
3353 assert_eq!(second["sequence"], 2);
3354
3355 let old = store
3357 .get(
3358 wm_core::Galaxy::Sessions,
3359 uuid::Uuid::parse_str(&first_id).unwrap(),
3360 )
3361 .unwrap()
3362 .unwrap();
3363 assert!(
3364 old.metadata
3365 .tags
3366 .iter()
3367 .any(|t| t.starts_with("superseded-by:")),
3368 "old turn must carry superseded-by: {:?}",
3369 old.metadata.tags
3370 );
3371
3372 let replay = SessionReplayTool::new(store.clone());
3373 let visible = replay
3374 .call(&mut ctx, json!({"mode": "full", "session_id": sid}))
3375 .await
3376 .unwrap();
3377 assert_eq!(visible["count"], 1, "default view shows the current story");
3378 assert_eq!(visible["turns"][0]["content"], "we chose B");
3379
3380 let full = replay
3381 .call(
3382 &mut ctx,
3383 json!({"mode": "full", "session_id": sid, "include_superseded": true}),
3384 )
3385 .await
3386 .unwrap();
3387 assert_eq!(full["count"], 2, "history stays queryable");
3388 }
3389
3390 #[tokio::test]
3391 async fn record_requires_content_and_valid_role() {
3392 let store = test_store();
3393 let sid = start_session(&store);
3394 let tool = SessionRecordTool::new(store);
3395 let mut ctx = Context::default();
3396 assert!(
3397 tool.call(&mut ctx, json!({"role": "system", "content": "x"}))
3398 .await
3399 .is_err()
3400 );
3401 for bad in [json!(""), json!(" "), json!("\n\t")] {
3404 let err = tool
3405 .call(&mut ctx, json!({"content": bad, "session_id": sid}))
3406 .await
3407 .unwrap_err();
3408 assert!(
3409 err.to_string().contains("content"),
3410 "content={bad:?}: {err}"
3411 );
3412 }
3413 let err = tool
3414 .call(
3415 &mut ctx,
3416 json!({"content": "x", "session_id": sid, "turn_type": "observation"}),
3417 )
3418 .await
3419 .unwrap_err();
3420 assert!(err.to_string().contains("turn_type"), "{err}");
3421 let v = tool
3423 .call(
3424 &mut ctx,
3425 json!({"content": "x", "session_id": sid, "turn_type": "context"}),
3426 )
3427 .await
3428 .unwrap();
3429 assert_eq!(v["status"], "success", "{v}");
3430 }
3431
3432 #[tokio::test]
3436 async fn record_rejects_out_of_range_importance() {
3437 let store = test_store();
3438 let sid = start_session(&store);
3439 let record = SessionRecordTool::new(store);
3440 let mut ctx = Context::default();
3441 for bad in [json!(1.5), json!(999), json!(-0.25), json!("2.0")] {
3442 let err = record
3443 .call(
3444 &mut ctx,
3445 json!({"content": "x", "session_id": sid, "importance": bad}),
3446 )
3447 .await
3448 .unwrap_err();
3449 assert!(
3450 err.to_string().contains("importance"),
3451 "importance={bad} must be rejected, got: {err}"
3452 );
3453 }
3454 for good in [json!(0.0), json!(1.0), json!(0.75), json!("0.4")] {
3457 let v = record
3458 .call(
3459 &mut ctx,
3460 json!({"content": "x", "session_id": sid, "importance": good}),
3461 )
3462 .await
3463 .unwrap();
3464 assert_eq!(v["status"], "success", "{v}");
3465 }
3466 }
3467
3468 #[tokio::test]
3472 async fn record_stamps_provenance_from_role() {
3473 let store = test_store();
3474 let sid = start_session(&store);
3475 let record = SessionRecordTool::new(store.clone());
3476 let mut ctx = Context::default();
3477 let ai = record
3478 .call(
3479 &mut ctx,
3480 json!({"role": "ai", "content": "agent turn", "session_id": sid}),
3481 )
3482 .await
3483 .unwrap();
3484 let user = record
3485 .call(
3486 &mut ctx,
3487 json!({"role": "user", "content": "human turn", "session_id": sid}),
3488 )
3489 .await
3490 .unwrap();
3491
3492 let ai_mem = store
3493 .get(
3494 Galaxy::Sessions,
3495 uuid::Uuid::parse_str(ai["memory_id"].as_str().unwrap()).unwrap(),
3496 )
3497 .expect("ai turn stored")
3498 .expect("ai turn present");
3499 assert_eq!(ai_mem.metadata.source, "agent");
3500 assert!((ai_mem.metadata.source_trust - 0.7).abs() < 1e-5);
3501
3502 let user_mem = store
3503 .get(
3504 Galaxy::Sessions,
3505 uuid::Uuid::parse_str(user["memory_id"].as_str().unwrap()).unwrap(),
3506 )
3507 .expect("user turn stored")
3508 .expect("user turn present");
3509 assert_eq!(user_mem.metadata.source, "user");
3510 assert!((user_mem.metadata.source_trust - 1.0).abs() < f32::EPSILON);
3511 }
3512
3513 #[tokio::test]
3514 async fn record_rejects_invalid_track_slug() {
3515 let store = test_store();
3516 let sid = start_session(&store);
3517 let record = SessionRecordTool::new(store);
3518 let mut ctx = Context::default();
3519 let mut bad = vec![
3520 String::new(),
3521 "Track-A".to_string(),
3522 "has space".to_string(),
3523 "-leading".to_string(),
3524 ];
3525 bad.push("x".repeat(TRACK_MAX_LEN + 1));
3526 for bad in &bad {
3527 let err = record
3528 .call(
3529 &mut ctx,
3530 json!({"content": "x", "session_id": sid, "track": bad}),
3531 )
3532 .await
3533 .unwrap_err()
3534 .to_string();
3535 assert!(err.contains("invalid track"), "slug {bad:?}: {err}");
3536 }
3537 let ok = record
3539 .call(
3540 &mut ctx,
3541 json!({"content": "x", "session_id": sid, "track": "wmv9/harness-2.1_a"}),
3542 )
3543 .await
3544 .unwrap();
3545 assert_eq!(ok["status"], "success");
3546 }
3547
3548 #[tokio::test]
3549 async fn track_log_reads_recorded_turns_per_track() {
3550 let store = test_store();
3551 let sid = start_session(&store);
3552 let record = SessionRecordTool::new(store.clone());
3553 let mut ctx = Context::default();
3554 for (track, content) in [
3555 ("harness-2", "harness: schema check landed"),
3556 ("safety-a", "safety: mask gate landed"),
3557 ("harness-2", "harness: manifest regen done"),
3558 ] {
3559 record
3560 .call(
3561 &mut ctx,
3562 json!({
3563 "role": "ai",
3564 "content": content,
3565 "session_id": sid,
3566 "track": track,
3567 "turn_type": "summary",
3568 }),
3569 )
3570 .await
3571 .unwrap();
3572 }
3573
3574 let log = SessionTrackLogTool::new(store);
3575 let v = log
3576 .call(&mut ctx, json!({"track": "harness-2"}))
3577 .await
3578 .unwrap();
3579 assert_eq!(v["mode"], "log");
3580 assert_eq!(v["tracks"][0]["track"], "harness-2");
3581 assert_eq!(v["tracks"][0]["entry_count"], 2);
3582 assert_eq!(v["tracks"][0]["returned"], 2);
3583 assert_eq!(
3584 v["tracks"][0]["entries"][0]["content"],
3585 "harness: schema check landed"
3586 );
3587 assert_eq!(
3588 v["tracks"][0]["entries"][1]["content"],
3589 "harness: manifest regen done"
3590 );
3591 assert!(v["tracks"][0]["age_seconds"].as_i64().unwrap() >= 0);
3592
3593 let all = log.call(&mut ctx, json!({})).await.unwrap();
3595 assert_eq!(all["mode"], "overview");
3596 assert_eq!(all["track_count"], 2);
3597 let names: Vec<&str> = all["tracks"]
3598 .as_array()
3599 .unwrap()
3600 .iter()
3601 .map(|t| t["track"].as_str().unwrap())
3602 .collect();
3603 assert_eq!(names, vec!["harness-2", "safety-a"]);
3604 assert!(all["tracks"][0]["entries"].is_null());
3605 assert_eq!(all["tracks"][0]["latest_entry"]["turn_type"], "summary");
3606 assert!(
3607 all["tracks"][0]["latest_entry"]["preview"]
3608 .as_str()
3609 .unwrap()
3610 .contains("manifest regen")
3611 );
3612
3613 let unknown = log
3615 .call(&mut ctx, json!({"track": "never-seen"}))
3616 .await
3617 .unwrap();
3618 assert_eq!(unknown["tracks"][0]["known"], false);
3619 assert_eq!(unknown["tracks"][0]["entry_count"], 0);
3620 }
3621
3622 #[tokio::test]
3623 async fn track_log_merges_related_tracks_and_filters_time() {
3624 let store = test_store();
3625 let sid = start_session(&store);
3626 let record = SessionRecordTool::new(store.clone());
3627 let mut ctx = Context::default();
3628 for (track, content) in [
3629 ("harness-2", "old harness note"),
3630 ("safety-a", "safety note"),
3631 ("harness-2", "new harness note"),
3632 ] {
3633 record
3634 .call(
3635 &mut ctx,
3636 json!({"role": "ai", "content": content, "session_id": sid, "track": track}),
3637 )
3638 .await
3639 .unwrap();
3640 }
3641 for mut mem in store.scan_all(Galaxy::Sessions).unwrap() {
3643 if mem.content.contains("old harness note") {
3644 mem.metadata.created_at = Utc::now() - chrono::Duration::days(3);
3645 store.put(Galaxy::Sessions, &mem).unwrap();
3646 }
3647 }
3648
3649 let log = SessionTrackLogTool::new(store);
3650 let cutoff = (Utc::now() - chrono::Duration::days(1))
3651 .format("%Y-%m-%d")
3652 .to_string();
3653 let v = log
3654 .call(
3655 &mut ctx,
3656 json!({"tracks": ["harness-2", "safety-a"], "since": cutoff}),
3657 )
3658 .await
3659 .unwrap();
3660 assert_eq!(v["track_count"], 2);
3661 let harness = &v["tracks"][0];
3662 assert_eq!(harness["track"], "harness-2");
3663 assert_eq!(harness["entry_count"], 1, "aged turn filtered: {v}");
3664 assert_eq!(harness["entries"][0]["content"], "new harness note");
3665 let safety = &v["tracks"][1];
3666 assert_eq!(safety["track"], "safety-a");
3667 assert_eq!(safety["entry_count"], 1);
3668
3669 let err = log
3671 .call(&mut ctx, json!({"tracks": ["harness-2", "Bad Slug"]}))
3672 .await
3673 .unwrap_err()
3674 .to_string();
3675 assert!(err.contains("invalid track"), "{err}");
3676 }
3677
3678 #[tokio::test]
3679 async fn track_log_hides_superseded_turns_by_default() {
3680 let store = test_store();
3681 let sid = start_session(&store);
3682 let record = SessionRecordTool::new(store.clone());
3683 let mut ctx = Context::default();
3684 let first = record
3685 .call(
3686 &mut ctx,
3687 json!({"role": "ai", "content": "v1 estimate", "session_id": sid,
3688 "track": "bench", "turn_type": "summary"}),
3689 )
3690 .await
3691 .unwrap();
3692 record
3693 .call(
3694 &mut ctx,
3695 json!({"role": "ai", "content": "v2 estimate", "session_id": sid,
3696 "track": "bench", "turn_type": "summary",
3697 "supersedes": first["memory_id"]}),
3698 )
3699 .await
3700 .unwrap();
3701
3702 let log = SessionTrackLogTool::new(store);
3703 let v = log.call(&mut ctx, json!({"track": "bench"})).await.unwrap();
3704 assert_eq!(v["tracks"][0]["entry_count"], 1);
3705 assert_eq!(v["tracks"][0]["entries"][0]["content"], "v2 estimate");
3706 let all = log
3707 .call(
3708 &mut ctx,
3709 json!({"track": "bench", "include_superseded": true}),
3710 )
3711 .await
3712 .unwrap();
3713 assert_eq!(all["tracks"][0]["entry_count"], 2);
3714 }
3715
3716 #[tokio::test]
3717 async fn checkpoint_track_lands_in_track_log() {
3718 let store = test_store();
3719 let sid = start_session(&store);
3720 let checkpoint = SessionCheckpointNodiscoveryTool::new(store.clone());
3721 let mut ctx = Context::default();
3722 let v = checkpoint
3723 .call(
3724 &mut ctx,
3725 json!({
3726 "session_id": sid,
3727 "track": "harness-2",
3728 "label": "slice2-done",
3729 "next_queue": ["regen manifest", "push after ceremony"],
3730 "open_flags": ["none"],
3731 }),
3732 )
3733 .await
3734 .unwrap();
3735 assert_eq!(v["status"], "success");
3736
3737 let log = SessionTrackLogTool::new(store);
3738 let out = log
3739 .call(&mut ctx, json!({"track": "harness-2"}))
3740 .await
3741 .unwrap();
3742 assert_eq!(out["tracks"][0]["entry_count"], 1);
3743 assert_eq!(out["tracks"][0]["entries"][0]["kind"], "checkpoint");
3744 assert_eq!(
3745 out["tracks"][0]["latest_checkpoint"]["entry"]["label"],
3746 "slice2-done"
3747 );
3748 assert_eq!(
3749 out["tracks"][0]["latest_checkpoint"]["entry"]["handoff"]["next_queue"][0],
3750 "regen manifest"
3751 );
3752
3753 let err = checkpoint
3755 .call(&mut ctx, json!({"session_id": sid, "track": "Bad"}))
3756 .await
3757 .unwrap_err()
3758 .to_string();
3759 assert!(err.contains("invalid track"), "{err}");
3760 }
3761
3762 #[tokio::test]
3763 async fn continuity_returns_previous_session_tail() {
3764 let store = test_store();
3765 let sid1 = start_session(&store);
3766 let record = SessionRecordTool::new(store.clone());
3767 let mut ctx = Context::default();
3768 for i in 0..5 {
3769 record
3770 .call(
3771 &mut ctx,
3772 json!({"role": "user", "content": format!("turn {i}"), "session_id": sid1}),
3773 )
3774 .await
3775 .unwrap();
3776 }
3777 let sid2 = start_session(&store);
3778 let continuity = SessionContinuityTool::new(store);
3779 let v = continuity
3780 .call(&mut ctx, json!({"current_session_id": sid2, "n": 2}))
3781 .await
3782 .unwrap();
3783 assert_eq!(v["previous_session"], sid1);
3784 assert_eq!(v["count"], 2);
3785 assert_eq!(v["turns"][1]["content"], "turn 4");
3786 assert!(
3787 v.get("hint").is_none(),
3788 "non-empty continuity must not carry the scoping hint"
3789 );
3790 assert!(
3791 v["checkpoint"].is_null() && v["checkpoint_id"].is_null(),
3792 "no checkpoint recorded — the fields must be present but null: {v}"
3793 );
3794 }
3795
3796 #[tokio::test]
3797 async fn continuity_surfaces_latest_checkpoint_handoff() {
3798 let store = test_store();
3804 let sid1 = start_session(&store);
3805 let record = SessionRecordTool::new(store.clone());
3806 let mut ctx = Context::default();
3807 record
3808 .call(
3809 &mut ctx,
3810 json!({"role": "ai", "turn_type": "summary", "importance": 0.9,
3811 "content": "seam work done", "session_id": sid1}),
3812 )
3813 .await
3814 .unwrap();
3815
3816 let seed_checkpoint = |handoff: Value, age_secs: i64| {
3817 let mut cp = Memory::new(
3818 Galaxy::Sessions,
3819 json!({
3820 "type": "checkpoint",
3821 "session_id": sid1,
3822 "label": "wrap",
3823 "data": {},
3824 "handoff": handoff,
3825 })
3826 .to_string(),
3827 );
3828 cp.metadata.tags = vec!["session".into(), "checkpoint".into()];
3829 cp.metadata.created_at = Utc::now() - chrono::Duration::seconds(age_secs);
3830 store.put(Galaxy::Sessions, &cp).unwrap();
3831 cp.metadata.id.to_string()
3832 };
3833 seed_checkpoint(
3834 json!({"next_queue": ["stale task"], "open_flags": ["stale flag"]}),
3835 120,
3836 );
3837 let newest_id = seed_checkpoint(
3838 json!({
3839 "git": {"commit": "abc1234", "branch": "main", "dirty_count": 0},
3840 "tests_green": true,
3841 "next_queue": ["fuzz malformed headers", "verify seal"],
3842 "open_flags": ["header length cap unresolved"],
3843 }),
3844 0,
3845 );
3846
3847 let sid2 = start_session(&store);
3848 let continuity = SessionContinuityTool::new(store);
3849 let v = continuity
3850 .call(&mut ctx, json!({"current_session_id": sid2}))
3851 .await
3852 .unwrap();
3853 assert_eq!(v["previous_session"], sid1);
3854 assert_eq!(
3855 v["checkpoint"]["next_queue"][0], "fuzz malformed headers",
3856 "the latest checkpoint must win over an older one: {v}"
3857 );
3858 assert_eq!(
3859 v["checkpoint"]["open_flags"][0], "header length cap unresolved",
3860 "open flags must survive the handoff through continuity: {v}"
3861 );
3862 assert_eq!(v["checkpoint"]["git"]["commit"], "abc1234");
3863 assert_eq!(v["checkpoint"]["tests_green"], true);
3864 assert_eq!(v["checkpoint_id"].as_str().unwrap(), newest_id);
3865 }
3866
3867 #[tokio::test]
3868 async fn continuity_empty_store_discloses_project_scoping() {
3869 let store = test_store();
3874 let continuity = SessionContinuityTool::new(store);
3875 let mut ctx = Context::default();
3876 let v = continuity.call(&mut ctx, json!({})).await.unwrap();
3877
3878 assert_eq!(v["status"], "success");
3879 assert_eq!(v["count"], 0);
3880 let hint = v["hint"].as_str().expect("hint present on empty store");
3881 assert!(hint.contains("project-scoped"), "got: {hint}");
3882 assert!(hint.contains("opencode config"), "got: {hint}");
3883 assert!(hint.contains("GET /status"), "got: {hint}");
3884 assert!(hint.contains("store "), "got: {hint}");
3886 }
3887
3888 #[tokio::test]
3889 async fn record_defaults_to_newest_start_by_time_not_key_order() {
3890 let store = test_store();
3895 for age in [50_000, 40_000, 30_000, 20_000, 10_000] {
3896 start_session_aged(&store, age);
3897 }
3898 let newest = start_session_aged(&store, 0);
3899
3900 let record = SessionRecordTool::new(store);
3901 let mut ctx = Context::default();
3902 let r = record
3903 .call(&mut ctx, json!({"role": "ai", "content": "latest turn"}))
3904 .await
3905 .unwrap();
3906 assert_eq!(
3907 r["session_id"], newest,
3908 "record without explicit session_id must target the newest start by created_at"
3909 );
3910 }
3911
3912 #[tokio::test]
3913 async fn continuity_picks_newest_prior_by_time_not_key_order() {
3914 let store = test_store();
3918 for age in [40_000, 30_000, 20_000] {
3919 start_session_aged(&store, age);
3920 }
3921 let newest_prior = start_session_aged(&store, 10);
3922 let current = start_session_aged(&store, 0);
3923
3924 let continuity = SessionContinuityTool::new(store);
3925 let mut ctx = Context::default();
3926 let v = continuity
3927 .call(&mut ctx, json!({"current_session_id": current, "n": 1}))
3928 .await
3929 .unwrap();
3930 assert_eq!(
3931 v["previous_session"], newest_prior,
3932 "continuity must select the newest prior session by created_at"
3933 );
3934 }
3935
3936 #[tokio::test]
3937 async fn handoff_transfer_accept_list() {
3938 let store = test_store();
3939 let sid = start_session(&store);
3940 let record = SessionRecordTool::new(store.clone());
3941 let mut ctx = Context::default();
3942 record
3943 .call(
3944 &mut ctx,
3945 json!({"role": "ai", "content": "context", "session_id": sid}),
3946 )
3947 .await
3948 .unwrap();
3949
3950 let handoff = SessionHandoffTool::new(store.clone());
3951 let t = handoff
3952 .call(
3953 &mut ctx,
3954 json!({"action": "transfer", "session_id": sid, "message": "take over"}),
3955 )
3956 .await
3957 .unwrap();
3958 assert_eq!(t["status"], "success");
3959 let hid = t["handoff_id"].as_str().unwrap().to_string();
3960
3961 let list = handoff
3962 .call(&mut ctx, json!({"action": "list"}))
3963 .await
3964 .unwrap();
3965 assert_eq!(list["pending_count"], 1);
3966
3967 let a = handoff
3968 .call(&mut ctx, json!({"action": "accept", "handoff_id": hid}))
3969 .await
3970 .unwrap();
3971 assert_eq!(a["status"], "success");
3972
3973 let list2 = handoff
3974 .call(&mut ctx, json!({"action": "list"}))
3975 .await
3976 .unwrap();
3977 assert_eq!(list2["pending_count"], 0);
3978 }
3979
3980 #[tokio::test]
3981 async fn lossless_replay_binds_explicit_session_chunks_and_detects_stale_view() {
3982 let store = test_store();
3983 let older = start_session_aged(&store, 60);
3984 let newer = start_session_aged(&store, 0);
3985 let content = format!("prefix {} DISTINCT-FACT-AFTER-120", "é".repeat(9000));
3986 let mut turn = Memory::new(Galaxy::Sessions, json!({"type":"session_turn","session_id":older,"sequence":1,"timestamp":1_i64,"content":content}).to_string());
3987 turn.metadata.tags = vec!["session".into(), "turn".into(), format!("session:{older}")];
3988 store.put(Galaxy::Sessions, &turn).unwrap();
3989 let mut other = Memory::new(Galaxy::Sessions, json!({"type":"session_turn","session_id":newer,"sequence":1,"timestamp":1_i64,"content":"newer-only"}).to_string());
3990 other.metadata.tags = vec!["session".into(), "turn".into(), format!("session:{newer}")];
3991 store.put(Galaxy::Sessions, &other).unwrap();
3992 let replay = SessionReplayTool::new(store.clone());
3993 let mut ctx = Context::default();
3994 let mut response = replay
3995 .call(
3996 &mut ctx,
3997 json!({"mode":"lossless","session_id":older,"page_size":1,"max_wire_bytes":2048}),
3998 )
3999 .await
4000 .unwrap();
4001 assert_eq!(response["session_id"], older);
4002 assert!(response["records"][0].get("chunk").is_some());
4003 let mut bytes = Vec::new();
4004 loop {
4005 assert!(serde_json::to_vec(&response).unwrap().len() <= 2048);
4006 for record in response["records"].as_array().unwrap() {
4007 if let Some(chunk) = record.get("chunk") {
4008 assert_eq!(chunk["byte_offset"].as_u64().unwrap() as usize, bytes.len());
4009 bytes.extend(base64_decode(chunk["data_b64"].as_str().unwrap()).unwrap());
4010 } else {
4011 bytes.extend(record["content"].as_str().unwrap().as_bytes());
4012 }
4013 }
4014 if response["complete"] == true {
4015 assert!(response["next_cursor"].is_null());
4016 break;
4017 }
4018 let cursor = response["next_cursor"].as_str().unwrap();
4019 response = replay.call(&mut ctx, json!({"mode":"lossless","session_id":older,"page_size":1,"max_wire_bytes":2048,"cursor":cursor})).await.unwrap();
4020 }
4021 assert_eq!(String::from_utf8(bytes).unwrap(), content);
4022 let first = replay
4023 .call(
4024 &mut ctx,
4025 json!({"mode":"lossless","session_id":older,"page_size":1,"max_wire_bytes":2048}),
4026 )
4027 .await
4028 .unwrap();
4029 let stale_cursor = first["next_cursor"].as_str().unwrap().to_string();
4030 let mut appended = Memory::new(Galaxy::Sessions, json!({"type":"session_turn","session_id":older,"sequence":2,"timestamp":2_i64,"content":"later"}).to_string());
4031 appended.metadata.tags = vec!["session".into(), "turn".into(), format!("session:{older}")];
4032 store.put(Galaxy::Sessions, &appended).unwrap();
4033 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"));
4034 }
4035
4036 fn lossless_fixture_turn(sid: &str, sequence: u64, text: &str) -> Memory {
4037 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());
4038 m.metadata.tags = vec!["turn".into(), format!("session:{sid}")];
4039 m
4040 }
4041
4042 #[tokio::test]
4043 async fn lossless_strict_args_and_cursor_prevalidation_before_corrupt_scan() {
4044 let store = test_store();
4045 let sid = start_session(&store);
4046 store
4047 .put_raw(
4048 Galaxy::Sessions,
4049 uuid::Uuid::new_v4().as_bytes(),
4050 b"not-a-memory",
4051 )
4052 .unwrap();
4053 let replay = SessionReplayTool::new(store);
4054 for (key, value) in [
4055 ("include_superseded", json!("false")),
4056 ("page_size", json!(null)),
4057 ("page_size", json!(-1)),
4058 ("max_wire_bytes", json!("2048")),
4059 ("cursor", json!(false)),
4060 ("cursor", json!("a€")),
4061 ("cursor", json!("🔑")),
4062 ] {
4063 let mut args = json!({"mode":"lossless","session_id":sid});
4064 args[key] = value;
4065 let error = replay
4066 .call(&mut Context::default(), args)
4067 .await
4068 .unwrap_err()
4069 .to_string();
4070 assert!(
4071 error.contains("invalid_args") || error.contains("invalid_cursor"),
4072 "{key}: {error}"
4073 );
4074 assert!(!error.contains("incomplete scan"));
4075 }
4076 let error = replay
4077 .call(
4078 &mut Context::default(),
4079 json!({"mode":"lossless","session_id":sid}),
4080 )
4081 .await
4082 .unwrap_err()
4083 .to_string();
4084 assert!(error.contains("refusing incomplete scan"));
4085 }
4086
4087 #[tokio::test]
4088 async fn lossless_metadata_staleness_settings_seek_and_partial_resume() {
4089 let store = test_store();
4090 let sid = start_session(&store);
4091 let mut turn = lossless_fixture_turn(&sid, 1, &"€\n".repeat(2000));
4092 store.put(Galaxy::Sessions, &turn).unwrap();
4093 let replay = SessionReplayTool::new(store.clone());
4094 let args = json!({"mode":"lossless","session_id":sid,"page_size":1,"max_wire_bytes":2048});
4095 let first = replay
4096 .call(&mut Context::default(), args.clone())
4097 .await
4098 .unwrap();
4099 let cursor = first["next_cursor"].as_str().unwrap();
4100 for (key, value) in [("page_size", json!(2)), ("max_wire_bytes", json!(4096))] {
4101 let mut next = args.clone();
4102 next["cursor"] = json!(cursor);
4103 next[key] = value;
4104 assert!(
4105 replay
4106 .call(&mut Context::default(), next)
4107 .await
4108 .unwrap_err()
4109 .to_string()
4110 .contains("invalid_cursor")
4111 );
4112 }
4113 let (view, _, _) = parse_lossless_cursor(cursor, &sid, false, 1, 2048).unwrap();
4114 let mut seek = args.clone();
4116 seek["cursor"] = json!(lossless_cursor(
4117 &sid,
4118 false,
4119 1,
4120 2048,
4121 &view,
4122 0,
4123 turn.content.len() - 2
4124 ));
4125 let inner = serde_json::from_str::<Value>(&turn.content).unwrap()["content"]
4127 .as_str()
4128 .unwrap()
4129 .to_string();
4130 seek["cursor"] = json!(lossless_cursor(
4131 &sid,
4132 false,
4133 1,
4134 2048,
4135 &view,
4136 0,
4137 inner.len() - 2
4138 ));
4139 let tail = replay.call(&mut Context::default(), seek).await.unwrap();
4140 assert!(tail["records"][0].get("content").is_none());
4141 assert_eq!(
4142 base64_decode(tail["records"][0]["chunk"]["data_b64"].as_str().unwrap()).unwrap(),
4143 inner.as_bytes()[inner.len() - 2..]
4144 );
4145 assert!(tail["next_cursor"].is_null());
4146 let mut changed = serde_json::from_str::<Value>(&turn.content).unwrap();
4147 changed["role"] = json!("human");
4148 turn.content = changed.to_string();
4149 store.put(Galaxy::Sessions, &turn).unwrap();
4150 let mut next = args;
4151 next["cursor"] = json!(cursor);
4152 assert!(
4153 replay
4154 .call(&mut Context::default(), next)
4155 .await
4156 .unwrap_err()
4157 .to_string()
4158 .contains("stale_view")
4159 );
4160 }
4161
4162 #[tokio::test]
4163 async fn lossless_visibility_empty_nonturn_ties_and_supersession() {
4164 let store = test_store();
4165 let sid = start_session(&store);
4166 let replay = SessionReplayTool::new(store.clone());
4167 let mut args = json!({"mode":"lossless","session_id":sid,"page_size":64});
4168 let empty = replay
4169 .call(&mut Context::default(), args.clone())
4170 .await
4171 .unwrap();
4172 assert_eq!(empty["records"], json!([]));
4173 assert!(empty["next_cursor"].is_null());
4174 assert_eq!(empty["complete"], true);
4175 let mut a = lossless_fixture_turn(&sid, 1, "");
4176 let mut b = lossless_fixture_turn(&sid, 1, "second");
4177 b.metadata.tags.clear();
4179 let mut hidden = lossless_fixture_turn(&sid, 2, "PRIVATE-SENTINEL");
4180 hidden.metadata.is_private = true;
4181 let mut excluded = lossless_fixture_turn(&sid, 3, "EXCLUDED-SENTINEL");
4182 excluded.metadata.model_exclude = true;
4183 a.metadata
4184 .tags
4185 .push(format!("supersedes:{}", hidden.metadata.id));
4186 let mut handoff = Memory::new(
4187 Galaxy::Sessions,
4188 json!({"type":"session_handoff"}).to_string(),
4189 );
4190 handoff.metadata.tags = vec![format!("session:{sid}")];
4191 store
4192 .put_batch(
4193 Galaxy::Sessions,
4194 &[a.clone(), b.clone(), hidden.clone(), excluded, handoff],
4195 )
4196 .unwrap();
4197 let result = replay
4198 .call(&mut Context::default(), args.clone())
4199 .await
4200 .unwrap();
4201 let wire = result.to_string();
4202 assert!(!wire.contains("SENTINEL"));
4203 assert!(!wire.contains(&hidden.metadata.id.to_string()));
4204 let mut ids = vec![a.metadata.id.to_string(), b.metadata.id.to_string()];
4205 ids.sort();
4206 assert_eq!(
4207 result["records"]
4208 .as_array()
4209 .unwrap()
4210 .iter()
4211 .map(|r| r["record_id"].as_str().unwrap().to_string())
4212 .collect::<Vec<_>>(),
4213 ids
4214 );
4215 a.metadata
4216 .tags
4217 .push(format!("superseded-by:{}", b.metadata.id));
4218 store.put(Galaxy::Sessions, &a).unwrap();
4219 assert_eq!(
4220 replay
4221 .call(&mut Context::default(), args.clone())
4222 .await
4223 .unwrap()["records"]
4224 .as_array()
4225 .unwrap()
4226 .len(),
4227 1
4228 );
4229 args["include_superseded"] = json!(true);
4230 assert_eq!(
4231 replay.call(&mut Context::default(), args).await.unwrap()["records"]
4232 .as_array()
4233 .unwrap()
4234 .len(),
4235 2
4236 );
4237 let mut start = store
4238 .get(Galaxy::Sessions, uuid::Uuid::parse_str(&sid).unwrap())
4239 .unwrap()
4240 .unwrap();
4241 start.metadata.model_exclude = true;
4242 store.put(Galaxy::Sessions, &start).unwrap();
4243 assert!(
4244 replay
4245 .call(
4246 &mut Context::default(),
4247 json!({"mode":"lossless","session_id":sid})
4248 )
4249 .await
4250 .unwrap_err()
4251 .to_string()
4252 .contains("not found")
4253 );
4254 }
4255
4256 #[tokio::test]
4257 async fn lossless_malformed_selected_turn_fails_and_schema_advertises_mode() {
4258 let store = test_store();
4259 let sid = start_session(&store);
4260 let mut bad = lossless_fixture_turn(&sid, 1, "good");
4261 bad.content = "{broken".into();
4262 store.put(Galaxy::Sessions, &bad).unwrap();
4263 let replay = SessionReplayTool::new(store);
4264 assert!(
4265 replay
4266 .call(
4267 &mut Context::default(),
4268 json!({"mode":"lossless","session_id":sid})
4269 )
4270 .await
4271 .unwrap_err()
4272 .to_string()
4273 .contains("malformed_selected_turn")
4274 );
4275 assert!(replay.input_schema().to_string().contains("lossless"));
4276 assert!(replay.description().contains("explicit session UUID"));
4277 }
4278
4279 #[test]
4280 fn lossless_arbitrary_cursor_text_is_panic_free() {
4281 use proptest::prelude::*;
4282 proptest!(|(text in any::<String>())| { let _ = hex_decode(&text); });
4283 assert!(hex_decode("AB").is_none());
4284 }
4285
4286 #[tokio::test]
4287 async fn lossless_selection_reaches_record_after_ten_thousand_sources() {
4288 let store = test_store();
4289 let sid = start_session(&store);
4290 let mut sources = Vec::new();
4291 for i in 1..=10_001_u128 {
4292 let mut m = Memory::new(Galaxy::Sessions, "unrelated".into());
4293 m.metadata.id = uuid::Uuid::from_u128(i);
4294 sources.push(m);
4295 }
4296 let mut target = lossless_fixture_turn(&sid, 1, "AFTER-TEN-THOUSAND");
4297 target.metadata.id = uuid::Uuid::from_u128(u128::MAX);
4298 sources.push(target);
4299 store.put_batch(Galaxy::Sessions, &sources).unwrap();
4300 let replay = SessionReplayTool::new(store);
4301 let result = replay
4302 .call(
4303 &mut Context::default(),
4304 json!({"mode":"lossless","session_id":sid}),
4305 )
4306 .await
4307 .unwrap();
4308 assert_eq!(result["records"][0]["content"], "AFTER-TEN-THOUSAND");
4309 }
4310
4311 #[tokio::test]
4312 async fn lossless_tag_payload_disagreement_refuses_and_missing_start_is_not_found() {
4313 let store = test_store();
4314 let sid = start_session(&store);
4315 let mut m = lossless_fixture_turn(&uuid::Uuid::new_v4().to_string(), 1, "contradiction");
4316 m.metadata.tags = vec![format!("session:{sid}")];
4317 store.put(Galaxy::Sessions, &m).unwrap();
4318 let replay = SessionReplayTool::new(store);
4319 assert!(
4320 replay
4321 .call(
4322 &mut Context::default(),
4323 json!({"mode":"lossless","session_id":sid})
4324 )
4325 .await
4326 .unwrap_err()
4327 .to_string()
4328 .contains("malformed_selected_turn")
4329 );
4330 assert!(
4331 replay
4332 .call(
4333 &mut Context::default(),
4334 json!({"mode":"lossless","session_id":uuid::Uuid::new_v4().to_string()})
4335 )
4336 .await
4337 .unwrap_err()
4338 .to_string()
4339 .contains("not found")
4340 );
4341 }
4342}