diff --git a/Cargo.lock b/Cargo.lock index f4f97ff..a15b555 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -21,9 +21,9 @@ dependencies = [ [[package]] name = "actix-http" -version = "3.13.5" +version = "3.13.6" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "86d62d1a48894ec9450bcde7ef1e3681205771ff6ea9c61cc5ae031145192787" +checksum = "3c4f07a021937af76abef4c458a0d6b25f63f82eec5394354deeb7de0551f523" dependencies = [ "actix-codec", "actix-service", @@ -49,7 +49,7 @@ dependencies = [ "mime", "percent-encoding", "pin-project-lite", - "rand 0.10.2", + "rand 0.10.3", "sha1 0.11.0", "smallvec", "tokio", @@ -60,12 +60,13 @@ dependencies = [ [[package]] name = "actix-macros" -version = "0.2.4" +version = "0.2.5" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e01ed3140b2f8d422c68afa1ed2e85d996ea619c988ac834d255db32138655cb" +checksum = "367f814ad4afbac74f07df5001214da65f65e185c90ef56c4dd8df23f8695b9b" dependencies = [ + "proc-macro2", "quote", - "syn 2.0.119", + "syn 3.0.6", ] [[package]] @@ -140,9 +141,9 @@ dependencies = [ [[package]] name = "actix-utils" -version = "3.0.1" +version = "3.0.2" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "88a1dcdff1466e3c2488e1cb5c36a71822750ad43839937f85d2f4d9f8b705d8" +checksum = "0128396dd7313f697ad05b21b1a7be7d4cbb81888704f55996e4a27db196bb4d" dependencies = [ "local-waker", "pin-project-lite", @@ -194,14 +195,14 @@ dependencies = [ [[package]] name = "actix-web-codegen" -version = "4.3.0" +version = "4.4.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "f591380e2e68490b5dfaf1dd1aa0ebe78d84ba7067078512b4ea6e4492d622b8" +checksum = "b96b09c4878563f8ab4a5fd0c59f9f0d6e0e9f60eb9b748526a0b9604fd89c50" dependencies = [ "actix-router", "proc-macro2", "quote", - "syn 2.0.119", + "syn 3.0.6", ] [[package]] @@ -390,7 +391,7 @@ dependencies = [ "proc-macro2", "quote", "syn 2.0.119", - "synstructure", + "synstructure 0.13.2", ] [[package]] @@ -412,7 +413,7 @@ checksum = "82f6aeea286b8eb4dd3431a1be1b59d290ace00f5bfd8e2a159bc2a05e2c1667" dependencies = [ "proc-macro2", "quote", - "syn 3.0.5", + "syn 3.0.6", ] [[package]] @@ -621,9 +622,9 @@ dependencies = [ [[package]] name = "cc" -version = "1.4.6" +version = "1.4.7" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "a3eb0f42d6c360dc3f8a821f6bf2fdea7f72bfd36b3076eb0e6d1e9e0752fff4" +checksum = "54413ede23c2daf518f35156dfde027feb2374004d63bd497f983c8db9c0e313" dependencies = [ "find-msvc-tools", "jobserver", @@ -633,9 +634,9 @@ dependencies = [ [[package]] name = "cfg-if" -version = "1.0.4" +version = "1.0.5" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "9330f8b2ff13f34540b44e946ef35111825727b38d33286ef986142615121801" +checksum = "4e7648175b45a9a48536d676f68d918270699102aa8dab5496df06904c914600" [[package]] name = "cfg_aliases" @@ -704,9 +705,9 @@ dependencies = [ [[package]] name = "clap" -version = "4.6.6" +version = "4.6.7" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "473c7e07f409a8d772161724aa8db6a765a2532a70f9667eeb7b49d3d02fbdca" +checksum = "aa8876b300ab35ba921adea3dfd70157a46249b33f95c9084ae5709785478946" dependencies = [ "clap_builder", "clap_derive", @@ -714,9 +715,9 @@ dependencies = [ [[package]] name = "clap_builder" -version = "4.6.6" +version = "4.6.7" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "7b48fea5a88e9ae728a2dcbedbfc0e730f7d60da42e1cb049a83c9fb8b789889" +checksum = "ec0797fb7aeb1406c84efac526901f7ec3ead2124f946b494e72879d4b54704d" dependencies = [ "anstream", "anstyle", @@ -726,21 +727,21 @@ dependencies = [ [[package]] name = "clap_derive" -version = "4.6.4" +version = "4.6.7" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "d012d2b9d65aca7f18f4d9878a045bc17899bba951561ba5ec3c2ba1eed9a061" +checksum = "f9c751b79415d4e559e3d1fcf128e09e720eb673a06d26cf6f392d37d75b66e0" dependencies = [ "heck", "proc-macro2", "quote", - "syn 3.0.5", + "syn 3.0.6", ] [[package]] name = "clap_lex" -version = "1.1.0" +version = "1.1.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "c8d4a3bb8b1e0c1050499d1815f5ab16d04f0959b233085fb31653fbfc9d98f9" +checksum = "1c133bc6a41be0d194c306b5506d15e6feeea7b1d6604bd3f8310dfb2ca96486" [[package]] name = "client" @@ -1064,7 +1065,7 @@ dependencies = [ "proc-macro2", "quote", "strsim", - "syn 3.0.5", + "syn 3.0.6", ] [[package]] @@ -1075,7 +1076,7 @@ checksum = "2ac7135c3ef02b2f7833bbeb1be5ba7f966dcde8a87c6b87f65a778d71a02785" dependencies = [ "darling_core", "quote", - "syn 3.0.5", + "syn 3.0.6", ] [[package]] @@ -1216,7 +1217,7 @@ checksum = "c6232dd377dcc64799954cbd3a9bb882e9cdc1308ccd87b1c098f1fb2eaf82a8" dependencies = [ "proc-macro2", "quote", - "syn 3.0.5", + "syn 3.0.6", ] [[package]] @@ -1406,15 +1407,15 @@ dependencies = [ [[package]] name = "find-msvc-tools" -version = "0.1.12" +version = "0.1.13" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "3e0f1c7c3a72c66fd80abe965175f7523475c0489a87d3ff9d6e8c87d87a9d2d" +checksum = "ef25905e51abafe4dcea6c15fec58c57b601cdbd0ee53d22ea1d3016c587d39b" [[package]] name = "finl_unicode" -version = "1.4.0" +version = "1.5.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "9844ddc3a6e533d62bba727eb6c28b5d360921d5175e9ff0f1e621a5c590a4d5" +checksum = "80bb028c8b4148c9ee0cca68fcd9add6044e81d3619f48577ddf13a263d047a2" [[package]] name = "fixedbitset" @@ -1526,7 +1527,7 @@ checksum = "9fb9654ba8355388abeb8dcb4fc62f511300867002afc858860463bdd9fe0c44" dependencies = [ "proc-macro2", "quote", - "syn 3.0.5", + "syn 3.0.6", ] [[package]] @@ -1879,9 +1880,9 @@ dependencies = [ [[package]] name = "hyper-rustls" -version = "0.27.9" +version = "0.27.10" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "33ca68d021ef39cf6463ab54c1d0f5daf03377b70561305bb89a8f83aab66e0f" +checksum = "dfa8e654703247911e29c23fbeaa261834bd9bb74efba2f9acddc37bfb127f53" dependencies = [ "http 1.5.0", "hyper", @@ -2063,9 +2064,9 @@ dependencies = [ [[package]] name = "impl-more" -version = "0.3.7" +version = "0.3.8" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "edaff2ce006342d4d0e00fae676f7082dada44a203560629ba18ea50a19277bb" +checksum = "30c0cddce6b7505483307994f60d97c4b156020e1504d0cf1f05a21b1b325e64" [[package]] name = "indexmap" @@ -2105,7 +2106,7 @@ dependencies = [ "indoc", "proc-macro2", "quote", - "syn 3.0.5", + "syn 3.0.6", ] [[package]] @@ -2698,9 +2699,9 @@ dependencies = [ [[package]] name = "lru-slab" -version = "0.1.2" +version = "0.1.3" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "112b39cec0b298b6c1999fee3e31427f74f676e4cb9879ed1a121b43661a4154" +checksum = "4050469837a6ff301cd14c1f8f24f88549e6d548f24f64e2148eb0f72cebc51f" [[package]] name = "lzma-rust2" @@ -2856,7 +2857,7 @@ dependencies = [ "mtp-common", "mtp-crypto", "mtp-native", - "rand 0.10.2", + "rand 0.10.3", "tokio", ] @@ -2870,7 +2871,7 @@ dependencies = [ "mtp-common", "mtp-crypto", "mtp-type-map", - "rand 0.10.2", + "rand 0.10.3", "thiserror 2.0.20", ] @@ -2879,7 +2880,7 @@ name = "mtp-common" version = "0.3.0" source = "git+https://git.methanium.net/Methanium/mtp.git?rev=bb0f682b735de5ebb36bb41dc699260578341828#bb0f682b735de5ebb36bb41dc699260578341828" dependencies = [ - "rand 0.10.2", + "rand 0.10.3", "thiserror 2.0.20", ] @@ -2908,7 +2909,7 @@ dependencies = [ "hkdf", "ml-dsa", "mlkem-tls", - "rand 0.10.2", + "rand 0.10.3", "rand_core 0.6.4", "rustls", "serde", @@ -2924,7 +2925,7 @@ version = "0.3.0" source = "git+https://git.methanium.net/Methanium/mtp.git?rev=bb0f682b735de5ebb36bb41dc699260578341828#bb0f682b735de5ebb36bb41dc699260578341828" dependencies = [ "mtp-crypto", - "rand 0.10.2", + "rand 0.10.3", "thiserror 2.0.20", "zeroize", ] @@ -2956,7 +2957,7 @@ dependencies = [ "mtp-core", "mtp-crypto", "mtp-native", - "rand 0.10.2", + "rand 0.10.3", "thiserror 2.0.20", "tokio", "tracing", @@ -2973,7 +2974,7 @@ dependencies = [ "mtp-common", "mtp-core", "mtp-crypto", - "rand 0.10.2", + "rand 0.10.3", "rcgen", "rustls", "rustls-native-certs", @@ -3041,7 +3042,7 @@ dependencies = [ "proc-macro2", "quote", "rustversion", - "syn 3.0.5", + "syn 3.0.6", ] [[package]] @@ -3630,9 +3631,9 @@ checksum = "439ee305def115ba05938db6eb1644ff94165c5ab5e9420d1c1bcedbba909391" [[package]] name = "ppmd-rust" -version = "1.4.1" +version = "1.5.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "9e9219bcb9d7aca6b2f63c83cf100cf78bcd619ac46e6ecbd0dd90869a39345d" +checksum = "196a7c80b9a7652aba7cc070827516c2abe4ccdf53d128e1944003cf5726cff1" [[package]] name = "ppv-lite86" @@ -3654,9 +3655,9 @@ dependencies = [ [[package]] name = "quinn" -version = "0.11.11" +version = "0.11.12" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "0c1a41e437b6bbd489372cd4971de128e85c855f56c57f283d20ff016cf7c0a8" +checksum = "4051e23e9185c255a7e33ef59cdbca87a22d359052eecd22fc6b901fb37d9d11" dependencies = [ "bytes", "cfg_aliases", @@ -3675,16 +3676,16 @@ dependencies = [ [[package]] name = "quinn-proto" -version = "0.11.17" +version = "0.11.18" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "04759210543be93709136e28212294a659ef5001836ff4eab4d663e4529bba83" +checksum = "a9746dbde176634f4f2f1faf2404e30a31b2bc1e9cafb5329c95d8177a18c9fc" dependencies = [ "aws-lc-rs", "bytes", "fastbloom", "getrandom 0.4.3", "lru-slab", - "rand 0.10.2", + "rand 0.10.3", "rand_pcg", "ring", "rustc-hash", @@ -3757,9 +3758,9 @@ dependencies = [ [[package]] name = "rand" -version = "0.10.2" +version = "0.10.3" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "c7f5fa3a058cd35567ef9bfa5e75732bee0f9e4c55fa90477bef2dfcdbc4be80" +checksum = "65c9fb96cbc91e3478eaae79a69fcd3f1ae4ad052e471fe6732fff548984b4af" dependencies = [ "chacha20 0.10.2", "getrandom 0.4.3", @@ -3803,7 +3804,7 @@ dependencies = [ [[package]] name = "ratatool" version = "0.1.0" -source = "git+https://git.methanium.net/methanium/ratatool.git#f94b8a25ffc5c196db2852ff576593df5ad96938" +source = "git+https://git.methanium.net/methanium/ratatool.git#ddd0ddea7235cfc1ca27f79cf166f40d3d1fdabe" dependencies = [ "crossterm", "futures-util", @@ -4084,9 +4085,9 @@ dependencies = [ [[package]] name = "rustix" -version = "1.1.4" +version = "1.1.5" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "b6fe4565b9518b83ef4f91bb47ce29620ca828bd32cb7e408f0062e9930ba190" +checksum = "891efababe418670775f199f0d233d84843c227a0949a883ce15b37c78d6629d" dependencies = [ "bitflags 2.13.2", "errno", @@ -4097,9 +4098,9 @@ dependencies = [ [[package]] name = "rustls" -version = "0.23.44" +version = "0.23.45" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "6725596c3f2c3a0aef021139e145d4eafe314a6623e4680ca83852b2c67ab2ba" +checksum = "0d41d731c7d2f962d1ccc364cec258de3c0e93b38c2fb3ba97ac74513048d634" dependencies = [ "aws-lc-rs", "log", @@ -4282,7 +4283,7 @@ checksum = "e7a5d71263a5a7d47b41f6b3f06ba276f10cc18b0931f1799f710578e2309348" dependencies = [ "proc-macro2", "quote", - "syn 3.0.5", + "syn 3.0.6", ] [[package]] @@ -4601,9 +4602,9 @@ dependencies = [ [[package]] name = "syn" -version = "3.0.5" +version = "3.0.6" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "12df2e0110f65b775f769bb17ef989067a1d931b2eb822bd4346631eeada89f9" +checksum = "8593e8e72159ed2257d083c7a454a85cbf854f37a0966d8d483aff8c8a3ebcee" dependencies = [ "proc-macro2", "quote", @@ -4630,6 +4631,17 @@ dependencies = [ "syn 2.0.119", ] +[[package]] +name = "synstructure" +version = "0.14.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "901704edd0dfe137f1987838ee4f259e4e063c31371bdb423f7ae38ec6f77f02" +dependencies = [ + "proc-macro2", + "quote", + "syn 3.0.6", +] + [[package]] name = "sysinfo" version = "0.38.4" @@ -4672,7 +4684,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "32497e9a4c7b38532efcdebeef879707aa9f794296a4f0244f6f69e9bc8574bd" dependencies = [ "fastrand", - "getrandom 0.3.4", + "getrandom 0.4.3", "once_cell", "rustix", "windows-sys 0.61.2", @@ -4791,7 +4803,7 @@ checksum = "bc04cd3e1236dd4a98afca4569f2deb3f120e5422a4023be2cb683f8486292af" dependencies = [ "proc-macro2", "quote", - "syn 3.0.5", + "syn 3.0.6", ] [[package]] @@ -4838,18 +4850,9 @@ dependencies = [ [[package]] name = "tinyvec" -version = "1.13.2" +version = "1.13.3" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "4cf0ded5c4e56918d8f8a339e1bb67d038d3bc6d144ac407904015ba2e4cde9b" -dependencies = [ - "tinyvec_macros", -] - -[[package]] -name = "tinyvec_macros" -version = "0.1.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1f3ccbac311fea05f86f61904b462b55fb3df8837a366dfc601a0161d0532f20" +checksum = "fd3ca314f692efd6c868f8408f53fe444634a845f96c028b97d35f6a1f79f0ee" [[package]] name = "tokio" @@ -4876,7 +4879,7 @@ checksum = "78773a2a397f451582ce068015985c33193cf6dea8b74d2a639fe457b2f07b0e" dependencies = [ "proc-macro2", "quote", - "syn 3.0.5", + "syn 3.0.6", ] [[package]] @@ -5064,9 +5067,9 @@ checksum = "5c1cb5db39152898a79168971543b1cb5020dff7fe43c8dc468b0885f5e29df5" [[package]] name = "unicode-ident" -version = "1.0.24" +version = "1.0.26" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e6e4313cd5fcd3dad5cafa179702e2b244f760991f45397d14d4ebf38247da75" +checksum = "d245f478577f809a851594d02313b640fb437e0bb33866753cff937863096954" [[package]] name = "unicode-normalization" @@ -5267,7 +5270,7 @@ dependencies = [ "bumpalo", "proc-macro2", "quote", - "syn 3.0.5", + "syn 3.0.6", "wasm-bindgen-shared", ] @@ -5757,14 +5760,14 @@ dependencies = [ [[package]] name = "yoke-derive" -version = "0.8.2" +version = "0.8.3" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "de844c262c8848816172cef550288e7dc6c7b7814b4ee56b3e1553f275f1858e" +checksum = "33811428bee40dbceb6d545e95754741d17a6aef9a4849f0fd62e2ba4f412a78" dependencies = [ "proc-macro2", "quote", - "syn 2.0.119", - "synstructure", + "syn 3.0.6", + "synstructure 0.14.0", ] [[package]] @@ -5798,14 +5801,14 @@ dependencies = [ [[package]] name = "zerofrom-derive" -version = "0.1.7" +version = "0.1.8" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "11532158c46691caf0f2593ea8358fed6bbf68a0315e80aae9bd41fbade684a1" +checksum = "f75b4683f6c7f45248d4d64056a24298c6281e0993356d7d1b4a1a962ef10d4a" dependencies = [ "proc-macro2", "quote", - "syn 2.0.119", - "synstructure", + "syn 3.0.6", + "synstructure 0.14.0", ] [[package]] @@ -5858,7 +5861,7 @@ checksum = "34df6fc39dbd26ddc9c10e6a2984476e13acce22e64e4487636ef494369225da" dependencies = [ "proc-macro2", "quote", - "syn 3.0.5", + "syn 3.0.6", ] [[package]] @@ -5890,9 +5893,9 @@ dependencies = [ [[package]] name = "zlib-rs" -version = "0.6.7" +version = "0.6.8" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "34b31d188d9d685a4f9c7b46d6e36631b07058d2cfe190267adce54dc230bf12" +checksum = "b268e58e7c693d7c271f93ffc4ba3b380412554231c85bf61ca7af91042a4112" [[package]] name = "zmij" diff --git a/iota-connection/src/message_handlers.rs b/iota-connection/src/message_handlers.rs index ce2b6d0..07983aa 100644 --- a/iota-connection/src/message_handlers.rs +++ b/iota-connection/src/message_handlers.rs @@ -31,6 +31,32 @@ pub struct MessageMutation { pub send_time: i64, } +#[derive(Clone, Debug)] +pub enum AppliedMessageMutation { + State { + send_time: i64, + state: MessageState, + }, + Edit { + send_time: i64, + content: String, + version: i64, + }, + Reaction { + send_time: i64, + reaction: String, + accepted: bool, + }, + Delete { + send_time: i64, + }, +} + +#[derive(Clone, Debug, Default)] +pub struct RelayApplicationResult { + pub mutation: Option, +} + #[derive(Debug)] pub struct SettingMutation { pub response: CommunicationValue, @@ -189,7 +215,7 @@ pub fn apply_verified_relay_content( accepted_at: i64, storage_owner: i64, sent_by_self: bool, -) -> Result<(), String> { +) -> Result { validate_relay_identity(context, &content.content)?; let sender_id = i64::try_from(context.signer_id) .map_err(|_| "Relay signer ID exceeds the local storage range".to_string())?; @@ -200,7 +226,7 @@ pub fn apply_verified_relay_content( match content.message_type { CommunicationType::MessageState => { - let partner_id = + let _partner_id = relay_number(&content.content, DataType::ChatPartnerId, &context.type_map) .and_then(|value| i64::try_from(value).ok()) .filter(|id| *id == recipient_id) @@ -218,31 +244,36 @@ pub fn apply_verified_relay_content( .map(MessageState::from_str) .filter(|state| matches!(state, MessageState::Received | MessageState::Read)) .ok_or_else(|| "Relay MessageState has an invalid state".to_string())?; - if sent_by_self { - return chat_files::change_message_state_by_relay_id( + let update = if sent_by_self { + chat_files::change_message_state_by_relay_id( storage_owner, recipient_id, recipient_principal, relay_message_id, state, ) - .map_err(|error| error.to_string()); - } - chat_files::record_message_receipt( - storage_owner, - recipient_id, - recipient_principal, - relay_message_id, - sender_id, - signer_principal, - &context.message_id, - state, - event_at, - now_millis_i64(), - ) - .map_err(|error| error.to_string())?; - let _ = partner_id; - Ok(()) + .map_err(|error| error.to_string())? + } else { + chat_files::record_message_receipt( + storage_owner, + recipient_id, + recipient_principal, + relay_message_id, + sender_id, + signer_principal, + &context.message_id, + state, + event_at, + now_millis_i64(), + ) + .map_err(|error| error.to_string())? + }; + Ok(RelayApplicationResult { + mutation: update.map(|update| AppliedMessageMutation::State { + send_time: update.send_time, + state: update.state, + }), + }) } CommunicationType::MessageSend => { let message = relay_string(&content.content, DataType::Content, &context.type_map) @@ -328,7 +359,7 @@ pub fn apply_verified_relay_content( }, }) .map_err(|error| error.to_string())?; - Ok(()) + Ok(RelayApplicationResult::default()) } CommunicationType::MessageEdit => { let message = relay_string(&content.content, DataType::Content, &context.type_map) @@ -336,6 +367,11 @@ pub fn apply_verified_relay_content( let send_time = relay_number(&content.content, DataType::SendTime, &context.type_map) .and_then(|value| i64::try_from(value).ok()) .ok_or_else(|| "Relay MessageEdit is missing SendTime".to_string())?; + let version = + relay_number(&content.content, DataType::VersionNumber, &context.type_map) + .and_then(|value| i64::try_from(value).ok()) + .filter(|value| *value > 0) + .ok_or_else(|| "Relay MessageEdit has an invalid VersionNumber".to_string())?; let external_user = if sent_by_self { recipient_id } else { @@ -365,7 +401,14 @@ pub fn apply_verified_relay_content( message, ) }; - result.map_err(|error| error.to_string()) + result.map_err(|error| error.to_string())?; + Ok(RelayApplicationResult { + mutation: Some(AppliedMessageMutation::Edit { + send_time, + content: message.to_string(), + version, + }), + }) } CommunicationType::MessageReactionAdd | CommunicationType::MessageReactionRemove => { let reaction = relay_string(&content.content, DataType::Reaction, &context.type_map) @@ -405,7 +448,14 @@ pub fn apply_verified_relay_content( reaction, ) }; - result.map_err(|error| error.to_string()) + result.map_err(|error| error.to_string())?; + Ok(RelayApplicationResult { + mutation: Some(AppliedMessageMutation::Reaction { + send_time, + reaction: reaction.to_string(), + accepted: content.message_type == CommunicationType::MessageReactionAdd, + }), + }) } CommunicationType::MessageDelete | CommunicationType::MessageDeleteLive => { let send_time = relay_number(&content.content, DataType::SendTime, &context.type_map) @@ -437,7 +487,10 @@ pub fn apply_verified_relay_content( sender_id, ) }; - result.map_err(|error| error.to_string()) + result.map_err(|error| error.to_string())?; + Ok(RelayApplicationResult { + mutation: Some(AppliedMessageMutation::Delete { send_time }), + }) } CommunicationType::SetChatSecret => { let frame = CommunicationValue::new(CommunicationType::SetChatSecret) @@ -526,7 +579,8 @@ pub fn apply_verified_relay_content( created_at, updated_at: now_millis_i64(), }) - .map_err(|error| error.to_string()) + .map_err(|error| error.to_string())?; + Ok(RelayApplicationResult::default()) } CommunicationType::AddConversation => { let other_id = @@ -560,9 +614,10 @@ pub fn apply_verified_relay_content( &context.type_map, ), ) - .map_err(|error| format!("AddConversation persistence failed: {error}")) + .map_err(|error| format!("AddConversation persistence failed: {error}"))?; + Ok(RelayApplicationResult::default()) } - _ => Ok(()), + _ => Ok(RelayApplicationResult::default()), } } @@ -3003,7 +3058,7 @@ pub fn apply_relay_application_content( context: &VerifiedRelayContext, content: &VerifiedRelayContent, accepted_at: i64, -) -> Result<(), String> { +) -> Result { if application.signer_principal.user_id != context.signer_id || application.recipient_principal.user_id != context.final_recipient_id { diff --git a/iota-connection/src/relay_service.rs b/iota-connection/src/relay_service.rs index 511156e..17118b5 100644 --- a/iota-connection/src/relay_service.rs +++ b/iota-connection/src/relay_service.rs @@ -457,6 +457,7 @@ impl RelayService { let signer_principal = normalized.signer; let recipient_principal = normalized.recipient; let accepted_at = crate::message_common::now_millis_i64(); + let mut source_mutation = None; let signer_id = match i64::try_from(verified.context.signer_id) { Ok(id) => id, Err(_) => { @@ -653,7 +654,7 @@ impl RelayService { return Self::relay_response(frame.id(), CommunicationType::ErrorInvalidData); } }; - if let Err(error) = message_handlers::apply_verified_relay_content( + let application = match message_handlers::apply_verified_relay_content( &verified.context, &content, signer_principal.handle, @@ -662,13 +663,17 @@ impl RelayService { owner, true, ) { - log!("Relay shared-Iota origin application failed: {}", error); - let _ = relay_replay::mark_rejected( - signer_principal.handle, - &verified.context.message_id, - ); - return Self::relay_response(frame.id(), CommunicationType::ErrorInvalidData); - } + Ok(application) => application, + Err(error) => { + log!("Relay shared-Iota origin application failed: {}", error); + let _ = relay_replay::mark_rejected( + signer_principal.handle, + &verified.context.message_id, + ); + return Self::relay_response(frame.id(), CommunicationType::ErrorInvalidData); + } + }; + source_mutation = application.mutation; if let Err(error) = chat_files::record_destination_iota_received( owner, signer_principal.handle, @@ -703,22 +708,29 @@ impl RelayService { ); } }; - if let Err(error) = message_handlers::apply_verified_relay_content( + let application = match message_handlers::apply_verified_relay_content( &verified.context, &content, signer_principal.handle, recipient_principal.handle, accepted_at, - i64::try_from(verified.context.signer_id).unwrap_or_default(), + signer_id, true, ) { - log!("Relay origin application failed: {}", error); - let _ = relay_replay::mark_rejected( - signer_principal.handle, - &verified.context.message_id, - ); - return Self::relay_response(frame.id(), CommunicationType::ErrorInvalidData); - } + Ok(application) => application, + Err(error) => { + log!("Relay origin application failed: {}", error); + let _ = relay_replay::mark_rejected( + signer_principal.handle, + &verified.context.message_id, + ); + return Self::relay_response( + frame.id(), + CommunicationType::ErrorInvalidData, + ); + } + }; + source_mutation = application.mutation; } let Some(route) = route_destination(&recipient_principal.home) else { log!("Relay recipient has no Iota route"); @@ -894,7 +906,18 @@ impl RelayService { .with_id(frame_id); return RelayOutcome::Accepted { ingress_response: response, - local_deliveries: Vec::new(), + local_deliveries: source_mutation + .as_ref() + .map(|mutation| { + vec![mutation_delivery( + mutation, + verified.context.signer_id, + verified.context.signer_id, + verified.context.final_recipient_id, + signer_principal.handle, + )] + }) + .unwrap_or_default(), }; } Ok(response) => { @@ -914,7 +937,18 @@ impl RelayService { &verified.context.message_id, accepted_at, true, - Vec::new(), + source_mutation + .as_ref() + .map(|mutation| { + vec![mutation_delivery( + mutation, + verified.context.signer_id, + verified.context.signer_id, + verified.context.final_recipient_id, + signer_principal.handle, + )] + }) + .unwrap_or_default(), ); } else if let Ok(signer_id) = i64::try_from(verified.context.signer_id) { if let Err(error) = @@ -971,7 +1005,18 @@ impl RelayService { &verified.context.message_id, accepted_at, true, - Vec::new(), + source_mutation + .as_ref() + .map(|mutation| { + vec![mutation_delivery( + mutation, + verified.context.signer_id, + verified.context.signer_id, + verified.context.final_recipient_id, + signer_principal.handle, + )] + }) + .unwrap_or_default(), ); } } @@ -1044,7 +1089,7 @@ impl RelayService { return Self::relay_response(frame.id(), CommunicationType::ErrorInvalidData); } }; - if let Err(error) = message_handlers::apply_verified_relay_content( + let application = match message_handlers::apply_verified_relay_content( &verified.context, &content, signer_principal.handle, @@ -1062,26 +1107,42 @@ impl RelayService { }, false, ) { - log!("Relay application dispatch failed: {}", error); - if let Err(queue_error) = - relay_queue::remove_for_frame(RouteTarget::User(destination), frame_id) - { - log!("Relay application queue cleanup failed: {}", queue_error); + Ok(application) => application, + Err(error) => { + log!("Relay application dispatch failed: {}", error); + if let Err(queue_error) = + relay_queue::remove_for_frame(RouteTarget::User(destination), frame_id) + { + log!("Relay application queue cleanup failed: {}", queue_error); + } + let _ = relay_replay::mark_rejected( + signer_principal.handle, + &verified.context.message_id, + ); + return Self::relay_response(frame.id(), CommunicationType::ErrorInvalidData); } - let _ = relay_replay::mark_rejected( - signer_principal.handle, - &verified.context.message_id, - ); - return Self::relay_response(frame.id(), CommunicationType::ErrorInvalidData); - } - client_event = client_event_from_relay( - content.message_type, - &content.content, - verified.context.signer_id, - destination, - &verified.context.message_id, - accepted_at, - ); + }; + client_event = application + .mutation + .as_ref() + .map(|mutation| { + client_event_from_mutation( + mutation, + verified.context.signer_id, + destination, + verified.context.signer_id, + ) + }) + .or_else(|| { + client_event_from_relay( + content.message_type, + &content.content, + verified.context.signer_id, + destination, + &verified.context.message_id, + accepted_at, + ) + }); if let Err(error) = relay_replay::mark_applied(signer_principal.handle, &verified.context.message_id) { @@ -1109,19 +1170,29 @@ impl RelayService { { log!("Relay queue state update failed: {}", error); } + let mut local_deliveries = Vec::new(); + if let Some(mutation) = source_mutation.as_ref() { + local_deliveries.push(mutation_delivery( + mutation, + verified.context.signer_id, + verified.context.signer_id, + verified.context.final_recipient_id, + signer_principal.handle, + )); + } + if let Some(frame) = client_event { + local_deliveries.push(ClientDelivery { + recipient: recipient_principal.handle, + frame, + }); + } Self::relay_success( frame.id(), local_iota_id, &verified.context.message_id, accepted_at, signer_is_local, - client_event - .into_iter() - .map(|frame| ClientDelivery { - recipient: recipient_principal.handle, - frame, - }) - .collect(), + local_deliveries, ) } @@ -1239,7 +1310,12 @@ impl RelayService { return Self::relay_response(Some(frame_id), CommunicationType::ErrorInvalidData); } relay_replay::RelayReservation::Existing { state, .. } if state == "delivered" => { - return federated_success(frame_id, &normalized.message_id, accepted_at); + return federated_success( + frame_id, + &normalized.message_id, + accepted_at, + Vec::new(), + ); } relay_replay::RelayReservation::Existing { state, .. } if state == "rejected" => { return Self::relay_response(Some(frame_id), CommunicationType::ErrorInvalidData); @@ -1253,18 +1329,59 @@ impl RelayService { ); if let Some(recipient) = recipient_local { - if !already_applied - && let Err(error) = message_handlers::apply_relay_application_content( - &application, + let mut source_mutation = None; + let mut destination_mutation = None; + if !already_applied { + if let Some(sender) = signer_local { + let source_application = message_handlers::apply_relay_application_content( + &message_handlers::RelayApplicationContext { + hosted_sender: Some(sender), + hosted_recipient: None, + ..application.clone() + }, + &context, + &content, + accepted_at, + ); + match source_application { + Ok(application) => source_mutation = application.mutation, + Err(error) => { + log!("Relay V2 shared-Iota origin application failed: {error}"); + let _ = relay_replay::mark_rejected( + normalized.signer.handle, + &normalized.message_id, + ); + return Self::relay_response( + Some(frame_id), + CommunicationType::ErrorInvalidData, + ); + } + } + } + let destination_application = message_handlers::apply_relay_application_content( + &message_handlers::RelayApplicationContext { + hosted_sender: None, + hosted_recipient: Some(recipient), + ..application.clone() + }, &context, &content, accepted_at, - ) - { - log!("Relay V2 destination application failed: {error}"); - let _ = - relay_replay::mark_rejected(normalized.signer.handle, &normalized.message_id); - return Self::relay_response(Some(frame_id), CommunicationType::ErrorInvalidData); + ); + match destination_application { + Ok(application) => destination_mutation = application.mutation, + Err(error) => { + log!("Relay V2 destination application failed: {error}"); + let _ = relay_replay::mark_rejected( + normalized.signer.handle, + &normalized.message_id, + ); + return Self::relay_response( + Some(frame_id), + CommunicationType::ErrorInvalidData, + ); + } + } } let relay_identity = relay_queue::RelayIdentity { signer: normalized.signer.handle, @@ -1285,45 +1402,86 @@ impl RelayService { return Self::relay_response(Some(frame_id), CommunicationType::ErrorInternal); } let _ = relay_replay::mark_queued(normalized.signer.handle, &normalized.message_id); - let delivery = client_event_from_relay( - message_type, - &content.content, - normalized.signer.principal.user_id, - normalized.recipient.principal.user_id, + let destination_event = destination_mutation + .as_ref() + .map(|mutation| { + client_event_from_mutation( + mutation, + normalized.signer.principal.user_id, + normalized.recipient.principal.user_id, + normalized.signer.principal.user_id, + ) + }) + .or_else(|| { + client_event_from_relay( + message_type, + &content.content, + normalized.signer.principal.user_id, + normalized.recipient.principal.user_id, + &normalized.message_id, + accepted_at, + ) + }); + let mut local_deliveries = Vec::new(); + if let Some(mutation) = source_mutation.as_ref() { + local_deliveries.push(mutation_delivery( + mutation, + normalized.signer.principal.user_id, + normalized.signer.principal.user_id, + normalized.recipient.principal.user_id, + normalized.signer.handle, + )); + } + if let Some(frame) = destination_event { + local_deliveries.push(ClientDelivery { + recipient: normalized.recipient.handle, + frame, + }); + } + return federated_success( + frame_id, &normalized.message_id, accepted_at, + local_deliveries, ); - return match federated_success(frame_id, &normalized.message_id, accepted_at) { - RelayOutcome::Accepted { - ingress_response, .. - } => RelayOutcome::Accepted { - ingress_response, - local_deliveries: delivery - .into_iter() - .map(|frame| ClientDelivery { - recipient: normalized.recipient.handle, - frame, - }) - .collect(), - }, - outcome => outcome, - }; } let Some(signer) = signer_local else { return Self::relay_response(Some(frame_id), CommunicationType::ErrorInvalidData); }; - if !already_applied - && let Err(error) = message_handlers::apply_relay_application_content( + let source_mutation = if already_applied { + None + } else { + match message_handlers::apply_relay_application_content( &application, &context, &content, accepted_at, - ) - { - log!("Relay V2 origin application failed: {error}"); - return Self::relay_response(Some(frame_id), CommunicationType::ErrorInvalidData); - } + ) { + Ok(application) => application.mutation, + Err(error) => { + log!("Relay V2 origin application failed: {error}"); + return Self::relay_response( + Some(frame_id), + CommunicationType::ErrorInvalidData, + ); + } + } + }; + let source_deliveries = || { + source_mutation + .as_ref() + .map(|mutation| { + vec![mutation_delivery( + mutation, + normalized.signer.principal.user_id, + normalized.signer.principal.user_id, + normalized.recipient.principal.user_id, + normalized.signer.handle, + )] + }) + .unwrap_or_default() + }; let Some(route) = route_destination(&normalized.recipient.home) else { return Self::relay_response(Some(frame_id), CommunicationType::ErrorNoIota); }; @@ -1368,7 +1526,12 @@ impl RelayService { log!("Relay V2 acknowledgement storage failed: {error}"); return Self::relay_response(Some(frame_id), CommunicationType::ErrorInternal); } - federated_success(frame_id, &relay_message_id, destination_accepted_at) + federated_success( + frame_id, + &relay_message_id, + destination_accepted_at, + source_deliveries(), + ) } Ok(RouteOutcome::Accepted { .. }) => { let _ = iota_storage::util::relay_queue::quarantine_for_frame( @@ -1380,9 +1543,12 @@ impl RelayService { Ok(RouteOutcome::Rejected { response_type }) => { Self::relay_response(Some(frame_id), response_type) } - Ok(RouteOutcome::Retryable { .. }) | Err(_) => { - federated_success(frame_id, &normalized.message_id, accepted_at) - } + Ok(RouteOutcome::Retryable { .. }) | Err(_) => federated_success( + frame_id, + &normalized.message_id, + accepted_at, + source_deliveries(), + ), } } @@ -1495,7 +1661,12 @@ impl RelayService { } } -fn federated_success(frame_id: u32, message_id: &str, accepted_at: i64) -> RelayOutcome { +fn federated_success( + frame_id: u32, + message_id: &str, + accepted_at: i64, + local_deliveries: Vec, +) -> RelayOutcome { RelayOutcome::Accepted { ingress_response: CommunicationValue::new(CommunicationType::Success) .with_id(frame_id) @@ -1507,7 +1678,7 @@ fn federated_success(frame_id: u32, message_id: &str, accepted_at: i64) -> Relay DataType::RelayAcceptedAt, DataValue::SignedNumber(accepted_at.into()), ), - local_deliveries: Vec::new(), + local_deliveries, } } @@ -1535,6 +1706,84 @@ fn record_origin_delivery_failure( } } +fn client_event_from_mutation( + mutation: &message_handlers::AppliedMessageMutation, + actor_id: u64, + account_id: u64, + chat_partner_id: u64, +) -> CommunicationValue { + let event = match mutation { + message_handlers::AppliedMessageMutation::State { send_time, state } => { + CommunicationValue::new(CommunicationType::MessageState) + .add_typed_default( + DataType::SendTime, + DataValue::SignedNumber((*send_time).into()), + ) + .add_typed_default( + DataType::MessageState, + DataValue::Str(state.as_str().to_string()), + ) + } + message_handlers::AppliedMessageMutation::Edit { + send_time, + content, + version, + } => CommunicationValue::new(CommunicationType::MessageEditLive) + .add_typed_default(DataType::Content, DataValue::Str(content.clone())) + .add_typed_default( + DataType::SendTime, + DataValue::SignedNumber((*send_time).into()), + ) + .add_typed_default( + DataType::VersionNumber, + DataValue::SignedNumber((*version).into()), + ), + message_handlers::AppliedMessageMutation::Reaction { + send_time, + reaction, + accepted, + } => CommunicationValue::new(CommunicationType::MessageReactionLive) + .add_typed_default(DataType::Reaction, DataValue::Str(reaction.clone())) + .add_typed_default( + DataType::SendTime, + DataValue::SignedNumber((*send_time).into()), + ) + .add_typed_default( + DataType::SenderId, + DataValue::UnsignedNumber(actor_id.into()), + ) + .add_typed_default(DataType::Accepted, DataValue::Bool(*accepted)), + message_handlers::AppliedMessageMutation::Delete { send_time } => { + CommunicationValue::new(CommunicationType::MessageDeleteLive).add_typed_default( + DataType::SendTime, + DataValue::SignedNumber((*send_time).into()), + ) + } + }; + + event + .with_id(next_client_event_id()) + .with_sender(actor_id) + .with_receiver(account_id) + .add_typed_default( + DataType::ChatPartnerId, + DataValue::UnsignedNumber(chat_partner_id.into()), + ) +} + +fn mutation_delivery( + mutation: &message_handlers::AppliedMessageMutation, + actor_id: u64, + account_id: u64, + chat_partner_id: u64, + recipient: iota_identity::PrincipalHandle, +) -> ClientDelivery { + ClientDelivery { + recipient, + frame: client_event_from_mutation(mutation, actor_id, account_id, chat_partner_id), + } +} + fn client_event_from_relay( message_type: CommunicationType, payload: &DataValue, @@ -1619,67 +1868,6 @@ fn client_event_from_relay( .add_typed_default(DataType::Payload, DataValue::Str("available".to_string())), ) } - CommunicationType::MessageEdit => { - let frame = CommunicationValue::new(CommunicationType::MessageEdit) - .with_payload(payload.clone()); - Some( - CommunicationValue::new(CommunicationType::MessageEditLive) - .with_id(event_id) - .with_sender(signer_id) - .with_receiver(receiver) - .add_typed_default( - DataType::Content, - frame.get_data(DataType::Content)?.clone(), - ) - .add_typed_default( - DataType::SendTime, - frame.get_data(DataType::SendTime)?.clone(), - ) - .add_typed_default( - DataType::VersionNumber, - frame.get_data(DataType::VersionNumber)?.clone(), - ) - .add_typed_default(DataType::ChatPartnerId, sender.clone()), - ) - } - CommunicationType::MessageReactionAdd | CommunicationType::MessageReactionRemove => { - let frame = CommunicationValue::new(message_type).with_payload(payload.clone()); - Some( - CommunicationValue::new(CommunicationType::MessageReactionLive) - .with_id(event_id) - .with_sender(signer_id) - .with_receiver(receiver) - .add_typed_default( - DataType::Reaction, - frame.get_data(DataType::Reaction)?.clone(), - ) - .add_typed_default( - DataType::SendTime, - frame.get_data(DataType::SendTime)?.clone(), - ) - .add_typed_default(DataType::ChatPartnerId, sender.clone()) - .add_typed_default(DataType::SenderId, sender) - .add_typed_default( - DataType::Accepted, - DataValue::Bool(message_type == CommunicationType::MessageReactionAdd), - ), - ) - } - CommunicationType::MessageDelete | CommunicationType::MessageDeleteLive => { - let frame = CommunicationValue::new(CommunicationType::MessageDelete) - .with_payload(payload.clone()); - Some( - CommunicationValue::new(CommunicationType::MessageDeleteLive) - .with_id(event_id) - .with_sender(signer_id) - .with_receiver(receiver) - .add_typed_default( - DataType::SendTime, - frame.get_data(DataType::SendTime)?.clone(), - ) - .add_typed_default(DataType::ChatPartnerId, sender), - ) - } CommunicationType::AddConversation if recipient_id < signer_id => { let chat_id = format!("{recipient_id}:{signer_id}"); Some( @@ -1872,4 +2060,24 @@ mod tests { assert!(recipient_block_policy_applies(true, 42, 43)); assert!(!recipient_block_policy_applies(false, 42, 43)); } + + #[test] + fn mutation_events_use_receiving_account_chat_partner() { + let event = client_event_from_mutation( + &message_handlers::AppliedMessageMutation::Reaction { + send_time: 12, + reaction: "thumbsup".into(), + accepted: true, + }, + 7, + 8, + 9, + ); + + assert!(event.is_type(CommunicationType::MessageReactionLive)); + assert_eq!(event.receiver(), Some(8)); + assert_eq!(event.sender(), Some(7)); + assert_eq!(event.get_data(DataType::ChatPartnerId).as_number(), Some(9)); + assert_eq!(event.get_data(DataType::SenderId).as_number(), Some(7)); + } } diff --git a/iota-storage/src/util/chat_files.rs b/iota-storage/src/util/chat_files.rs index a9d5664..c697c47 100644 --- a/iota-storage/src/util/chat_files.rs +++ b/iota-storage/src/util/chat_files.rs @@ -14,6 +14,12 @@ pub enum MessageState { Sending, } +#[derive(Clone, Debug)] +pub struct MessageStateUpdate { + pub send_time: i64, + pub state: MessageState, +} + impl MessageState { pub fn as_str(&self) -> &'static str { match self { @@ -979,23 +985,29 @@ pub fn change_message_state_by_relay_id( relay_signer_principal: iota_identity::PrincipalHandle, relay_message_id: &str, new_state: MessageState, -) -> Result<(), StorageError> { +) -> Result, StorageError> { db::with_db(|conn| { let tx = conn.unchecked_transaction()?; - let Some((msg_id, current)) = tx + let Some((msg_id, current, message_time)) = tx .query_row( - "SELECT id, message_state FROM messages WHERE storage_owner = ?1 AND relay_signer_principal = ?2 AND relay_message_id = ?3", + "SELECT id, message_state, message_time FROM messages WHERE storage_owner = ?1 AND relay_signer_principal = ?2 AND relay_message_id = ?3", params![storage_owner, relay_signer_principal.0, relay_message_id], - |row| Ok((row.get::<_, i64>(0)?, row.get::<_, String>(1)?)), + |row| { + Ok(( + row.get::<_, i64>(0)?, + row.get::<_, String>(1)?, + row.get::<_, i64>(2)?, + )) + }, ) .optional()? else { - return Ok(()); + return Ok(None); }; - let state = MessageState::from_str(¤t).upgrade(new_state).as_str(); + let final_state = MessageState::from_str(¤t).upgrade(new_state); tx.execute( "UPDATE messages SET message_state = ?1 WHERE id = ?2", - params![state, msg_id], + params![final_state.as_str(), msg_id], )?; sync::record_event( &tx, @@ -1005,7 +1017,10 @@ pub fn change_message_state_by_relay_id( Operation::Upsert, )?; tx.commit()?; - Ok(()) + Ok(Some(MessageStateUpdate { + send_time: message_time, + state: final_state, + })) }) } @@ -1020,7 +1035,7 @@ pub fn record_message_receipt( receipt_type: MessageState, event_at: i64, recorded_at: i64, -) -> Result<(), StorageError> { +) -> Result, StorageError> { let receipt_type = match receipt_type { MessageState::Received => "received", MessageState::Read => "read", @@ -1028,11 +1043,19 @@ pub fn record_message_receipt( }; db::with_db(|conn| { let tx = conn.unchecked_transaction()?; - let Some((message_id, external_principal, authored_at, history_deleted)) = tx + let Some((message_id, message_time, external_principal, authored_at, history_deleted)) = tx .query_row( - "SELECT id, external_principal, authored_at, history_deleted FROM messages WHERE storage_owner = ?1 AND relay_signer_principal = ?2 AND relay_message_id = ?3", + "SELECT id, message_time, external_principal, authored_at, history_deleted FROM messages WHERE storage_owner = ?1 AND relay_signer_principal = ?2 AND relay_message_id = ?3", params![storage_owner, target_signer_principal.0, target_message_id], - |row| Ok((row.get::<_, i64>(0)?, row.get::<_, i64>(1)?, row.get::<_, Option>(2)?, row.get::<_, i64>(3)?)), + |row| { + Ok(( + row.get::<_, i64>(0)?, + row.get::<_, i64>(1)?, + row.get::<_, i64>(2)?, + row.get::<_, Option>(3)?, + row.get::<_, i64>(4)?, + )) + }, ) .optional()? else { @@ -1051,7 +1074,7 @@ pub fn record_message_receipt( )); } if history_deleted != 0 { - return Ok(()); + return Ok(None); } tx.execute( "INSERT OR IGNORE INTO message_receipts (storage_owner, target_signer_id, target_signer_principal, target_message_id, receipt_signer_id, receipt_signer_principal, receipt_message_id, receipt_type, event_at, recorded_at) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10)", @@ -1062,17 +1085,15 @@ pub fn record_message_receipt( } else { ("client_received_at", "client_received_recorded_at") }; - let state = MessageState::from_str(&tx.query_row( + let final_state = MessageState::from_str(&tx.query_row( "SELECT message_state FROM messages WHERE id = ?1", [message_id], |row| row.get::<_, String>(0), )?) - .upgrade(MessageState::from_str(receipt_type)) - .as_str() - .to_string(); + .upgrade(MessageState::from_str(receipt_type)); tx.execute( &format!("UPDATE messages SET {state_column} = COALESCE({state_column}, ?1), {recorded_column} = COALESCE({recorded_column}, ?2), message_state = ?3 WHERE id = ?4"), - params![event_at, recorded_at, state, message_id], + params![event_at, recorded_at, final_state.as_str(), message_id], )?; sync::record_event( &tx, @@ -1082,7 +1103,10 @@ pub fn record_message_receipt( Operation::Upsert, )?; tx.commit()?; - Ok(()) + Ok(Some(MessageStateUpdate { + send_time: message_time, + state: final_state, + })) }) } @@ -1579,6 +1603,10 @@ mod tests { MessageState::Received.upgrade(MessageState::Read), MessageState::Read ); + assert_eq!( + MessageState::Read.upgrade(MessageState::Received), + MessageState::Read + ); } #[test] diff --git a/iota-storage/tests/message_offsets.rs b/iota-storage/tests/message_offsets.rs new file mode 100644 index 0000000..102ed41 --- /dev/null +++ b/iota-storage/tests/message_offsets.rs @@ -0,0 +1,84 @@ +use iota_identity::PrincipalHandle; +use iota_storage::util::chat_files::{self, MessageState, NewMessage}; + +#[test] +fn message_offset_counts_visible_newer_messages_in_one_chat() { + let storage = tempfile::tempdir().unwrap(); + iota_util::file_util::configure_storage_directory(storage.path().to_owned()); + iota_storage::util::db::initialize_database().unwrap(); + + iota_storage::util::db::with_db(|connection| { + connection.execute( + "INSERT INTO principals (authority_kind, authority_id, remote_user_id, descriptor_revision, last_resolved_at) VALUES ('iota', 'iota:peer', 2, 1, 1)", + [], + )?; + let principal = connection.last_insert_rowid(); + connection.execute( + "INSERT INTO contacts (storage_owner, user_id, principal_handle, created_at) VALUES (1, 2, ?1, 1)", + [principal], + )?; + connection.execute( + "INSERT INTO principals (authority_kind, authority_id, remote_user_id, descriptor_revision, last_resolved_at) VALUES ('iota', 'iota:other', 3, 1, 1)", + [], + )?; + let other_principal = connection.last_insert_rowid(); + connection.execute( + "INSERT INTO contacts (storage_owner, user_id, principal_handle, created_at) VALUES (1, 3, ?1, 1)", + [other_principal], + )?; + Ok(()) + }) + .unwrap(); + + for (send_time, relay_message_id, external_user, external_principal) in [ + (10, "first", 2, 1), + (20, "second", 2, 1), + (30, "hidden", 2, 1), + (40, "other-chat", 3, 2), + ] { + chat_files::add_message(NewMessage { + relay_signer_id: external_user, + relay_signer_principal: PrincipalHandle(external_principal), + relay_message_id, + authored_at: send_time, + send_time, + storage_owner: 1, + external_user, + external_principal: PrincipalHandle(external_principal), + sent_by_self: false, + content: relay_message_id, + height: 0, + key_version: 1, + reply_to: None, + origin_iota_received_at: None, + destination_iota_received_at: Some(send_time), + initial_state: MessageState::Sent, + }) + .unwrap(); + } + + iota_storage::util::db::with_db(|connection| { + connection.execute( + "UPDATE messages SET history_deleted = 1 WHERE relay_message_id = 'hidden'", + [], + )?; + Ok(()) + }) + .unwrap(); + + let newest = chat_files::get_message_with_offset(1, 2, 20) + .unwrap() + .unwrap(); + assert_eq!(newest.1, 0); + + let first = chat_files::get_message_with_offset(1, 2, 10) + .unwrap() + .unwrap(); + assert_eq!(first.1, 1); + + assert!( + chat_files::get_message_with_offset(1, 3, 10) + .unwrap() + .is_none() + ); +}