pub const ENVELOPE_FORMAT_VERSION: u32 = 2;
pub const ENVELOPE_KEY: &str = "wm_envelope";
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
pub struct EnvelopeHeader {
pub format_version: u32,
pub kind: String,
pub created_at: String,
pub count: usize,
pub generator: String,
}
impl EnvelopeHeader {
#[must_use]
pub fn new(kind: &str, count: usize) -> Self {
Self {
format_version: ENVELOPE_FORMAT_VERSION,
kind: kind.to_string(),
created_at: chrono::Utc::now().to_rfc3339(),
count,
generator: format!("wm {}", env!("CARGO_PKG_VERSION")),
}
}
#[must_use]
pub fn header_line(&self) -> String {
let mut v = serde_json::to_value(self).unwrap_or_else(|_| serde_json::json!({}));
let inner = v.take();
serde_json::to_string(&serde_json::json!({ ENVELOPE_KEY: inner }))
.unwrap_or_else(|_| format!("{{\"{ENVELOPE_KEY}\":null}}"))
}
pub fn check_version(&self) -> Result<(), String> {
if self.format_version > ENVELOPE_FORMAT_VERSION {
return Err(format!(
"envelope format_version {} is newer than this build supports ({}); \
upgrade `wm` to read this stream — refusing to import it partially",
self.format_version, ENVELOPE_FORMAT_VERSION
));
}
Ok(())
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum HeaderRead {
NotAHeader,
Header(EnvelopeHeader),
Refused(String),
}
#[must_use]
pub fn read_header_line(line: &str) -> HeaderRead {
let trimmed = line.trim();
let Ok(v) = serde_json::from_str::<serde_json::Value>(trimmed) else {
return HeaderRead::NotAHeader;
};
let Some(inner) = v.get(ENVELOPE_KEY) else {
return HeaderRead::NotAHeader;
};
if inner.is_null() {
return HeaderRead::Refused(format!(
"envelope header present but malformed (null {ENVELOPE_KEY})"
));
}
match serde_json::from_value::<EnvelopeHeader>(inner.clone()) {
Ok(h) => match h.check_version() {
Ok(()) => HeaderRead::Header(h),
Err(msg) => HeaderRead::Refused(msg),
},
Err(e) => HeaderRead::Refused(format!(
"envelope header failed to parse: {e} (required: format_version, kind, \
created_at, count, generator)"
)),
}
}
#[derive(Debug, Default, PartialEq, Eq)]
pub struct StreamScan {
pub header: Option<EnvelopeHeader>,
pub warnings: Vec<String>,
pub records: usize,
pub unparseable_lines: Vec<usize>,
}
#[must_use]
pub fn scan_stream(payload: &str) -> StreamScan {
let mut scan = StreamScan::default();
let mut header_checked = false;
for (idx, line) in payload.lines().enumerate() {
if line.trim().is_empty() {
continue;
}
if !header_checked {
header_checked = true;
match read_header_line(line) {
HeaderRead::NotAHeader => { }
HeaderRead::Header(h) => {
scan.header = Some(h);
continue;
}
HeaderRead::Refused(msg) => {
scan.warnings.push(msg);
continue;
}
}
}
match serde_json::from_str::<serde_json::Value>(line.trim()) {
Ok(v) if v.is_object() => scan.records += 1,
_ => scan.unparseable_lines.push(idx + 1),
}
}
if let Some(h) = &scan.header {
if h.count != scan.records {
scan.warnings.push(format!(
"envelope declares count {} but stream carries {} records",
h.count, scan.records
));
}
}
if !scan.unparseable_lines.is_empty() {
scan.warnings.push(format!(
"{} line(s) skipped as unparseable JSON at lines {:?}",
scan.unparseable_lines.len(),
scan.unparseable_lines
));
}
scan
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn header_line_roundtrips_through_read() {
let h = EnvelopeHeader::new("session_export", 5);
let line = h.header_line();
assert!(line.contains(ENVELOPE_KEY));
match read_header_line(&line) {
HeaderRead::Header(parsed) => {
assert_eq!(parsed, h);
assert_eq!(parsed.format_version, ENVELOPE_FORMAT_VERSION);
assert_eq!(parsed.kind, "session_export");
assert_eq!(parsed.count, 5);
}
other => panic!("expected header, got {other:?}"),
}
}
#[test]
fn bare_record_line_is_not_a_header() {
let mem = serde_json::json!({
"metadata": {"id": "00000000-0000-0000-0000-000000000000"},
"content": "old memory",
"embedding": null
});
let line = serde_json::to_string(&mem).unwrap();
assert_eq!(read_header_line(&line), HeaderRead::NotAHeader);
}
#[test]
fn non_json_line_is_not_a_header() {
assert_eq!(read_header_line("not json at all"), HeaderRead::NotAHeader);
}
#[test]
fn newer_format_version_is_refused() {
let h = EnvelopeHeader {
format_version: ENVELOPE_FORMAT_VERSION + 1,
kind: "session_export".into(),
created_at: chrono::Utc::now().to_rfc3339(),
count: 1,
generator: "wm 99.0.0".into(),
};
match read_header_line(&h.header_line()) {
HeaderRead::Refused(msg) => {
assert!(msg.contains("newer than this build supports"), "{msg}");
}
other => panic!("expected refusal, got {other:?}"),
}
}
#[test]
fn malformed_header_is_refused_with_field_names() {
let line = serde_json::to_string(&serde_json::json!({
ENVELOPE_KEY: {"format_version": 2}
}))
.unwrap();
match read_header_line(&line) {
HeaderRead::Refused(msg) => assert!(msg.contains("required"), "{msg}"),
other => panic!("expected refusal, got {other:?}"),
}
}
#[test]
fn scan_stream_v2_payload() {
let rec1 = serde_json::json!({"metadata": {}, "content": "a"}).to_string();
let rec2 = serde_json::json!({"metadata": {}, "content": "b"}).to_string();
let payload = format!(
"{}\n{rec1}\n{rec2}\n",
EnvelopeHeader::new("session_export", 2).header_line()
);
let scan = scan_stream(&payload);
assert!(scan.header.is_some());
assert_eq!(scan.records, 2);
assert!(scan.warnings.is_empty(), "{:?}", scan.warnings);
assert!(scan.unparseable_lines.is_empty());
}
#[test]
fn scan_stream_bare_v1_payload_has_no_header() {
let rec1 = serde_json::json!({"metadata": {}, "content": "a"}).to_string();
let payload = format!("{rec1}\n");
let scan = scan_stream(&payload);
assert!(scan.header.is_none());
assert_eq!(scan.records, 1);
assert!(scan.warnings.is_empty());
}
#[test]
fn scan_stream_flags_count_mismatch_and_bad_lines() {
let good = serde_json::json!({"metadata": {}, "content": "a"}).to_string();
let payload = format!(
"{}\n{good}\n{{\"broken\":\nnot json\n",
EnvelopeHeader::new("session_export", 3).header_line()
);
let scan = scan_stream(&payload);
assert_eq!(scan.records, 1);
assert!(
scan.warnings
.iter()
.any(|w| w.contains("declares count 3 but stream carries 1"))
);
assert!(scan.warnings.iter().any(|w| w.contains("unparseable JSON")));
assert_eq!(scan.unparseable_lines, vec![3, 4]);
}
#[test]
fn header_survives_backup_json_roundtrip() {
let h = EnvelopeHeader::new("store_backup", 42);
let json = serde_json::to_string_pretty(&h).unwrap();
let parsed: EnvelopeHeader = serde_json::from_str(&json).unwrap();
assert_eq!(parsed, h);
}
}