Compare commits

...
Author SHA1 Message Date
9b1cb4ad7e chore(deps): update rust crate livekit-api to 0.6.0
Some checks failed
renovate/artifacts Artifact file update failure
renovate/stability-days Updates have met minimum release age requirement
2026-08-28 18:00:49 +03:00
Alex Emmet
233b944007
[Add] Replyjumps 2026-08-28 16:21:11 +02:00
10 changed files with 79 additions and 17 deletions

View file

@ -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"

@ -1 +1 @@
Subproject commit 486541b9483356ff49ff3ec7016f87d3ecbeaa0e
Subproject commit f3c5037b0a099ae5389486eec77471bd5addab4a

View file

@ -129,6 +129,20 @@ where
}
}
fn parse_positive_or_default<T>(name: &'static str, default: T) -> Result<T, ConfigError>
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<Option<LiveKitConfig>, ConfigError> {
livekit_from_values(
env::var("LIVEKIT_HOSTNAME").ok(),

View file

@ -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) => {

View file

@ -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,

View file

@ -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,

View file

@ -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;
}

View file

@ -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,

View file

@ -140,7 +140,7 @@ pub async fn route_relay(
state: &Arc<AppState>,
source: RelaySource,
frame: CommunicationValue,
) -> Result<(), RelayRouteError> {
) -> Result<CommunicationValue, 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<(), RelayRouteError> {
) -> Result<CommunicationValue, 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<(), RelayRouteError> {
) -> Result<CommunicationValue, 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_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<CommunicationValue, RelayRouteError> {
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());
}
}

View file

@ -328,6 +328,16 @@ impl RhoConnection {
.await
}
pub async fn await_peer_control(
&self,
cv: &CommunicationValue,
) -> Result<CommunicationValue, String> {
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)