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