From 945b58178dbf947f89c15864d5de954954424bd0 Mon Sep 17 00:00:00 2001 From: TroyHernandez Date: Fri, 2 Oct 2026 18:20:32 -0500 Subject: [PATCH 1/3] Add MatrixRTC calls over LiveKit mx_call_join() finds the LiveKit JWT service, trades an OpenID token for a media token, publishes this device's org.matrix.msc3401.call.member state event, and makes a media key that goes Olm-encrypted to every device in the call as io.element.call.encryption_keys. mx_call_handle() applies each sync: peers' keys, membership changes with the Element/FluffyChat rotation policy (leaver rotates, joiner within 10 s gets the current key), and the hourly membership refresh. mx_call_connect() joins the LiveKit room through livekitr with per-participant HKDF keys. mx_call_leave() clears the membership. Underneath: mx_crypto_process_sync() now returns other decrypted Olm to-device events in to_device, and mx_send_to_device_encrypted() sends an Olm-encrypted to-device event of any type. --- DESCRIPTION | 3 +- NAMESPACE | 12 + R/call.R | 590 +++++++++++++++++++++++++++++ R/e2ee.R | 19 +- R/to-device.R | 163 ++++++++ inst/tinytest/test_call.R | 384 +++++++++++++++++++ inst/tinytest/test_to_device.R | 174 +++++++++ man/mx_call_connect.Rd | 23 ++ man/mx_call_handle.Rd | 30 ++ man/mx_call_join.Rd | 64 ++++ man/mx_call_key_parse.Rd | 29 ++ man/mx_call_key_plan.Rd | 38 ++ man/mx_call_leave.Rd | 16 + man/mx_call_members.Rd | 35 ++ man/mx_call_poll.Rd | 24 ++ man/mx_call_service_url.Rd | 27 ++ man/mx_crypto_encrypt_to_device.Rd | 53 +++ man/mx_crypto_process_sync.Rd | 7 +- man/mx_send_to_device_encrypted.Rd | 50 +++ 19 files changed, 1735 insertions(+), 6 deletions(-) create mode 100644 R/call.R create mode 100644 R/to-device.R create mode 100644 inst/tinytest/test_call.R create mode 100644 inst/tinytest/test_to_device.R create mode 100644 man/mx_call_connect.Rd create mode 100644 man/mx_call_handle.Rd create mode 100644 man/mx_call_join.Rd create mode 100644 man/mx_call_key_parse.Rd create mode 100644 man/mx_call_key_plan.Rd create mode 100644 man/mx_call_leave.Rd create mode 100644 man/mx_call_members.Rd create mode 100644 man/mx_call_poll.Rd create mode 100644 man/mx_call_service_url.Rd create mode 100644 man/mx_crypto_encrypt_to_device.Rd create mode 100644 man/mx_send_to_device_encrypted.Rd diff --git a/DESCRIPTION b/DESCRIPTION index 9f94afb..6d2171b 100644 --- a/DESCRIPTION +++ b/DESCRIPTION @@ -22,11 +22,12 @@ Depends: R (>= 4.0) Imports: jsonlite, - mx.api (>= 0.3.0.2), + mx.api (>= 0.3.1.1), stats, tools, utils Suggests: + livekitr, mx.crypto (>= 0.2.2), simplermarkdown, tinytest diff --git a/NAMESPACE b/NAMESPACE index c96d777..0daa144 100644 --- a/NAMESPACE +++ b/NAMESPACE @@ -1,6 +1,15 @@ # tinyrox says don't edit this manually, but it can't stop you! export(mx_accept_invites) +export(mx_call_connect) +export(mx_call_handle) +export(mx_call_join) +export(mx_call_key_parse) +export(mx_call_key_plan) +export(mx_call_leave) +export(mx_call_members) +export(mx_call_poll) +export(mx_call_service_url) export(mx_client_config_path) export(mx_client_configure) export(mx_client_from_config) @@ -18,6 +27,7 @@ export(mx_crypto_decrypt_event) export(mx_crypto_device_keys) export(mx_crypto_encrypt_event) export(mx_crypto_encrypt_for_devices) +export(mx_crypto_encrypt_to_device) export(mx_crypto_handle_to_device) export(mx_crypto_inbound_session) export(mx_crypto_known_devices) @@ -58,10 +68,12 @@ export(mx_send_encrypted) export(mx_send_media) export(mx_send_table) export(mx_send_text) +export(mx_send_to_device_encrypted) export(mx_set_displayname) export(mx_sync_update) export(mx_table_html) export(mx_verify_console) export(mx_with_relogin) +S3method(print,mx_call) S3method(print,mx_client_config) diff --git a/R/call.R b/R/call.R new file mode 100644 index 0000000..a3a35e3 --- /dev/null +++ b/R/call.R @@ -0,0 +1,590 @@ +# MatrixRTC calls over LiveKit, as Element Call and FluffyChat run them: +# a per-device membership state event, a media token from the LiveKit +# JWT service, and per-participant media keys exchanged as Olm-encrypted +# to-device events. Media itself goes through the livekitr package. +# +# Wire format (both clients read, 2026-10): membership is a state event +# org.matrix.msc3401.call.member with state key "___m.call", +# keys are io.element.call.encryption_keys to-device events, and the +# LiveKit participant identity is ":". + +MX_CALL_MEMBER <- "org.matrix.msc3401.call.member" +MX_CALL_KEYS <- "io.element.call.encryption_keys" +MX_CALL_EXPIRES_MS <- 4 * 60 * 60 * 1000 +MX_CALL_REFRESH_S <- 60 * 60 +# A joiner inside this window gets the current key; later, the key rotates. +MX_CALL_KEY_GRACE_S <- 10 +# Key indexes cycle below 255: FluffyChat's key ring has 255 slots. +MX_CALL_KEY_INDEXES <- 255L + +mx_call_now_ms <- function(now = Sys.time()) { + floor(as.numeric(now) * 1000) +} + +mx_call_state_key <- function(user_id, device_id) { + paste0("_", user_id, "_", device_id, "_m.call") +} + +mx_call_identity <- function(user_id, device_id) { + paste(user_id, device_id, sep = ":") +} + +# Content of this device's membership event. created_ts is what Element +# reads, created_at what FluffyChat reads; both fall back to the event's +# origin_server_ts, so sending both only keeps the expiry exact. +mx_call_member_content <- function(user_id, device_id, service_url, room_id, + intent = "voice", now = Sys.time()) { + stopifnot(intent %in% c("voice", "video")) + ts <- mx_call_now_ms(now) + list( + application = "m.call", + call_id = "", + scope = "m.room", + device_id = device_id, + membershipID = mx_call_identity(user_id, device_id), + expires = MX_CALL_EXPIRES_MS, + `m.call.intent` = intent, + focus_active = list(type = "livekit", + focus_selection = "oldest_membership"), + foci_preferred = list(list(type = "livekit", + livekit_service_url = service_url, + livekit_alias = room_id)), + created_ts = ts, + created_at = ts + ) +} + +#' Active call memberships in a room's state +#' +#' Reads the \code{org.matrix.msc3401.call.member} state events out of a +#' room's full state (\code{mx.api::mx_room_state()}) and keeps the ones +#' that count as in the call: non-empty content, a LiveKit focus, and an +#' expiry (\code{created_ts}, or \code{created_at}, or the event's own +#' timestamp, plus \code{expires}, default 4 hours) still in the future. +#' +#' @param state List of state events. +#' @param now The current time. +#' @return A list of members, each \code{list(user_id, device_id, +#' identity, membership_id, service_urls, expires_at)}, where +#' \code{identity} is the member's LiveKit participant identity +#' \code{":"} and \code{service_urls} the LiveKit JWT +#' services it prefers. +#' @examples +#' state <- list(list(type = "org.matrix.msc3401.call.member", +#' sender = "@alice:example.org", origin_server_ts = 1e12, +#' content = list(application = "m.call", device_id = "PHONE", +#' focus_active = list(type = "livekit"), +#' foci_preferred = list(list(type = "livekit", +#' livekit_service_url = "https://jwt.example.org"))))) +#' mx_call_members(state, now = as.POSIXct(1e9, origin = "1970-01-01")) +#' @export +mx_call_members <- function(state, now = Sys.time()) { + now_ms <- mx_call_now_ms(now) + members <- list() + for (ev in state) { + if (!identical(ev$type, MX_CALL_MEMBER)) next + c <- ev$content + if (!is.list(c) || !length(c)) next + if (!identical(c$focus_active$type, "livekit")) next + if (!is.character(c$device_id) || !nzchar(c$device_id)) next + if (!is.character(ev$sender) || !nzchar(ev$sender)) next + created <- c$created_ts %||% c$created_at %||% ev$origin_server_ts + expires <- c$expires %||% MX_CALL_EXPIRES_MS + if (!is.numeric(created) || !is.numeric(expires)) next + expires_at <- created + expires + if (expires_at <= now_ms) next + urls <- character() + for (focus in c$foci_preferred %||% list()) { + if (identical(focus$type, "livekit") && + is.character(focus$livekit_service_url)) { + urls <- c(urls, focus$livekit_service_url) + } + } + members[[length(members) + 1L]] <- list( + user_id = ev$sender, + device_id = c$device_id, + identity = mx_call_identity(ev$sender, c$device_id), + membership_id = c$membershipID %||% + mx_call_identity(ev$sender, c$device_id), + service_urls = urls, + expires_at = expires_at) + } + members +} + +#' Find the LiveKit JWT service for a room's call +#' +#' Tries, in order, the services the current call members advertise in +#' their memberships, the homeserver's RTC transports +#' (\code{mx.api::mx_rtc_transports()}), and the +#' \code{org.matrix.msc4143.rtc_foci} entry of the server's client +#' well-known file. +#' +#' @param client Matrix client config. +#' @param members Current members from \code{\link{mx_call_members}}. +#' @return The service URL, a string. +#' @examples +#' \dontrun{ +#' mx_call_service_url(client, members) +#' } +#' @export +mx_call_service_url <- function(client, members = list()) { + for (member in members) { + if (length(member$service_urls)) { + return(member$service_urls[[1L]]) + } + } + s <- mx_client_session(client) + for (transport in mx.api::mx_rtc_transports(s)) { + if (identical(transport$type, "livekit") && + is.character(transport$livekit_service_url)) { + return(transport$livekit_service_url) + } + } + server_name <- sub("^@[^:]+:", "", client$user_id) + well_known <- mx.api::mx_well_known_client(server_name) + foci <- well_known[["org.matrix.msc4143.rtc_foci"]] + if (is.list(foci) && !is.null(foci$livekit_service_url)) { + foci <- list(foci) + } + for (focus in foci %||% list()) { + if (is.character(focus$livekit_service_url)) { + return(focus$livekit_service_url) + } + } + stop("no LiveKit JWT service found: nobody in the call advertises one, ", + "the homeserver lists no RTC transports, and ", server_name, + " has no rtc_foci in its well-known file", call. = FALSE) +} + +# ---- media keys ------------------------------------------------------- + +mx_call_b64 <- function(bytes) { + jsonlite::base64_enc(bytes) +} + +# Accepts standard or URL-safe alphabets, padded or not, as the Element +# decoder does. +mx_call_b64_decode <- function(x) { + x <- gsub("[\r\n= ]", "", chartr("-_", "+/", x)) + pad <- (4L - nchar(x) %% 4L) %% 4L + jsonlite::base64_dec(paste0(x, strrep("=", pad))) +} + +# Content of a key event: the key the receiver should use for this +# device's media. FluffyChat casts every field but sent_ts, so all of them +# are always present. +mx_call_key_content <- function(key, index, user_id, device_id, room_id, + now = Sys.time()) { + list( + keys = list(index = as.integer(index), key = mx_call_b64(key)), + member = list(id = mx_call_identity(user_id, device_id), + claimed_device_id = device_id), + room_id = room_id, + session = list(application = "m.call", call_id = "", scope = "m.room"), + sent_ts = mx_call_now_ms(now) + ) +} + +#' Read a call key event +#' +#' Parses a decrypted \code{io.element.call.encryption_keys} to-device +#' event (from the \code{to_device} list of +#' \code{\link{mx_crypto_process_sync}}) into the LiveKit identity it +#' belongs to and the key to set for it. +#' +#' @param event List with \code{type}, \code{content} and \code{sender}. +#' @param room_id The call's room; keys for other rooms are ignored. +#' @return \code{list(identity, key, index)}, with \code{key} a raw +#' vector, or NULL when the event is not a usable key for this room. +#' @examples +#' ev <- list(type = "io.element.call.encryption_keys", +#' sender = "@alice:example.org", +#' content = list(keys = list(index = 3, key = "AAECAwQFBgcICQoLDA0ODw=="), +#' member = list(claimed_device_id = "PHONE"), room_id = "!r:example.org")) +#' mx_call_key_parse(ev, "!r:example.org") +#' @export +mx_call_key_parse <- function(event, room_id) { + if (!identical(event$type, MX_CALL_KEYS)) return(NULL) + c <- event$content + if (!identical(c$room_id, room_id)) return(NULL) + device <- c$member$claimed_device_id + index <- c$keys$index + key <- c$keys$key + if (!is.character(event$sender) || !is.character(device) || + !nzchar(device) || !is.numeric(index) || length(index) != 1L || + index < 0 || index != floor(index) || !is.character(key)) { + return(NULL) + } + bytes <- tryCatch(mx_call_b64_decode(key), error = function(e) NULL) + if (is.null(bytes) || !length(bytes)) return(NULL) + list(identity = mx_call_identity(event$sender, device), key = bytes, + index = as.integer(index)) +} + +# The key state of one call: our key and index, when the key was made, +# whom it was sent to, and the peers' keys. +mx_call_keys_new <- function() { + keys <- new.env(parent = emptyenv()) + keys$key <- NULL + keys$index <- -1L + keys$created <- NULL + keys$shared_with <- character() + keys$peers <- list() + keys +} + +mx_call_keys_rotate <- function(keys, now = Sys.time()) { + keys$key <- mx_crypto_random_bytes(16L) + keys$index <- (keys$index + 1L) %% MX_CALL_KEY_INDEXES + keys$created <- now + keys$shared_with <- character() + invisible(keys) +} + +#' Decide what to do with the call key when membership changes +#' +#' Implements the rotation policy Element Call and FluffyChat follow: a +#' leaver forces a new key for everyone; a joiner within the grace period +#' after the key was made receives the current key; a joiner after it +#' gets a new key, as does everyone else. +#' +#' @param keys Key state from \code{mx_call_keys_new()} (an environment; +#' it is not modified). +#' @param identities LiveKit identities of the other members now in the +#' call. +#' @param now The current time. +#' @param grace Seconds after a key is made during which it is still +#' handed to joiners instead of rotated. +#' @return \code{list(rotate, targets)}: whether to make a new key, and +#' the identities to send the (new or current) key to. +#' @examples +#' keys <- mx.client:::mx_call_keys_new() +#' mx_call_key_plan(keys, c("@a:ex:D1", "@b:ex:D2")) +#' @export +mx_call_key_plan <- function(keys, identities, now = Sys.time(), + grace = MX_CALL_KEY_GRACE_S) { + identities <- unique(identities) + if (is.null(keys$key)) { + return(list(rotate = TRUE, targets = identities)) + } + left <- setdiff(keys$shared_with, identities) + joined <- setdiff(identities, keys$shared_with) + if (length(left)) { + return(list(rotate = TRUE, targets = identities)) + } + if (!length(joined)) { + return(list(rotate = FALSE, targets = character())) + } + age <- as.numeric(difftime(now, keys$created, units = "secs")) + if (age < grace) { + list(rotate = FALSE, targets = joined) + } else { + list(rotate = TRUE, targets = identities) + } +} + +# Devices of the given identities, with verified keys. +mx_call_devices <- function(client, identities) { + if (!length(identities)) return(list()) + parts <- regmatches(identities, regexec("^(.*):([^:]+)$", identities)) + user_ids <- unique(vapply(parts, `[`, "", 2L)) + known <- mx_crypto_known_devices(client, user_ids, strict = TRUE) + Filter(function(d) { + mx_call_identity(d$user_id, d$device_id) %in% identities + }, known) +} + +# ---- the call object -------------------------------------------------- + +mx_require_livekitr <- function() { + if (!requireNamespace("livekitr", quietly = TRUE)) { + stop("joining the media of a call needs the livekitr package", + call. = FALSE) + } + invisible(TRUE) +} + +# Apply a peer's key to the media session, remembering it for a session +# that is not connected yet. +mx_call_apply_peer_key <- function(call, parsed) { + call$keys$peers[[parsed$identity]] <- parsed[c("key", "index")] + if (!is.null(call$session)) { + livekitr::lk_set_e2ee_key(call$session, parsed$key, + identity = parsed$identity, + key_index = parsed$index) + } + invisible(NULL) +} + +# Send our current key to the given identities, then make sure our own +# media uses it. +mx_call_send_key <- function(call, targets) { + if (length(targets)) { + devices <- mx_call_devices(call$client, targets) + content <- mx_call_key_content(call$keys$key, call$keys$index, + call$client$user_id, + call$client$device_id, call$room_id) + res <- mx_send_to_device_encrypted(call$client, call$account, + call$sessions, MX_CALL_KEYS, content, + devices, call$store_dir) + call$sessions <- res$sessions + reached <- vapply(res$sent, function(p) { + mx_call_identity(p$user_id, p$device_id) + }, "") + call$keys$shared_with <- union(call$keys$shared_with, + intersect(targets, reached)) + } + if (!is.null(call$session)) { + livekitr::lk_set_e2ee_key(call$session, call$keys$key, + identity = call$identity, + key_index = call$keys$index) + } + invisible(NULL) +} + +# Reconcile the key state with the current members. +mx_call_update_keys <- function(call, members, now = Sys.time()) { + others <- setdiff(vapply(members, `[[`, "", "identity"), call$identity) + plan <- mx_call_key_plan(call$keys, others, now) + if (plan$rotate) { + mx_call_keys_rotate(call$keys, now) + } + mx_call_send_key(call, plan$targets) + call$members <- members + invisible(NULL) +} + +mx_call_send_membership <- function(call, now = Sys.time()) { + content <- mx_call_member_content(call$client$user_id, + call$client$device_id, call$service_url, + call$room_id, call$intent, now) + mx.api::mx_set_state(mx_client_session(call$client), call$room_id, + MX_CALL_MEMBER, content, + mx_call_state_key(call$client$user_id, call$client$device_id)) + call$membership_sent <- now + invisible(NULL) +} + +#' Join a MatrixRTC call +#' +#' Joins the call of a room the way Element Call and FluffyChat do: +#' finds the LiveKit JWT service, trades an OpenID token for a media +#' token, announces this device's membership as a state event, makes a +#' media key and sends it Olm-encrypted to every device already in the +#' call, and, with \code{connect = TRUE}, joins the LiveKit room through +#' \pkg{livekitr} with end-to-end encryption on. The returned call must +#' then be driven with \code{\link{mx_call_handle}} (or +#' \code{\link{mx_call_poll}}) and ended with \code{\link{mx_call_leave}}. +#' +#' @param client Matrix client config. +#' @param account An mx.crypto account handle. +#' @param sessions A session set. +#' @param room_id The room whose call to join. +#' @param intent \code{"voice"} or \code{"video"}, what the membership +#' announces. +#' @param store_dir Character or NULL. Where the updated crypto sessions +#' are saved after each key send. +#' @param connect Join the LiveKit room now? \code{FALSE} does everything +#' on the Matrix side only. +#' @param service_url Character or NULL. The LiveKit JWT service to use +#' instead of discovering one. +#' @return An object of class \code{"mx_call"}: an environment holding +#' the \code{client} and \code{sessions} (both updated as the call +#' runs), the LiveKit \code{session} from \pkg{livekitr} (or NULL), this +#' device's LiveKit \code{identity}, the current \code{members}, and +#' the media \code{token}. +#' @examples +#' \dontrun{ +#' call <- mx_call_join(client, acct, sessions, "!room:example.org", +#' store_dir = store) +#' livekitr::lk_on_audio(call$session, function(pcm, info) { +#' cat(info$identity, "spoke\n") +#' }) +#' repeat mx_call_poll(call) +#' } +#' @export +mx_call_join <- function(client, account, sessions, room_id, + intent = "voice", store_dir = NULL, connect = TRUE, + service_url = NULL) { + mx_require_crypto() + call <- new.env(parent = emptyenv()) + call$client <- client + call$account <- account + call$sessions <- sessions + call$store_dir <- store_dir + call$room_id <- room_id + call$intent <- intent + call$identity <- mx_call_identity(client$user_id, client$device_id) + call$keys <- mx_call_keys_new() + call$session <- NULL + call$members <- list() + class(call) <- "mx_call" + + s <- mx_client_session(client) + members <- mx_call_members(mx.api::mx_room_state(s, room_id)) + call$service_url <- service_url %||% mx_call_service_url(client, members) + openid <- mx.api::mx_openid_token(s) + call$token <- mx.api::mx_rtc_livekit_token(call$service_url, room_id, + openid, client$device_id) + mx_call_send_membership(call) + mx_call_update_keys(call, members) + if (connect) { + mx_call_connect(call) + } + call +} + +#' Join the LiveKit room of a call +#' +#' Connects the media side of a call made with \code{mx_call_join(connect +#' = FALSE)}: joins the LiveKit room with the media token, with +#' end-to-end encryption configured as Element Call and FluffyChat +#' expect (per-participant HKDF keys, a 256-slot key ring), and sets our +#' key and every peer key received so far. +#' +#' @param call An \code{"mx_call"}. +#' @param ... Further arguments to \code{livekitr::lk_connect()}, such +#' as \code{opts}. +#' @return \code{call}, invisibly. +#' @export +mx_call_connect <- function(call, ...) { + mx_require_livekitr() + if (!is.null(call$session)) { + return(invisible(call)) + } + call$session <- livekitr::lk_connect(call$token$url, call$token$jwt, + e2ee = mx_call_e2ee_options(), ...) + if (!identical(call$session$identity, call$identity)) { + warning("the LiveKit identity ", call$session$identity, + " differs from the Matrix identity ", call$identity, + "; other clients will not match this device's keys", + call. = FALSE) + } + livekitr::lk_set_e2ee_key(call$session, call$keys$key, + identity = call$identity, + key_index = call$keys$index) + for (identity in names(call$keys$peers)) { + peer <- call$keys$peers[[identity]] + livekitr::lk_set_e2ee_key(call$session, peer$key, identity = identity, + key_index = peer$index) + } + invisible(call) +} + +# How Element Call and FluffyChat configure the frame cryptor: keys are +# used as given (HKDF, not the shared-key PBKDF2 path), one slot per key +# index, and a few failed frames are tolerated while keys are in flight. +mx_call_e2ee_options <- function() { + list(kdf = "hkdf", key_ring_size = 256L, ratchet_window_size = 10L, + failure_tolerance = 10L) +} + +#' Feed a sync response to a call +#' +#' Applies what a \code{/sync} response means for the call: peers' media +#' keys from the decrypted to-device events, membership changes (which +#' may rotate and resend our key), and the hourly refresh of our own +#' membership. Call it with every sync while in the call. When the +#' application already runs \code{\link{mx_crypto_process_sync}} on the +#' response, pass its result as \code{processed} so to-device events +#' are not decrypted twice. +#' +#' @param call An \code{"mx_call"}. +#' @param sync A parsed \code{/sync} response. +#' @param processed The result of \code{mx_crypto_process_sync()} on +#' \code{sync}, or NULL to have it run here (its \code{sessions} are +#' then kept in \code{call$sessions}). +#' @return A list of what changed: \code{keys}, the identities whose +#' keys were received, and \code{members}, the current members when +#' membership changed (NULL otherwise). +#' @export +mx_call_handle <- function(call, sync, processed = NULL) { + if (is.null(processed)) { + identity <- mx.crypto::mxc_account_identity_keys(call$account) + processed <- mx_crypto_process_sync(call$account, call$sessions, + sync, identity$curve25519, self_id = call$client$user_id, + self_device_id = call$client$device_id) + call$sessions <- processed$sessions + } + received <- character() + for (ev in processed$to_device %||% list()) { + parsed <- mx_call_key_parse(ev, call$room_id) + if (!is.null(parsed)) { + mx_call_apply_peer_key(call, parsed) + received <- c(received, parsed$identity) + } + } + members <- NULL + room <- sync$rooms$join[[call$room_id]] + changed <- any(vapply(c(room$state$events %||% list(), + room$timeline$events %||% list()), + function(ev) identical(ev$type, MX_CALL_MEMBER), + logical(1))) + now <- Sys.time() + if (changed) { + members <- mx_call_members( + mx.api::mx_room_state(mx_client_session(call$client), + call$room_id), now) + mx_call_update_keys(call, members, now) + } + if (as.numeric(difftime(now, call$membership_sent, units = "secs")) >= + MX_CALL_REFRESH_S) { + mx_call_send_membership(call, now) + } + list(keys = received, members = members) +} + +#' Sync once and poll the media of a call +#' +#' One iteration of a call loop for programs with no sync loop of their +#' own: syncs with \code{\link{mx_sync_update}}, applies the response with +#' \code{\link{mx_call_handle}}, and polls the LiveKit session with +#' \code{livekitr::lk_poll()}, which is where audio callbacks run. +#' +#' @param call An \code{"mx_call"}. +#' @param timeout Seconds to wait for the sync long poll. +#' @param media_timeout Seconds to wait in the media poll. +#' @return The LiveKit events from \code{livekitr::lk_poll()}, or an +#' empty list when the call has no media session. +#' @export +mx_call_poll <- function(call, timeout = 1, media_timeout = 0.1) { + res <- mx_sync_update(call$client, timeout = as.integer(timeout * 1000), + save = !is.null(attr(call$client, "path"))) + call$client <- res$client + mx_call_handle(call, res$sync) + if (is.null(call$session)) { + return(list()) + } + livekitr::lk_poll(call$session, timeout = media_timeout) +} + +#' Leave a MatrixRTC call +#' +#' Leaves the LiveKit room and clears this device's membership event. +#' +#' @param call An \code{"mx_call"}. +#' @return \code{NULL}, invisibly. +#' @export +mx_call_leave <- function(call) { + if (!is.null(call$session)) { + livekitr::lk_disconnect(call$session) + call$session <- NULL + } + # An empty content object is how a device leaves; the state key stays. + mx.api::mx_set_state(mx_client_session(call$client), call$room_id, + MX_CALL_MEMBER, stats::setNames(list(), character()), + mx_call_state_key(call$client$user_id, call$client$device_id)) + invisible(NULL) +} + +#' @export +print.mx_call <- function(x, ...) { + cat(" room ", x$room_id, " as ", x$identity, + if (is.null(x$session)) " (media not connected)" else " (connected)", + "\n", sep = "") + cat(" JWT service: ", x$service_url, "\n", sep = "") + cat(" members in call: ", length(x$members), ", key index: ", + x$keys$index, ", peer keys: ", length(x$keys$peers), "\n", sep = "") + invisible(x) +} diff --git a/R/e2ee.R b/R/e2ee.R index cd2b7fc..f1dca78 100644 --- a/R/e2ee.R +++ b/R/e2ee.R @@ -286,8 +286,11 @@ mx_crypto_encrypt_for_devices <- function(account, sessions, room_id, #' \code{verification_events} (original verification envelopes, separated #' from chat messages; no handshake or network side effect is performed), #' \code{sessions}, unsent \code{key_requests}, matching -#' \code{key_request_cancellations}, and \code{incoming_key_requests} -#' for a policy-aware sharing layer to inspect. +#' \code{key_request_cancellations}, \code{incoming_key_requests} +#' for a policy-aware sharing layer to inspect, and \code{to_device}: +#' every other Olm-encrypted to-device event that decrypted and whose +#' claimed sender matches the envelope, each as \code{list(type, +#' content, sender, sender_bound)}. Call encryption keys arrive here. #' @examples #' \donttest{ #' if (requireNamespace("mx.crypto", quietly = TRUE)) { @@ -308,6 +311,7 @@ mx_crypto_process_sync <- function(account, sessions, sync_resp, cancellations <- list() incoming_requests <- list() verification_events <- list() + to_device <- list() # 1. To-device: recover shared room keys. for (ev in sync_resp$to_device$events %||% list()) { @@ -425,6 +429,14 @@ mx_crypto_process_sync <- function(account, sessions, sync_resp, cancellations[[length(cancellations) + 1L]] <- mx_crypto_key_request_cancellation(pending) sessions$key_requests[[key]] <- NULL + } else if (identical(decoded$sender, ev$sender)) { + # Any other Olm-encrypted to-device event (call encryption + # keys, for one) is handed to the caller. The envelope sender + # is the homeserver's word and the plaintext sender is the + # device's claim; they have to agree. + to_device[[length(to_device) + 1L]] <- list( + type = decoded$type, content = decoded$content, + sender = decoded$sender, sender_bound = chk$bound) } } @@ -540,5 +552,6 @@ mx_crypto_process_sync <- function(account, sessions, sync_resp, verification_events = verification_events, key_requests = key_requests, key_request_cancellations = cancellations, - incoming_key_requests = incoming_requests) + incoming_key_requests = incoming_requests, + to_device = to_device) } diff --git a/R/to-device.R b/R/to-device.R new file mode 100644 index 0000000..48ba20b --- /dev/null +++ b/R/to-device.R @@ -0,0 +1,163 @@ +# Olm-encrypted to-device events of any type. Room keys use the dedicated +# path in e2ee.R; this is for everything else a device sends another +# device directly, such as the encryption keys of a MatrixRTC call. + +# The Olm payload for one recipient device, with the sender and recipient +# blocks that let the receiver attribute it (see mx_crypto_room_key_payload). +mx_crypto_olm_payload <- function(olm_session, event_type, content, + sender_user_id, sender_curve25519, + sender_ed25519, recipient_user_id, + recipient_curve25519, recipient_ed25519) { + payload <- list(type = event_type, content = content, + sender = sender_user_id, recipient = recipient_user_id, + recipient_keys = list(ed25519 = recipient_ed25519), + keys = list(ed25519 = sender_ed25519)) + plaintext <- mx.api::mx_canonical_json(payload) + ct <- mx.crypto::mxc_olm_encrypt(olm_session, charToRaw(plaintext)) + list( + algorithm = MX_OLM, + sender_key = sender_curve25519, + ciphertext = stats::setNames( + list(list(type = ct$type, body = ct$body)), + recipient_curve25519) + ) +} + +#' Encrypt a to-device event for a set of devices +#' +#' Produces one Olm-encrypted \code{m.room.encrypted} payload per +#' recipient device, opening an Olm session where none exists yet. Unlike +#' Megolm room keys, which are shared once per session, the event is +#' encrypted for every recipient on every call. +#' +#' @param account An mx.crypto account handle. +#' @param sessions A session set. +#' @param event_type The inner event type, e.g. +#' \code{"io.element.call.encryption_keys"}. +#' @param content Named list. Plaintext event content. +#' @param recipients List of recipient devices as returned by +#' \code{mx_crypto_known_devices()}: \code{user_id}, \code{device_id}, +#' \code{curve25519}, \code{ed25519}, and \code{otk} for a device +#' without an Olm session yet. +#' @param sender_user_id Character. This user's Matrix id. +#' @return List with \code{to_device} (per-device payloads, each +#' \code{list(user_id, device_id, content)}) and the updated +#' \code{sessions}. +#' @examples +#' \donttest{ +#' if (requireNamespace("mx.crypto", quietly = TRUE)) { +#' acct <- mx.crypto::mxc_account_new() +#' out <- mx_crypto_encrypt_to_device(acct, mx_crypto_sessions_new(), +#' "org.example.ping", list(n = 1), recipients = list(), +#' sender_user_id = "@me:ex") +#' length(out$to_device) +#' } +#' } +#' @export +mx_crypto_encrypt_to_device <- function(account, sessions, event_type, + content, recipients, sender_user_id) { + mx_require_crypto() + identity <- mx.crypto::mxc_account_identity_keys(account) + to_device <- list() + for (r in recipients) { + peer <- r$curve25519 + if (is.null(peer) || !nzchar(peer) || is.null(r$ed25519) || + !nzchar(r$ed25519)) { + stop("recipient ", r$user_id %||% "", "/", + r$device_id %||% "", " has no verified keys; ", + "recipients come from mx_crypto_known_devices()", + call. = FALSE) + } + olm <- sessions$olm[[peer]] + if (is.null(olm)) { + if (is.null(r$otk)) { + stop("no Olm session for ", peer, + " and no one-time key supplied to open one", + call. = FALSE) + } + olm <- mx.crypto::mxc_olm_create_outbound( + account, peer_curve25519 = peer, peer_otk = r$otk) + sessions$olm[[peer]] <- olm + } + to_device[[length(to_device) + 1L]] <- list( + user_id = r$user_id, device_id = r$device_id, + content = mx_crypto_olm_payload(olm, event_type, content, + sender_user_id, identity$curve25519, identity$ed25519, + r$user_id, peer, r$ed25519)) + } + list(to_device = to_device, sessions = sessions) +} + +# Devices ready to receive Olm: those with a session as they are, the +# rest with a freshly claimed one-time key. Devices whose claimed key did +# not verify are dropped with a warning. +mx_crypto_olm_recipients <- function(client, sessions, devices) { + need <- Filter(function(d) is.null(sessions$olm[[d$curve25519]]), devices) + have <- Filter(function(d) !is.null(sessions$olm[[d$curve25519]]), devices) + claimed <- if (length(need)) { + mx_crypto_claim_otks(client, need, strict = TRUE) + } else { + list() + } + usable <- Filter(function(d) !is.null(d$otk), claimed) + if (length(usable) < length(claimed)) { + warning("mx.client: ", length(claimed) - length(usable), " of ", + length(claimed), " devices had no usable one-time key and ", + "were skipped", call. = FALSE) + } + c(usable, have) +} + +#' Send an Olm-encrypted to-device event +#' +#' Encrypts \code{content} for each device in \code{devices}, claiming +#' one-time keys where no Olm session exists, and delivers it with +#' \code{mx.api::mx_send_to_device()} in one request. This device itself +#' is never a recipient. +#' +#' @param client Matrix client config. +#' @param account An mx.crypto account handle. +#' @param sessions A session set. +#' @param event_type The inner event type. +#' @param content Named list. Plaintext event content. +#' @param devices List of target devices from +#' \code{mx_crypto_known_devices()}, already narrowed to the devices +#' that should receive the event. +#' @param store_dir Character or NULL. Where to persist the updated +#' sessions; NULL leaves saving to the caller. +#' @return List with \code{sessions} (updated) and \code{sent}, the +#' \code{list(user_id, device_id)} pairs the event went to. +#' @examples +#' \dontrun{ +#' devs <- mx_crypto_known_devices(client, "@bob:example.org") +#' mx_send_to_device_encrypted(client, acct, sessions, "org.example.ping", +#' list(n = 1), devs, store_dir) +#' } +#' @export +mx_send_to_device_encrypted <- function(client, account, sessions, + event_type, content, devices, + store_dir = NULL) { + mx_require_crypto() + devices <- Filter(function(d) { + !(identical(d$user_id, client$user_id) && + identical(d$device_id, client$device_id)) + }, devices) + recipients <- mx_crypto_olm_recipients(client, sessions, devices) + out <- mx_crypto_encrypt_to_device(account, sessions, event_type, + content, recipients, client$user_id) + if (length(out$to_device)) { + messages <- list() + for (p in out$to_device) { + messages[[p$user_id]][[p$device_id]] <- p$content + } + mx.api::mx_send_to_device(mx_client_session(client), "m.room.encrypted", + messages) + } + if (!is.null(store_dir)) { + mx_crypto_sessions_save(out$sessions, store_dir) + } + list(sessions = out$sessions, + sent = lapply(out$to_device, function(p) { + list(user_id = p$user_id, device_id = p$device_id) + })) +} diff --git a/inst/tinytest/test_call.R b/inst/tinytest/test_call.R new file mode 100644 index 0000000..9ce4548 --- /dev/null +++ b/inst/tinytest/test_call.R @@ -0,0 +1,384 @@ +# MatrixRTC call layer: membership events, key events, the rotation +# policy and service discovery. No homeserver and no LiveKit: HTTP is +# replaced below where a function reaches for it. + +library(tinytest) +library(mx.client) + +ns <- asNamespace("mx.client") +ROOM <- "!call:example.org" +T0 <- as.POSIXct(1.7e9, origin = "1970-01-01", tz = "UTC") + +with_api <- function(name, handler, code) { + api <- asNamespace("mx.api") + original <- get(name, envir = api, inherits = FALSE) + assignInNamespace(name, handler, ns = "mx.api") + on.exit(assignInNamespace(name, original, ns = "mx.api"), add = TRUE) + force(code) +} + +client <- list(server = "https://example.org", token = "t", + user_id = "@bot:example.org", device_id = "BOTDEV") + +# ---- identity and membership content --------------------------------- + +expect_identical(ns$mx_call_identity("@a:ex", "D"), "@a:ex:D") +expect_identical(ns$mx_call_state_key("@a:ex", "D"), "_@a:ex_D_m.call") + +content <- ns$mx_call_member_content("@bot:example.org", "BOTDEV", + "https://jwt.example.org", ROOM, now = T0) +expect_identical(content$application, "m.call") +expect_identical(content$call_id, "") +expect_identical(content$scope, "m.room") +expect_identical(content$device_id, "BOTDEV") +expect_identical(content$membershipID, "@bot:example.org:BOTDEV") +expect_equal(content$expires, 4 * 3600 * 1000) +expect_identical(content$`m.call.intent`, "voice") +expect_identical(content$focus_active$type, "livekit") +expect_identical(content$foci_preferred[[1]]$livekit_service_url, + "https://jwt.example.org") +expect_identical(content$foci_preferred[[1]]$livekit_alias, ROOM) +expect_equal(content$created_ts, 1.7e12) +expect_equal(content$created_at, 1.7e12) +expect_error(ns$mx_call_member_content("@a:ex", "D", "u", ROOM, intent = "x")) +# It survives JSON as the clients read it: expires and timestamps as numbers +js <- jsonlite::fromJSON(jsonlite::toJSON(content, auto_unbox = TRUE, + digits = NA), simplifyVector = FALSE) +expect_equal(js$created_ts, 1.7e12) +expect_equal(js$expires, 14400000) + +# ---- reading memberships from room state ----------------------------- + +member_event <- function(sender, device, created, expires = NULL, + focus = "livekit", content_extra = list(), + origin_ts = NULL) { + c <- list(application = "m.call", device_id = device, + focus_active = list(type = focus), + foci_preferred = list(list(type = "livekit", + livekit_service_url = "https://jwt.a"))) + if (!is.null(created)) c$created_ts <- created + if (!is.null(expires)) c$expires <- expires + c <- c(c, content_extra) + ev <- list(type = "org.matrix.msc3401.call.member", sender = sender, + state_key = ns$mx_call_state_key(sender, device), content = c) + if (!is.null(origin_ts)) ev$origin_server_ts <- origin_ts + ev +} +now_ms <- ns$mx_call_now_ms(T0) +state <- list( + member_event("@alice:example.org", "PHONE", now_ms - 1000), + member_event("@carol:example.org", "OLD", now_ms - 5 * 3600 * 1000), + member_event("@dave:example.org", "SHORT", now_ms - 2000, expires = 1000), + list(type = "org.matrix.msc3401.call.member", sender = "@erin:example.org", + state_key = "_@erin:example.org_GONE_m.call", + content = stats::setNames(list(), character())), + member_event("@frank:example.org", "JITSI", now_ms, focus = "jitsi"), + member_event("@bob:example.org", "LAPTOP", NULL, origin_ts = now_ms - 60000, + content_extra = list(membershipID = "custom")), + list(type = "m.room.member", sender = "@alice:example.org", + state_key = "@alice:example.org", content = list(membership = "join")) +) +members <- mx_call_members(state, now = T0) +expect_equal(length(members), 2L) +expect_identical(vapply(members, `[[`, "", "identity"), + c("@alice:example.org:PHONE", "@bob:example.org:LAPTOP")) +expect_identical(members[[1]]$service_urls, "https://jwt.a") +expect_identical(members[[1]]$membership_id, "@alice:example.org:PHONE") +expect_identical(members[[2]]$membership_id, "custom") +expect_equal(members[[2]]$expires_at, now_ms - 60000 + 4 * 3600 * 1000) +expect_equal(length(mx_call_members(list(), now = T0)), 0L) +# created_at (FluffyChat's name) is read too +fc <- member_event("@g:ex", "FC", NULL) +fc$content$created_at <- now_ms +expect_equal(length(mx_call_members(list(fc), now = T0)), 1L) + +# ---- service discovery ------------------------------------------------ + +expect_identical(mx_call_service_url(client, members), "https://jwt.a") +transports <- list(list(type = "livekit", + livekit_service_url = "https://jwt.transport")) +expect_identical( + with_api("mx_rtc_transports", function(session) transports, + mx_call_service_url(client, list())), + "https://jwt.transport") +expect_identical( + with_api("mx_rtc_transports", function(session) list(), + with_api("mx_well_known_client", function(server_name) { + expect_identical(server_name, "example.org") + list(`org.matrix.msc4143.rtc_foci` = list(list( + type = "livekit", + livekit_service_url = "https://jwt.wellknown"))) + }, mx_call_service_url(client, list()))), + "https://jwt.wellknown") +expect_error( + with_api("mx_rtc_transports", function(session) list(), + with_api("mx_well_known_client", function(server_name) NULL, + mx_call_service_url(client, list()))), + "no LiveKit JWT service found") + +# ---- key events ------------------------------------------------------- + +key <- as.raw(0:15) +kc <- ns$mx_call_key_content(key, 3L, "@bot:example.org", "BOTDEV", ROOM, + now = T0) +expect_identical(kc$keys$index, 3L) +expect_identical(kc$keys$key, "AAECAwQFBgcICQoLDA0ODw==") +expect_identical(kc$member$id, "@bot:example.org:BOTDEV") +expect_identical(kc$member$claimed_device_id, "BOTDEV") +expect_identical(kc$room_id, ROOM) +expect_identical(kc$session$application, "m.call") +expect_equal(kc$sent_ts, 1.7e12) + +parsed <- mx_call_key_parse(list(type = "io.element.call.encryption_keys", + sender = "@bot:example.org", content = kc), + ROOM) +expect_identical(parsed$identity, "@bot:example.org:BOTDEV") +expect_identical(parsed$key, key) +expect_identical(parsed$index, 3L) +# Unpadded and URL-safe base64 decode too +kc2 <- kc +kc2$keys$key <- "AAECAwQFBgcICQoLDA0ODw" +expect_identical(mx_call_key_parse(list(type = "io.element.call.encryption_keys", + sender = "@bot:example.org", content = kc2), ROOM)$key, key) +kc3 <- kc +kc3$keys$key <- jsonlite::base64url_enc(as.raw(c(0xfb, 0xff, 0xfe, 1))) +expect_identical(mx_call_key_parse(list(type = "io.element.call.encryption_keys", + sender = "@bot:example.org", content = kc3), ROOM)$key, + as.raw(c(0xfb, 0xff, 0xfe, 1))) +# Rejections: other room, other type, bad index, missing device +expect_null(mx_call_key_parse(list(type = "io.element.call.encryption_keys", + sender = "@bot:example.org", content = kc), "!other:example.org")) +expect_null(mx_call_key_parse(list(type = "m.room_key", sender = "@b:ex", + content = kc), ROOM)) +kc4 <- kc +kc4$keys$index <- -1 +expect_null(mx_call_key_parse(list(type = "io.element.call.encryption_keys", + sender = "@bot:example.org", content = kc4), ROOM)) +kc5 <- kc +kc5$member$claimed_device_id <- NULL +expect_null(mx_call_key_parse(list(type = "io.element.call.encryption_keys", + sender = "@bot:example.org", content = kc5), ROOM)) + +# ---- rotation policy -------------------------------------------------- + +keys <- ns$mx_call_keys_new() +expect_identical(keys$index, -1L) +# No key yet: make one, send to everyone +plan <- mx_call_key_plan(keys, c("@a:ex:A", "@b:ex:B"), now = T0) +expect_true(plan$rotate) +expect_identical(plan$targets, c("@a:ex:A", "@b:ex:B")) +ns$mx_call_keys_rotate(keys, T0) +expect_identical(keys$index, 0L) +expect_equal(length(keys$key), 16L) +keys$shared_with <- c("@a:ex:A", "@b:ex:B") +# Same members: nothing to do +plan <- mx_call_key_plan(keys, c("@b:ex:B", "@a:ex:A"), now = T0 + 5) +expect_false(plan$rotate) +expect_identical(plan$targets, character()) +# Joiner inside the grace period: current key to the joiner only +plan <- mx_call_key_plan(keys, c("@a:ex:A", "@b:ex:B", "@c:ex:C"), now = T0 + 5) +expect_false(plan$rotate) +expect_identical(plan$targets, "@c:ex:C") +# Joiner after the grace period: rotate, send to all +plan <- mx_call_key_plan(keys, c("@a:ex:A", "@b:ex:B", "@c:ex:C"), now = T0 + 30) +expect_true(plan$rotate) +expect_identical(plan$targets, c("@a:ex:A", "@b:ex:B", "@c:ex:C")) +# Leaver: rotate for the rest, even inside the grace period +plan <- mx_call_key_plan(keys, "@a:ex:A", now = T0 + 1) +expect_true(plan$rotate) +expect_identical(plan$targets, "@a:ex:A") +# Everyone left: rotate, nobody to send to +plan <- mx_call_key_plan(keys, character(), now = T0 + 1) +expect_true(plan$rotate) +expect_identical(plan$targets, character()) +# Indexes wrap below 255 +keys$index <- 254L +ns$mx_call_keys_rotate(keys, T0) +expect_identical(keys$index, 0L) +expect_identical(keys$shared_with, character()) + +# ---- update_keys drives send_key with the plan ----------------------- + +local({ + sends <- list() + orig <- get("mx_call_send_key", envir = ns) + assignInNamespace("mx_call_send_key", function(call, targets) { + sends[[length(sends) + 1L]] <<- list(index = call$keys$index, + targets = targets) + call$keys$shared_with <- union(call$keys$shared_with, targets) + invisible(NULL) + }, ns = "mx.client") + on.exit(assignInNamespace("mx_call_send_key", orig, ns = "mx.client"), + add = TRUE) + call <- new.env() + call$identity <- "@bot:example.org:BOTDEV" + call$keys <- ns$mx_call_keys_new() + call$session <- NULL + me <- list(identity = "@bot:example.org:BOTDEV") + a <- list(identity = "@alice:example.org:PHONE") + b <- list(identity = "@bob:example.org:LAPTOP") + ns$mx_call_update_keys(call, list(me, a), now = T0) + expect_identical(sends[[1]], list(index = 0L, targets = a$identity)) + ns$mx_call_update_keys(call, list(me, a, b), now = T0 + 2) + expect_identical(sends[[2]], list(index = 0L, targets = b$identity)) + ns$mx_call_update_keys(call, list(me, b), now = T0 + 3) + expect_identical(sends[[3]], list(index = 1L, targets = b$identity)) + expect_equal(length(call$members), 2L) +}) + +# ---- mx_call_handle: keys from to_device, membership changes ---------- + +local({ + applied <- list() + orig_apply <- get("mx_call_apply_peer_key", envir = ns) + orig_update <- get("mx_call_update_keys", envir = ns) + assignInNamespace("mx_call_apply_peer_key", function(call, parsed) { + applied[[length(applied) + 1L]] <<- parsed$identity + call$keys$peers[[parsed$identity]] <- parsed[c("key", "index")] + }, ns = "mx.client") + assignInNamespace("mx_call_update_keys", function(call, members, now) { + call$members <- members + }, ns = "mx.client") + on.exit({ + assignInNamespace("mx_call_apply_peer_key", orig_apply, ns = "mx.client") + assignInNamespace("mx_call_update_keys", orig_update, ns = "mx.client") + }, add = TRUE) + + call <- new.env() + call$client <- client + call$room_id <- ROOM + call$identity <- "@bot:example.org:BOTDEV" + call$keys <- ns$mx_call_keys_new() + call$session <- NULL + call$members <- list() + call$membership_sent <- Sys.time() + + peer_key <- ns$mx_call_key_content(as.raw(1:16), 2L, "@alice:example.org", + "PHONE", ROOM) + other_room <- ns$mx_call_key_content(as.raw(1:16), 2L, "@alice:example.org", + "PHONE", "!other:example.org") + processed <- list(to_device = list( + list(type = "io.element.call.encryption_keys", + sender = "@alice:example.org", content = peer_key), + list(type = "io.element.call.encryption_keys", + sender = "@alice:example.org", content = other_room), + list(type = "org.example.other", sender = "@x:ex", content = list()))) + quiet <- list(rooms = list(join = list())) + res <- mx_call_handle(call, quiet, processed) + expect_identical(res$keys, "@alice:example.org:PHONE") + expect_null(res$members) + expect_identical(applied, list("@alice:example.org:PHONE")) + expect_identical(call$keys$peers[["@alice:example.org:PHONE"]]$index, 2L) + + # A membership event in the room's timeline triggers a state refresh + with_state <- list(rooms = list(join = stats::setNames(list(list( + timeline = list(events = list(list( + type = "org.matrix.msc3401.call.member", + sender = "@alice:example.org"))))), ROOM))) + res <- with_api("mx_room_state", function(session, room_id) { + expect_identical(room_id, ROOM) + list(member_event("@alice:example.org", "PHONE", + ns$mx_call_now_ms(Sys.time()))) + }, mx_call_handle(call, with_state, list(to_device = list()))) + expect_equal(length(res$members), 1L) + expect_identical(call$members[[1]]$identity, "@alice:example.org:PHONE") + + # Membership is re-sent once an hour + sent <- NULL + call$membership_sent <- Sys.time() - 3601 + call$service_url <- "https://jwt.a" + call$intent <- "voice" + with_api("mx_set_state", function(session, room_id, event_type, content, + state_key = "") { + sent <<- list(type = event_type, state_key = state_key, + content = content) + list(event_id = "$e") + }, mx_call_handle(call, quiet, list(to_device = list()))) + expect_identical(sent$type, "org.matrix.msc3401.call.member") + expect_identical(sent$state_key, "_@bot:example.org_BOTDEV_m.call") + expect_identical(sent$content$device_id, "BOTDEV") + expect_true(as.numeric(Sys.time()) - as.numeric(call$membership_sent) < 5) +}) + +# ---- mx_call_leave clears the membership ----------------------------- + +local({ + sent <- NULL + call <- new.env() + call$client <- client + call$room_id <- ROOM + call$session <- NULL + with_api("mx_set_state", function(session, room_id, event_type, content, + state_key = "") { + sent <<- list(room_id = room_id, type = event_type, content = content, + state_key = state_key) + list(event_id = "$e") + }, mx_call_leave(call)) + expect_identical(sent$room_id, ROOM) + expect_identical(sent$state_key, "_@bot:example.org_BOTDEV_m.call") + expect_equal(length(sent$content), 0L) + expect_identical(as.character(jsonlite::toJSON(sent$content, + auto_unbox = TRUE)), "{}") +}) + +# ---- mx_call_join on the Matrix side only ---------------------------- + +local({ + calls <- list() + record <- function(what) { + function(...) { + calls[[length(calls) + 1L]] <<- what + switch(what, + room_state = list(member_event("@alice:example.org", "PHONE", + ns$mx_call_now_ms(Sys.time()))), + openid = list(access_token = "oid", token_type = "Bearer", + matrix_server_name = "example.org", + expires_in = 3600L), + livekit_token = list(url = "wss://sfu.example.org", + jwt = "jwt"), + set_state = list(event_id = "$e")) + } + } + orig_send <- get("mx_call_send_key", envir = ns) + orig_req <- get("mx_require_crypto", envir = ns) + assignInNamespace("mx_call_send_key", function(call, targets) { + calls[[length(calls) + 1L]] <<- paste("send_key", paste(targets, + collapse = ",")) + call$keys$shared_with <- targets + }, ns = "mx.client") + assignInNamespace("mx_require_crypto", function() invisible(TRUE), + ns = "mx.client") + on.exit({ + assignInNamespace("mx_call_send_key", orig_send, ns = "mx.client") + assignInNamespace("mx_require_crypto", orig_req, ns = "mx.client") + }, add = TRUE) + call <- with_api("mx_room_state", record("room_state"), + with_api("mx_openid_token", record("openid"), + with_api("mx_rtc_livekit_token", record("livekit_token"), + with_api("mx_set_state", record("set_state"), + mx_call_join(client, NULL, mx_crypto_sessions_new(), ROOM, + connect = FALSE))))) + expect_inherits(call, "mx_call") + expect_identical(calls, list("room_state", "openid", "livekit_token", + "set_state", + "send_key @alice:example.org:PHONE")) + expect_identical(call$service_url, "https://jwt.a") + expect_identical(call$token$url, "wss://sfu.example.org") + expect_identical(call$identity, "@bot:example.org:BOTDEV") + expect_equal(length(call$members), 1L) + expect_identical(call$keys$index, 0L) + expect_null(call$session) + expect_stdout(print(call), "media not connected") + # An explicit service URL skips discovery + call2 <- with_api("mx_room_state", function(...) list(), + with_api("mx_openid_token", record("openid"), + with_api("mx_rtc_livekit_token", function(service_url, ...) { + expect_identical(service_url, "https://jwt.explicit") + list(url = "wss://x", jwt = "j") + }, + with_api("mx_set_state", record("set_state"), + mx_call_join(client, NULL, mx_crypto_sessions_new(), ROOM, + connect = FALSE, service_url = "https://jwt.explicit"))))) + expect_identical(call2$service_url, "https://jwt.explicit") + expect_equal(length(call2$members), 0L) +}) diff --git a/inst/tinytest/test_to_device.R b/inst/tinytest/test_to_device.R new file mode 100644 index 0000000..8a6db13 --- /dev/null +++ b/inst/tinytest/test_to_device.R @@ -0,0 +1,174 @@ +# Olm to-device events of arbitrary type: Alice encrypts one for Bob's +# device, Bob's sync processing decrypts it and hands it back in +# `to_device`; the shape of the payload and the sender checks. + +library(tinytest) + +if (!requireNamespace("mx.crypto", quietly = TRUE)) { + exit_file("mx.crypto not available (needs a Rust toolchain)") +} +library(mx.client) + +alice <- mx.crypto::mxc_account_new() +bob <- mx.crypto::mxc_account_new() +alice_keys <- mx.crypto::mxc_account_identity_keys(alice) +bob_keys <- mx.crypto::mxc_account_identity_keys(bob) +mx.crypto::mxc_account_generate_one_time_keys(bob, 2L) +bob_otk <- mx.crypto::mxc_account_one_time_keys(bob)[[1]] + +bob_device <- list(user_id = "@bob:example.org", device_id = "BOBDEV", + curve25519 = bob_keys$curve25519, + ed25519 = bob_keys$ed25519, otk = bob_otk) + +content <- list(keys = list(index = 3L, key = "AAECAwQFBgcICQoLDA0ODw=="), + member = list(id = "@alice:example.org:ALICEDEV", + claimed_device_id = "ALICEDEV"), + room_id = "!call:example.org") + +# Encrypt: one payload per recipient, Olm session opened and retained +a_sess <- mx_crypto_sessions_new() +out <- mx_crypto_encrypt_to_device(alice, a_sess, + "io.element.call.encryption_keys", content, + recipients = list(bob_device), + sender_user_id = "@alice:example.org") +expect_equal(length(out$to_device), 1L) +expect_identical(out$to_device[[1]]$user_id, "@bob:example.org") +expect_identical(out$to_device[[1]]$device_id, "BOBDEV") +payload <- out$to_device[[1]]$content +expect_identical(payload$algorithm, "m.olm.v1.curve25519-aes-sha2") +expect_identical(payload$sender_key, alice_keys$curve25519) +expect_identical(names(payload$ciphertext), bob_keys$curve25519) +expect_identical(payload$ciphertext[[1]]$type, 0L) # prekey message +expect_true(!is.null(out$sessions$olm[[bob_keys$curve25519]])) + +# No recipients: nothing to send, no error +expect_equal(length(mx_crypto_encrypt_to_device(alice, a_sess, "x", list(), + list(), "@alice:example.org")$to_device), + 0L) + +# A recipient without verified keys is refused +expect_error(mx_crypto_encrypt_to_device(alice, a_sess, "x", list(), + list(list(user_id = "@c:ex", device_id = "C", curve25519 = "k")), "@a:ex"), + "no verified keys") + +# Decrypt on Bob's side through process_sync: the event comes back in +# to_device with its type, content and sender +b_sess <- mx_crypto_sessions_new() +sync <- list(to_device = list(events = list(list( + type = "m.room.encrypted", sender = "@alice:example.org", + content = payload))), rooms = list(join = list())) +res <- mx_crypto_process_sync(bob, b_sess, sync, bob_keys$curve25519, + self_id = "@bob:example.org", + self_device_id = "BOBDEV") +expect_equal(length(res$to_device), 1L) +ev <- res$to_device[[1]] +expect_identical(ev$type, "io.element.call.encryption_keys") +expect_identical(ev$sender, "@alice:example.org") +expect_identical(ev$content$keys$key, content$keys$key) +expect_identical(ev$content$member$claimed_device_id, "ALICEDEV") +expect_false(ev$sender_bound) # no device list given, so a claim only +expect_equal(length(res$events), 0L) +expect_equal(length(res$verification_events), 0L) + +# With Alice's device known, the sender is bound +alice_device <- list(user_id = "@alice:example.org", device_id = "ALICEDEV", + curve25519 = alice_keys$curve25519, + ed25519 = alice_keys$ed25519) +out2 <- mx_crypto_encrypt_to_device(alice, out$sessions, "org.example.ping", + list(n = 2L), list(bob_device), + "@alice:example.org") +expect_identical(out2$to_device[[1]]$content$ciphertext[[1]]$type, 0L) +sync2 <- list(to_device = list(events = list(list( + type = "m.room.encrypted", sender = "@alice:example.org", + content = out2$to_device[[1]]$content))), rooms = list(join = list())) +res2 <- mx_crypto_process_sync(bob, res$sessions, sync2, bob_keys$curve25519, + self_id = "@bob:example.org", + devices = list(alice_device), + self_device_id = "BOBDEV") +expect_equal(length(res2$to_device), 1L) +expect_true(res2$to_device[[1]]$sender_bound) +expect_identical(res2$to_device[[1]]$content$n, 2L) + +# An envelope whose sender disagrees with the plaintext is dropped +out3 <- mx_crypto_encrypt_to_device(alice, out2$sessions, "org.example.ping", + list(n = 3L), list(bob_device), + "@alice:example.org") +sync3 <- list(to_device = list(events = list(list( + type = "m.room.encrypted", sender = "@mallory:example.org", + content = out3$to_device[[1]]$content))), rooms = list(join = list())) +res3 <- mx_crypto_process_sync(bob, res2$sessions, sync3, bob_keys$curve25519, + self_id = "@bob:example.org", + self_device_id = "BOBDEV") +expect_equal(length(res3$to_device), 0L) + +# Room keys still take their own path, not the generic one +room_out <- mx_crypto_encrypt_for_devices(alice, out3$sessions, "!r:example.org", + list(msgtype = "m.text", body = "hi"), alice_keys$curve25519, "ALICEDEV", + recipients = list(bob_device), sender_user_id = "@alice:example.org") +sync4 <- list(to_device = list(events = list(list( + type = "m.room.encrypted", sender = "@alice:example.org", + content = room_out$to_device[[1]]$content))), rooms = list(join = list())) +res4 <- mx_crypto_process_sync(bob, res3$sessions, sync4, bob_keys$curve25519, + self_id = "@bob:example.org", + self_device_id = "BOBDEV") +expect_equal(length(res4$to_device), 0L) +expect_equal(length(res4$sessions$megolm_in), 1L) + +# mx_send_to_device_encrypted: groups payloads by user and device, never +# sends to this device, and persists sessions +local({ + bob2 <- mx.crypto::mxc_account_new() + mx.crypto::mxc_account_generate_one_time_keys(bob2, 1L) + bob_otk2 <- mx.crypto::mxc_account_one_time_keys(bob2)[[1]] + bob2_keys <- mx.crypto::mxc_account_identity_keys(bob2) + + ns <- asNamespace("mx.client") + sent <- NULL + orig_claim <- get("mx_crypto_claim_otks", envir = ns) + assignInNamespace("mx_crypto_claim_otks", function(client, devices, strict) { + lapply(devices, function(d) { d$otk <- bob_otk2; d }) + }, ns = "mx.client") + on.exit(assignInNamespace("mx_crypto_claim_otks", orig_claim, + ns = "mx.client"), add = TRUE) + api <- asNamespace("mx.api") + orig_api_send <- get("mx_send_to_device", envir = api) + assignInNamespace("mx_send_to_device", function(session, event_type, + messages, txn_id = NULL) { + sent <<- list(event_type = event_type, messages = messages) + list() + }, ns = "mx.api") + on.exit(assignInNamespace("mx_send_to_device", orig_api_send, ns = "mx.api"), + add = TRUE) + + bob2_device <- list(user_id = "@bob:example.org", device_id = "BOBTWO", + curve25519 = bob2_keys$curve25519, + ed25519 = bob2_keys$ed25519) + self_device <- list(user_id = "@alice:example.org", device_id = "ALICEDEV", + curve25519 = alice_keys$curve25519, + ed25519 = alice_keys$ed25519) + client <- list(server = "https://example.org", token = "t", + user_id = "@alice:example.org", device_id = "ALICEDEV") + store <- tempfile("td-store") + res <- mx_send_to_device_encrypted(client, alice, room_out$sessions, + "org.example.ping", list(n = 4L), + list(bob_device, bob2_device, self_device), + store_dir = store) + expect_identical(sent$event_type, "m.room.encrypted") + expect_identical(names(sent$messages), "@bob:example.org") + expect_identical(sort(names(sent$messages[["@bob:example.org"]])), + c("BOBDEV", "BOBTWO")) + expect_equal(length(res$sent), 2L) + expect_true(!is.null(res$sessions$olm[[bob2_keys$curve25519]])) + expect_true(file.exists(store)) + # Bob's second device opens the session from the prekey message + msg <- sent$messages[["@bob:example.org"]][["BOBTWO"]] + syncb <- list(to_device = list(events = list(list( + type = "m.room.encrypted", sender = "@alice:example.org", + content = msg))), rooms = list(join = list())) + rb <- mx_crypto_process_sync(bob2, mx_crypto_sessions_new(), syncb, + bob2_keys$curve25519, + self_id = "@bob:example.org", + self_device_id = "BOBTWO") + expect_identical(rb$to_device[[1]]$content$n, 4L) + unlink(store, recursive = TRUE) +}) diff --git a/man/mx_call_connect.Rd b/man/mx_call_connect.Rd new file mode 100644 index 0000000..a6b90dd --- /dev/null +++ b/man/mx_call_connect.Rd @@ -0,0 +1,23 @@ +% tinyrox says don't edit this manually, but it can't stop you! +\name{mx_call_connect} +\alias{mx_call_connect} +\title{Join the LiveKit room of a call} +\usage{ +mx_call_connect(call, ...) +} +\arguments{ +\item{call}{An \code{"mx_call"}.} + +\item{...}{Further arguments to \code{livekitr::lk_connect()}, such +as \code{opts}.} +} +\value{ +\code{call}, invisibly. +} +\description{ +Connects the media side of a call made with \code{mx_call_join(connect += FALSE)}: joins the LiveKit room with the media token, with +end-to-end encryption configured as Element Call and FluffyChat +expect (per-participant HKDF keys, a 256-slot key ring), and sets our +key and every peer key received so far. +} diff --git a/man/mx_call_handle.Rd b/man/mx_call_handle.Rd new file mode 100644 index 0000000..1a2a1a7 --- /dev/null +++ b/man/mx_call_handle.Rd @@ -0,0 +1,30 @@ +% tinyrox says don't edit this manually, but it can't stop you! +\name{mx_call_handle} +\alias{mx_call_handle} +\title{Feed a sync response to a call} +\usage{ +mx_call_handle(call, sync, processed = NULL) +} +\arguments{ +\item{call}{An \code{"mx_call"}.} + +\item{sync}{A parsed \code{/sync} response.} + +\item{processed}{The result of \code{mx_crypto_process_sync()} on +\code{sync}, or NULL to have it run here (its \code{sessions} are +then kept in \code{call$sessions}).} +} +\value{ +A list of what changed: \code{keys}, the identities whose + keys were received, and \code{members}, the current members when + membership changed (NULL otherwise). +} +\description{ +Applies what a \code{/sync} response means for the call: peers' media +keys from the decrypted to-device events, membership changes (which +may rotate and resend our key), and the hourly refresh of our own +membership. Call it with every sync while in the call. When the +application already runs \code{\link{mx_crypto_process_sync}} on the +response, pass its result as \code{processed} so to-device events +are not decrypted twice. +} diff --git a/man/mx_call_join.Rd b/man/mx_call_join.Rd new file mode 100644 index 0000000..cd98d0a --- /dev/null +++ b/man/mx_call_join.Rd @@ -0,0 +1,64 @@ +% tinyrox says don't edit this manually, but it can't stop you! +\name{mx_call_join} +\alias{mx_call_join} +\title{Join a MatrixRTC call} +\usage{ +mx_call_join( + client, + account, + sessions, + room_id, + intent = "voice", + store_dir = NULL, + connect = TRUE, + service_url = NULL +) +} +\arguments{ +\item{client}{Matrix client config.} + +\item{account}{An mx.crypto account handle.} + +\item{sessions}{A session set.} + +\item{room_id}{The room whose call to join.} + +\item{intent}{\code{"voice"} or \code{"video"}, what the membership +announces.} + +\item{store_dir}{Character or NULL. Where the updated crypto sessions +are saved after each key send.} + +\item{connect}{Join the LiveKit room now? \code{FALSE} does everything +on the Matrix side only.} + +\item{service_url}{Character or NULL. The LiveKit JWT service to use +instead of discovering one.} +} +\value{ +An object of class \code{"mx_call"}: an environment holding + the \code{client} and \code{sessions} (both updated as the call + runs), the LiveKit \code{session} from \pkg{livekitr} (or NULL), this + device's LiveKit \code{identity}, the current \code{members}, and + the media \code{token}. +} +\description{ +Joins the call of a room the way Element Call and FluffyChat do: +finds the LiveKit JWT service, trades an OpenID token for a media +token, announces this device's membership as a state event, makes a +media key and sends it Olm-encrypted to every device already in the +call, and, with \code{connect = TRUE}, joins the LiveKit room through +\pkg{livekitr} with end-to-end encryption on. The returned call must +then be driven with \code{\link{mx_call_handle}} (or +\code{\link{mx_call_poll}}) and ended with \code{\link{mx_call_leave}}. +} +\examples{ +\dontrun{ +call <- mx_call_join(client, acct, sessions, "!room:example.org", + store_dir = store) +livekitr::lk_on_audio(call$session, function(pcm, info) { + cat(info$identity, "spoke\n") +}) +repeat mx_call_poll(call) +} +} diff --git a/man/mx_call_key_parse.Rd b/man/mx_call_key_parse.Rd new file mode 100644 index 0000000..eea6ed0 --- /dev/null +++ b/man/mx_call_key_parse.Rd @@ -0,0 +1,29 @@ +% tinyrox says don't edit this manually, but it can't stop you! +\name{mx_call_key_parse} +\alias{mx_call_key_parse} +\title{Read a call key event} +\usage{ +mx_call_key_parse(event, room_id) +} +\arguments{ +\item{event}{List with \code{type}, \code{content} and \code{sender}.} + +\item{room_id}{The call's room; keys for other rooms are ignored.} +} +\value{ +\code{list(identity, key, index)}, with \code{key} a raw + vector, or NULL when the event is not a usable key for this room. +} +\description{ +Parses a decrypted \code{io.element.call.encryption_keys} to-device +event (from the \code{to_device} list of +\code{\link{mx_crypto_process_sync}}) into the LiveKit identity it +belongs to and the key to set for it. +} +\examples{ +ev <- list(type = "io.element.call.encryption_keys", + sender = "@alice:example.org", + content = list(keys = list(index = 3, key = "AAECAwQFBgcICQoLDA0ODw=="), + member = list(claimed_device_id = "PHONE"), room_id = "!r:example.org")) +mx_call_key_parse(ev, "!r:example.org") +} diff --git a/man/mx_call_key_plan.Rd b/man/mx_call_key_plan.Rd new file mode 100644 index 0000000..2c9c1cc --- /dev/null +++ b/man/mx_call_key_plan.Rd @@ -0,0 +1,38 @@ +% tinyrox says don't edit this manually, but it can't stop you! +\name{mx_call_key_plan} +\alias{mx_call_key_plan} +\title{Decide what to do with the call key when membership changes} +\usage{ +mx_call_key_plan( + keys, + identities, + now = Sys.time(), + grace = MX_CALL_KEY_GRACE_S +) +} +\arguments{ +\item{keys}{Key state from \code{mx_call_keys_new()} (an environment; +it is not modified).} + +\item{identities}{LiveKit identities of the other members now in the +call.} + +\item{now}{The current time.} + +\item{grace}{Seconds after a key is made during which it is still +handed to joiners instead of rotated.} +} +\value{ +\code{list(rotate, targets)}: whether to make a new key, and + the identities to send the (new or current) key to. +} +\description{ +Implements the rotation policy Element Call and FluffyChat follow: a +leaver forces a new key for everyone; a joiner within the grace period +after the key was made receives the current key; a joiner after it +gets a new key, as does everyone else. +} +\examples{ +keys <- mx.client:::mx_call_keys_new() +mx_call_key_plan(keys, c("@a:ex:D1", "@b:ex:D2")) +} diff --git a/man/mx_call_leave.Rd b/man/mx_call_leave.Rd new file mode 100644 index 0000000..a6e732c --- /dev/null +++ b/man/mx_call_leave.Rd @@ -0,0 +1,16 @@ +% tinyrox says don't edit this manually, but it can't stop you! +\name{mx_call_leave} +\alias{mx_call_leave} +\title{Leave a MatrixRTC call} +\usage{ +mx_call_leave(call) +} +\arguments{ +\item{call}{An \code{"mx_call"}.} +} +\value{ +\code{NULL}, invisibly. +} +\description{ +Leaves the LiveKit room and clears this device's membership event. +} diff --git a/man/mx_call_members.Rd b/man/mx_call_members.Rd new file mode 100644 index 0000000..e89da82 --- /dev/null +++ b/man/mx_call_members.Rd @@ -0,0 +1,35 @@ +% tinyrox says don't edit this manually, but it can't stop you! +\name{mx_call_members} +\alias{mx_call_members} +\title{Active call memberships in a room's state} +\usage{ +mx_call_members(state, now = Sys.time()) +} +\arguments{ +\item{state}{List of state events.} + +\item{now}{The current time.} +} +\value{ +A list of members, each \code{list(user_id, device_id, + identity, membership_id, service_urls, expires_at)}, where + \code{identity} is the member's LiveKit participant identity + \code{":"} and \code{service_urls} the LiveKit JWT + services it prefers. +} +\description{ +Reads the \code{org.matrix.msc3401.call.member} state events out of a +room's full state (\code{mx.api::mx_room_state()}) and keeps the ones +that count as in the call: non-empty content, a LiveKit focus, and an +expiry (\code{created_ts}, or \code{created_at}, or the event's own +timestamp, plus \code{expires}, default 4 hours) still in the future. +} +\examples{ +state <- list(list(type = "org.matrix.msc3401.call.member", + sender = "@alice:example.org", origin_server_ts = 1e12, + content = list(application = "m.call", device_id = "PHONE", + focus_active = list(type = "livekit"), + foci_preferred = list(list(type = "livekit", + livekit_service_url = "https://jwt.example.org"))))) +mx_call_members(state, now = as.POSIXct(1e9, origin = "1970-01-01")) +} diff --git a/man/mx_call_poll.Rd b/man/mx_call_poll.Rd new file mode 100644 index 0000000..783afa2 --- /dev/null +++ b/man/mx_call_poll.Rd @@ -0,0 +1,24 @@ +% tinyrox says don't edit this manually, but it can't stop you! +\name{mx_call_poll} +\alias{mx_call_poll} +\title{Sync once and poll the media of a call} +\usage{ +mx_call_poll(call, timeout = 1, media_timeout = 0.1) +} +\arguments{ +\item{call}{An \code{"mx_call"}.} + +\item{timeout}{Seconds to wait for the sync long poll.} + +\item{media_timeout}{Seconds to wait in the media poll.} +} +\value{ +The LiveKit events from \code{livekitr::lk_poll()}, or an + empty list when the call has no media session. +} +\description{ +One iteration of a call loop for programs with no sync loop of their +own: syncs with \code{\link{mx_sync_update}}, applies the response with +\code{\link{mx_call_handle}}, and polls the LiveKit session with +\code{livekitr::lk_poll()}, which is where audio callbacks run. +} diff --git a/man/mx_call_service_url.Rd b/man/mx_call_service_url.Rd new file mode 100644 index 0000000..dfae8bc --- /dev/null +++ b/man/mx_call_service_url.Rd @@ -0,0 +1,27 @@ +% tinyrox says don't edit this manually, but it can't stop you! +\name{mx_call_service_url} +\alias{mx_call_service_url} +\title{Find the LiveKit JWT service for a room's call} +\usage{ +mx_call_service_url(client, members = list()) +} +\arguments{ +\item{client}{Matrix client config.} + +\item{members}{Current members from \code{\link{mx_call_members}}.} +} +\value{ +The service URL, a string. +} +\description{ +Tries, in order, the services the current call members advertise in +their memberships, the homeserver's RTC transports +(\code{mx.api::mx_rtc_transports()}), and the +\code{org.matrix.msc4143.rtc_foci} entry of the server's client +well-known file. +} +\examples{ +\dontrun{ +mx_call_service_url(client, members) +} +} diff --git a/man/mx_crypto_encrypt_to_device.Rd b/man/mx_crypto_encrypt_to_device.Rd new file mode 100644 index 0000000..8201e19 --- /dev/null +++ b/man/mx_crypto_encrypt_to_device.Rd @@ -0,0 +1,53 @@ +% tinyrox says don't edit this manually, but it can't stop you! +\name{mx_crypto_encrypt_to_device} +\alias{mx_crypto_encrypt_to_device} +\title{Encrypt a to-device event for a set of devices} +\usage{ +mx_crypto_encrypt_to_device( + account, + sessions, + event_type, + content, + recipients, + sender_user_id +) +} +\arguments{ +\item{account}{An mx.crypto account handle.} + +\item{sessions}{A session set.} + +\item{event_type}{The inner event type, e.g. +\code{"io.element.call.encryption_keys"}.} + +\item{content}{Named list. Plaintext event content.} + +\item{recipients}{List of recipient devices as returned by +\code{mx_crypto_known_devices()}: \code{user_id}, \code{device_id}, +\code{curve25519}, \code{ed25519}, and \code{otk} for a device +without an Olm session yet.} + +\item{sender_user_id}{Character. This user's Matrix id.} +} +\value{ +List with \code{to_device} (per-device payloads, each + \code{list(user_id, device_id, content)}) and the updated + \code{sessions}. +} +\description{ +Produces one Olm-encrypted \code{m.room.encrypted} payload per +recipient device, opening an Olm session where none exists yet. Unlike +Megolm room keys, which are shared once per session, the event is +encrypted for every recipient on every call. +} +\examples{ +\donttest{ +if (requireNamespace("mx.crypto", quietly = TRUE)) { + acct <- mx.crypto::mxc_account_new() + out <- mx_crypto_encrypt_to_device(acct, mx_crypto_sessions_new(), + "org.example.ping", list(n = 1), recipients = list(), + sender_user_id = "@me:ex") + length(out$to_device) +} +} +} diff --git a/man/mx_crypto_process_sync.Rd b/man/mx_crypto_process_sync.Rd index a5753e7..2dce28a 100644 --- a/man/mx_crypto_process_sync.Rd +++ b/man/mx_crypto_process_sync.Rd @@ -44,8 +44,11 @@ List with \code{events} (decrypted, normalized), updated \code{verification_events} (original verification envelopes, separated from chat messages; no handshake or network side effect is performed), \code{sessions}, unsent \code{key_requests}, matching - \code{key_request_cancellations}, and \code{incoming_key_requests} - for a policy-aware sharing layer to inspect. + \code{key_request_cancellations}, \code{incoming_key_requests} + for a policy-aware sharing layer to inspect, and \code{to_device}: + every other Olm-encrypted to-device event that decrypted and whose + claimed sender matches the envelope, each as \code{list(type, + content, sender, sender_bound)}. Call encryption keys arrive here. } \description{ Handles inbound to-device \code{m.room.encrypted} (Olm) messages, diff --git a/man/mx_send_to_device_encrypted.Rd b/man/mx_send_to_device_encrypted.Rd new file mode 100644 index 0000000..f860644 --- /dev/null +++ b/man/mx_send_to_device_encrypted.Rd @@ -0,0 +1,50 @@ +% tinyrox says don't edit this manually, but it can't stop you! +\name{mx_send_to_device_encrypted} +\alias{mx_send_to_device_encrypted} +\title{Send an Olm-encrypted to-device event} +\usage{ +mx_send_to_device_encrypted( + client, + account, + sessions, + event_type, + content, + devices, + store_dir = NULL +) +} +\arguments{ +\item{client}{Matrix client config.} + +\item{account}{An mx.crypto account handle.} + +\item{sessions}{A session set.} + +\item{event_type}{The inner event type.} + +\item{content}{Named list. Plaintext event content.} + +\item{devices}{List of target devices from +\code{mx_crypto_known_devices()}, already narrowed to the devices +that should receive the event.} + +\item{store_dir}{Character or NULL. Where to persist the updated +sessions; NULL leaves saving to the caller.} +} +\value{ +List with \code{sessions} (updated) and \code{sent}, the + \code{list(user_id, device_id)} pairs the event went to. +} +\description{ +Encrypts \code{content} for each device in \code{devices}, claiming +one-time keys where no Olm session exists, and delivers it with +\code{mx.api::mx_send_to_device()} in one request. This device itself +is never a recipient. +} +\examples{ +\dontrun{ +devs <- mx_crypto_known_devices(client, "@bob:example.org") +mx_send_to_device_encrypted(client, acct, sessions, "org.example.ping", + list(n = 1), devs, store_dir) +} +} From 09314e6bed0c959ee588c06ee2c6847af8244ae8 Mon Sep 17 00:00:00 2001 From: TroyHernandez Date: Fri, 2 Oct 2026 18:20:51 -0500 Subject: [PATCH 2/3] Bump version to 0.2.1.1 --- DESCRIPTION | 4 ++-- NEWS.md | 20 ++++++++++++++++++++ 2 files changed, 22 insertions(+), 2 deletions(-) diff --git a/DESCRIPTION b/DESCRIPTION index 6d2171b..2c3ebba 100644 --- a/DESCRIPTION +++ b/DESCRIPTION @@ -1,8 +1,8 @@ Package: mx.client Type: Package Title: Stateful Matrix Client Helpers -Version: 0.2.1 -Date: 2026-09-11 +Version: 0.2.1.1 +Date: 2026-10-02 Authors@R: c( person("Troy", "Hernandez", role = c("aut", "cre"), email = "troy@cornball.ai", diff --git a/NEWS.md b/NEWS.md index c7291c6..357cb3d 100644 --- a/NEWS.md +++ b/NEWS.md @@ -1,3 +1,23 @@ +# mx.client 0.2.1.1 + +## MatrixRTC calls + +* New: `mx_call_join()`, `mx_call_handle()`, `mx_call_poll()`, + `mx_call_connect()` and `mx_call_leave()` run a MatrixRTC call over + LiveKit the way Element Call and FluffyChat do: a per-device + `org.matrix.msc3401.call.member` state event, a media token from the + LiveKit JWT service, and per-participant media keys exchanged as + Olm-encrypted `io.element.call.encryption_keys` to-device events, with + their key rotation policy. Media goes through the `livekitr` package + (Suggests). `mx_call_members()`, `mx_call_service_url()`, + `mx_call_key_parse()` and `mx_call_key_plan()` expose the pieces. +* New: `mx_send_to_device_encrypted()` and `mx_crypto_encrypt_to_device()` + send an Olm-encrypted to-device event of any type. +* `mx_crypto_process_sync()` returns decrypted Olm to-device events other + than room keys and verification messages in a new `to_device` element. +* Requires mx.api 0.3.1.1 for the OpenID, RTC transport, LiveKit token and + room-state endpoints. + # mx.client 0.2.1 ## Store compatibility and dependencies From f48612b385f9f365d6a9dfde5e10713be0b5e609 Mon Sep 17 00:00:00 2001 From: TroyHernandez Date: Fri, 2 Oct 2026 18:31:58 -0500 Subject: [PATCH 3/3] Poll media before sync; explain a power-level refusal of call membership Found running two devices through a call on a local Tuwunel 1.9.3 with lk-jwt-service 0.7.0 and livekit-server 1.13.7: a room made without a call client's power levels refuses org.matrix.msc3401.call.member at the default level 50, and Tuwunel holds an empty incremental sync for 5 s whatever the timeout, so mx_call_poll() runs the media poll first and documents the hold. --- R/call.R | 51 +++++++++++++++++++++++++++++++++------------ man/mx_call_poll.Rd | 27 +++++++++++++++++------- 2 files changed, 58 insertions(+), 20 deletions(-) diff --git a/R/call.R b/R/call.R index a3a35e3..3cec94d 100644 --- a/R/call.R +++ b/R/call.R @@ -359,9 +359,21 @@ mx_call_send_membership <- function(call, now = Sys.time()) { content <- mx_call_member_content(call$client$user_id, call$client$device_id, call$service_url, call$room_id, call$intent, now) - mx.api::mx_set_state(mx_client_session(call$client), call$room_id, - MX_CALL_MEMBER, content, - mx_call_state_key(call$client$user_id, call$client$device_id)) + tryCatch( + mx.api::mx_set_state(mx_client_session(call$client), call$room_id, + MX_CALL_MEMBER, content, + mx_call_state_key(call$client$user_id, + call$client$device_id)), + mx_error_M_FORBIDDEN = function(e) { + # Rooms default to power level 50 for state events. Element + # and FluffyChat set this event type to 0 when they create a + # room with calls; a room made another way needs the same. + stop(call$client$user_id, " may not send ", MX_CALL_MEMBER, + " in ", call$room_id, ": ", conditionMessage(e), + ". The room's m.room.power_levels needs events[\"", + MX_CALL_MEMBER, "\"] low enough for every participant ", + "(call clients set it to 0).", call. = FALSE) + }) call$membership_sent <- now invisible(NULL) } @@ -535,28 +547,41 @@ mx_call_handle <- function(call, sync, processed = NULL) { list(keys = received, members = members) } -#' Sync once and poll the media of a call +#' Poll the media of a call, then sync once #' #' One iteration of a call loop for programs with no sync loop of their -#' own: syncs with \code{\link{mx_sync_update}}, applies the response with -#' \code{\link{mx_call_handle}}, and polls the LiveKit session with -#' \code{livekitr::lk_poll()}, which is where audio callbacks run. +#' own: polls the LiveKit session with \code{livekitr::lk_poll()}, which +#' is where audio callbacks run, for up to \code{media_timeout} seconds +#' (it returns as soon as media events arrive, so the wait is short while +#' someone speaks), then syncs with \code{\link{mx_sync_update}} and +#' applies the response with \code{\link{mx_call_handle}}. +#' +#' Audio callbacks only run while this process polls, so the sync must +#' not block for long: the default \code{timeout = 0} asks the homeserver +#' to answer at once. Not every homeserver complies: Tuwunel 1.9 holds an +#' incremental sync that has nothing new for 5 s whatever the timeout, +#' and returns at once only when something happened. Frames that arrive +#' meanwhile wait in \pkg{livekitr}'s native queue, so nothing is lost, +#' but they reach the callback late. A program that needs prompt audio +#' on such a server should sync from another process or thread. #' #' @param call An \code{"mx_call"}. -#' @param timeout Seconds to wait for the sync long poll. #' @param media_timeout Seconds to wait in the media poll. +#' @param timeout Seconds to let the sync long poll wait. #' @return The LiveKit events from \code{livekitr::lk_poll()}, or an #' empty list when the call has no media session. #' @export -mx_call_poll <- function(call, timeout = 1, media_timeout = 0.1) { +mx_call_poll <- function(call, media_timeout = 0.5, timeout = 0) { + events <- if (is.null(call$session)) { + list() + } else { + livekitr::lk_poll(call$session, timeout = media_timeout) + } res <- mx_sync_update(call$client, timeout = as.integer(timeout * 1000), save = !is.null(attr(call$client, "path"))) call$client <- res$client mx_call_handle(call, res$sync) - if (is.null(call$session)) { - return(list()) - } - livekitr::lk_poll(call$session, timeout = media_timeout) + events } #' Leave a MatrixRTC call diff --git a/man/mx_call_poll.Rd b/man/mx_call_poll.Rd index 783afa2..fe4b622 100644 --- a/man/mx_call_poll.Rd +++ b/man/mx_call_poll.Rd @@ -1,16 +1,16 @@ % tinyrox says don't edit this manually, but it can't stop you! \name{mx_call_poll} \alias{mx_call_poll} -\title{Sync once and poll the media of a call} +\title{Poll the media of a call, then sync once} \usage{ -mx_call_poll(call, timeout = 1, media_timeout = 0.1) +mx_call_poll(call, media_timeout = 0.5, timeout = 0) } \arguments{ \item{call}{An \code{"mx_call"}.} -\item{timeout}{Seconds to wait for the sync long poll.} - \item{media_timeout}{Seconds to wait in the media poll.} + +\item{timeout}{Seconds to let the sync long poll wait.} } \value{ The LiveKit events from \code{livekitr::lk_poll()}, or an @@ -18,7 +18,20 @@ The LiveKit events from \code{livekitr::lk_poll()}, or an } \description{ One iteration of a call loop for programs with no sync loop of their -own: syncs with \code{\link{mx_sync_update}}, applies the response with -\code{\link{mx_call_handle}}, and polls the LiveKit session with -\code{livekitr::lk_poll()}, which is where audio callbacks run. +own: polls the LiveKit session with \code{livekitr::lk_poll()}, which +is where audio callbacks run, for up to \code{media_timeout} seconds +(it returns as soon as media events arrive, so the wait is short while +someone speaks), then syncs with \code{\link{mx_sync_update}} and +applies the response with \code{\link{mx_call_handle}}. +} +\details{ +Audio callbacks only run while this process polls, so the sync must +not block for long: the default \code{timeout = 0} asks the homeserver +to answer at once. Not every homeserver complies: Tuwunel 1.9 holds an +incremental sync that has nothing new for 5 s whatever the timeout, +and returns at once only when something happened. Frames that arrive +meanwhile wait in \pkg{livekitr}'s native queue, so nothing is lost, +but they reach the callback late. A program that needs prompt audio +on such a server should sync from another process or thread. + }