1pub const ENVELOPE_FORMAT_VERSION: u32 = 2;
32
33pub const ENVELOPE_KEY: &str = "wm_envelope";
35
36#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
38pub struct EnvelopeHeader {
39 pub format_version: u32,
41 pub kind: String,
43 pub created_at: String,
45 pub count: usize,
47 pub generator: String,
49}
50
51impl EnvelopeHeader {
52 #[must_use]
54 pub fn new(kind: &str, count: usize) -> Self {
55 Self {
56 format_version: ENVELOPE_FORMAT_VERSION,
57 kind: kind.to_string(),
58 created_at: chrono::Utc::now().to_rfc3339(),
59 count,
60 generator: format!("wm {}", env!("CARGO_PKG_VERSION")),
61 }
62 }
63
64 #[must_use]
66 pub fn header_line(&self) -> String {
67 let mut v = serde_json::to_value(self).unwrap_or_else(|_| serde_json::json!({}));
68 let inner = v.take();
71 serde_json::to_string(&serde_json::json!({ ENVELOPE_KEY: inner }))
72 .unwrap_or_else(|_| format!("{{\"{ENVELOPE_KEY}\":null}}"))
73 }
74
75 pub fn check_version(&self) -> Result<(), String> {
81 if self.format_version > ENVELOPE_FORMAT_VERSION {
82 return Err(format!(
83 "envelope format_version {} is newer than this build supports ({}); \
84 upgrade `wm` to read this stream — refusing to import it partially",
85 self.format_version, ENVELOPE_FORMAT_VERSION
86 ));
87 }
88 Ok(())
89 }
90}
91
92#[derive(Debug, Clone, PartialEq, Eq)]
94pub enum HeaderRead {
95 NotAHeader,
97 Header(EnvelopeHeader),
99 Refused(String),
101}
102
103#[must_use]
108pub fn read_header_line(line: &str) -> HeaderRead {
109 let trimmed = line.trim();
110 let Ok(v) = serde_json::from_str::<serde_json::Value>(trimmed) else {
111 return HeaderRead::NotAHeader;
112 };
113 let Some(inner) = v.get(ENVELOPE_KEY) else {
114 return HeaderRead::NotAHeader;
115 };
116 if inner.is_null() {
117 return HeaderRead::Refused(format!(
118 "envelope header present but malformed (null {ENVELOPE_KEY})"
119 ));
120 }
121 match serde_json::from_value::<EnvelopeHeader>(inner.clone()) {
122 Ok(h) => match h.check_version() {
123 Ok(()) => HeaderRead::Header(h),
124 Err(msg) => HeaderRead::Refused(msg),
125 },
126 Err(e) => HeaderRead::Refused(format!(
127 "envelope header failed to parse: {e} (required: format_version, kind, \
128 created_at, count, generator)"
129 )),
130 }
131}
132
133#[derive(Debug, Default, PartialEq, Eq)]
135pub struct StreamScan {
136 pub header: Option<EnvelopeHeader>,
138 pub warnings: Vec<String>,
140 pub records: usize,
142 pub unparseable_lines: Vec<usize>,
145}
146
147#[must_use]
153pub fn scan_stream(payload: &str) -> StreamScan {
154 let mut scan = StreamScan::default();
155 let mut header_checked = false;
156 for (idx, line) in payload.lines().enumerate() {
157 if line.trim().is_empty() {
158 continue;
159 }
160 if !header_checked {
161 header_checked = true;
162 match read_header_line(line) {
163 HeaderRead::NotAHeader => { }
164 HeaderRead::Header(h) => {
165 scan.header = Some(h);
166 continue;
167 }
168 HeaderRead::Refused(msg) => {
169 scan.warnings.push(msg);
170 continue;
171 }
172 }
173 }
174 match serde_json::from_str::<serde_json::Value>(line.trim()) {
175 Ok(v) if v.is_object() => scan.records += 1,
176 _ => scan.unparseable_lines.push(idx + 1),
177 }
178 }
179 if let Some(h) = &scan.header {
180 if h.count != scan.records {
181 scan.warnings.push(format!(
182 "envelope declares count {} but stream carries {} records",
183 h.count, scan.records
184 ));
185 }
186 }
187 if !scan.unparseable_lines.is_empty() {
188 scan.warnings.push(format!(
189 "{} line(s) skipped as unparseable JSON at lines {:?}",
190 scan.unparseable_lines.len(),
191 scan.unparseable_lines
192 ));
193 }
194 scan
195}
196
197#[cfg(test)]
198mod tests {
199 use super::*;
200
201 #[test]
202 fn header_line_roundtrips_through_read() {
203 let h = EnvelopeHeader::new("session_export", 5);
204 let line = h.header_line();
205 assert!(line.contains(ENVELOPE_KEY));
206 match read_header_line(&line) {
207 HeaderRead::Header(parsed) => {
208 assert_eq!(parsed, h);
209 assert_eq!(parsed.format_version, ENVELOPE_FORMAT_VERSION);
210 assert_eq!(parsed.kind, "session_export");
211 assert_eq!(parsed.count, 5);
212 }
213 other => panic!("expected header, got {other:?}"),
214 }
215 }
216
217 #[test]
218 fn bare_record_line_is_not_a_header() {
219 let mem = serde_json::json!({
220 "metadata": {"id": "00000000-0000-0000-0000-000000000000"},
221 "content": "old memory",
222 "embedding": null
223 });
224 let line = serde_json::to_string(&mem).unwrap();
225 assert_eq!(read_header_line(&line), HeaderRead::NotAHeader);
226 }
227
228 #[test]
229 fn non_json_line_is_not_a_header() {
230 assert_eq!(read_header_line("not json at all"), HeaderRead::NotAHeader);
231 }
232
233 #[test]
234 fn newer_format_version_is_refused() {
235 let h = EnvelopeHeader {
236 format_version: ENVELOPE_FORMAT_VERSION + 1,
237 kind: "session_export".into(),
238 created_at: chrono::Utc::now().to_rfc3339(),
239 count: 1,
240 generator: "wm 99.0.0".into(),
241 };
242 match read_header_line(&h.header_line()) {
243 HeaderRead::Refused(msg) => {
244 assert!(msg.contains("newer than this build supports"), "{msg}");
245 }
246 other => panic!("expected refusal, got {other:?}"),
247 }
248 }
249
250 #[test]
251 fn malformed_header_is_refused_with_field_names() {
252 let line = serde_json::to_string(&serde_json::json!({
253 ENVELOPE_KEY: {"format_version": 2}
254 }))
255 .unwrap();
256 match read_header_line(&line) {
257 HeaderRead::Refused(msg) => assert!(msg.contains("required"), "{msg}"),
258 other => panic!("expected refusal, got {other:?}"),
259 }
260 }
261
262 #[test]
263 fn scan_stream_v2_payload() {
264 let rec1 = serde_json::json!({"metadata": {}, "content": "a"}).to_string();
265 let rec2 = serde_json::json!({"metadata": {}, "content": "b"}).to_string();
266 let payload = format!(
267 "{}\n{rec1}\n{rec2}\n",
268 EnvelopeHeader::new("session_export", 2).header_line()
269 );
270 let scan = scan_stream(&payload);
271 assert!(scan.header.is_some());
272 assert_eq!(scan.records, 2);
273 assert!(scan.warnings.is_empty(), "{:?}", scan.warnings);
274 assert!(scan.unparseable_lines.is_empty());
275 }
276
277 #[test]
278 fn scan_stream_bare_v1_payload_has_no_header() {
279 let rec1 = serde_json::json!({"metadata": {}, "content": "a"}).to_string();
280 let payload = format!("{rec1}\n");
281 let scan = scan_stream(&payload);
282 assert!(scan.header.is_none());
283 assert_eq!(scan.records, 1);
284 assert!(scan.warnings.is_empty());
285 }
286
287 #[test]
288 fn scan_stream_flags_count_mismatch_and_bad_lines() {
289 let good = serde_json::json!({"metadata": {}, "content": "a"}).to_string();
290 let payload = format!(
291 "{}\n{good}\n{{\"broken\":\nnot json\n",
292 EnvelopeHeader::new("session_export", 3).header_line()
293 );
294 let scan = scan_stream(&payload);
295 assert_eq!(scan.records, 1);
296 assert!(
297 scan.warnings
298 .iter()
299 .any(|w| w.contains("declares count 3 but stream carries 1"))
300 );
301 assert!(scan.warnings.iter().any(|w| w.contains("unparseable JSON")));
302 assert_eq!(scan.unparseable_lines, vec![3, 4]);
303 }
304
305 #[test]
306 fn header_survives_backup_json_roundtrip() {
307 let h = EnvelopeHeader::new("store_backup", 42);
310 let json = serde_json::to_string_pretty(&h).unwrap();
311 let parsed: EnvelopeHeader = serde_json::from_str(&json).unwrap();
312 assert_eq!(parsed, h);
313 }
314}