Compare commits

...
Author SHA1 Message Date
bf96891092 chore(deps): update rust crate livekit-protocol to v0.7.12
All checks were successful
renovate/stability-days Updates have met minimum release age requirement
2026-08-28 18:00:44 +03:00
Alex Emmet
233b944007
[Add] Replyjumps 2026-08-28 16:21:11 +02:00
11 changed files with 88 additions and 26 deletions

18
Cargo.lock generated
View file

@ -643,7 +643,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "39cab71617ae0d63f51a36d69f866391735b51691dbda63cf6f96d042b63efeb"
dependencies = [
"libc",
"windows-sys 0.59.0",
"windows-sys 0.61.2",
]
[[package]]
@ -1486,9 +1486,9 @@ dependencies = [
[[package]]
name = "livekit-protocol"
version = "0.7.10"
version = "0.7.12"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "4d26880e94e2f9bab298445e7d86a3794453d211a12ddbd051bd9991a343f9ff"
checksum = "526f22bddf409e5f15449d55cf341647d7c16a1f43df9db9b05be106e4208e1c"
dependencies = [
"pbjson",
"pbjson-types",
@ -2365,7 +2365,7 @@ dependencies = [
"once_cell",
"socket2",
"tracing",
"windows-sys 0.59.0",
"windows-sys 0.61.2",
]
[[package]]
@ -2619,7 +2619,7 @@ dependencies = [
"errno",
"libc",
"linux-raw-sys",
"windows-sys 0.59.0",
"windows-sys 0.61.2",
]
[[package]]
@ -2678,7 +2678,7 @@ dependencies = [
"security-framework",
"security-framework-sys",
"webpki-root-certs",
"windows-sys 0.59.0",
"windows-sys 0.61.2",
]
[[package]]
@ -3054,10 +3054,10 @@ 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.59.0",
"windows-sys 0.61.2",
]
[[package]]
@ -3618,7 +3618,7 @@ version = "0.1.11"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "c2a7b1c03c876122aa43f3020e6c3c3ee5c05081c9a00739faf7503aeba10d22"
dependencies = [
"windows-sys 0.59.0",
"windows-sys 0.61.2",
]
[[package]]

View file

@ -29,7 +29,7 @@ log = "0.4"
dotenv = "0.15.0"
strum_macros = "0.28.0"
livekit-api = { version = "0.5.6", features = ["rustls-tls-native-roots"] }
livekit-protocol = "=0.7.10"
livekit-protocol = "=0.7.12"
thiserror = "2.0.19"
trust-dns-resolver = "0.23.2"
serde_json = "1.0.151"

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