diff --git a/src/openhuman/memory_sync/canonicalize/chat.rs b/src/openhuman/memory_sync/canonicalize/chat.rs index cd565942e..7936605a0 100644 --- a/src/openhuman/memory_sync/canonicalize/chat.rs +++ b/src/openhuman/memory_sync/canonicalize/chat.rs @@ -25,8 +25,11 @@ use crate::openhuman::memory_store::chunks::types::{Metadata, SourceKind}; pub struct ChatMessage { /// Author display name or id. pub author: String, - /// When the message was sent. - #[serde(with = "chrono::serde::ts_milliseconds")] + /// When the message was sent (epoch-ms integer or RFC 3339 string). + #[serde( + serialize_with = "chrono::serde::ts_milliseconds::serialize", + deserialize_with = "super::deserialize_flexible_timestamp" + )] pub timestamp: DateTime, /// Plain text / markdown body. pub text: String, @@ -221,4 +224,39 @@ mod tests { .unwrap(); assert!(out.metadata.source_ref.is_none()); } + + // ── Serde regression tests (CORE-2K / #3568) ──────────────────────────── + + #[test] + fn timestamp_epoch_ms_integer_still_works() { + let json = r#"{ + "author": "alice", + "timestamp": 1700000000000, + "text": "hello" + }"#; + let msg: ChatMessage = serde_json::from_str(json).expect("epoch-ms integer should parse"); + assert_eq!(msg.timestamp.timestamp_millis(), 1_700_000_000_000); + } + + #[test] + fn timestamp_iso8601_string_accepted() { + let json = r#"{ + "author": "alice", + "timestamp": "2026-05-17T19:30:00Z", + "text": "hello" + }"#; + let msg: ChatMessage = serde_json::from_str(json).expect("ISO-8601 string should parse"); + assert_eq!(msg.timestamp.timestamp(), 1_779_046_200); + } + + #[test] + fn timestamp_numeric_string_accepted() { + let json = r#"{ + "author": "alice", + "timestamp": "1700000000000", + "text": "hello" + }"#; + let msg: ChatMessage = serde_json::from_str(json).expect("numeric string should parse"); + assert_eq!(msg.timestamp.timestamp_millis(), 1_700_000_000_000); + } } diff --git a/src/openhuman/memory_sync/canonicalize/document.rs b/src/openhuman/memory_sync/canonicalize/document.rs index 78b8abf6f..434122975 100644 --- a/src/openhuman/memory_sync/canonicalize/document.rs +++ b/src/openhuman/memory_sync/canonicalize/document.rs @@ -6,7 +6,7 @@ //! is kept verbatim. use chrono::{DateTime, Utc}; -use serde::{Deserialize, Deserializer, Serialize}; +use serde::{Deserialize, Serialize}; use super::{normalize_source_ref, CanonicalisedSource}; use crate::openhuman::memory_store::chunks::types::{Metadata, SourceKind}; @@ -21,56 +21,6 @@ fn now_utc() -> DateTime { Utc::now() } -/// Deserialise a `DateTime` from either: -/// - a JSON integer = epoch **milliseconds** (legacy callers — back-compat), -/// - a JSON string = RFC 3339 / ISO-8601 (e.g. `"2026-05-17T19:30:00Z"`), or -/// a decimal string containing epoch milliseconds. -/// -/// On an unparseable string a serde error is returned (no silent default). -fn deserialize_flexible_timestamp<'de, D>(deserializer: D) -> Result, D::Error> -where - D: Deserializer<'de>, -{ - /// Untagged helper so serde tries each variant in order. - #[derive(Deserialize)] - #[serde(untagged)] - enum RawTs { - Millis(i64), - Text(String), - } - - let raw = RawTs::deserialize(deserializer)?; - match raw { - RawTs::Millis(ms) => { - tracing::debug!("[memory][document] parsed modified_at as epoch-ms: {ms}"); - chrono::TimeZone::timestamp_millis_opt(&Utc, ms) - .single() - .ok_or_else(|| serde::de::Error::custom(format!("invalid epoch-ms: {ms}"))) - } - RawTs::Text(s) => { - // Try RFC 3339 / ISO-8601 first. - if let Ok(dt) = DateTime::parse_from_rfc3339(&s) { - tracing::debug!("[memory][document] parsed modified_at as ISO-8601 string: {s}"); - return Ok(dt.with_timezone(&Utc)); - } - // Fall back: numeric string = epoch milliseconds. - if let Ok(ms) = s.parse::() { - tracing::debug!( - "[memory][document] parsed modified_at as numeric-string epoch-ms: {s}" - ); - return chrono::TimeZone::timestamp_millis_opt(&Utc, ms) - .single() - .ok_or_else(|| { - serde::de::Error::custom(format!("invalid epoch-ms string: {s}")) - }); - } - Err(serde::de::Error::custom(format!( - "modified_at: cannot parse '{s}' as RFC 3339 or epoch-ms" - ))) - } - } -} - // ── Input struct ────────────────────────────────────────────────────────────── /// Adapter input for a single document. @@ -90,7 +40,7 @@ pub struct DocumentInput { /// string (fixes CORE-2K), or absent → `Utc::now()` (fixes CORE-2J). #[serde( default = "now_utc", - deserialize_with = "deserialize_flexible_timestamp" + deserialize_with = "super::deserialize_flexible_timestamp" )] pub modified_at: DateTime, /// Optional pointer back to source (URL, file path, Notion page id). diff --git a/src/openhuman/memory_sync/canonicalize/email.rs b/src/openhuman/memory_sync/canonicalize/email.rs index 5d005346e..f9420fec5 100644 --- a/src/openhuman/memory_sync/canonicalize/email.rs +++ b/src/openhuman/memory_sync/canonicalize/email.rs @@ -22,7 +22,11 @@ pub struct EmailMessage { #[serde(default)] pub cc: Vec, pub subject: String, - #[serde(with = "chrono::serde::ts_milliseconds")] + /// When the message was sent (epoch-ms integer or RFC 3339 string). + #[serde( + serialize_with = "chrono::serde::ts_milliseconds::serialize", + deserialize_with = "super::deserialize_flexible_timestamp" + )] pub sent_at: DateTime, /// Plain-text or markdown body. pub body: String, @@ -255,4 +259,42 @@ mod tests { let out = canonicalise("gmail:t1", "a", &[], t).unwrap().unwrap(); assert!(out.metadata.source_ref.is_none()); } + + // ── Serde regression tests (CORE-2K / #3568) ──────────────────────────── + + #[test] + fn sent_at_epoch_ms_integer_still_works() { + let json = r#"{ + "from": "alice@example.com", + "subject": "Launch", + "sent_at": 1700000000000, + "body": "content" + }"#; + let msg: EmailMessage = serde_json::from_str(json).expect("epoch-ms integer should parse"); + assert_eq!(msg.sent_at.timestamp_millis(), 1_700_000_000_000); + } + + #[test] + fn sent_at_iso8601_string_accepted() { + let json = r#"{ + "from": "alice@example.com", + "subject": "Launch", + "sent_at": "2026-05-17T19:30:00Z", + "body": "content" + }"#; + let msg: EmailMessage = serde_json::from_str(json).expect("ISO-8601 string should parse"); + assert_eq!(msg.sent_at.timestamp(), 1_779_046_200); + } + + #[test] + fn sent_at_numeric_string_accepted() { + let json = r#"{ + "from": "alice@example.com", + "subject": "Launch", + "sent_at": "1700000000000", + "body": "content" + }"#; + let msg: EmailMessage = serde_json::from_str(json).expect("numeric string should parse"); + assert_eq!(msg.sent_at.timestamp_millis(), 1_700_000_000_000); + } } diff --git a/src/openhuman/memory_sync/canonicalize/mod.rs b/src/openhuman/memory_sync/canonicalize/mod.rs index 7ef754939..e6d636a8f 100644 --- a/src/openhuman/memory_sync/canonicalize/mod.rs +++ b/src/openhuman/memory_sync/canonicalize/mod.rs @@ -14,10 +14,61 @@ pub mod document; pub mod email; pub mod email_clean; -use serde::{Deserialize, Serialize}; +use chrono::{DateTime, Utc}; +use serde::{Deserialize, Deserializer, Serialize}; use crate::openhuman::memory_store::chunks::types::{Metadata, SourceRef}; +/// Deserialise a `DateTime` from either: +/// - a JSON integer = epoch **milliseconds** (legacy callers — back-compat), +/// - a JSON string = RFC 3339 / ISO-8601 (e.g. `"2026-05-17T19:30:00Z"`), or +/// a decimal string containing epoch milliseconds. +/// +/// On an unparseable string a serde error is returned (no silent default). +/// Shared across chat, email, and document canonicalisers. +pub(crate) fn deserialize_flexible_timestamp<'de, D>( + deserializer: D, +) -> Result, D::Error> +where + D: Deserializer<'de>, +{ + #[derive(Deserialize)] + #[serde(untagged)] + enum RawTs { + Millis(i64), + Text(String), + } + + let raw = RawTs::deserialize(deserializer)?; + match raw { + RawTs::Millis(ms) => { + tracing::debug!("[memory][canonicalize] parsed timestamp as epoch-ms: {ms}"); + chrono::TimeZone::timestamp_millis_opt(&Utc, ms) + .single() + .ok_or_else(|| serde::de::Error::custom(format!("invalid epoch-ms: {ms}"))) + } + RawTs::Text(s) => { + if let Ok(dt) = DateTime::parse_from_rfc3339(&s) { + tracing::debug!("[memory][canonicalize] parsed timestamp as ISO-8601 string: {s}"); + return Ok(dt.with_timezone(&Utc)); + } + if let Ok(ms) = s.parse::() { + tracing::debug!( + "[memory][canonicalize] parsed timestamp as numeric-string epoch-ms: {s}" + ); + return chrono::TimeZone::timestamp_millis_opt(&Utc, ms) + .single() + .ok_or_else(|| { + serde::de::Error::custom(format!("invalid epoch-ms string: {s}")) + }); + } + Err(serde::de::Error::custom(format!( + "cannot parse '{s}' as RFC 3339 or epoch-ms" + ))) + } + } +} + /// Output of a canonicaliser — one per logical source record /// (a chat batch, an email, a document). #[derive(Clone, Debug)] diff --git a/src/openhuman/memory_tree/tree/rpc.rs b/src/openhuman/memory_tree/tree/rpc.rs index 329b24d34..de3ac39b3 100644 --- a/src/openhuman/memory_tree/tree/rpc.rs +++ b/src/openhuman/memory_tree/tree/rpc.rs @@ -905,6 +905,73 @@ mod tests { assert_eq!(listed.len(), first.chunks_written); } + /// Regression #3568 / CORE-2K: chat payloads with RFC-3339 timestamps must + /// be accepted — not rejected with "expected unix timestamp in milliseconds". + #[tokio::test] + async fn ingest_chat_accepts_rfc3339_timestamps() { + let (_tmp, cfg) = test_config(); + let outcome = ingest_rpc( + &cfg, + IngestRequest { + source_kind: SourceKind::Chat, + source_id: "slack:#rfc3339-test".into(), + owner: "alice".into(), + tags: vec![], + payload: json!({ + "platform": "slack", + "channel_label": "#eng", + "messages": [ + { + "author": "alice", + "timestamp": "2026-05-17T19:30:00Z", + "text": "planning the launch" + }, + { + "author": "bob", + "timestamp": 1779046260000_i64, + "text": "confirmed" + } + ] + }), + }, + ) + .await + .unwrap(); + assert!(!outcome.value.chunk_ids.is_empty()); + } + + /// Regression #3568 / CORE-2K: email payloads with RFC-3339 timestamps must + /// be accepted. + #[tokio::test] + async fn ingest_email_accepts_rfc3339_timestamps() { + let (_tmp, cfg) = test_config(); + let outcome = ingest_rpc( + &cfg, + IngestRequest { + source_kind: SourceKind::Email, + source_id: "gmail:rfc3339-test".into(), + owner: "alice@example.com".into(), + tags: vec![], + payload: json!({ + "provider": "gmail", + "thread_subject": "Launch", + "messages": [ + { + "from": "bob@example.com", + "to": ["alice@example.com"], + "subject": "Launch", + "sent_at": "2026-05-17T19:30:00Z", + "body": "Let's ship this." + } + ] + }), + }, + ) + .await + .unwrap(); + assert!(!outcome.value.chunk_ids.is_empty()); + } + #[tokio::test] async fn ingest_rpc_rejects_invalid_document_payload() { let (_tmp, cfg) = test_config();