1#![forbid(unsafe_code)]
9
10use async_trait::async_trait;
11
12use chrono::{DateTime, NaiveDate, TimeZone, Utc};
13use serde_json::{Value, json};
14use sha2::{Digest, Sha256};
15use std::fmt::Write as _;
16use std::sync::Arc;
17use wm_core::{Context, EffectRow, Galaxy, Gana, Resource, Tool, ToolStats};
18use wm_memory::{Memory, MemoryStore};
19
20fn parse_time_bound(v: &Value, end_of_day: bool) -> Option<DateTime<Utc>> {
24 if let Some(secs) = v
25 .as_i64()
26 .or_else(|| v.as_u64().and_then(|u| i64::try_from(u).ok()))
27 {
28 return Utc.timestamp_opt(secs, 0).single();
29 }
30 let s = v.as_str()?;
31 if let Ok(dt) = DateTime::parse_from_rfc3339(s) {
32 return Some(dt.with_timezone(&Utc));
33 }
34 let day = NaiveDate::parse_from_str(s, "%Y-%m-%d").ok()?;
35 let naive = if end_of_day {
36 day.and_hms_opt(23, 59, 59)?
37 } else {
38 day.and_hms_opt(0, 0, 0)?
39 };
40 Some(naive.and_utc())
41}
42
43fn filter_by_time(
46 turns: Vec<(Memory, Value)>,
47 args: &Value,
48) -> wm_core::Result<Vec<(Memory, Value)>> {
49 let since = match args.get("since") {
50 Some(v) if !v.is_null() => Some(parse_time_bound(v, false).ok_or_else(|| {
51 wm_core::CoreError::InvalidArgs(
52 "invalid 'since' — use epoch seconds, RFC 3339, or YYYY-MM-DD".into(),
53 )
54 })?),
55 _ => None,
56 };
57 let until = match args.get("until") {
58 Some(v) if !v.is_null() => Some(parse_time_bound(v, true).ok_or_else(|| {
59 wm_core::CoreError::InvalidArgs(
60 "invalid 'until' — use epoch seconds, RFC 3339, or YYYY-MM-DD".into(),
61 )
62 })?),
63 _ => None,
64 };
65 Ok(turns
66 .into_iter()
67 .filter(|(m, _)| {
68 since.is_none_or(|t| m.metadata.created_at >= t)
69 && until.is_none_or(|t| m.metadata.created_at <= t)
70 })
71 .collect())
72}
73
74fn turn_json(mem: &Memory) -> Option<Value> {
75 let v: Value = serde_json::from_str(&mem.content).ok()?;
76 if v.get("type").and_then(Value::as_str) == Some("session_turn") {
77 Some(v)
78 } else {
79 None
80 }
81}
82
83fn load_turns(
89 store: &MemoryStore,
90 session_id: Option<&str>,
91 limit: usize,
92 include_superseded: bool,
93) -> wm_core::Result<Vec<(Memory, Value)>> {
94 let memories = store.scan_all(Galaxy::Sessions)?;
95 let mut turns: Vec<(Memory, Value)> = memories
96 .iter()
97 .filter(|m| {
98 include_superseded
99 || !m
100 .metadata
101 .tags
102 .iter()
103 .any(|t| t.starts_with("superseded-by:"))
104 })
105 .filter_map(|m| turn_json(m).map(|v| (m.clone(), v)))
106 .filter(|(_, v)| {
107 session_id.is_none_or(|sid| v.get("session_id").and_then(Value::as_str) == Some(sid))
108 })
109 .collect();
110 turns.sort_by_key(|(_, v)| {
111 (
112 v.get("sequence").and_then(Value::as_u64).unwrap_or(0),
113 v.get("timestamp").and_then(Value::as_i64).unwrap_or(0),
114 )
115 });
116 turns.truncate(limit);
117 Ok(turns)
118}
119
120fn latest_checkpoint_handoff(
125 store: &MemoryStore,
126 session_id: &str,
127) -> wm_core::Result<Option<(String, DateTime<Utc>, Value)>> {
128 Ok(store
129 .scan_all(Galaxy::Sessions)?
130 .iter()
131 .filter(|m| {
132 m.metadata.tags.contains(&"checkpoint".to_string()) && m.content.contains(session_id)
133 })
134 .filter_map(|m| {
135 let parsed: Value = serde_json::from_str(&m.content).ok()?;
136 parsed
137 .get("handoff")
138 .filter(|h| !h.is_null())
139 .cloned()
140 .map(|h| (m.metadata.id.to_string(), m.metadata.created_at, h))
141 })
142 .max_by_key(|(_, created_at, _)| *created_at))
143}
144
145fn format_turn(v: &Value, full: bool) -> Value {
146 let role = v.get("role").and_then(Value::as_str).unwrap_or("?");
147 let content = v.get("content").and_then(Value::as_str).unwrap_or("");
148 if full {
149 json!({
150 "session_id": v.get("session_id"),
151 "sequence": v.get("sequence"),
152 "role": role,
153 "turn_type": v.get("turn_type"),
154 "importance": v.get("importance"),
155 "content": content,
156 })
157 } else {
158 json!({
159 "sequence": v.get("sequence"),
160 "role": role,
161 "turn_type": v.get("turn_type"),
162 "preview": content.chars().take(120).collect::<String>(),
163 })
164 }
165}
166
167const LOSSLESS_MAX_PAGE_SIZE: usize = 64;
168const LOSSLESS_DEFAULT_PAGE_SIZE: usize = 16;
169const LOSSLESS_MIN_WIRE_BYTES: usize = 1024;
170const LOSSLESS_MAX_WIRE_BYTES: usize = 49_152;
171
172fn hex_encode(bytes: &[u8]) -> String {
173 use std::fmt::Write as _;
174 bytes
175 .iter()
176 .fold(String::with_capacity(bytes.len() * 2), |mut out, b| {
177 let _ = write!(out, "{b:02x}");
178 out
179 })
180}
181
182fn hex_decode(value: &str) -> Option<Vec<u8>> {
183 if value.len() > 4096
184 || value.len() % 2 != 0
185 || !value
186 .bytes()
187 .all(|b| b.is_ascii_digit() || (b'a'..=b'f').contains(&b))
188 {
189 return None;
190 }
191 (0..value.len())
192 .step_by(2)
193 .map(|i| u8::from_str_radix(&value[i..i + 2], 16).ok())
194 .collect()
195}
196
197fn base64_encode(bytes: &[u8]) -> String {
198 const TABLE: &[u8; 64] = b"ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz0123456789+/";
199 let mut out = String::with_capacity(bytes.len().div_ceil(3) * 4);
200 for chunk in bytes.chunks(3) {
201 let n = u32::from(chunk[0]) << 16
202 | u32::from(*chunk.get(1).unwrap_or(&0)) << 8
203 | u32::from(*chunk.get(2).unwrap_or(&0));
204 out.push(char::from(TABLE[((n >> 18) & 63) as usize]));
205 out.push(char::from(TABLE[((n >> 12) & 63) as usize]));
206 out.push(if chunk.len() > 1 {
207 char::from(TABLE[((n >> 6) & 63) as usize])
208 } else {
209 '='
210 });
211 out.push(if chunk.len() > 2 {
212 char::from(TABLE[(n & 63) as usize])
213 } else {
214 '='
215 });
216 }
217 out
218}
219
220#[cfg(test)]
221fn base64_decode(value: &str) -> Option<Vec<u8>> {
222 const fn digit(byte: u8) -> Option<u8> {
223 match byte {
224 b'A'..=b'Z' => Some(byte - b'A'),
225 b'a'..=b'z' => Some(byte - b'a' + 26),
226 b'0'..=b'9' => Some(byte - b'0' + 52),
227 b'+' => Some(62),
228 b'/' => Some(63),
229 _ => None,
230 }
231 }
232 if value.len() % 4 != 0 {
233 return None;
234 }
235 let mut out = Vec::new();
236 for chunk in value.as_bytes().chunks_exact(4) {
237 let a = digit(chunk[0])?;
238 let b = digit(chunk[1])?;
239 let c = if chunk[2] == b'=' {
240 0
241 } else {
242 digit(chunk[2])?
243 };
244 let d = if chunk[3] == b'=' {
245 0
246 } else {
247 digit(chunk[3])?
248 };
249 out.push((a << 2) | (b >> 4));
250 if chunk[2] != b'=' {
251 out.push((b << 4) | (c >> 2));
252 }
253 if chunk[3] != b'=' {
254 out.push((c << 6) | d);
255 }
256 }
257 Some(out)
258}
259
260#[derive(Clone)]
261struct LosslessTurn {
262 memory: Memory,
263 turn: Value,
264 content: String,
265 content_hash: String,
266}
267
268fn lossless_error(kind: &str) -> wm_core::CoreError {
269 wm_core::CoreError::InvalidArgs(format!("lossless_{kind}"))
270}
271
272fn lossless_cursor(
273 session_id: &str,
274 include_superseded: bool,
275 page_size: usize,
276 max_wire: usize,
277 view: &str,
278 index: usize,
279 offset: usize,
280) -> String {
281 let value = json!({"v":1,"session_id":session_id,"include_superseded":include_superseded,"page_size":page_size,"max_wire_bytes":max_wire,"view":view,"index":index,"offset":offset});
284 hex_encode(value.to_string().as_bytes())
285}
286
287fn parse_lossless_cursor(
288 cursor: &str,
289 session_id: &str,
290 include_superseded: bool,
291 page_size: usize,
292 max_wire: usize,
293) -> wm_core::Result<(String, usize, usize)> {
294 let bytes = hex_decode(cursor).ok_or_else(|| lossless_error("invalid_cursor"))?;
295 let text = String::from_utf8(bytes).map_err(|_| lossless_error("invalid_cursor"))?;
296 let value: Value = serde_json::from_str(&text).map_err(|_| lossless_error("invalid_cursor"))?;
297 let canonical = json!({"v":value.get("v"),"session_id":value.get("session_id"),"include_superseded":value.get("include_superseded"),"page_size":value.get("page_size"),"max_wire_bytes":value.get("max_wire_bytes"),"view":value.get("view"),"index":value.get("index"),"offset":value.get("offset")});
298 let canonical_text =
299 serde_json::to_string(&canonical).map_err(|_| lossless_error("invalid_cursor"))?;
300 if canonical_text != text
301 || value.get("v").and_then(Value::as_u64) != Some(1)
302 || value.get("session_id").and_then(Value::as_str) != Some(session_id)
303 || value.get("include_superseded").and_then(Value::as_bool) != Some(include_superseded)
304 || value.get("page_size").and_then(Value::as_u64) != Some(page_size as u64)
305 || value.get("max_wire_bytes").and_then(Value::as_u64) != Some(max_wire as u64)
306 {
307 return Err(lossless_error("invalid_cursor"));
308 }
309 let view = value
310 .get("view")
311 .and_then(Value::as_str)
312 .filter(|v| {
313 v.len() == 64
314 && v.bytes()
315 .all(|b| b.is_ascii_digit() || (b'a'..=b'f').contains(&b))
316 })
317 .ok_or_else(|| lossless_error("invalid_cursor"))?;
318 let index = value
319 .get("index")
320 .and_then(Value::as_u64)
321 .and_then(|v| usize::try_from(v).ok())
322 .ok_or_else(|| lossless_error("invalid_cursor"))?;
323 let offset = value
324 .get("offset")
325 .and_then(Value::as_u64)
326 .and_then(|v| usize::try_from(v).ok())
327 .ok_or_else(|| lossless_error("invalid_cursor"))?;
328 Ok((view.to_string(), index, offset))
329}
330
331pub const TURN_TYPES: &[&str] = &[
336 "message",
337 "decision",
338 "breakthrough",
339 "question",
340 "answer",
341 "code_change",
342 "error",
343 "summary",
344 "context",
345];
346
347pub struct SessionRecordTool {
349 store: Arc<MemoryStore>,
350 stats: ToolStats,
351 effects: EffectRow,
352 search: Option<Arc<wm_memory::SearchEngine>>,
353}
354
355impl SessionRecordTool {
356 #[must_use]
357 pub fn new(store: Arc<MemoryStore>) -> Self {
358 Self {
359 store,
360 stats: ToolStats::default(),
361 effects: EffectRow {
362 writes: vec![Resource::Galaxy("sessions".into())],
363 ..Default::default()
364 },
365 search: None,
366 }
367 }
368
369 #[must_use]
373 pub fn with_search(mut self, search: Option<Arc<wm_memory::SearchEngine>>) -> Self {
374 self.search = search;
375 self
376 }
377}
378
379#[async_trait]
380impl Tool for SessionRecordTool {
381 fn name(&self) -> &str {
382 "session.record"
383 }
384 fn gana(&self) -> Gana {
385 Gana::StraddlingLegs
386 }
387 fn effects(&self) -> &EffectRow {
388 &self.effects
389 }
390 fn input_schema(&self) -> Value {
391 super::common::schema(
392 &json!({
393 "content": super::common::str_prop("Turn content"),
394 "role": super::common::str_prop("user | ai (default user)"),
395 "turn_type": json!({
396 "type": "string",
397 "enum": TURN_TYPES,
398 "description": "Turn type (default message)",
399 }),
400 "importance": super::common::bounded_num_prop("0-1 importance (default 0.5)", 0.0, 1.0),
401 "session_id": super::common::str_prop("Target session (default: most recent session)"),
402 "supersedes": super::common::str_prop("Memory id of an earlier turn this record corrects/replaces (amend-with-supersede)"),
403 }),
404 &["content"],
405 )
406 }
407 fn description(&self) -> &str {
408 "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)."
409 }
410 async fn call(&self, _ctx: &mut Context, args: Value) -> wm_core::Result<Value> {
411 let role = args.get("role").and_then(Value::as_str).unwrap_or("user");
412 if !matches!(role, "user" | "ai") {
413 return Err(wm_core::CoreError::InvalidArgs(
414 "role must be 'user' or 'ai'".into(),
415 ));
416 }
417 let content = args
420 .get("content")
421 .and_then(Value::as_str)
422 .filter(|s| !s.trim().is_empty())
423 .ok_or_else(|| {
424 wm_core::CoreError::InvalidArgs("content is required and must not be blank".into())
425 })?;
426 let turn_type = args
430 .get("turn_type")
431 .and_then(Value::as_str)
432 .unwrap_or("message");
433 if !TURN_TYPES.contains(&turn_type) {
434 return Err(wm_core::CoreError::InvalidArgs(format!(
435 "turn_type must be one of: {}",
436 TURN_TYPES.join(", ")
437 )));
438 }
439 let importance = wm_dispatch::write_gate::parse_importance_value(args.get("importance"))
444 .map_err(wm_core::CoreError::InvalidArgs)?
445 .unwrap_or(0.5);
446 let session_id = args.get("session_id").and_then(Value::as_str);
447
448 let session_id: String = if let Some(sid) = session_id {
454 sid.to_string()
455 } else {
456 self.store
457 .scan_all(Galaxy::Sessions)?
458 .iter()
459 .filter(|m| m.metadata.tags.contains(&"start".to_string()))
460 .max_by_key(|m| m.metadata.created_at)
461 .map(|m| m.metadata.id.to_string())
462 .ok_or_else(|| {
463 wm_core::CoreError::Tool("no session found — run session.start first".into())
464 })?
465 };
466
467 let sequence = load_turns(&self.store, Some(&session_id), 10_000, true)?.len() as u64 + 1;
470
471 let mut mem = Memory::new(
472 Galaxy::Sessions,
473 json!({
474 "type": "session_turn",
475 "session_id": session_id,
476 "sequence": sequence,
477 "role": role,
478 "turn_type": turn_type,
479 "importance": importance,
480 "content": content,
481 "timestamp": wm_core::time::now_unix_millis(),
482 })
483 .to_string(),
484 );
485 mem.metadata.tags = vec![
486 "session".into(),
487 "turn".into(),
488 role.into(),
489 turn_type.into(),
490 format!("session:{session_id}"),
491 ];
492 let (source, trust): (&str, f32) = if role == "user" {
499 ("user", 1.0)
500 } else {
501 ("agent", 0.7)
502 };
503 mem.metadata.source = source.to_string();
504 mem.metadata.source_trust = trust;
505
506 if let Some(old_id_str) = args.get("supersedes").and_then(Value::as_str) {
510 let old_id = uuid::Uuid::parse_str(old_id_str).map_err(|e| {
511 wm_core::CoreError::InvalidArgs(format!("invalid 'supersedes' id: {e}"))
512 })?;
513 let mut old = self.store.get(Galaxy::Sessions, old_id)?.ok_or_else(|| {
514 wm_core::CoreError::NotFound(format!("superseded turn {old_id} not found"))
515 })?;
516 old.metadata
517 .tags
518 .push(format!("superseded-by:{}", mem.metadata.id));
519 self.store.put(Galaxy::Sessions, &old)?;
520 super::common::index_memory(self.search.as_deref(), &old);
521 mem.metadata.tags.push(format!("supersedes:{old_id}"));
522 }
523
524 mem.metadata.importance = importance as f32;
525 self.store.put(Galaxy::Sessions, &mem)?;
526 super::common::index_memory(self.search.as_deref(), &mem);
527 Ok(json!({
528 "status": "success",
529 "session_id": session_id,
530 "sequence": sequence,
531 "memory_id": mem.metadata.id.to_string(),
532 }))
533 }
534 fn stats(&self) -> &ToolStats {
535 &self.stats
536 }
537}
538
539pub struct SessionReplayTool {
541 store: Arc<MemoryStore>,
542 stats: ToolStats,
543 effects: EffectRow,
544}
545
546impl SessionReplayTool {
547 #[must_use]
548 pub fn new(store: Arc<MemoryStore>) -> Self {
549 Self {
550 store,
551 stats: ToolStats::default(),
552 effects: EffectRow::read_only(vec![Resource::Galaxy("sessions".into())]),
553 }
554 }
555
556 fn lossless(&self, args: &Value) -> wm_core::Result<Value> {
557 const ALLOWED: &[&str] = &[
558 "mode",
559 "session_id",
560 "include_superseded",
561 "page_size",
562 "max_wire_bytes",
563 "cursor",
564 ];
565 let object = args
566 .as_object()
567 .ok_or_else(|| lossless_error("invalid_args"))?;
568 if object.keys().any(|key| !ALLOWED.contains(&key.as_str())) {
569 return Err(lossless_error("unsupported_selection_args"));
570 }
571 let session_id = args
572 .get("session_id")
573 .and_then(Value::as_str)
574 .filter(|s| !s.is_empty())
575 .ok_or_else(|| lossless_error("session_id_required"))?;
576 uuid::Uuid::parse_str(session_id).map_err(|_| lossless_error("invalid_session_id"))?;
577 let include_superseded = match args.get("include_superseded") {
578 None => false,
579 Some(v) => v.as_bool().ok_or_else(|| lossless_error("invalid_args"))?,
580 };
581 let size_arg = |name: &str, default: usize| -> wm_core::Result<usize> {
582 match args.get(name) {
583 None => Ok(default),
584 Some(v) => v
585 .as_u64()
586 .and_then(|v| usize::try_from(v).ok())
587 .ok_or_else(|| lossless_error("invalid_args")),
588 }
589 };
590 let page_size = size_arg("page_size", LOSSLESS_DEFAULT_PAGE_SIZE)?;
591 let max_wire = size_arg("max_wire_bytes", LOSSLESS_MAX_WIRE_BYTES)?;
592 if !(1..=LOSSLESS_MAX_PAGE_SIZE).contains(&page_size)
593 || !(LOSSLESS_MIN_WIRE_BYTES..=LOSSLESS_MAX_WIRE_BYTES).contains(&max_wire)
594 {
595 return Err(lossless_error("invalid_args"));
596 }
597
598 let placement = match args.get("cursor") {
600 None => None,
601 Some(v) => Some(parse_lossless_cursor(
602 v.as_str().ok_or_else(|| lossless_error("invalid_cursor"))?,
603 session_id,
604 include_superseded,
605 page_size,
606 max_wire,
607 )?),
608 };
609 let memories = self.store.scan_all_strict(Galaxy::Sessions)?;
610 let start = memories.iter().find(|m| {
611 m.metadata.id.to_string() == session_id
612 && m.metadata.tags.contains(&"start".to_string())
613 && !m.metadata.is_private
614 && !m.metadata.model_exclude
615 });
616 if start.is_none() {
617 return Err(wm_core::CoreError::NotFound("session not found".into()));
618 }
619 let mut turns = Vec::new();
620 for memory in memories {
621 if memory.metadata.is_private
622 || memory.metadata.model_exclude
623 || (!include_superseded
624 && memory
625 .metadata
626 .tags
627 .iter()
628 .any(|t| t.starts_with("superseded-by:")))
629 {
630 continue;
631 }
632 let tagged = memory
633 .metadata
634 .tags
635 .contains(&format!("session:{session_id}"));
636 let turn = match serde_json::from_str::<Value>(&memory.content) {
637 Ok(v) => v,
638 Err(_) if tagged => return Err(lossless_error("malformed_selected_turn")),
639 Err(_) => continue,
640 };
641 if turn.get("type").and_then(Value::as_str) != Some("session_turn") {
642 if tagged && memory.metadata.tags.iter().any(|tag| tag == "turn") {
643 return Err(lossless_error("malformed_selected_turn"));
644 }
645 continue;
646 }
647 if !tagged && turn.get("session_id").and_then(Value::as_str) != Some(session_id) {
648 continue;
649 }
650 let content = turn
651 .get("content")
652 .and_then(Value::as_str)
653 .ok_or_else(|| lossless_error("malformed_selected_turn"))?
654 .to_string();
655 if turn.get("session_id").and_then(Value::as_str) != Some(session_id)
658 || turn.get("sequence").and_then(Value::as_u64).is_none()
659 || turn.get("timestamp").and_then(Value::as_i64).is_none()
660 {
661 return Err(lossless_error("malformed_selected_turn"));
662 }
663 turns.push(LosslessTurn {
664 content_hash: hex_encode(&Sha256::digest(content.as_bytes())),
665 memory,
666 turn,
667 content,
668 });
669 }
670 turns.sort_by_key(|turn| {
671 (
672 turn.turn["sequence"].as_u64().unwrap(),
673 turn.turn["timestamp"].as_i64().unwrap(),
674 turn.memory.metadata.id,
675 )
676 });
677 let visible_ids: std::collections::HashSet<String> = turns
678 .iter()
679 .map(|t| t.memory.metadata.id.to_string())
680 .collect();
681 let metadata: Vec<Value> = turns.iter().map(|t| {
682 let relationships: Vec<&str> = t.memory.metadata.tags.iter().filter_map(|tag| {
683 let (_, id) = tag.split_once(':')?;
684 ((tag.starts_with("superseded-by:") || tag.starts_with("supersedes:")) && visible_ids.contains(id)).then_some(tag.as_str())
685 }).collect();
686 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})
687 }).collect();
688 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();
689 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()));
690 let (mut index, mut offset) = match placement {
691 Some((token_view, index, offset)) => {
692 if token_view != view {
693 return Err(lossless_error("stale_view"));
694 }
695 (index, offset)
696 }
697 None => (0, 0),
698 };
699 if index > turns.len() || (index == turns.len() && offset != 0) {
700 return Err(lossless_error("invalid_placement"));
701 }
702 if index < turns.len() && offset != 0 && offset >= turns[index].content.len() {
703 return Err(lossless_error("invalid_placement"));
704 }
705 let mut records = Vec::new();
706 while index < turns.len() && records.len() < page_size {
707 let turn = &turns[index];
708 let bytes = turn.content.as_bytes();
709 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});
710 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()});
711 if offset == 0 && serde_json::to_vec(&candidate).unwrap().len() <= max_wire {
712 records.push(whole);
713 index += 1;
714 offset = 0;
715 continue;
716 }
717 if !records.is_empty() {
718 break;
719 }
720 let mut take = bytes.len().saturating_sub(offset);
721 while take > 0 {
722 let end = offset + take;
723 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()}});
724 let next = if end == bytes.len() {
725 lossless_cursor(
726 session_id,
727 include_superseded,
728 page_size,
729 max_wire,
730 &view,
731 index + 1,
732 0,
733 )
734 } else {
735 lossless_cursor(
736 session_id,
737 include_superseded,
738 page_size,
739 max_wire,
740 &view,
741 index,
742 end,
743 )
744 };
745 let final_chunk = end == bytes.len() && index + 1 == turns.len();
746 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});
747 if serde_json::to_vec(&candidate).unwrap().len() <= max_wire {
748 records.push(chunk);
749 if end == bytes.len() {
750 index += 1;
751 offset = 0;
752 } else {
753 offset = end;
754 }
755 break;
756 }
757 take /= 2;
758 }
759 if records.is_empty() {
760 return Err(lossless_error("wire_ceiling_too_small"));
761 }
762 break;
763 }
764 let complete = index == turns.len() && offset == 0;
765 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});
766 if serde_json::to_vec(&response).unwrap().len() > max_wire {
767 return Err(lossless_error("wire_ceiling_too_small"));
768 }
769 Ok(response)
770 }
771}
772
773#[async_trait]
774impl Tool for SessionReplayTool {
775 fn name(&self) -> &str {
776 "session.replay"
777 }
778 fn gana(&self) -> Gana {
779 Gana::StraddlingLegs
780 }
781 fn effects(&self) -> &EffectRow {
782 &self.effects
783 }
784 fn input_schema(&self) -> Value {
785 super::common::schema(
786 &json!({
787 "mode": super::common::str_prop("full | selective | progressive | lossless (default full)"),
788 "session_id": super::common::str_prop("Target session (required explicit UUID for lossless; otherwise default most recent)"),
789 "n": super::common::int_prop("Maximum turns (default 50)"),
790 "since": super::common::str_prop("Time-range floor: epoch seconds, RFC 3339, or YYYY-MM-DD"),
791 "until": super::common::str_prop("Time-range ceiling: epoch seconds, RFC 3339, or YYYY-MM-DD"),
792 "include_superseded": {
793 "type": "boolean",
794 "description": "Also return turns replaced via supersedes (default false)."
795 },
796 "turn_types": super::common::str_array_prop("Selective mode: turn types to keep"),
797 "min_importance": super::common::num_prop("Selective mode floor (default 0.7)"),
798 "token_budget": super::common::int_prop("Progressive mode token budget (default 2000)"),
799 "page_size": super::common::int_prop("Lossless mode records per page (1-64, default 16)"),
800 "max_wire_bytes": super::common::int_prop("Lossless mode serialized JSON ceiling (1024-49152)"),
801 "cursor": super::common::str_prop("Lossless mode opaque placement cursor"),
802 }),
803 &[],
804 )
805 }
806 fn description(&self) -> &str {
807 "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."
808 }
809 async fn call(&self, _ctx: &mut Context, args: Value) -> wm_core::Result<Value> {
810 let mode = args.get("mode").and_then(Value::as_str).unwrap_or("full");
811 if mode == "lossless" {
812 return self.lossless(&args);
813 }
814 let requested_session_id = args
815 .get("session_id")
816 .and_then(Value::as_str)
817 .filter(|sid| !sid.is_empty());
818 let session_id = match requested_session_id {
823 Some(sid) => Some(sid.to_string()),
824 None => self
825 .store
826 .scan_all(Galaxy::Sessions)?
827 .iter()
828 .filter(|m| m.metadata.tags.contains(&"start".to_string()))
829 .max_by_key(|m| m.metadata.created_at)
830 .map(|m| m.metadata.id.to_string()),
831 };
832 let n = args.get("n").and_then(Value::as_u64).unwrap_or(50) as usize;
833 let include_superseded = args
834 .get("include_superseded")
835 .and_then(Value::as_bool)
836 .unwrap_or(false);
837 let loaded_turns = match session_id.as_deref() {
841 Some(sid) => load_turns(&self.store, Some(sid), 10_000, include_superseded)?,
842 None => Vec::new(),
843 };
844 let turns = filter_by_time(loaded_turns, &args)?;
845
846 if requested_session_id.is_some() && turns.is_empty() {
849 return Err(wm_core::CoreError::InvalidArgs(format!(
850 "no session found with id {requested_session_id:?}"
851 )));
852 }
853
854 let selected: Vec<(Memory, Value)> = match mode {
855 "selective" => {
856 let min_importance = args
857 .get("min_importance")
858 .and_then(Value::as_f64)
859 .unwrap_or(0.7);
860 let turn_types: Vec<String> = args
861 .get("turn_types")
862 .and_then(Value::as_array)
863 .map_or_else(
864 || vec!["decision".into(), "breakthrough".into(), "answer".into()],
865 |a| {
866 a.iter()
867 .filter_map(Value::as_str)
868 .map(str::to_string)
869 .collect()
870 },
871 );
872 turns
873 .into_iter()
874 .filter(|(_, v)| {
875 v.get("importance").and_then(Value::as_f64).unwrap_or(0.0) >= min_importance
876 && v.get("turn_type")
877 .and_then(Value::as_str)
878 .is_some_and(|t| turn_types.contains(&t.to_string()))
879 })
880 .collect()
881 }
882 "progressive" => {
883 let budget = args
884 .get("token_budget")
885 .and_then(Value::as_u64)
886 .unwrap_or(2000) as usize;
887 let mut used = 0usize;
888 let mut out = Vec::new();
889 for (m, v) in turns.into_iter().rev() {
890 let approx = v
891 .get("content")
892 .and_then(Value::as_str)
893 .map_or(0, |c| c.len() / 4);
894 if used + approx > budget {
895 break;
896 }
897 used += approx;
898 out.push((m, v));
899 }
900 out.reverse();
901 out
902 }
903 _ => turns
904 .into_iter()
905 .rev()
906 .take(n)
907 .collect::<Vec<_>>()
908 .into_iter()
909 .rev()
910 .collect(),
911 };
912
913 let full = mode != "progressive";
914 let formatted: Vec<Value> = selected.iter().map(|(_, v)| format_turn(v, full)).collect();
915 Ok(json!({
916 "status": "success",
917 "mode": mode,
918 "count": formatted.len(),
919 "session_id": session_id,
920 "turns": formatted,
921 }))
922 }
923 fn stats(&self) -> &ToolStats {
924 &self.stats
925 }
926}
927
928pub struct SessionContinuityTool {
930 store: Arc<MemoryStore>,
931 stats: ToolStats,
932 effects: EffectRow,
933}
934
935impl SessionContinuityTool {
936 #[must_use]
937 pub fn new(store: Arc<MemoryStore>) -> Self {
938 Self {
939 store,
940 stats: ToolStats::default(),
941 effects: EffectRow::read_only(vec![Resource::Galaxy("sessions".into())]),
942 }
943 }
944}
945
946#[async_trait]
947impl Tool for SessionContinuityTool {
948 fn name(&self) -> &str {
949 "session.continuity"
950 }
951 fn gana(&self) -> Gana {
952 Gana::StraddlingLegs
953 }
954 fn effects(&self) -> &EffectRow {
955 &self.effects
956 }
957 fn input_schema(&self) -> Value {
958 super::common::schema(
959 &json!({
960 "current_session_id": super::common::str_prop("Session to exclude (optional)"),
961 "n": super::common::int_prop("Number of prior turns (default 10)"),
962 "since": super::common::str_prop("Time-range floor: epoch seconds, RFC 3339, or YYYY-MM-DD"),
963 "until": super::common::str_prop("Time-range ceiling: epoch seconds, RFC 3339, or YYYY-MM-DD"),
964 }),
965 &[],
966 )
967 }
968 fn description(&self) -> &str {
969 "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)."
970 }
971 async fn call(&self, _ctx: &mut Context, args: Value) -> wm_core::Result<Value> {
972 let current = args
973 .get("current_session_id")
974 .or_else(|| args.get("session_id"))
975 .and_then(Value::as_str);
976 let n = args.get("n").and_then(Value::as_u64).unwrap_or(10) as usize;
977
978 let memories = self.store.scan_all(Galaxy::Sessions)?;
989 let mut starts: Vec<_> = memories
990 .iter()
991 .filter(|m| {
992 m.metadata.tags.contains(&"start".to_string())
993 && current.is_none_or(|c| m.metadata.id.to_string() != c)
994 })
995 .collect();
996 starts.sort_by_key(|m| std::cmp::Reverse(m.metadata.created_at));
997 let previous = starts
998 .iter()
999 .find(|m| {
1000 let sid = m.metadata.id.to_string();
1001 load_turns(&self.store, Some(&sid), 1, false).is_ok_and(|turns| !turns.is_empty())
1002 })
1003 .copied()
1004 .or_else(|| starts.first().copied());
1005
1006 let Some(prev) = previous else {
1007 let project = std::env::var("WM_PROJECT").ok().filter(|s| !s.is_empty());
1014 let store = self.store.path().display().to_string();
1015 let scope = project.map_or_else(
1016 || format!("store {store}"),
1017 |p| format!("store {store}, project '{p}'"),
1018 );
1019 return Ok(json!({
1020 "status": "success",
1021 "previous_session": null,
1022 "turns": [],
1023 "count": 0,
1024 "message": "no previous session found",
1025 "hint": format!(
1026 "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."
1027 ),
1028 }));
1029 };
1030
1031 let prev_id = prev.metadata.id.to_string();
1032 let mut turns = filter_by_time(
1033 load_turns(&self.store, Some(&prev_id), 10_000, false)?,
1034 &args,
1035 )?;
1036 let total = turns.len();
1037 let tail: Vec<Value> = turns
1038 .split_off(total.saturating_sub(n))
1039 .iter()
1040 .map(|(_, v)| format_turn(v, true))
1041 .collect();
1042 let (checkpoint, checkpoint_id) = match latest_checkpoint_handoff(&self.store, &prev_id)? {
1049 Some((id, _, handoff)) => (handoff, Value::String(id)),
1050 None => (Value::Null, Value::Null),
1051 };
1052 Ok(json!({
1053 "status": "success",
1054 "previous_session": prev_id,
1055 "count": tail.len(),
1056 "turns": tail,
1057 "checkpoint": checkpoint,
1058 "checkpoint_id": checkpoint_id,
1059 }))
1060 }
1061 fn stats(&self) -> &ToolStats {
1062 &self.stats
1063 }
1064}
1065
1066pub struct SessionDigestTool {
1073 store: Arc<MemoryStore>,
1074 stats: ToolStats,
1075 effects: EffectRow,
1076}
1077
1078impl SessionDigestTool {
1079 #[must_use]
1080 pub fn new(store: Arc<MemoryStore>) -> Self {
1081 Self {
1082 store,
1083 stats: ToolStats::default(),
1084 effects: EffectRow::read_only(vec![Resource::Galaxy("sessions".into())]),
1085 }
1086 }
1087}
1088
1089const DIGEST_SECTION_ORDER: &[&str] = &["decision", "breakthrough", "error", "summary"];
1091
1092#[async_trait]
1093impl Tool for SessionDigestTool {
1094 fn name(&self) -> &str {
1095 "session.digest"
1096 }
1097 fn gana(&self) -> Gana {
1098 Gana::StraddlingLegs
1099 }
1100 fn effects(&self) -> &EffectRow {
1101 &self.effects
1102 }
1103 fn input_schema(&self) -> Value {
1104 super::common::schema(
1105 &json!({
1106 "session_id": super::common::str_prop("Session to digest (default: most recent)"),
1107 "min_importance": super::common::num_prop("Importance floor (default 0.5)"),
1108 "include_checkpoint": {
1109 "type": "boolean",
1110 "description": "Append the latest checkpoint's git/handoff state (default true)."
1111 },
1112 "since": super::common::str_prop("Time-range floor: epoch seconds, RFC 3339, or YYYY-MM-DD"),
1113 "until": super::common::str_prop("Time-range ceiling: epoch seconds, RFC 3339, or YYYY-MM-DD"),
1114 }),
1115 &[],
1116 )
1117 }
1118 fn description(&self) -> &str {
1119 "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."
1120 }
1121 async fn call(&self, _ctx: &mut Context, args: Value) -> wm_core::Result<Value> {
1122 let session_id = match args.get("session_id").and_then(Value::as_str) {
1123 Some(sid) if !sid.is_empty() => sid.to_string(),
1124 _ => self
1125 .store
1126 .scan_all(Galaxy::Sessions)?
1127 .iter()
1128 .filter(|m| m.metadata.tags.contains(&"start".to_string()))
1129 .max_by_key(|m| m.metadata.created_at)
1130 .map(|m| m.metadata.id.to_string())
1131 .ok_or_else(|| {
1132 wm_core::CoreError::Tool("no session found — run session.start first".into())
1133 })?,
1134 };
1135 let min_importance = args
1136 .get("min_importance")
1137 .and_then(Value::as_f64)
1138 .unwrap_or(0.5);
1139 let include_checkpoint = args
1140 .get("include_checkpoint")
1141 .and_then(Value::as_bool)
1142 .unwrap_or(true);
1143
1144 let mut turns: Vec<_> = filter_by_time(
1145 load_turns(&self.store, Some(&session_id), 10_000, false)?,
1146 &args,
1147 )?
1148 .into_iter()
1149 .filter(|(_, v)| {
1150 v.get("importance").and_then(Value::as_f64).unwrap_or(0.0) >= min_importance
1151 })
1152 .collect();
1153 turns.sort_by(|a, b| {
1154 b.1.get("importance")
1155 .and_then(Value::as_f64)
1156 .unwrap_or(0.0)
1157 .total_cmp(&a.1.get("importance").and_then(Value::as_f64).unwrap_or(0.0))
1158 });
1159
1160 let mut groups: Vec<(String, Vec<&Value>)> = Vec::new();
1162 for (_, v) in &turns {
1163 let t = v
1164 .get("turn_type")
1165 .and_then(Value::as_str)
1166 .unwrap_or("message")
1167 .to_string();
1168 match groups.iter_mut().find(|(name, _)| *name == t) {
1169 Some((_, list)) => list.push(v),
1170 None => groups.push((t, vec![v])),
1171 }
1172 }
1173 groups.sort_by_key(|(name, _)| {
1174 (
1175 DIGEST_SECTION_ORDER
1176 .iter()
1177 .position(|k| k == name)
1178 .unwrap_or(DIGEST_SECTION_ORDER.len()),
1179 name.clone(),
1180 )
1181 });
1182
1183 let mut digest = format!("# Session handoff — {session_id}\n");
1184 let mut included = 0usize;
1185 for (turn_type, items) in &groups {
1186 writeln!(
1187 digest,
1188 "\n## {} ({})",
1189 capitalize(&pluralize(turn_type)),
1190 items.len()
1191 )
1192 .expect("write to String cannot fail");
1193 for v in items {
1194 let importance = v.get("importance").and_then(Value::as_f64).unwrap_or(0.0);
1195 let content = v.get("content").and_then(Value::as_str).unwrap_or("");
1196 writeln!(digest, "- ({importance:.2}) {content}")
1197 .expect("write to String cannot fail");
1198 included += 1;
1199 }
1200 }
1201
1202 let mut checkpoint_state = Value::Null;
1204 if include_checkpoint {
1205 if let Some((_, _, cp)) = latest_checkpoint_handoff(&self.store, &session_id)? {
1206 digest.push_str("\n## Checkpoint state\n");
1207 if let Some(git) = cp.get("git") {
1208 writeln!(
1209 digest,
1210 "- commit `{}` on `{}` ({} dirty files)",
1211 git.get("commit").and_then(Value::as_str).unwrap_or("?"),
1212 git.get("branch").and_then(Value::as_str).unwrap_or("?"),
1213 git.get("dirty_count").and_then(Value::as_i64).unwrap_or(0)
1214 )
1215 .expect("write to String cannot fail");
1216 }
1217 if let Some(q) = cp.get("next_queue").and_then(Value::as_array) {
1218 if !q.is_empty() {
1219 writeln!(
1220 digest,
1221 "- next queue: {}",
1222 q.iter()
1223 .filter_map(Value::as_str)
1224 .collect::<Vec<_>>()
1225 .join(" → ")
1226 )
1227 .expect("write to String cannot fail");
1228 }
1229 }
1230 if let Some(f) = cp.get("open_flags").and_then(Value::as_array) {
1231 if !f.is_empty() {
1232 writeln!(
1233 digest,
1234 "- open flags: {}",
1235 f.iter()
1236 .filter_map(Value::as_str)
1237 .collect::<Vec<_>>()
1238 .join("; ")
1239 )
1240 .expect("write to String cannot fail");
1241 }
1242 }
1243 if let Some(tg) = cp.get("tests_green") {
1244 writeln!(digest, "- tests green: {tg}").expect("write to String cannot fail");
1245 }
1246 checkpoint_state = cp;
1247 }
1248 }
1249
1250 Ok(json!({
1251 "status": "success",
1252 "session_id": session_id,
1253 "digest": digest,
1254 "turns_included": included,
1255 "turns_total_scanned": turns.len(),
1256 "sections": groups.iter().map(|(t, items)| json!({"type": t, "count": items.len()})).collect::<Vec<_>>(),
1257 "checkpoint": checkpoint_state,
1258 }))
1259 }
1260 fn stats(&self) -> &ToolStats {
1261 &self.stats
1262 }
1263}
1264
1265fn capitalize(s: &str) -> String {
1267 let mut chars = s.chars();
1268 match chars.next() {
1269 Some(first) => first.to_uppercase().collect::<String>() + chars.as_str(),
1270 None => String::new(),
1271 }
1272}
1273
1274fn pluralize(s: &str) -> String {
1277 if let Some(stem) = s.strip_suffix('y') {
1278 format!("{stem}ies")
1279 } else {
1280 format!("{s}s")
1281 }
1282}
1283
1284pub struct SessionHandoffTool {
1286 store: Arc<MemoryStore>,
1287 stats: ToolStats,
1288 effects: EffectRow,
1289 search: Option<Arc<wm_memory::SearchEngine>>,
1290}
1291
1292impl SessionHandoffTool {
1293 #[must_use]
1294 pub fn new(store: Arc<MemoryStore>) -> Self {
1295 Self {
1296 store,
1297 stats: ToolStats::default(),
1298 effects: EffectRow {
1299 writes: vec![Resource::Galaxy("sessions".into())],
1300 ..Default::default()
1301 },
1302 search: None,
1303 }
1304 }
1305
1306 #[must_use]
1308 pub fn with_search(mut self, search: Option<Arc<wm_memory::SearchEngine>>) -> Self {
1309 self.search = search;
1310 self
1311 }
1312}
1313
1314#[async_trait]
1315impl Tool for SessionHandoffTool {
1316 fn name(&self) -> &str {
1317 "session.handoff"
1318 }
1319 fn gana(&self) -> Gana {
1320 Gana::StraddlingLegs
1321 }
1322 fn effects(&self) -> &EffectRow {
1323 &self.effects
1324 }
1325 fn input_schema(&self) -> Value {
1326 super::common::schema(
1327 &json!({
1328 "action": super::common::str_prop("transfer | accept | list"),
1329 "session_id": super::common::str_prop("transfer: session to hand off"),
1330 "message": super::common::str_prop("transfer: handoff note"),
1331 "handoff_id": super::common::str_prop("accept: handoff to accept"),
1332 }),
1333 &["action"],
1334 )
1335 }
1336 fn description(&self) -> &str {
1337 "Transfer or resume a session across devices (actions: transfer, accept, list). transfer: session_id (required) + message; accept: handoff_id; list: pending handoffs."
1338 }
1339 async fn call(&self, _ctx: &mut Context, args: Value) -> wm_core::Result<Value> {
1340 let action = args.get("action").and_then(Value::as_str).unwrap_or("list");
1341 match action {
1342 "transfer" => {
1343 let session_id =
1344 args.get("session_id")
1345 .and_then(Value::as_str)
1346 .ok_or_else(|| {
1347 wm_core::CoreError::InvalidArgs(
1348 "session_id required for transfer".into(),
1349 )
1350 })?;
1351 let message = args.get("message").and_then(Value::as_str).unwrap_or("");
1352 let turns = load_turns(&self.store, Some(session_id), 10_000, false)?;
1353 if turns.is_empty() {
1354 return Err(wm_core::CoreError::Tool(format!(
1355 "session {session_id} has no recorded turns"
1356 )));
1357 }
1358 let summary: Vec<Value> =
1359 turns.iter().map(|(_, v)| format_turn(v, false)).collect();
1360 let handoff_id = format!("handoff-{}", uuid::Uuid::new_v4());
1361 let mut mem = Memory::new(
1362 Galaxy::Sessions,
1363 json!({
1364 "type": "session_handoff",
1365 "handoff_id": handoff_id,
1366 "session_id": session_id,
1367 "message": message,
1368 "status": "pending",
1369 "turn_count": summary.len(),
1370 "summary": summary,
1371 "created_at": wm_core::time::now_unix_millis(),
1372 })
1373 .to_string(),
1374 );
1375 mem.metadata.tags = vec![
1376 "session".into(),
1377 "handoff".into(),
1378 format!("session:{session_id}"),
1379 ];
1380 mem.metadata.importance = 0.8;
1381 self.store.put(Galaxy::Sessions, &mem)?;
1382 super::common::index_memory(self.search.as_deref(), &mem);
1383 Ok(json!({
1384 "status": "success",
1385 "action": "transfer",
1386 "handoff_id": handoff_id,
1387 "session_id": session_id,
1388 "turn_count": summary.len(),
1389 }))
1390 }
1391 "accept" => {
1392 let handoff_id =
1393 args.get("handoff_id")
1394 .and_then(Value::as_str)
1395 .ok_or_else(|| {
1396 wm_core::CoreError::InvalidArgs("handoff_id required for accept".into())
1397 })?;
1398 let memories = self.store.scan_all(Galaxy::Sessions)?;
1399 let found = memories.iter().find(|m| {
1400 m.metadata.tags.contains(&"handoff".to_string())
1401 && m.content.contains(handoff_id)
1402 });
1403 let Some(mem) = found else {
1404 return Err(wm_core::CoreError::Tool(format!(
1405 "handoff {handoff_id} not found"
1406 )));
1407 };
1408 let mut updated = mem.clone();
1409 if let Ok(mut v) = serde_json::from_str::<Value>(&updated.content) {
1410 v["status"] = json!("accepted");
1411 updated.content = v.to_string();
1412 }
1413 self.store.put(Galaxy::Sessions, &updated)?;
1414 super::common::index_memory(self.search.as_deref(), &updated);
1415 Ok(json!({
1416 "status": "success",
1417 "action": "accept",
1418 "handoff_id": handoff_id,
1419 }))
1420 }
1421 "list" => {
1422 let memories = self.store.scan_all(Galaxy::Sessions)?;
1423 let handoffs: Vec<Value> = memories
1424 .iter()
1425 .filter(|m| m.metadata.tags.contains(&"handoff".to_string()))
1426 .filter_map(|m| serde_json::from_str::<Value>(&m.content).ok())
1427 .filter(|v| v.get("status").and_then(Value::as_str) == Some("pending"))
1428 .map(|v| {
1429 json!({
1430 "handoff_id": v.get("handoff_id"),
1431 "session_id": v.get("session_id"),
1432 "message": v.get("message"),
1433 "turn_count": v.get("turn_count"),
1434 })
1435 })
1436 .collect();
1437 Ok(json!({
1438 "status": "success",
1439 "action": "list",
1440 "pending_count": handoffs.len(),
1441 "handoffs": handoffs,
1442 }))
1443 }
1444 other => Err(wm_core::CoreError::InvalidArgs(format!(
1445 "unknown session.handoff action: {other}"
1446 ))),
1447 }
1448 }
1449 fn stats(&self) -> &ToolStats {
1450 &self.stats
1451 }
1452}
1453
1454#[must_use]
1456pub struct SessionExportTool {
1462 store: Arc<MemoryStore>,
1463 stats: ToolStats,
1464 effects: EffectRow,
1465}
1466
1467impl SessionExportTool {
1468 pub fn new(store: Arc<MemoryStore>) -> Self {
1469 Self {
1470 store,
1471 stats: ToolStats::default(),
1472 effects: EffectRow::read_only(vec![Resource::Galaxy("sessions".into())]),
1473 }
1474 }
1475}
1476
1477#[async_trait]
1478impl Tool for SessionExportTool {
1479 fn name(&self) -> &str {
1480 "session.export"
1481 }
1482 fn gana(&self) -> Gana {
1483 Gana::StraddlingLegs
1484 }
1485 fn effects(&self) -> &EffectRow {
1486 &self.effects
1487 }
1488 fn input_schema(&self) -> Value {
1489 super::common::schema(
1490 &json!({
1491 "session_id": super::common::str_prop("Session to export (default: most recent)"),
1492 "path": super::common::str_prop("Write JSONL to this file instead of returning inline"),
1493 }),
1494 &[],
1495 )
1496 }
1497 fn description(&self) -> &str {
1498 "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)."
1499 }
1500 async fn call(&self, _ctx: &mut Context, args: Value) -> wm_core::Result<Value> {
1501 let session_id = match args.get("session_id").and_then(Value::as_str) {
1502 Some(sid) if !sid.is_empty() => sid.to_string(),
1503 _ => self
1504 .store
1505 .scan_all(Galaxy::Sessions)?
1506 .iter()
1507 .filter(|m| m.metadata.tags.contains(&"start".to_string()))
1508 .max_by_key(|m| m.metadata.created_at)
1509 .map(|m| m.metadata.id.to_string())
1510 .ok_or_else(|| {
1511 wm_core::CoreError::Tool("no session found — run session.start first".into())
1512 })?,
1513 };
1514
1515 let mut members: Vec<Memory> = self
1519 .store
1520 .scan_all(Galaxy::Sessions)?
1521 .into_iter()
1522 .filter(|m| m.metadata.id.to_string() == session_id || m.content.contains(&session_id))
1523 .collect();
1524 members.sort_by_key(|m| m.metadata.created_at);
1525
1526 let mut jsonl = String::new();
1527 let header = wm_memory::envelope::EnvelopeHeader::new("session_export", members.len());
1533 jsonl.push_str(&header.header_line());
1534 jsonl.push('\n');
1535 for m in &members {
1536 let line = serde_json::to_string(m)
1537 .map_err(|e| wm_core::CoreError::Tool(format!("export serialize: {e}")))?;
1538 jsonl.push_str(&line);
1539 jsonl.push('\n');
1540 }
1541
1542 let path_arg = args
1543 .get("path")
1544 .and_then(Value::as_str)
1545 .filter(|s| !s.is_empty());
1546 if let Some(dest) = path_arg {
1547 std::fs::write(dest, &jsonl)
1548 .map_err(|e| wm_core::CoreError::Tool(format!("export write {dest}: {e}")))?;
1549 Ok(json!({
1550 "status": "success",
1551 "session_id": session_id,
1552 "records": members.len(),
1553 "path": dest,
1554 }))
1555 } else {
1556 Ok(json!({
1557 "status": "success",
1558 "session_id": session_id,
1559 "records": members.len(),
1560 "jsonl": jsonl,
1561 }))
1562 }
1563 }
1564 fn stats(&self) -> &ToolStats {
1565 &self.stats
1566 }
1567}
1568
1569pub struct SessionImportTool {
1580 store: Arc<MemoryStore>,
1581 search: Option<Arc<wm_memory::SearchEngine>>,
1582 stats: ToolStats,
1583 effects: EffectRow,
1584}
1585
1586impl SessionImportTool {
1587 pub fn new(store: Arc<MemoryStore>, search: Option<Arc<wm_memory::SearchEngine>>) -> Self {
1588 Self {
1589 store,
1590 search,
1591 stats: ToolStats::default(),
1592 effects: EffectRow {
1593 writes: vec![Resource::Galaxy("sessions".into())],
1594 ..Default::default()
1595 },
1596 }
1597 }
1598}
1599
1600#[async_trait]
1601impl Tool for SessionImportTool {
1602 fn name(&self) -> &str {
1603 "session.import"
1604 }
1605 fn gana(&self) -> Gana {
1606 Gana::StraddlingLegs
1607 }
1608 fn effects(&self) -> &EffectRow {
1609 &self.effects
1610 }
1611 fn input_schema(&self) -> Value {
1612 super::common::schema(
1613 &json!({
1614 "path": super::common::str_prop("Read JSONL from this file"),
1615 "jsonl": super::common::str_prop("Or pass the export payload inline"),
1616 }),
1617 &[],
1618 )
1619 }
1620 fn description(&self) -> &str {
1621 "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."
1622 }
1623 async fn call(&self, _ctx: &mut Context, args: Value) -> wm_core::Result<Value> {
1624 let payload = match args
1625 .get("path")
1626 .and_then(Value::as_str)
1627 .filter(|s| !s.is_empty())
1628 {
1629 Some(path) => std::fs::read_to_string(path)
1630 .map_err(|e| wm_core::CoreError::Tool(format!("import read {path}: {e}")))?,
1631 None => args
1632 .get("jsonl")
1633 .and_then(Value::as_str)
1634 .filter(|s| !s.is_empty())
1635 .ok_or_else(|| {
1636 wm_core::CoreError::InvalidArgs(
1637 "provide either 'path' or inline 'jsonl'".into(),
1638 )
1639 })?
1640 .to_string(),
1641 };
1642
1643 let mut envelope: Option<wm_memory::envelope::EnvelopeHeader> = None;
1648 let mut record_lines: Vec<&str> = Vec::new();
1649 let mut header_consumed = false;
1650 for line in payload.lines() {
1651 if line.trim().is_empty() {
1652 continue;
1653 }
1654 if !header_consumed {
1655 header_consumed = true;
1656 match wm_memory::envelope::read_header_line(line) {
1657 wm_memory::envelope::HeaderRead::Header(h) => {
1658 envelope = Some(h);
1659 continue;
1660 }
1661 wm_memory::envelope::HeaderRead::Refused(msg) => {
1662 return Err(wm_core::CoreError::Tool(msg));
1663 }
1664 wm_memory::envelope::HeaderRead::NotAHeader => {}
1665 }
1666 }
1667 record_lines.push(line);
1668 }
1669
1670 let readonly_engine = self.search.as_ref().is_some_and(|s| s.is_readonly());
1675 if readonly_engine {
1676 tracing::warn!(
1677 "session.import running against a read-only search engine — records land \
1678 in LMDB unindexed; they become searchable at the next writable startup \
1679 (heal_index_drift)"
1680 );
1681 }
1682 let mut writer_slot = match (&self.search, readonly_engine) {
1683 (Some(s), false) => s.writer().ok(),
1684 _ => None,
1685 };
1686
1687 let mut imported = 0usize;
1688 let mut indexed = 0usize;
1689 let mut skipped = 0usize;
1690 let mut session_ids: Vec<String> = Vec::new();
1691 for (lineno, line) in record_lines.iter().enumerate() {
1692 let mem: Memory = match serde_json::from_str(line) {
1693 Ok(m) => m,
1694 Err(e) => {
1695 skipped += 1;
1696 tracing::warn!(line = lineno + 1, error = %e, "skipping unparseable export line");
1697 continue;
1698 }
1699 };
1700 if let Ok(parsed) = serde_json::from_str::<Value>(&mem.content) {
1701 if let Some(sid) = parsed.get("session_id").and_then(Value::as_str) {
1702 if !session_ids.iter().any(|s| s == sid) {
1703 session_ids.push(sid.to_string());
1704 }
1705 }
1706 }
1707 if let (Some(search), Some(writer)) = (&self.search, writer_slot.as_mut()) {
1711 let id_str = mem.metadata.id.to_string();
1712 let _ = search.delete_document(writer, &id_str);
1713 match search.add_document(
1714 writer,
1715 &id_str,
1716 mem.metadata.galaxy.db_name(),
1717 &mem.content,
1718 &mem.metadata.tags,
1719 mem.metadata.created_at.timestamp(),
1720 ) {
1721 Ok(()) => indexed += 1,
1722 Err(e) => {
1723 tracing::warn!(id = %id_str, error = %e, "import index add failed (LMDB record kept)");
1724 }
1725 }
1726 }
1727 self.store.put(Galaxy::Sessions, &mem)?;
1728 imported += 1;
1729 }
1730
1731 if let Some(search) = &self.search {
1732 if let Some(mut writer) = writer_slot {
1733 search
1734 .commit(&mut writer)
1735 .map_err(|e| wm_core::CoreError::Tool(format!("import index commit: {e}")))?;
1736 }
1737 }
1738
1739 let mut warnings: Vec<String> = Vec::new();
1740 if let Some(h) = &envelope {
1741 if h.count != imported {
1742 let msg = format!(
1743 "envelope declares count {} but {} records imported",
1744 h.count, imported
1745 );
1746 tracing::warn!("{msg}");
1747 warnings.push(msg);
1748 }
1749 }
1750
1751 let envelope_info = envelope.as_ref().map(|h| {
1752 json!({
1753 "format_version": h.format_version,
1754 "kind": h.kind,
1755 "generator": h.generator,
1756 "created_at": h.created_at,
1757 "declared_count": h.count,
1758 })
1759 });
1760
1761 Ok(json!({
1762 "status": "success",
1763 "imported": imported,
1764 "skipped": skipped,
1765 "session_ids": session_ids,
1766 "indexed": indexed,
1767 "envelope": envelope_info,
1768 "warnings": warnings,
1769 }))
1770 }
1771 fn stats(&self) -> &ToolStats {
1772 &self.stats
1773 }
1774}
1775
1776pub fn register_session_ops(
1777 registry: &wm_dispatch::ToolRegistry,
1778 store: &Arc<MemoryStore>,
1779 search: Option<Arc<wm_memory::SearchEngine>>,
1780) -> wm_dispatch::ToolRegistry {
1781 registry
1782 .register(Arc::new(
1783 SessionRecordTool::new(store.clone()).with_search(search.clone()),
1784 ))
1785 .register(Arc::new(SessionReplayTool::new(store.clone())))
1786 .register(Arc::new(SessionContinuityTool::new(store.clone())))
1787 .register(Arc::new(
1788 SessionHandoffTool::new(store.clone()).with_search(search.clone()),
1789 ))
1790 .register(Arc::new(SessionExportTool::new(store.clone())))
1791 .register(Arc::new(SessionImportTool::new(store.clone(), search)))
1792}
1793
1794#[cfg(test)]
1795mod tests {
1796 use super::*;
1797
1798 fn test_store() -> Arc<MemoryStore> {
1799 let dir = tempfile::tempdir().unwrap();
1800 let path = dir.path().join("lmdb");
1801 std::fs::create_dir_all(&path).unwrap();
1802 Arc::new(MemoryStore::open_default(path).unwrap())
1803 }
1804
1805 fn start_session(store: &MemoryStore) -> String {
1806 let mut mem = Memory::new(
1807 Galaxy::Sessions,
1808 json!({"type": "session_start"}).to_string(),
1809 );
1810 mem.metadata.tags = vec!["session".into(), "start".into()];
1811 store.put(Galaxy::Sessions, &mem).unwrap();
1812 mem.metadata.id.to_string()
1813 }
1814
1815 fn start_session_aged(store: &MemoryStore, age_secs: i64) -> String {
1818 let mut mem = Memory::new(
1819 Galaxy::Sessions,
1820 json!({"type": "session_start"}).to_string(),
1821 );
1822 mem.metadata.tags = vec!["session".into(), "start".into()];
1823 mem.metadata.created_at = chrono::Utc::now() - chrono::Duration::seconds(age_secs);
1824 store.put(Galaxy::Sessions, &mem).unwrap();
1825 mem.metadata.id.to_string()
1826 }
1827
1828 fn record_aged_turn(store: &MemoryStore, sid: &str, age_days: u32, content: &str) {
1830 let mut mem = Memory::new(
1831 Galaxy::Sessions,
1832 json!({
1833 "type": "session_turn",
1834 "session_id": sid,
1835 "role": "ai",
1836 "turn_type": "decision",
1837 "importance": 0.9,
1838 "content": content,
1839 })
1840 .to_string(),
1841 );
1842 mem.metadata.tags = vec!["session".into(), "turn".into(), format!("session:{sid}")];
1843 mem.metadata.created_at = Utc::now() - chrono::Duration::days(i64::from(age_days));
1844 store.put(Galaxy::Sessions, &mem).unwrap();
1845 }
1846
1847 #[tokio::test]
1848 async fn replay_time_filters_since_and_until() {
1849 let store = test_store();
1850 let sid = start_session(&store);
1851 record_aged_turn(&store, &sid, 3, "three days ago");
1852 record_aged_turn(&store, &sid, 2, "two days ago");
1853 record_aged_turn(&store, &sid, 1, "yesterday");
1854 record_aged_turn(&store, &sid, 0, "today");
1855
1856 let replay = SessionReplayTool::new(store);
1857 let mut ctx = Context::default();
1858
1859 let two_days_ago = (Utc::now() - chrono::Duration::days(2)).format("%Y-%m-%d");
1861 let v = replay
1862 .call(
1863 &mut ctx,
1864 json!({"session_id": sid, "since": two_days_ago.to_string()}),
1865 )
1866 .await
1867 .unwrap();
1868 assert_eq!(v["count"], 3, "since=date keeps day-of + later: {v}");
1869
1870 let until_epoch = (Utc::now() - chrono::Duration::hours(23)).timestamp();
1872 let v = replay
1873 .call(&mut ctx, json!({"session_id": sid, "until": until_epoch}))
1874 .await
1875 .unwrap();
1876 assert_eq!(
1877 v["count"], 3,
1878 "until=epoch(23h ago) keeps the three older turns: {v}"
1879 );
1880
1881 let since = (Utc::now() - chrono::Duration::hours(60)).to_rfc3339();
1883 let until = (Utc::now() - chrono::Duration::hours(12)).to_rfc3339();
1884 let v = replay
1885 .call(
1886 &mut ctx,
1887 json!({"session_id": sid, "since": since, "until": until}),
1888 )
1889 .await
1890 .unwrap();
1891 assert_eq!(v["count"], 2, "window keeps two/two-days-ago turns: {v}");
1892 for turn in v["turns"].as_array().unwrap() {
1893 assert_ne!(
1894 turn["content"], "today",
1895 "time filters must exclude out-of-window turns"
1896 );
1897 }
1898
1899 assert!(
1901 replay
1902 .call(&mut ctx, json!({"session_id": sid, "since": "not-a-date"}))
1903 .await
1904 .is_err()
1905 );
1906 }
1907
1908 #[tokio::test]
1909 async fn replay_omitted_session_id_uses_latest_session_start() {
1910 let store = test_store();
1911 let older = start_session_aged(&store, 120);
1915 record_aged_turn(&store, &older, 0, "older-session-only");
1916 let latest = start_session_aged(&store, 60);
1917 record_aged_turn(&store, &latest, 0, "latest-session-only");
1918
1919 let replay = SessionReplayTool::new(store);
1920 let mut ctx = Context::default();
1921
1922 let omitted = replay
1923 .call(&mut ctx, json!({"mode": "full"}))
1924 .await
1925 .unwrap();
1926 assert_eq!(omitted["session_id"], latest);
1927 assert_eq!(
1928 omitted["count"], 1,
1929 "omitted id must not combine sessions: {omitted}"
1930 );
1931 assert_eq!(omitted["turns"][0]["content"], "latest-session-only");
1932
1933 let explicit_older = replay
1934 .call(&mut ctx, json!({"session_id": older, "mode": "full"}))
1935 .await
1936 .unwrap();
1937 assert_eq!(explicit_older["session_id"], older);
1938 assert_eq!(explicit_older["count"], 1);
1939 assert_eq!(explicit_older["turns"][0]["content"], "older-session-only");
1940 }
1941
1942 #[tokio::test]
1943 async fn replay_without_any_session_remains_truthfully_empty() {
1944 let replay = SessionReplayTool::new(test_store());
1945 let mut ctx = Context::default();
1946
1947 let value = replay.call(&mut ctx, json!({})).await.unwrap();
1948 assert_eq!(value["status"], "success");
1949 assert_eq!(value["session_id"], Value::Null);
1950 assert_eq!(value["count"], 0);
1951 assert_eq!(value["turns"], json!([]));
1952 }
1953
1954 #[tokio::test]
1955 async fn replay_omitted_id_does_not_combine_orphan_turns_without_a_start() {
1956 let store = test_store();
1957 record_aged_turn(&store, "orphan-a", 0, "orphan-a-only");
1961 record_aged_turn(&store, "orphan-b", 0, "orphan-b-only");
1962
1963 let replay = SessionReplayTool::new(store);
1964 let mut ctx = Context::default();
1965
1966 let omitted = replay.call(&mut ctx, json!({})).await.unwrap();
1967 assert_eq!(omitted["session_id"], Value::Null);
1968 assert_eq!(
1969 omitted["count"], 0,
1970 "omitted id must not combine orphans: {omitted}"
1971 );
1972 assert_eq!(omitted["turns"], json!([]));
1973
1974 let explicit = replay
1975 .call(&mut ctx, json!({"session_id": "orphan-a"}))
1976 .await
1977 .unwrap();
1978 assert_eq!(explicit["session_id"], "orphan-a");
1979 assert_eq!(explicit["count"], 1);
1980 assert_eq!(explicit["turns"][0]["content"], "orphan-a-only");
1981 }
1982
1983 #[tokio::test]
1984 async fn continuity_respects_since_filter() {
1985 let store = test_store();
1986 let sid1 = start_session(&store);
1987 record_aged_turn(&store, &sid1, 5, "ancient decision");
1988 record_aged_turn(&store, &sid1, 0, "fresh decision");
1989 let sid2 = start_session(&store);
1990
1991 let continuity = SessionContinuityTool::new(store);
1992 let mut ctx = Context::default();
1993 let cutoff = (Utc::now() - chrono::Duration::days(1))
1994 .format("%Y-%m-%d")
1995 .to_string();
1996 let v = continuity
1997 .call(
1998 &mut ctx,
1999 json!({"current_session_id": sid2, "since": cutoff, "n": 10}),
2000 )
2001 .await
2002 .unwrap();
2003 assert_eq!(v["count"], 1, "only the fresh turn is in range: {v}");
2004 assert_eq!(v["turns"][0]["content"], "fresh decision");
2005
2006 let all = continuity
2008 .call(&mut ctx, json!({"current_session_id": sid2, "n": 10}))
2009 .await
2010 .unwrap();
2011 assert_eq!(all["count"], 2);
2012 }
2013
2014 #[tokio::test]
2015 async fn continuity_skips_empty_newest_session() {
2016 let store = test_store();
2020 let sid1 = start_session(&store);
2021 record_aged_turn(
2022 &store,
2023 &sid1,
2024 0,
2025 "Decision: use SQLite for the report cache",
2026 );
2027 let sid2 = start_session(&store); let continuity = SessionContinuityTool::new(store);
2030 let mut ctx = Context::default();
2031 let v = continuity.call(&mut ctx, json!({"n": 5})).await.unwrap();
2032 assert_ne!(
2033 v["previous_session"], sid2,
2034 "empty newest session must be skipped: {v}"
2035 );
2036 assert_eq!(v["previous_session"], sid1);
2037 assert_eq!(v["count"], 1);
2038 assert!(
2039 v["turns"][0]["content"]
2040 .as_str()
2041 .unwrap()
2042 .contains("SQLite")
2043 );
2044 }
2045
2046 #[tokio::test]
2047 async fn digest_groups_by_type_and_respects_importance_floor() {
2048 let store = test_store();
2049 let sid = start_session(&store);
2050 for (turn_type, importance, content) in [
2052 ("summary", 0.6, "wrapped up"),
2053 ("decision", 0.9, "picked architecture CO over alternatives"),
2054 ("error", 0.95, "startForce root cause found"),
2055 ("breakthrough", 0.85, "watch resolution insight"),
2056 ("message", 0.3, "low-value chatter"),
2057 ] {
2058 let mut mem = Memory::new(
2059 Galaxy::Sessions,
2060 json!({
2061 "type": "session_turn",
2062 "session_id": sid,
2063 "role": "ai",
2064 "turn_type": turn_type,
2065 "importance": importance,
2066 "content": content,
2067 })
2068 .to_string(),
2069 );
2070 mem.metadata.tags = vec!["session".into(), "turn".into()];
2071 store.put(Galaxy::Sessions, &mem).unwrap();
2072 }
2073
2074 let tool = SessionDigestTool::new(store);
2075 let mut ctx = Context::default();
2076 let v = tool
2077 .call(&mut ctx, json!({"session_id": sid, "min_importance": 0.5}))
2078 .await
2079 .unwrap();
2080
2081 assert_eq!(v["status"], "success");
2082 let digest = v["digest"].as_str().unwrap();
2083 for expected in [
2085 "startForce root cause found",
2086 "picked architecture CO over alternatives",
2087 "watch resolution insight",
2088 ] {
2089 assert!(
2090 digest.contains(expected),
2091 "digest must contain '{expected}': {digest}"
2092 );
2093 }
2094 assert!(!digest.contains("low-value chatter"), "got: {digest}");
2096 let d = digest.find("## Decisions").unwrap();
2098 let b = digest.find("## Breakthroughs").unwrap();
2099 let e = digest.find("## Errors").unwrap();
2100 let s = digest.find("## Summaries").unwrap();
2101 assert!(
2102 d < b && b < e && e < s,
2103 "sections must follow canonical order: {digest}"
2104 );
2105 assert_eq!(v["turns_included"], 4);
2106
2107 let strict = tool
2109 .call(&mut ctx, json!({"session_id": sid, "min_importance": 0.9}))
2110 .await
2111 .unwrap();
2112 let strict_digest = strict["digest"].as_str().unwrap();
2113 assert!(strict_digest.contains("startForce"));
2114 assert!(!strict_digest.contains("watch resolution insight"));
2115 }
2116
2117 #[tokio::test]
2118 async fn digest_appends_checkpoint_state() {
2119 let store = test_store();
2120 let sid = start_session(&store);
2121
2122 let mut cp = Memory::new(
2124 Galaxy::Sessions,
2125 json!({
2126 "type": "checkpoint",
2127 "session_id": sid,
2128 "label": "wrap",
2129 "data": {},
2130 "handoff": {
2131 "git": {
2132 "commit": "abc1234",
2133 "branch": "main",
2134 "dirty_count": 2
2135 },
2136 "tests_green": true,
2137 "next_queue": ["first task", "second task"],
2138 "open_flags": ["flaky probe"]
2139 }
2140 })
2141 .to_string(),
2142 );
2143 cp.metadata.tags = vec!["session".into(), "checkpoint".into()];
2144 store.put(Galaxy::Sessions, &cp).unwrap();
2145
2146 let tool = SessionDigestTool::new(store);
2147 let mut ctx = Context::default();
2148 let v = tool
2149 .call(&mut ctx, json!({"session_id": sid}))
2150 .await
2151 .unwrap();
2152
2153 let digest = v["digest"].as_str().unwrap();
2154 assert!(digest.contains("## Checkpoint state"), "got: {digest}");
2155 assert!(digest.contains("abc1234"));
2156 assert!(digest.contains("first task → second task"));
2157 assert!(digest.contains("flaky probe"));
2158 assert_eq!(v["checkpoint"]["git"]["branch"], "main");
2159 }
2160
2161 #[tokio::test]
2162 async fn supersedes_hides_old_turn_until_requested() {
2163 let store = test_store();
2167 let sid = start_session(&store);
2168 let record = SessionRecordTool::new(store.clone());
2169 let mut ctx = Context::default();
2170
2171 let first = record
2172 .call(
2173 &mut ctx,
2174 json!({"role": "ai", "turn_type": "decision", "importance": 0.9,
2175 "content": "perf: 240ms", "session_id": sid}),
2176 )
2177 .await
2178 .unwrap();
2179 let old_id = first["memory_id"].as_str().unwrap().to_string();
2180
2181 let second = record
2182 .call(
2183 &mut ctx,
2184 json!({"role": "ai", "turn_type": "decision", "importance": 0.9,
2185 "content": "perf revised: 180ms after warm cache",
2186 "session_id": sid, "supersedes": old_id}),
2187 )
2188 .await
2189 .unwrap();
2190 assert_eq!(second["status"], "success");
2191 let _new_id = second["memory_id"].as_str().unwrap().to_string();
2192
2193 let replay = SessionReplayTool::new(store.clone());
2195 let v = replay
2196 .call(&mut ctx, json!({"session_id": sid}))
2197 .await
2198 .unwrap();
2199 assert_eq!(
2200 v["count"], 1,
2201 "superseded turn must be hidden by default: {v}"
2202 );
2203 assert_eq!(
2204 v["turns"][0]["content"],
2205 "perf revised: 180ms after warm cache"
2206 );
2207
2208 let with_history = replay
2210 .call(
2211 &mut ctx,
2212 json!({"session_id": sid, "include_superseded": true}),
2213 )
2214 .await
2215 .unwrap();
2216 assert_eq!(with_history["count"], 2, "got: {with_history}");
2217
2218 let sid2 = start_session(&store);
2220 let continuity = SessionContinuityTool::new(store.clone());
2221 let c = continuity
2222 .call(&mut ctx, json!({"current_session_id": sid2}))
2223 .await
2224 .unwrap();
2225 assert_eq!(c["count"], 1, "continuity must skip superseded turns: {c}");
2226 assert_eq!(
2227 c["turns"][0]["content"],
2228 "perf revised: 180ms after warm cache"
2229 );
2230
2231 let digest = SessionDigestTool::new(store);
2232 let d = digest
2233 .call(&mut ctx, json!({"session_id": sid, "min_importance": 0.5}))
2234 .await
2235 .unwrap();
2236 let digest_text = d["digest"].as_str().unwrap();
2237 assert!(digest_text.contains("180ms"), "got: {digest_text}");
2238 assert!(
2239 !digest_text.contains("240ms"),
2240 "superseded claim must not leak: {digest_text}"
2241 );
2242 }
2243
2244 #[tokio::test]
2245 async fn export_import_roundtrip_preserves_history() {
2246 let store_a = test_store();
2249 let sid = start_session(&store_a);
2250 let record = SessionRecordTool::new(store_a.clone());
2251 let mut ctx = Context::default();
2252 record_aged_turn(&store_a, &sid, 2, "day-one decision");
2253 record_aged_turn(&store_a, &sid, 1, "day-two decision");
2254 let first = record
2255 .call(
2256 &mut ctx,
2257 json!({"role": "ai", "turn_type": "decision", "importance": 0.9,
2258 "content": "original claim", "session_id": sid}),
2259 )
2260 .await
2261 .unwrap();
2262 record
2263 .call(
2264 &mut ctx,
2265 json!({"role": "ai", "turn_type": "decision", "importance": 0.9,
2266 "content": "corrected claim",
2267 "session_id": sid,
2268 "supersedes": first["memory_id"].as_str().unwrap()}),
2269 )
2270 .await
2271 .unwrap();
2272
2273 let export = SessionExportTool::new(store_a.clone());
2275 {
2277 let mut marker: Memory = store_a
2278 .scan_all(Galaxy::Sessions)
2279 .unwrap()
2280 .into_iter()
2281 .find(|m| m.metadata.tags.contains(&"start".to_string()))
2282 .unwrap();
2283 marker.metadata.title = Some("The Big Decision".to_string());
2284 marker.metadata.topic = Some("v8-slices".to_string());
2285 store_a.put(Galaxy::Sessions, &marker).unwrap();
2286 }
2287 let exported = export
2288 .call(&mut ctx, json!({"session_id": sid}))
2289 .await
2290 .unwrap();
2291 assert_eq!(exported["status"], "success");
2292 let jsonl = exported["jsonl"].as_str().unwrap();
2293 assert_eq!(
2295 exported["records"], 5,
2296 "start + 2 aged + 2 claims: {exported}"
2297 );
2298 assert_eq!(jsonl.lines().count(), 6);
2300 let header_line = jsonl.lines().next().unwrap();
2301 match wm_memory::envelope::read_header_line(header_line) {
2302 wm_memory::envelope::HeaderRead::Header(h) => {
2303 assert_eq!(h.kind, "session_export");
2304 assert_eq!(h.count, 5);
2305 }
2306 other => panic!("first line must be the envelope header, got {other:?}"),
2307 }
2308
2309 let store_b = test_store();
2311 let import = SessionImportTool::new(store_b.clone(), None);
2312 let imported = import
2313 .call(&mut ctx, json!({"jsonl": jsonl}))
2314 .await
2315 .unwrap();
2316 assert_eq!(imported["imported"], 5, "got: {imported}");
2317 assert_eq!(imported["skipped"], 0);
2318 assert_eq!(imported["session_ids"], json!([sid]));
2319 assert_eq!(imported["envelope"]["format_version"], 2);
2321 assert_eq!(imported["envelope"]["declared_count"], 5);
2322 assert_eq!(imported["warnings"], json!([]));
2323 let marker_b = store_b
2325 .scan_all(Galaxy::Sessions)
2326 .unwrap()
2327 .into_iter()
2328 .find(|m| m.metadata.tags.contains(&"start".to_string()))
2329 .unwrap();
2330 assert_eq!(marker_b.metadata.title.as_deref(), Some("The Big Decision"));
2331 assert_eq!(marker_b.metadata.topic.as_deref(), Some("v8-slices"));
2332
2333 let replay_b = SessionReplayTool::new(store_b.clone());
2335 let v = replay_b
2336 .call(&mut ctx, json!({"session_id": sid}))
2337 .await
2338 .unwrap();
2339 assert_eq!(v["count"], 3, "two aged turns + correction: {v}");
2340 let contents: Vec<&str> = v["turns"]
2341 .as_array()
2342 .unwrap()
2343 .iter()
2344 .filter_map(|t| t["content"].as_str())
2345 .collect();
2346 assert!(contents.contains(&"day-one decision"));
2347 assert!(contents.contains(&"corrected claim"));
2348 assert!(!contents.contains(&"original claim"));
2349
2350 let full = replay_b
2352 .call(
2353 &mut ctx,
2354 json!({"session_id": sid, "include_superseded": true}),
2355 )
2356 .await
2357 .unwrap();
2358 assert_eq!(full["count"], 4);
2359
2360 let new_sid = start_session(&store_b);
2362 let continuity = SessionContinuityTool::new(store_b);
2363 let c = continuity
2364 .call(
2365 &mut ctx,
2366 json!({"current_session_id": new_sid, "since":
2367 (Utc::now() - chrono::Duration::days(3)).format("%Y-%m-%d").to_string()}),
2368 )
2369 .await
2370 .unwrap();
2371 assert_eq!(
2372 c["previous_session"], sid,
2373 "import must preserve created_at so recency resolution works"
2374 );
2375 assert_eq!(c["count"], 3);
2376 }
2377
2378 #[tokio::test]
2379 async fn import_rejects_missing_payload() {
2380 let store = test_store();
2381 let tool = SessionImportTool::new(store, None);
2382 let mut ctx = Context::default();
2383 assert!(tool.call(&mut ctx, json!({})).await.is_err());
2384 }
2385
2386 #[tokio::test]
2387 async fn import_refuses_newer_envelope_format() {
2388 let store = test_store();
2389 let mut ctx = Context::default();
2390 let header = wm_memory::envelope::EnvelopeHeader {
2391 format_version: wm_memory::envelope::ENVELOPE_FORMAT_VERSION + 1,
2392 kind: "session_export".into(),
2393 created_at: chrono::Utc::now().to_rfc3339(),
2394 count: 1,
2395 generator: "wm 99.0.0".into(),
2396 };
2397 let record =
2398 serde_json::to_string(&Memory::new(Galaxy::Sessions, "future".into())).unwrap();
2399 let payload = format!("{}\n{record}\n", header.header_line());
2400 let tool = SessionImportTool::new(store, None);
2401 let result = tool.call(&mut ctx, json!({"jsonl": payload})).await;
2402 let err = format!("{:?}", result.unwrap_err());
2403 assert!(err.contains("newer than this build supports"), "{err}");
2404 }
2405
2406 #[tokio::test]
2410 async fn session_record_indexes_at_write_time() {
2411 let dir = tempfile::tempdir().unwrap();
2412 let lmdb = dir.path().join("lmdb");
2413 std::fs::create_dir_all(&lmdb).unwrap();
2414 let store = Arc::new(MemoryStore::open_default(&lmdb).unwrap());
2415 let tantivy = dir.path().join("tantivy");
2416 std::fs::create_dir_all(&tantivy).unwrap();
2417 let search = Arc::new(wm_memory::SearchEngine::open(&tantivy).unwrap());
2418
2419 let sid = start_session(&store);
2420 let mut ctx = Context::default();
2421 SessionRecordTool::new(store.clone())
2422 .with_search(Some(search.clone()))
2423 .call(
2424 &mut ctx,
2425 json!({"role": "ai", "turn_type": "decision", "importance": 0.7,
2426 "content": "amber lighthouse protocol engaged", "session_id": sid}),
2427 )
2428 .await
2429 .unwrap();
2430
2431 let docs = search.count_docs_in_galaxy("sessions").unwrap();
2432 assert!(
2433 docs >= 1,
2434 "session.record must index its write immediately (docs={docs})"
2435 );
2436 }
2437
2438 #[tokio::test]
2442 async fn import_indexes_tantivy_no_drift_even_on_reimport() {
2443 let dir = tempfile::tempdir().unwrap();
2444 let lmdb = dir.path().join("lmdb");
2445 std::fs::create_dir_all(&lmdb).unwrap();
2446 let store = Arc::new(MemoryStore::open_default(&lmdb).unwrap());
2447 let tantivy = dir.path().join("tantivy");
2448 std::fs::create_dir_all(&tantivy).unwrap();
2449 let search = Arc::new(wm_memory::SearchEngine::open(&tantivy).unwrap());
2450
2451 let store_a = test_store();
2453 let sid = start_session(&store_a);
2454 let mut ctx = Context::default();
2455 SessionRecordTool::new(store_a.clone())
2456 .call(
2457 &mut ctx,
2458 json!({"role": "ai", "turn_type": "decision", "importance": 0.9,
2459 "content": "kumquat governance ratchet engaged", "session_id": sid}),
2460 )
2461 .await
2462 .unwrap();
2463 let exported = SessionExportTool::new(store_a.clone())
2464 .call(&mut ctx, json!({"session_id": sid}))
2465 .await
2466 .unwrap();
2467 let jsonl = exported["jsonl"].as_str().unwrap().to_string();
2468
2469 let import = SessionImportTool::new(store.clone(), Some(search.clone()));
2471 for round in 1..=2 {
2472 let r = import
2473 .call(&mut ctx, json!({"jsonl": jsonl}))
2474 .await
2475 .unwrap();
2476 assert_eq!(r["status"], "success", "round {round}: {r}");
2477 assert_eq!(r["skipped"], 0);
2478 assert_eq!(
2479 r["indexed"], r["imported"],
2480 "round {round}: every record indexed: {r}"
2481 );
2482 }
2483
2484 let report = wm_memory::reindex::check_consistency(&store, &search);
2487 let drifted: Vec<_> = report
2488 .galaxies
2489 .iter()
2490 .filter(|g| g.drift)
2491 .map(|g| g.galaxy.clone())
2492 .collect();
2493 assert!(
2494 drifted.is_empty(),
2495 "import must leave zero index drift, drifted: {drifted:?}"
2496 );
2497
2498 let needle_id = store
2500 .scan_all(Galaxy::Sessions)
2501 .unwrap()
2502 .iter()
2503 .find(|m| m.content.contains("kumquat"))
2504 .unwrap()
2505 .metadata
2506 .id
2507 .to_string();
2508 let hits = search.search("kumquat governance ratchet", 10).unwrap();
2509 assert!(
2510 hits.iter().any(|h| h.memory_id == needle_id),
2511 "imported record must be searchable via the index: {hits:?}"
2512 );
2513 }
2514
2515 #[tokio::test]
2516 async fn record_then_replay_full() {
2517 let store = test_store();
2518 let sid = start_session(&store);
2519 let record = SessionRecordTool::new(store.clone());
2520 let mut ctx = Context::default();
2521 for (i, role) in [("user", "hello"), ("ai", "hi there")].iter().enumerate() {
2522 let r = record
2523 .call(
2524 &mut ctx,
2525 json!({"role": role.0, "content": role.1, "session_id": sid}),
2526 )
2527 .await
2528 .unwrap();
2529 assert_eq!(r["sequence"], i as u64 + 1);
2530 }
2531
2532 let replay = SessionReplayTool::new(store.clone());
2533 let v = replay
2534 .call(&mut ctx, json!({"mode": "full", "session_id": sid}))
2535 .await
2536 .unwrap();
2537 assert_eq!(v["count"], 2);
2538 assert_eq!(v["turns"][0]["content"], "hello");
2539 assert_eq!(v["turns"][1]["role"], "ai");
2540 }
2541
2542 #[tokio::test]
2543 async fn record_requires_content_and_valid_role() {
2544 let store = test_store();
2545 let sid = start_session(&store);
2546 let tool = SessionRecordTool::new(store);
2547 let mut ctx = Context::default();
2548 assert!(tool.call(&mut ctx, json!({})).await.is_err());
2549 assert!(
2550 tool.call(&mut ctx, json!({"role": "system", "content": "x"}))
2551 .await
2552 .is_err()
2553 );
2554 for bad in [json!(""), json!(" "), json!("\n\t")] {
2557 let err = tool
2558 .call(&mut ctx, json!({"content": bad, "session_id": sid}))
2559 .await
2560 .unwrap_err();
2561 assert!(
2562 err.to_string().contains("content"),
2563 "content={bad:?}: {err}"
2564 );
2565 }
2566 let err = tool
2567 .call(
2568 &mut ctx,
2569 json!({"content": "x", "session_id": sid, "turn_type": "observation"}),
2570 )
2571 .await
2572 .unwrap_err();
2573 assert!(err.to_string().contains("turn_type"), "{err}");
2574 let v = tool
2576 .call(
2577 &mut ctx,
2578 json!({"content": "x", "session_id": sid, "turn_type": "context"}),
2579 )
2580 .await
2581 .unwrap();
2582 assert_eq!(v["status"], "success", "{v}");
2583 }
2584
2585 #[tokio::test]
2589 async fn record_rejects_out_of_range_importance() {
2590 let store = test_store();
2591 let sid = start_session(&store);
2592 let record = SessionRecordTool::new(store);
2593 let mut ctx = Context::default();
2594 for bad in [json!(1.5), json!(999), json!(-0.25), json!("2.0")] {
2595 let err = record
2596 .call(
2597 &mut ctx,
2598 json!({"content": "x", "session_id": sid, "importance": bad}),
2599 )
2600 .await
2601 .unwrap_err();
2602 assert!(
2603 err.to_string().contains("importance"),
2604 "importance={bad} must be rejected, got: {err}"
2605 );
2606 }
2607 for good in [json!(0.0), json!(1.0), json!(0.75), json!("0.4")] {
2610 let v = record
2611 .call(
2612 &mut ctx,
2613 json!({"content": "x", "session_id": sid, "importance": good}),
2614 )
2615 .await
2616 .unwrap();
2617 assert_eq!(v["status"], "success", "{v}");
2618 }
2619 }
2620
2621 #[tokio::test]
2625 async fn record_stamps_provenance_from_role() {
2626 let store = test_store();
2627 let sid = start_session(&store);
2628 let record = SessionRecordTool::new(store.clone());
2629 let mut ctx = Context::default();
2630 let ai = record
2631 .call(
2632 &mut ctx,
2633 json!({"role": "ai", "content": "agent turn", "session_id": sid}),
2634 )
2635 .await
2636 .unwrap();
2637 let user = record
2638 .call(
2639 &mut ctx,
2640 json!({"role": "user", "content": "human turn", "session_id": sid}),
2641 )
2642 .await
2643 .unwrap();
2644
2645 let ai_mem = store
2646 .get(
2647 Galaxy::Sessions,
2648 uuid::Uuid::parse_str(ai["memory_id"].as_str().unwrap()).unwrap(),
2649 )
2650 .expect("ai turn stored")
2651 .expect("ai turn present");
2652 assert_eq!(ai_mem.metadata.source, "agent");
2653 assert!((ai_mem.metadata.source_trust - 0.7).abs() < 1e-5);
2654
2655 let user_mem = store
2656 .get(
2657 Galaxy::Sessions,
2658 uuid::Uuid::parse_str(user["memory_id"].as_str().unwrap()).unwrap(),
2659 )
2660 .expect("user turn stored")
2661 .expect("user turn present");
2662 assert_eq!(user_mem.metadata.source, "user");
2663 assert!((user_mem.metadata.source_trust - 1.0).abs() < f32::EPSILON);
2664 }
2665
2666 #[tokio::test]
2667 async fn continuity_returns_previous_session_tail() {
2668 let store = test_store();
2669 let sid1 = start_session(&store);
2670 let record = SessionRecordTool::new(store.clone());
2671 let mut ctx = Context::default();
2672 for i in 0..5 {
2673 record
2674 .call(
2675 &mut ctx,
2676 json!({"role": "user", "content": format!("turn {i}"), "session_id": sid1}),
2677 )
2678 .await
2679 .unwrap();
2680 }
2681 let sid2 = start_session(&store);
2682 let continuity = SessionContinuityTool::new(store);
2683 let v = continuity
2684 .call(&mut ctx, json!({"current_session_id": sid2, "n": 2}))
2685 .await
2686 .unwrap();
2687 assert_eq!(v["previous_session"], sid1);
2688 assert_eq!(v["count"], 2);
2689 assert_eq!(v["turns"][1]["content"], "turn 4");
2690 assert!(
2691 v.get("hint").is_none(),
2692 "non-empty continuity must not carry the scoping hint"
2693 );
2694 assert!(
2695 v["checkpoint"].is_null() && v["checkpoint_id"].is_null(),
2696 "no checkpoint recorded — the fields must be present but null: {v}"
2697 );
2698 }
2699
2700 #[tokio::test]
2701 async fn continuity_surfaces_latest_checkpoint_handoff() {
2702 let store = test_store();
2708 let sid1 = start_session(&store);
2709 let record = SessionRecordTool::new(store.clone());
2710 let mut ctx = Context::default();
2711 record
2712 .call(
2713 &mut ctx,
2714 json!({"role": "ai", "turn_type": "summary", "importance": 0.9,
2715 "content": "seam work done", "session_id": sid1}),
2716 )
2717 .await
2718 .unwrap();
2719
2720 let seed_checkpoint = |handoff: Value, age_secs: i64| {
2721 let mut cp = Memory::new(
2722 Galaxy::Sessions,
2723 json!({
2724 "type": "checkpoint",
2725 "session_id": sid1,
2726 "label": "wrap",
2727 "data": {},
2728 "handoff": handoff,
2729 })
2730 .to_string(),
2731 );
2732 cp.metadata.tags = vec!["session".into(), "checkpoint".into()];
2733 cp.metadata.created_at = Utc::now() - chrono::Duration::seconds(age_secs);
2734 store.put(Galaxy::Sessions, &cp).unwrap();
2735 cp.metadata.id.to_string()
2736 };
2737 seed_checkpoint(
2738 json!({"next_queue": ["stale task"], "open_flags": ["stale flag"]}),
2739 120,
2740 );
2741 let newest_id = seed_checkpoint(
2742 json!({
2743 "git": {"commit": "abc1234", "branch": "main", "dirty_count": 0},
2744 "tests_green": true,
2745 "next_queue": ["fuzz malformed headers", "verify seal"],
2746 "open_flags": ["header length cap unresolved"],
2747 }),
2748 0,
2749 );
2750
2751 let sid2 = start_session(&store);
2752 let continuity = SessionContinuityTool::new(store);
2753 let v = continuity
2754 .call(&mut ctx, json!({"current_session_id": sid2}))
2755 .await
2756 .unwrap();
2757 assert_eq!(v["previous_session"], sid1);
2758 assert_eq!(
2759 v["checkpoint"]["next_queue"][0], "fuzz malformed headers",
2760 "the latest checkpoint must win over an older one: {v}"
2761 );
2762 assert_eq!(
2763 v["checkpoint"]["open_flags"][0], "header length cap unresolved",
2764 "open flags must survive the handoff through continuity: {v}"
2765 );
2766 assert_eq!(v["checkpoint"]["git"]["commit"], "abc1234");
2767 assert_eq!(v["checkpoint"]["tests_green"], true);
2768 assert_eq!(v["checkpoint_id"].as_str().unwrap(), newest_id);
2769 }
2770
2771 #[tokio::test]
2772 async fn continuity_empty_store_discloses_project_scoping() {
2773 let store = test_store();
2778 let continuity = SessionContinuityTool::new(store);
2779 let mut ctx = Context::default();
2780 let v = continuity.call(&mut ctx, json!({})).await.unwrap();
2781
2782 assert_eq!(v["status"], "success");
2783 assert_eq!(v["count"], 0);
2784 let hint = v["hint"].as_str().expect("hint present on empty store");
2785 assert!(hint.contains("project-scoped"), "got: {hint}");
2786 assert!(hint.contains("opencode config"), "got: {hint}");
2787 assert!(hint.contains("GET /status"), "got: {hint}");
2788 assert!(hint.contains("store "), "got: {hint}");
2790 }
2791
2792 #[tokio::test]
2793 async fn record_defaults_to_newest_start_by_time_not_key_order() {
2794 let store = test_store();
2799 for age in [50_000, 40_000, 30_000, 20_000, 10_000] {
2800 start_session_aged(&store, age);
2801 }
2802 let newest = start_session_aged(&store, 0);
2803
2804 let record = SessionRecordTool::new(store);
2805 let mut ctx = Context::default();
2806 let r = record
2807 .call(&mut ctx, json!({"role": "ai", "content": "latest turn"}))
2808 .await
2809 .unwrap();
2810 assert_eq!(
2811 r["session_id"], newest,
2812 "record without explicit session_id must target the newest start by created_at"
2813 );
2814 }
2815
2816 #[tokio::test]
2817 async fn continuity_picks_newest_prior_by_time_not_key_order() {
2818 let store = test_store();
2822 for age in [40_000, 30_000, 20_000] {
2823 start_session_aged(&store, age);
2824 }
2825 let newest_prior = start_session_aged(&store, 10);
2826 let current = start_session_aged(&store, 0);
2827
2828 let continuity = SessionContinuityTool::new(store);
2829 let mut ctx = Context::default();
2830 let v = continuity
2831 .call(&mut ctx, json!({"current_session_id": current, "n": 1}))
2832 .await
2833 .unwrap();
2834 assert_eq!(
2835 v["previous_session"], newest_prior,
2836 "continuity must select the newest prior session by created_at"
2837 );
2838 }
2839
2840 #[tokio::test]
2841 async fn handoff_transfer_accept_list() {
2842 let store = test_store();
2843 let sid = start_session(&store);
2844 let record = SessionRecordTool::new(store.clone());
2845 let mut ctx = Context::default();
2846 record
2847 .call(
2848 &mut ctx,
2849 json!({"role": "ai", "content": "context", "session_id": sid}),
2850 )
2851 .await
2852 .unwrap();
2853
2854 let handoff = SessionHandoffTool::new(store.clone());
2855 let t = handoff
2856 .call(
2857 &mut ctx,
2858 json!({"action": "transfer", "session_id": sid, "message": "take over"}),
2859 )
2860 .await
2861 .unwrap();
2862 assert_eq!(t["status"], "success");
2863 let hid = t["handoff_id"].as_str().unwrap().to_string();
2864
2865 let list = handoff
2866 .call(&mut ctx, json!({"action": "list"}))
2867 .await
2868 .unwrap();
2869 assert_eq!(list["pending_count"], 1);
2870
2871 let a = handoff
2872 .call(&mut ctx, json!({"action": "accept", "handoff_id": hid}))
2873 .await
2874 .unwrap();
2875 assert_eq!(a["status"], "success");
2876
2877 let list2 = handoff
2878 .call(&mut ctx, json!({"action": "list"}))
2879 .await
2880 .unwrap();
2881 assert_eq!(list2["pending_count"], 0);
2882 }
2883
2884 #[tokio::test]
2885 async fn lossless_replay_binds_explicit_session_chunks_and_detects_stale_view() {
2886 let store = test_store();
2887 let older = start_session_aged(&store, 60);
2888 let newer = start_session_aged(&store, 0);
2889 let content = format!("prefix {} DISTINCT-FACT-AFTER-120", "é".repeat(9000));
2890 let mut turn = Memory::new(Galaxy::Sessions, json!({"type":"session_turn","session_id":older,"sequence":1,"timestamp":1_i64,"content":content}).to_string());
2891 turn.metadata.tags = vec!["session".into(), "turn".into(), format!("session:{older}")];
2892 store.put(Galaxy::Sessions, &turn).unwrap();
2893 let mut other = Memory::new(Galaxy::Sessions, json!({"type":"session_turn","session_id":newer,"sequence":1,"timestamp":1_i64,"content":"newer-only"}).to_string());
2894 other.metadata.tags = vec!["session".into(), "turn".into(), format!("session:{newer}")];
2895 store.put(Galaxy::Sessions, &other).unwrap();
2896 let replay = SessionReplayTool::new(store.clone());
2897 let mut ctx = Context::default();
2898 let mut response = replay
2899 .call(
2900 &mut ctx,
2901 json!({"mode":"lossless","session_id":older,"page_size":1,"max_wire_bytes":2048}),
2902 )
2903 .await
2904 .unwrap();
2905 assert_eq!(response["session_id"], older);
2906 assert!(response["records"][0].get("chunk").is_some());
2907 let mut bytes = Vec::new();
2908 loop {
2909 assert!(serde_json::to_vec(&response).unwrap().len() <= 2048);
2910 for record in response["records"].as_array().unwrap() {
2911 if let Some(chunk) = record.get("chunk") {
2912 assert_eq!(chunk["byte_offset"].as_u64().unwrap() as usize, bytes.len());
2913 bytes.extend(base64_decode(chunk["data_b64"].as_str().unwrap()).unwrap());
2914 } else {
2915 bytes.extend(record["content"].as_str().unwrap().as_bytes());
2916 }
2917 }
2918 if response["complete"] == true {
2919 assert!(response["next_cursor"].is_null());
2920 break;
2921 }
2922 let cursor = response["next_cursor"].as_str().unwrap();
2923 response = replay.call(&mut ctx, json!({"mode":"lossless","session_id":older,"page_size":1,"max_wire_bytes":2048,"cursor":cursor})).await.unwrap();
2924 }
2925 assert_eq!(String::from_utf8(bytes).unwrap(), content);
2926 let first = replay
2927 .call(
2928 &mut ctx,
2929 json!({"mode":"lossless","session_id":older,"page_size":1,"max_wire_bytes":2048}),
2930 )
2931 .await
2932 .unwrap();
2933 let stale_cursor = first["next_cursor"].as_str().unwrap().to_string();
2934 let mut appended = Memory::new(Galaxy::Sessions, json!({"type":"session_turn","session_id":older,"sequence":2,"timestamp":2_i64,"content":"later"}).to_string());
2935 appended.metadata.tags = vec!["session".into(), "turn".into(), format!("session:{older}")];
2936 store.put(Galaxy::Sessions, &appended).unwrap();
2937 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"));
2938 }
2939
2940 fn lossless_fixture_turn(sid: &str, sequence: u64, text: &str) -> Memory {
2941 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());
2942 m.metadata.tags = vec!["turn".into(), format!("session:{sid}")];
2943 m
2944 }
2945
2946 #[tokio::test]
2947 async fn lossless_strict_args_and_cursor_prevalidation_before_corrupt_scan() {
2948 let store = test_store();
2949 let sid = start_session(&store);
2950 store
2951 .put_raw(
2952 Galaxy::Sessions,
2953 uuid::Uuid::new_v4().as_bytes(),
2954 b"not-a-memory",
2955 )
2956 .unwrap();
2957 let replay = SessionReplayTool::new(store);
2958 for (key, value) in [
2959 ("include_superseded", json!("false")),
2960 ("page_size", json!(null)),
2961 ("page_size", json!(-1)),
2962 ("max_wire_bytes", json!("2048")),
2963 ("cursor", json!(false)),
2964 ("cursor", json!("a€")),
2965 ("cursor", json!("🔑")),
2966 ] {
2967 let mut args = json!({"mode":"lossless","session_id":sid});
2968 args[key] = value;
2969 let error = replay
2970 .call(&mut Context::default(), args)
2971 .await
2972 .unwrap_err()
2973 .to_string();
2974 assert!(
2975 error.contains("invalid_args") || error.contains("invalid_cursor"),
2976 "{key}: {error}"
2977 );
2978 assert!(!error.contains("incomplete scan"));
2979 }
2980 let error = replay
2981 .call(
2982 &mut Context::default(),
2983 json!({"mode":"lossless","session_id":sid}),
2984 )
2985 .await
2986 .unwrap_err()
2987 .to_string();
2988 assert!(error.contains("refusing incomplete scan"));
2989 }
2990
2991 #[tokio::test]
2992 async fn lossless_metadata_staleness_settings_seek_and_partial_resume() {
2993 let store = test_store();
2994 let sid = start_session(&store);
2995 let mut turn = lossless_fixture_turn(&sid, 1, &"€\n".repeat(2000));
2996 store.put(Galaxy::Sessions, &turn).unwrap();
2997 let replay = SessionReplayTool::new(store.clone());
2998 let args = json!({"mode":"lossless","session_id":sid,"page_size":1,"max_wire_bytes":2048});
2999 let first = replay
3000 .call(&mut Context::default(), args.clone())
3001 .await
3002 .unwrap();
3003 let cursor = first["next_cursor"].as_str().unwrap();
3004 for (key, value) in [("page_size", json!(2)), ("max_wire_bytes", json!(4096))] {
3005 let mut next = args.clone();
3006 next["cursor"] = json!(cursor);
3007 next[key] = value;
3008 assert!(
3009 replay
3010 .call(&mut Context::default(), next)
3011 .await
3012 .unwrap_err()
3013 .to_string()
3014 .contains("invalid_cursor")
3015 );
3016 }
3017 let (view, _, _) = parse_lossless_cursor(cursor, &sid, false, 1, 2048).unwrap();
3018 let mut seek = args.clone();
3020 seek["cursor"] = json!(lossless_cursor(
3021 &sid,
3022 false,
3023 1,
3024 2048,
3025 &view,
3026 0,
3027 turn.content.len() - 2
3028 ));
3029 let inner = serde_json::from_str::<Value>(&turn.content).unwrap()["content"]
3031 .as_str()
3032 .unwrap()
3033 .to_string();
3034 seek["cursor"] = json!(lossless_cursor(
3035 &sid,
3036 false,
3037 1,
3038 2048,
3039 &view,
3040 0,
3041 inner.len() - 2
3042 ));
3043 let tail = replay.call(&mut Context::default(), seek).await.unwrap();
3044 assert!(tail["records"][0].get("content").is_none());
3045 assert_eq!(
3046 base64_decode(tail["records"][0]["chunk"]["data_b64"].as_str().unwrap()).unwrap(),
3047 inner.as_bytes()[inner.len() - 2..]
3048 );
3049 assert!(tail["next_cursor"].is_null());
3050 let mut changed = serde_json::from_str::<Value>(&turn.content).unwrap();
3051 changed["role"] = json!("human");
3052 turn.content = changed.to_string();
3053 store.put(Galaxy::Sessions, &turn).unwrap();
3054 let mut next = args;
3055 next["cursor"] = json!(cursor);
3056 assert!(
3057 replay
3058 .call(&mut Context::default(), next)
3059 .await
3060 .unwrap_err()
3061 .to_string()
3062 .contains("stale_view")
3063 );
3064 }
3065
3066 #[tokio::test]
3067 async fn lossless_visibility_empty_nonturn_ties_and_supersession() {
3068 let store = test_store();
3069 let sid = start_session(&store);
3070 let replay = SessionReplayTool::new(store.clone());
3071 let mut args = json!({"mode":"lossless","session_id":sid,"page_size":64});
3072 let empty = replay
3073 .call(&mut Context::default(), args.clone())
3074 .await
3075 .unwrap();
3076 assert_eq!(empty["records"], json!([]));
3077 assert!(empty["next_cursor"].is_null());
3078 assert_eq!(empty["complete"], true);
3079 let mut a = lossless_fixture_turn(&sid, 1, "");
3080 let mut b = lossless_fixture_turn(&sid, 1, "second");
3081 b.metadata.tags.clear();
3083 let mut hidden = lossless_fixture_turn(&sid, 2, "PRIVATE-SENTINEL");
3084 hidden.metadata.is_private = true;
3085 let mut excluded = lossless_fixture_turn(&sid, 3, "EXCLUDED-SENTINEL");
3086 excluded.metadata.model_exclude = true;
3087 a.metadata
3088 .tags
3089 .push(format!("supersedes:{}", hidden.metadata.id));
3090 let mut handoff = Memory::new(
3091 Galaxy::Sessions,
3092 json!({"type":"session_handoff"}).to_string(),
3093 );
3094 handoff.metadata.tags = vec![format!("session:{sid}")];
3095 store
3096 .put_batch(
3097 Galaxy::Sessions,
3098 &[a.clone(), b.clone(), hidden.clone(), excluded, handoff],
3099 )
3100 .unwrap();
3101 let result = replay
3102 .call(&mut Context::default(), args.clone())
3103 .await
3104 .unwrap();
3105 let wire = result.to_string();
3106 assert!(!wire.contains("SENTINEL"));
3107 assert!(!wire.contains(&hidden.metadata.id.to_string()));
3108 let mut ids = vec![a.metadata.id.to_string(), b.metadata.id.to_string()];
3109 ids.sort();
3110 assert_eq!(
3111 result["records"]
3112 .as_array()
3113 .unwrap()
3114 .iter()
3115 .map(|r| r["record_id"].as_str().unwrap().to_string())
3116 .collect::<Vec<_>>(),
3117 ids
3118 );
3119 a.metadata
3120 .tags
3121 .push(format!("superseded-by:{}", b.metadata.id));
3122 store.put(Galaxy::Sessions, &a).unwrap();
3123 assert_eq!(
3124 replay
3125 .call(&mut Context::default(), args.clone())
3126 .await
3127 .unwrap()["records"]
3128 .as_array()
3129 .unwrap()
3130 .len(),
3131 1
3132 );
3133 args["include_superseded"] = json!(true);
3134 assert_eq!(
3135 replay.call(&mut Context::default(), args).await.unwrap()["records"]
3136 .as_array()
3137 .unwrap()
3138 .len(),
3139 2
3140 );
3141 let mut start = store
3142 .get(Galaxy::Sessions, uuid::Uuid::parse_str(&sid).unwrap())
3143 .unwrap()
3144 .unwrap();
3145 start.metadata.model_exclude = true;
3146 store.put(Galaxy::Sessions, &start).unwrap();
3147 assert!(
3148 replay
3149 .call(
3150 &mut Context::default(),
3151 json!({"mode":"lossless","session_id":sid})
3152 )
3153 .await
3154 .unwrap_err()
3155 .to_string()
3156 .contains("not found")
3157 );
3158 }
3159
3160 #[tokio::test]
3161 async fn lossless_malformed_selected_turn_fails_and_schema_advertises_mode() {
3162 let store = test_store();
3163 let sid = start_session(&store);
3164 let mut bad = lossless_fixture_turn(&sid, 1, "good");
3165 bad.content = "{broken".into();
3166 store.put(Galaxy::Sessions, &bad).unwrap();
3167 let replay = SessionReplayTool::new(store);
3168 assert!(
3169 replay
3170 .call(
3171 &mut Context::default(),
3172 json!({"mode":"lossless","session_id":sid})
3173 )
3174 .await
3175 .unwrap_err()
3176 .to_string()
3177 .contains("malformed_selected_turn")
3178 );
3179 assert!(replay.input_schema().to_string().contains("lossless"));
3180 assert!(replay.description().contains("explicit session UUID"));
3181 }
3182
3183 #[test]
3184 fn lossless_arbitrary_cursor_text_is_panic_free() {
3185 use proptest::prelude::*;
3186 proptest!(|(text in any::<String>())| { let _ = hex_decode(&text); });
3187 assert!(hex_decode("AB").is_none());
3188 }
3189
3190 #[tokio::test]
3191 async fn lossless_selection_reaches_record_after_ten_thousand_sources() {
3192 let store = test_store();
3193 let sid = start_session(&store);
3194 let mut sources = Vec::new();
3195 for i in 1..=10_001_u128 {
3196 let mut m = Memory::new(Galaxy::Sessions, "unrelated".into());
3197 m.metadata.id = uuid::Uuid::from_u128(i);
3198 sources.push(m);
3199 }
3200 let mut target = lossless_fixture_turn(&sid, 1, "AFTER-TEN-THOUSAND");
3201 target.metadata.id = uuid::Uuid::from_u128(u128::MAX);
3202 sources.push(target);
3203 store.put_batch(Galaxy::Sessions, &sources).unwrap();
3204 let replay = SessionReplayTool::new(store);
3205 let result = replay
3206 .call(
3207 &mut Context::default(),
3208 json!({"mode":"lossless","session_id":sid}),
3209 )
3210 .await
3211 .unwrap();
3212 assert_eq!(result["records"][0]["content"], "AFTER-TEN-THOUSAND");
3213 }
3214
3215 #[tokio::test]
3216 async fn lossless_tag_payload_disagreement_refuses_and_missing_start_is_not_found() {
3217 let store = test_store();
3218 let sid = start_session(&store);
3219 let mut m = lossless_fixture_turn(&uuid::Uuid::new_v4().to_string(), 1, "contradiction");
3220 m.metadata.tags = vec![format!("session:{sid}")];
3221 store.put(Galaxy::Sessions, &m).unwrap();
3222 let replay = SessionReplayTool::new(store);
3223 assert!(
3224 replay
3225 .call(
3226 &mut Context::default(),
3227 json!({"mode":"lossless","session_id":sid})
3228 )
3229 .await
3230 .unwrap_err()
3231 .to_string()
3232 .contains("malformed_selected_turn")
3233 );
3234 assert!(
3235 replay
3236 .call(
3237 &mut Context::default(),
3238 json!({"mode":"lossless","session_id":uuid::Uuid::new_v4().to_string()})
3239 )
3240 .await
3241 .unwrap_err()
3242 .to_string()
3243 .contains("not found")
3244 );
3245 }
3246}