From 3bfec968485d82e329f1d53b03e1c96c0de7a2b7 Mon Sep 17 00:00:00 2001 From: Alex Emmet <111742636+Alex-Emmet@users.noreply.github.com> Date: Fri, 28 Aug 2026 17:00:44 +0200 Subject: [PATCH] [Add] timestamping --- iota-connection/src/message_handlers.rs | 3 + iota-storage/src/util/relay_replay.rs | 37 +++--- iota-storage/src/util/sync.rs | 2 +- omikron-connector/src/omikron_connection.rs | 137 +++++++++++++++----- 4 files changed, 133 insertions(+), 46 deletions(-) diff --git a/iota-connection/src/message_handlers.rs b/iota-connection/src/message_handlers.rs index 1f247aa..d692be4 100644 --- a/iota-connection/src/message_handlers.rs +++ b/iota-connection/src/message_handlers.rs @@ -185,6 +185,9 @@ pub fn apply_verified_relay_content( &context.type_map, ) .ok_or_else(|| "Relay MessageSend is missing RelayMessageId".to_string())?; + if relay_message_id != context.message_id { + return Err("Relay MessageSend identity does not match its protected message ID".into()); + } chat_files::add_message(chat_files::NewMessage { relay_signer_id: sender_id, relay_message_id, diff --git a/iota-storage/src/util/relay_replay.rs b/iota-storage/src/util/relay_replay.rs index db3cf0e..0c5bddc 100644 --- a/iota-storage/src/util/relay_replay.rs +++ b/iota-storage/src/util/relay_replay.rs @@ -81,13 +81,12 @@ pub fn mark_delivered_for_frame(destination_id: u64, frame_id: u32) -> Result<() }) } -pub fn mark_state(signer_id: u64, message_id: &str, state: &str) -> Result<(), StorageError> { - if !matches!( - state, - "received" | "applied" | "queued" | "delivered" | "rejected" - ) { - return Err(StorageError::Other("invalid relay inbox state".into())); - } +fn mark_transition( + signer_id: u64, + message_id: &str, + state: &str, + column: &str, +) -> Result<(), StorageError> { let signer_id = i64::try_from(signer_id) .map_err(|_| StorageError::Other("relay signer ID exceeds SQLite range".into()))?; let timestamp = std::time::SystemTime::now() @@ -95,14 +94,6 @@ pub fn mark_state(signer_id: u64, message_id: &str, state: &str) -> Result<(), S .unwrap_or_default() .as_millis() as i64; db::with_db(|connection| { - let column = match state { - "applied" => "applied_at", - "queued" => "queued_at", - "delivered" => "downstream_acked_at", - "rejected" => "rejected_at", - "received" => "accepted_at", - _ => return Err(StorageError::Other("invalid relay inbox state".into())), - }; connection.execute( &format!("UPDATE relay_inbox SET state = ?3, {column} = COALESCE({column}, ?4) WHERE signer_id = ?1 AND message_id = ?2"), params![signer_id, message_id, state, timestamp], @@ -111,6 +102,22 @@ pub fn mark_state(signer_id: u64, message_id: &str, state: &str) -> Result<(), S }) } +pub fn mark_applied(signer_id: u64, message_id: &str) -> Result<(), StorageError> { + mark_transition(signer_id, message_id, "applied", "applied_at") +} + +pub fn mark_queued(signer_id: u64, message_id: &str) -> Result<(), StorageError> { + mark_transition(signer_id, message_id, "queued", "queued_at") +} + +pub fn mark_downstream_acked(signer_id: u64, message_id: &str) -> Result<(), StorageError> { + mark_transition(signer_id, message_id, "delivered", "downstream_acked_at") +} + +pub fn mark_rejected(signer_id: u64, message_id: &str) -> Result<(), StorageError> { + mark_transition(signer_id, message_id, "rejected", "rejected_at") +} + pub fn prune_completed(before_terminal_at: i64) -> Result<(), StorageError> { db::with_db(|connection| { connection.execute( diff --git a/iota-storage/src/util/sync.rs b/iota-storage/src/util/sync.rs index b84c94a..a358a1d 100644 --- a/iota-storage/src/util/sync.rs +++ b/iota-storage/src/util/sync.rs @@ -4,7 +4,7 @@ use crate::util::db; use rusqlite::{Transaction, params}; use std::collections::BTreeMap; -pub const CACHE_SCHEMA_VERSION: i64 = 1; +pub const CACHE_SCHEMA_VERSION: i64 = 2; #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub enum EntityType { diff --git a/omikron-connector/src/omikron_connection.rs b/omikron-connector/src/omikron_connection.rs index 0a3a6a3..0c37013 100644 --- a/omikron-connector/src/omikron_connection.rs +++ b/omikron-connector/src/omikron_connection.rs @@ -893,6 +893,7 @@ impl OmikronConnection { iota_id: u64, relay_message_id: &str, accepted_at: i64, + include_origin_timestamp: bool, ) { let Some(frame_id) = frame_id else { return }; let response = CommunicationValue::new(CommunicationType::Success) @@ -906,6 +907,19 @@ impl OmikronConnection { DataType::RelayAcceptedAt, DataValue::SignedNumber(accepted_at.into()), ); + let response = if include_origin_timestamp { + response + .add_typed_default( + DataType::OriginIotaReceivedAt, + DataValue::SignedNumber(accepted_at.into()), + ) + .add_typed_default( + DataType::DestinationIotaReceivedAt, + DataValue::SignedNumber(accepted_at.into()), + ) + } else { + response + }; if let Err(error) = self.send_message(&response).await { log!("Relay response could not be sent: {}", error); } @@ -1039,6 +1053,63 @@ impl OmikronConnection { relay_replay::RelayReservation::Existing { .. } => false, }; + /* A shared Iota owns both independent replicas before delivering to its + * local recipient. The destination path below writes the recipient copy. */ + if signer_is_local && recipient_is_local && !already_applied { + let content = match open_verified_relay_content( + &verified, + &[&keyring], + verified.context.signer_id, + ) { + Ok(value) => value, + Err(error) => { + log!("Relay shared-Iota origin content verification failed: {}", error); + let _ = relay_replay::mark_rejected( + verified.context.signer_id, + &verified.context.message_id, + ); + self.send_relay_response(frame.id(), CommunicationType::ErrorInvalidData) + .await; + return; + } + }; + let owner = match i64::try_from(verified.context.signer_id) { + Ok(value) => value, + Err(_) => { + self.send_relay_response(frame.id(), CommunicationType::ErrorInvalidData) + .await; + return; + } + }; + if let Err(error) = message_handlers::apply_verified_relay_content( + &verified.context, + &content, + accepted_at, + owner, + true, + ) { + log!("Relay shared-Iota origin application failed: {}", error); + let _ = relay_replay::mark_rejected( + verified.context.signer_id, + &verified.context.message_id, + ); + self.send_relay_response(frame.id(), CommunicationType::ErrorInvalidData) + .await; + return; + } + if let Err(error) = chat_files::record_destination_iota_received( + owner, + owner, + &verified.context.message_id, + accepted_at, + ) { + log!("Relay shared-Iota destination timestamp storage failed: {}", error); + self.send_relay_response(frame.id(), CommunicationType::ErrorInternal) + .await; + return; + } + } + if signer_is_local && !recipient_is_local { if !already_applied { let content = match open_verified_relay_content( @@ -1049,10 +1120,9 @@ impl OmikronConnection { Ok(value) => value, Err(error) => { log!("Relay origin content verification failed: {}", error); - let _ = relay_replay::mark_state( + let _ = relay_replay::mark_rejected( verified.context.signer_id, &verified.context.message_id, - "rejected", ); self.send_relay_response(frame.id(), CommunicationType::ErrorInvalidData) .await; @@ -1067,10 +1137,9 @@ impl OmikronConnection { true, ) { log!("Relay origin application failed: {}", error); - let _ = relay_replay::mark_state( + let _ = relay_replay::mark_rejected( verified.context.signer_id, &verified.context.message_id, - "rejected", ); self.send_relay_response(frame.id(), CommunicationType::ErrorInvalidData) .await; @@ -1119,10 +1188,9 @@ impl OmikronConnection { .await; return; } - if let Err(error) = relay_replay::mark_state( + if let Err(error) = relay_replay::mark_queued( verified.context.signer_id, &verified.context.message_id, - "queued", ) { log!("Relay origin state update failed: {}", error); } @@ -1134,7 +1202,7 @@ impl OmikronConnection { let returned_id = response .get_data(DataType::RelayMessageId) .as_str(); - let accepted_at = response + let destination_accepted_at = response .get_data(DataType::RelayAcceptedAt) .as_number() .and_then(|value| i64::try_from(value).ok()); @@ -1144,19 +1212,23 @@ impl OmikronConnection { .await; return; } - if let Some(accepted_at) = accepted_at { - if let (Ok(owner), Ok(signer)) = ( - i64::try_from(verified.context.signer_id), - i64::try_from(verified.context.signer_id), + let Some(destination_accepted_at) = destination_accepted_at else { + log!("Relay acknowledgement is missing RelayAcceptedAt"); + self.send_relay_response(frame.id(), CommunicationType::ErrorInvalidData) + .await; + return; + }; + if let (Ok(owner), Ok(signer)) = ( + i64::try_from(verified.context.signer_id), + i64::try_from(verified.context.signer_id), + ) { + if let Err(error) = chat_files::record_destination_iota_received( + owner, + signer, + &verified.context.message_id, + destination_accepted_at, ) { - if let Err(error) = chat_files::record_destination_iota_received( - owner, - signer, - &verified.context.message_id, - accepted_at, - ) { - log!("Relay destination acknowledgement storage failed: {}", error); - } + log!("Relay destination acknowledgement storage failed: {}", error); } } if let Err(error) = relay_queue::acknowledge_iota(router, frame_id) { @@ -1165,14 +1237,22 @@ impl OmikronConnection { error ); } - if let Err(error) = relay_replay::mark_state( + if let Err(error) = relay_replay::mark_downstream_acked( verified.context.signer_id, &verified.context.message_id, - "delivered", ) { log!("Relay origin delivery state update failed: {}", error); } - let response = response.with_id(frame_id); + let response = response + .add_typed_default( + DataType::OriginIotaReceivedAt, + DataValue::SignedNumber(accepted_at.into()), + ) + .add_typed_default( + DataType::DestinationIotaReceivedAt, + DataValue::SignedNumber(destination_accepted_at.into()), + ) + .with_id(frame_id); if let Err(error) = self.send_message(&response).await { log!("Relay response could not be sent: {}", error); } @@ -1248,10 +1328,9 @@ impl OmikronConnection { queue_error ); } - let _ = relay_replay::mark_state( + let _ = relay_replay::mark_rejected( verified.context.signer_id, &verified.context.message_id, - "rejected", ); self.send_relay_response(frame.id(), CommunicationType::ErrorInvalidData) .await; @@ -1284,19 +1363,17 @@ impl OmikronConnection { { log!("Relay application queue cleanup failed: {}", queue_error); } - let _ = relay_replay::mark_state( + let _ = relay_replay::mark_rejected( verified.context.signer_id, &verified.context.message_id, - "rejected", ); self.send_relay_response(frame.id(), CommunicationType::ErrorInvalidData) .await; return; } - if let Err(error) = relay_replay::mark_state( + if let Err(error) = relay_replay::mark_applied( verified.context.signer_id, &verified.context.message_id, - "applied", ) { log!("Relay application state update failed: {}", error); self.send_relay_response(frame.id(), CommunicationType::ErrorInternal) @@ -1305,10 +1382,9 @@ impl OmikronConnection { } } - if let Err(error) = relay_replay::mark_state( + if let Err(error) = relay_replay::mark_queued( verified.context.signer_id, &verified.context.message_id, - "queued", ) { log!("Relay queue state update failed: {}", error); } @@ -1317,6 +1393,7 @@ impl OmikronConnection { local_iota_id, &verified.context.message_id, accepted_at, + signer_is_local, ) .await; if let Err(error) = self.send_message(&forwarded).await {