diff --git a/mtp-type-maps b/mtp-type-maps index f3c5037..486541b 160000 --- a/mtp-type-maps +++ b/mtp-type-maps @@ -1 +1 @@ -Subproject commit f3c5037b0a099ae5389486eec77471bd5addab4a +Subproject commit 486541b9483356ff49ff3ec7016f87d3ecbeaa0e diff --git a/src/config.rs b/src/config.rs index e9ef9c8..e69f383 100644 --- a/src/config.rs +++ b/src/config.rs @@ -129,20 +129,6 @@ 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 ffc3a70..7202ca2 100644 --- a/src/main.rs +++ b/src/main.rs @@ -57,17 +57,6 @@ 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 f785aed..3e7755c 100644 --- a/src/omega/omega_connection.rs +++ b/src/omega/omega_connection.rs @@ -657,7 +657,8 @@ impl OmegaConnection { } }; let response = match relay_router::route_from_omega(&self.rho, cv).await { - Ok(response) => response.with_id(request_id), + Ok(()) => CommunicationValue::new(CommunicationType::Success) + .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 3c07452..c6af5e9 100644 --- a/src/rho/app_connection.rs +++ b/src/rho/app_connection.rs @@ -135,7 +135,9 @@ impl AppConnection { None => Err(relay_router::RelayRouteError::DestinationIotaNotLocal), }; let response = match result { - Ok(response) => response.with_id(message_id), + Ok(()) => { + CommunicationValue::new(CommunicationType::Success).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 d87eea1..5a25ec0 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(response) => { - response.with_id(message_id) + Ok(()) => { + CommunicationValue::new(CommunicationType::Success).with_id(message_id) } Err(error) => { log_err!( @@ -201,25 +201,6 @@ 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. @@ -903,7 +884,7 @@ impl ClientConnection { let response = CommunicationValue::new(CommunicationType::LoadTxtRecord) .with_id(message_id) - .add_typed_default(DataType::Content, DataValue::Str(record_text)); + .add_typed_default(DataType::AppContent, 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 51d89fd..bd822cb 100755 --- a/src/rho/iota_connection.rs +++ b/src/rho/iota_connection.rs @@ -276,7 +276,7 @@ impl IotaConnection { ) .await { - Ok(response) => response.with_id(message_id), + Ok(()) => CommunicationValue::new(CommunicationType::Success).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 97a5384..bd0f0ab 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 { +) -> Result<(), RelayRouteError> { 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 { +) -> Result<(), RelayRouteError> { 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 { +) -> Result<(), RelayRouteError> { 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(|_| CommunicationValue::new(CommunicationType::Success)).map_err(|error| { + target.send_relay_to_client(&frame).await.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 { +fn route_response(response: CommunicationValue) -> Result<(), RelayRouteError> { if response.is_type(CommunicationType::Success) { - Ok(response) + Ok(()) } else { Err(RelayRouteError::Send(format!( "next Relay hop rejected the frame with {}", @@ -407,15 +407,4 @@ 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 c933acf..f4aa658 100644 --- a/src/rho/rho_connection.rs +++ b/src/rho/rho_connection.rs @@ -328,16 +328,6 @@ 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)