#[cfg(test)]
mod tests {
use std::{collections::BTreeSet, fs};
use rusqlite::{Connection, types::Value};
use crate::{
error::{message::MessageError, table::TableError},
tables::{
capabilities::Capabilities,
messages::{Message, ParsedBody},
table::Table,
},
test_support::schema_db,
util::query_context::QueryContext,
};
fn body_db() -> Connection {
let db = schema_db(true, true, true);
db.execute_batch(
"ALTER TABLE message ADD COLUMN attributedBody BLOB;
ALTER TABLE message ADD COLUMN message_summary_info BLOB;
ALTER TABLE message ADD COLUMN payload_data BLOB;",
)
.unwrap();
db
}
fn insert_body(db: &Connection, id: i32, body: Value, text: Option<&str>) {
db.execute(
"INSERT INTO message (rowid, guid, date, is_from_me, attributedBody, text)
VALUES (?1, ?2, ?1, 1, ?3, ?4)",
rusqlite::params![id, id.to_string(), body, text],
)
.unwrap();
}
fn assert_same_body(
borrowed: Result<ParsedBody, MessageError>,
lazy: Result<ParsedBody, MessageError>,
) {
match (borrowed, lazy) {
(Ok(borrowed), Ok(lazy)) => {
assert_eq!(borrowed.text, lazy.text);
assert_eq!(borrowed.components, lazy.components);
assert_eq!(borrowed.edited_parts, lazy.edited_parts);
assert_eq!(borrowed.balloon_bundle_id, lazy.balloon_bundle_id);
}
(Err(borrowed), Err(lazy)) => {
assert_eq!(format!("{borrowed:?}"), format!("{lazy:?}"));
}
(borrowed, lazy) => panic!("body results differ: {borrowed:?}, {lazy:?}"),
}
}
#[test]
fn borrowed_bodies_match_lazy_reads_for_every_typedstream_fixture() {
let db = body_db();
let mut count = 0;
for entry in fs::read_dir("test_data/typedstream").unwrap() {
let path = entry.unwrap().path();
if path.is_file() {
count += 1;
insert_body(&db, count, Value::Blob(fs::read(path).unwrap()), None);
}
}
assert!(count > 0);
let capabilities = Capabilities::determine(&db).unwrap();
assert!(capabilities.attributed_body);
let mut visited = 0;
Message::stream_with_body(
&db,
&capabilities,
&QueryContext::default(),
|message, bytes| {
assert_eq!(bytes.map(<[u8]>::to_vec), message.attributed_body(&db));
assert_same_body(
message.parse_body_from_bytes(&db, bytes),
message.parse_body(&db),
);
visited += 1;
Ok::<(), TableError>(())
},
)
.unwrap();
assert_eq!(visited, count);
}
#[test]
fn borrowed_bodies_preserve_null_numeric_text_and_malformed_values() {
let db = body_db();
for (idx, body) in [
Value::Null,
Value::Integer(123),
Value::Real(1.5),
Value::Text("invalid body".to_string()),
Value::Text(String::new()),
Value::Blob(vec![]),
Value::Blob(vec![0, 1, 2]),
]
.into_iter()
.enumerate()
{
insert_body(&db, idx as i32 * 2, body.clone(), None);
insert_body(&db, idx as i32 * 2 + 1, body, Some("plain text fallback"));
}
let capabilities = Capabilities::determine(&db).unwrap();
Message::stream_with_body(
&db,
&capabilities,
&QueryContext::default(),
|message, bytes| {
assert_eq!(bytes.map(<[u8]>::to_vec), message.attributed_body(&db));
let borrowed = message.parse_body_from_bytes(&db, bytes);
if message.rowid % 2 == 1 {
assert_eq!(
borrowed.as_ref().unwrap().text.as_deref(),
Some("plain text fallback")
);
}
assert_same_body(borrowed, message.parse_body(&db));
Ok::<(), TableError>(())
},
)
.unwrap();
}
#[test]
fn borrowed_text_preserves_utf16_storage_bytes() {
for encoding in ["UTF-16le", "UTF-16be"] {
let db = Connection::open_in_memory().unwrap();
db.execute_batch(&format!(
"PRAGMA encoding = '{encoding}';
CREATE TABLE message (
rowid INTEGER PRIMARY KEY, guid TEXT, date INTEGER,
is_from_me INTEGER, attributedBody BLOB, text TEXT
);
CREATE TABLE chat_message_join (chat_id INTEGER, message_id INTEGER);
CREATE TABLE message_attachment_join (message_id INTEGER, attachment_id INTEGER);"
))
.unwrap();
insert_body(&db, 1, Value::Text("café".to_string()), None);
insert_body(&db, 2, Value::Blob(vec![0, 128, 255]), None);
insert_body(&db, 3, Value::Integer(123), None);
let capabilities = Capabilities::determine(&db).unwrap();
let mut visited = 0;
Message::stream_with_body(
&db,
&capabilities,
&QueryContext::default(),
|message, bytes| {
let lazy = message.attributed_body(&db);
if message.rowid == 1 {
assert_ne!(lazy.as_deref(), Some("café".as_bytes()));
}
assert_eq!(bytes.map(<[u8]>::to_vec), lazy);
visited += 1;
Ok::<(), TableError>(())
},
)
.unwrap();
assert_eq!(visited, 3);
}
}
#[test]
fn borrowed_body_recovers_legacy_text_after_typedstream_failure() {
let db = body_db();
insert_body(
&db,
1,
Value::Blob(fs::read("test_data/typedstream/ExtraData").unwrap()),
None,
);
let capabilities = Capabilities::determine(&db).unwrap();
Message::stream_with_body(
&db,
&capabilities,
&QueryContext::default(),
|message, bytes| {
let body = message.parse_body_from_bytes(&db, bytes).unwrap();
assert!(body.text.as_deref().unwrap().starts_with("This is parsing"));
assert_same_body(Ok(body), message.parse_body(&db));
Ok::<(), TableError>(())
},
)
.unwrap();
}
#[test]
fn missing_body_column_preserves_plain_text_fallback() {
let db = schema_db(false, false, false);
db.execute_batch(
"INSERT INTO message (rowid, guid, date, is_from_me, text)
VALUES (1, 'legacy', 1, 1, 'legacy text');",
)
.unwrap();
let capabilities = Capabilities::determine(&db).unwrap();
assert!(!capabilities.attributed_body);
let mut visited = 0;
Message::stream_with_body(
&db,
&capabilities,
&QueryContext::default(),
|message, bytes| {
assert!(bytes.is_none());
let body = message.parse_body_from_bytes(&db, bytes).unwrap();
assert_eq!(body.text.as_deref(), Some("legacy text"));
assert_same_body(Ok(body), message.parse_body(&db));
visited += 1;
Ok::<(), TableError>(())
},
)
.unwrap();
assert_eq!(visited, 1);
}
#[test]
fn borrowed_body_preserves_edits_without_remaining_text() {
let db = body_db();
insert_body(&db, 1, Value::Null, None);
let summary = fs::read("test_data/edited_message/EditedAndUnsent.plist").unwrap();
db.execute(
"UPDATE message SET date_edited = 1, message_summary_info = ?1 WHERE rowid = 1",
[summary],
)
.unwrap();
let capabilities = Capabilities::determine(&db).unwrap();
Message::stream_with_body(
&db,
&capabilities,
&QueryContext::default(),
|message, bytes| {
let body = message.parse_body_from_bytes(&db, bytes).unwrap();
assert!(body.text.is_none());
assert!(body.edited_parts.is_some());
assert_same_body(Ok(body), message.parse_body(&db));
Ok::<(), TableError>(())
},
)
.unwrap();
}
#[test]
fn borrowed_body_preserves_url_payload_presence_check() {
let db = body_db();
let bytes = fs::read("test_data/typedstream/SingleLink").unwrap();
for id in 1..=2 {
insert_body(&db, id, Value::Blob(bytes.clone()), None);
}
db.execute_batch("UPDATE message SET payload_data = 123 WHERE rowid = 2")
.unwrap();
let capabilities = Capabilities::determine(&db).unwrap();
Message::stream_with_body(
&db,
&capabilities,
&QueryContext::default(),
|message, bytes| {
let body = message.parse_body_from_bytes(&db, bytes).unwrap();
assert_eq!(
body.balloon_bundle_id.as_deref(),
(message.rowid == 2).then_some("com.apple.messages.URLBalloonProvider")
);
assert_same_body(Ok(body), message.parse_body(&db));
Ok::<(), TableError>(())
},
)
.unwrap();
}
#[test]
fn body_stream_preserves_filters_and_duplicate_associations() {
let db = body_db();
for id in 1..=3 {
insert_body(&db, id, Value::Null, Some("text"));
}
db.execute_batch(
"INSERT INTO chat_message_join (chat_id, message_id)
VALUES (9, 1), (9, 1), (10, 2), (11, 3);
INSERT INTO chat_recoverable_message_join (chat_id, message_id) VALUES (9, 2);",
)
.unwrap();
let capabilities = Capabilities::determine(&db).unwrap();
let mut context = QueryContext::default();
context.set_selected_chat_ids(BTreeSet::from([9]));
let mut statement = Message::stream_rows(&db, &capabilities, &context).unwrap();
let lazy: Vec<_> = Message::rows(&mut statement, [])
.unwrap()
.map(|message| {
let message = message.unwrap();
(message.rowid, message.chat_id, message.deleted_from)
})
.collect();
let mut borrowed = vec![];
Message::stream_with_body(&db, &capabilities, &context, |message, _| {
borrowed.push((message.rowid, message.chat_id, message.deleted_from));
Ok::<(), TableError>(())
})
.unwrap();
assert_eq!(borrowed, lazy);
assert_eq!(borrowed.len(), 3);
assert_eq!(
borrowed.iter().map(|row| row.0).collect::<Vec<_>>(),
[1, 1, 2]
);
}
}