From 233b9440079bcf0ff0008e5a048be9de8767c5cf Mon Sep 17 00:00:00 2001 From: Alex Emmet <111742636+Alex-Emmet@users.noreply.github.com> Date: Fri, 28 Aug 2026 16:21:11 +0200 Subject: [PATCH 1/2] [Add] Replyjumps --- mtp-type-maps | 2 +- src/config.rs | 14 ++++++++++++++ src/main.rs | 11 +++++++++++ src/omega/omega_connection.rs | 3 +-- src/rho/app_connection.rs | 4 +--- src/rho/client_connection.rs | 25 ++++++++++++++++++++++--- src/rho/iota_connection.rs | 2 +- src/rho/relay_router.rs | 23 +++++++++++++++++------ src/rho/rho_connection.rs | 10 ++++++++++ 9 files changed, 78 insertions(+), 16 deletions(-) diff --git a/mtp-type-maps b/mtp-type-maps index 486541b..f3c5037 160000 --- a/mtp-type-maps +++ b/mtp-type-maps @@ -1 +1 @@ -Subproject commit 486541b9483356ff49ff3ec7016f87d3ecbeaa0e +Subproject commit f3c5037b0a099ae5389486eec77471bd5addab4a diff --git a/src/config.rs b/src/config.rs index e69f383..e9ef9c8 100644 --- a/src/config.rs +++ b/src/config.rs @@ -129,6 +129,20 @@ where } } +fn parse_positive_or_default(name: &'static str, default: T) -> Result +where + T: std::str::FromStr + PartialEq + Default, +{ + let value = parse_or_default(name, default)?; + if value == T::default() { + return Err(ConfigError::InvalidValue { + name, + kind: "positive number", + }); + } + Ok(value) +} + fn livekit_from_environment() -> Result, ConfigError> { livekit_from_values( env::var("LIVEKIT_HOSTNAME").ok(), diff --git a/src/main.rs b/src/main.rs index 7202ca2..ffc3a70 100644 --- a/src/main.rs +++ b/src/main.rs @@ -57,6 +57,17 @@ async fn main() { return; } }; + let keyring = match load_or_create_keyring( + &identity_secret, + KEYRING_PATH, + PUBLIC_KEY_PATH, + ) { + Ok(keyring) => keyring, + Err(error) => { + eprintln!("Unable to load Omikron keyring: {error}"); + return; + } + }; let public_key = match omega_database_public_key(&keyring) { Ok(public_key) => public_key, Err(error) => { diff --git a/src/omega/omega_connection.rs b/src/omega/omega_connection.rs index 3e7755c..f785aed 100644 --- a/src/omega/omega_connection.rs +++ b/src/omega/omega_connection.rs @@ -657,8 +657,7 @@ impl OmegaConnection { } }; let response = match relay_router::route_from_omega(&self.rho, cv).await { - Ok(()) => CommunicationValue::new(CommunicationType::Success) - .with_id(request_id), + Ok(response) => response.with_id(request_id), Err(error) => { log_err!( self.omikron_id as i64, diff --git a/src/rho/app_connection.rs b/src/rho/app_connection.rs index c6af5e9..3c07452 100644 --- a/src/rho/app_connection.rs +++ b/src/rho/app_connection.rs @@ -135,9 +135,7 @@ impl AppConnection { None => Err(relay_router::RelayRouteError::DestinationIotaNotLocal), }; let response = match result { - Ok(()) => { - CommunicationValue::new(CommunicationType::Success).with_id(message_id) - } + Ok(response) => response.with_id(message_id), Err(error) => { log_err!( self.user_id as i64, diff --git a/src/rho/client_connection.rs b/src/rho/client_connection.rs index 5a25ec0..d87eea1 100644 --- a/src/rho/client_connection.rs +++ b/src/rho/client_connection.rs @@ -166,8 +166,8 @@ impl ClientConnection { None => Err(relay_router::RelayRouteError::DestinationIotaNotLocal), }; let response = match result { - Ok(()) => { - CommunicationValue::new(CommunicationType::Success).with_id(message_id) + Ok(response) => { + response.with_id(message_id) } Err(error) => { log_err!( @@ -201,6 +201,25 @@ impl ClientConnection { return; } + if matches!( + relay_router::message_security_class(&cv), + relay_router::MessageSecurityClass::AuthenticatedPeerControl + ) { + let Some(rho) = self.get_rho_connection().await else { + self.send_error_response(message_id, CommunicationType::ErrorNoIota) + .await; + return; + }; + match rho.await_peer_control(&cv.with_sender(self.user_id)).await { + Ok(response) => self.send_message(&response).await, + Err(_) => { + self.send_error_response(message_id, CommunicationType::ErrorInternal) + .await; + } + } + return; + } + // Compatibility for clients predating SetUserState. The target // user fields, if present, are deliberately ignored: an // authenticated connection may only change its own state. @@ -884,7 +903,7 @@ impl ClientConnection { let response = CommunicationValue::new(CommunicationType::LoadTxtRecord) .with_id(message_id) - .add_typed_default(DataType::AppContent, DataValue::Str(record_text)); + .add_typed_default(DataType::Content, DataValue::Str(record_text)); self.send_message(&response).await; return; } diff --git a/src/rho/iota_connection.rs b/src/rho/iota_connection.rs index bd822cb..51d89fd 100755 --- a/src/rho/iota_connection.rs +++ b/src/rho/iota_connection.rs @@ -276,7 +276,7 @@ impl IotaConnection { ) .await { - Ok(()) => CommunicationValue::new(CommunicationType::Success).with_id(message_id), + Ok(response) => response.with_id(message_id), Err(error) => { log_err!( self.iota_id as i64, diff --git a/src/rho/relay_router.rs b/src/rho/relay_router.rs index bd0f0ab..97a5384 100644 --- a/src/rho/relay_router.rs +++ b/src/rho/relay_router.rs @@ -140,7 +140,7 @@ pub async fn route_relay( state: &Arc, source: RelaySource, frame: CommunicationValue, -) -> Result<(), RelayRouteError> { +) -> Result { let (next_hop, frame) = prepare_frame(ensure_relay_frame_id(frame))?; validate_source_next_hop(source, next_hop)?; @@ -169,7 +169,7 @@ pub async fn route_relay( pub async fn route_from_omega( rho: &RhoManager, frame: CommunicationValue, -) -> Result<(), RelayRouteError> { +) -> Result { let (next_hop, frame) = prepare_frame(ensure_relay_frame_id(frame))?; let RouteTarget::Iota(iota_id) = next_hop else { return Err(RelayRouteError::InvalidRouteTarget(next_hop.id())); @@ -193,7 +193,7 @@ async fn route_from_iota( source_iota_id: u64, next_hop: RouteTarget, frame: CommunicationValue, -) -> Result<(), RelayRouteError> { +) -> Result { match next_hop { RouteTarget::User(user_id) => { let target = rho @@ -206,7 +206,7 @@ async fn route_from_iota( if !target.has_local_client(user_id).await { return Err(RelayRouteError::ClientOffline); } - target.send_relay_to_client(&frame).await.map_err(|error| { + target.send_relay_to_client(&frame).await.map(|_| CommunicationValue::new(CommunicationType::Success)).map_err(|error| { if error == "client offline" { RelayRouteError::ClientOffline } else { @@ -240,9 +240,9 @@ async fn route_from_iota( } } -fn route_response(response: CommunicationValue) -> Result<(), RelayRouteError> { +fn route_response(response: CommunicationValue) -> Result { if response.is_type(CommunicationType::Success) { - Ok(()) + Ok(response) } else { Err(RelayRouteError::Send(format!( "next Relay hop rejected the frame with {}", @@ -407,4 +407,15 @@ mod tests { ); assert_eq!(RouteTarget::from_wire_id(7), None); } + + #[test] + fn successful_route_response_is_not_reconstructed() { + let response = CommunicationValue::new(CommunicationType::Success) + .with_payload(DataValue::Bytes(vec![9, 8, 7])); + let routed = match super::route_response(response.clone()) { + Ok(value) => value, + Err(_) => return, + }; + assert_eq!(routed.payload(), response.payload()); + } } diff --git a/src/rho/rho_connection.rs b/src/rho/rho_connection.rs index f4aa658..c933acf 100644 --- a/src/rho/rho_connection.rs +++ b/src/rho/rho_connection.rs @@ -328,6 +328,16 @@ impl RhoConnection { .await } + pub async fn await_peer_control( + &self, + cv: &CommunicationValue, + ) -> Result { + self.iota_connection + .clone() + .await_response(cv, Some(std::time::Duration::from_secs(20))) + .await + } + pub async fn forward_relay_ack(&self, user_id: u64, frame_id: u32) { let acknowledgement = CommunicationValue::new(CommunicationType::Success) .with_id(frame_id) From 9b1cb4ad7ecfa604ffe830ff3ee7159e40a8b5ad Mon Sep 17 00:00:00 2001 From: Rasensprenger Date: Fri, 28 Aug 2026 18:00:49 +0300 Subject: [PATCH 2/2] chore(deps): update rust crate livekit-api to 0.6.0 --- Cargo.toml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/Cargo.toml b/Cargo.toml index b758cd8..101bb0d 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -28,7 +28,7 @@ tokio = { version = "1.53.0", features = ["full"] } log = "0.4" dotenv = "0.15.0" strum_macros = "0.28.0" -livekit-api = { version = "0.5.6", features = ["rustls-tls-native-roots"] } +livekit-api = { version = "0.6.0", features = ["rustls-tls-native-roots"] } livekit-protocol = "=0.7.10" thiserror = "2.0.19" trust-dns-resolver = "0.23.2"