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