Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions .changes/rust-net-transport
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
patch type="changed" "Register the SDK's own WebSocket/URLSession stack as livekit-net's transport, so Rust-side signalling and HTTP reuse the tuned session (call-signaling QoS, multipath handover, TLS workarounds) instead of a second parallel client; a rejected WebSocket upgrade now keeps its HTTP status"
2 changes: 1 addition & 1 deletion Package.swift
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,7 @@ let package = Package(
dependencies: [
// LK-Prefixed Dynamic WebRTC XCFramework
.package(url: "https://github.com/livekit/webrtc-xcframework.git", exact: "150.7871.02"),
.package(url: "https://github.com/livekit/livekit-uniffi-xcframework.git", exact: "0.1.9"),
.package(url: "https://github.com/livekit/livekit-uniffi-xcframework.git", exact: "0.1.12"),
// Test-only: conformance oracle for the nanopb facades.
.package(url: "https://github.com/apple/swift-protobuf.git", from: "1.31.0"),
// Only used for DocC generation
Expand Down
2 changes: 1 addition & 1 deletion Package@swift-6.2.swift
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,7 @@ let package = Package(
dependencies: [
// LK-Prefixed Dynamic WebRTC XCFramework
.package(url: "https://github.com/livekit/webrtc-xcframework.git", exact: "150.7871.02"),
.package(url: "https://github.com/livekit/livekit-uniffi-xcframework.git", exact: "0.1.9"),
.package(url: "https://github.com/livekit/livekit-uniffi-xcframework.git", exact: "0.1.12"),
// Test-only: conformance oracle for the nanopb facades.
.package(url: "https://github.com/apple/swift-protobuf.git", from: "1.31.0"),
// Only used for DocC generation
Expand Down
5 changes: 5 additions & 0 deletions Sources/LiveKit/Core/RoomDependencies.swift
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,11 @@ final class ConnectionDependencies: Sendable {
let e2ee: StateSync<E2EEManager?>

init(room: Room, roomOptions: RoomOptions) {
// Process-wide rather than connection-scoped — livekit-net keeps the first
// registration — but this is the earliest construction on the connect path, so
// no signalling can outrun it.
_ = rustTransport

dataTracks = DataTracks(room: room)
let manager: E2EEManager? = if let e2eeOptions = roomOptions.e2eeOptions {
E2EEManager(e2eeOptions: e2eeOptions)
Expand Down
4 changes: 2 additions & 2 deletions Sources/LiveKit/Core/SignalClient.swift
Original file line number Diff line number Diff line change
Expand Up @@ -189,8 +189,8 @@ actor SignalClient: Loggable {

do {
let socket = try await WebSocket(url: url,
token: token,
connectOptions: connectOptions)
headers: ["Authorization": "Bearer \(token)"],
timeoutInterval: connectOptions?.socketConnectTimeoutInterval ?? .defaultSocketConnect)
connectSpan?.record("ws_open")

startDataTrackResponses()
Expand Down
35 changes: 26 additions & 9 deletions Sources/LiveKit/Support/Network/HTTP.swift
Original file line number Diff line number Diff line change
Expand Up @@ -23,19 +23,36 @@ class HTTP: NSObject {
delegate: nil,
delegateQueue: operationQueue)

static func requestValidation(from url: URL, token: String) async throws {
var request = URLRequest(url: url,
cachePolicy: .reloadIgnoringLocalAndRemoteCacheData,
timeoutInterval: .defaultHTTPConnect)
// Attach token to header
request.addValue("Bearer \(token)", forHTTPHeaderField: "Authorization")

// Make the data request
/// Perform one request on the shared session.
///
/// Both parameters are applied to `request`, overriding whatever it carries, so the
/// transport owns them rather than each call site.
///
/// - Parameter cachePolicy: Defaults to bypassing the cache: the shared session uses
/// `URLCache.shared`, and a cached validation or region response is a stale one.
/// - Parameter timeoutInterval: Defaults to ``TimeInterval/defaultHTTPConnect``.
static func request(_ request: URLRequest,
cachePolicy: URLRequest.CachePolicy = .reloadIgnoringLocalAndRemoteCacheData,
timeoutInterval: TimeInterval = .defaultHTTPConnect) async throws -> (Data, HTTPURLResponse)
{
var request = request
request.cachePolicy = cachePolicy
request.timeoutInterval = timeoutInterval
let (data, response) = try await session.data(for: request)

guard let httpResponse = response as? HTTPURLResponse else {
throw URLError(.badServerResponse)
}
return (data, httpResponse)
}

static func requestValidation(from url: URL, token: String) async throws {
var request = URLRequest(url: url)
// Attach token to header
request.addValue("Bearer \(token)", forHTTPHeaderField: "Authorization")

let (data, httpResponse) = try await Self.request(request,
cachePolicy: .reloadIgnoringLocalAndRemoteCacheData,
timeoutInterval: .defaultHTTPConnect)

guard (200 ..< 300).contains(httpResponse.statusCode) else {
let statusCode = httpResponse.statusCode
Expand Down
183 changes: 183 additions & 0 deletions Sources/LiveKit/Support/Network/RustTransport.swift
Original file line number Diff line number Diff line change
@@ -0,0 +1,183 @@
/*
* Copyright 2026 LiveKit
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/

internal import LiveKitUniFFI
import Foundation

/// Registers the SDK's own network stack as `livekit-net`'s transport, so anything the
/// Rust side sends goes out over the same ``WebSocket``/``HTTP`` session the Swift path
/// uses — one place tuning `networkServiceType`, multipath, TLS workarounds and logging.
///
/// Lazily initialized on first access, like ``sharedLogger``: `livekit-net` keeps its
/// clients in a `OnceLock`, so this runs once and later accesses are free.
let rustTransport: Void = {
setWsClient(c: WsClientAdapter())
setHttpClient(c: HttpClientAdapter())
}()

// MARK: - WebSocket

final class WsClientAdapter: WsClient {
func connect(url: String, headers: [Header], timeoutMs: UInt64) async throws -> WsConnectResult {
guard let url = URL(string: url) else {
throw TransportError.Other("invalid url: \(url)")
}
do {
let socket = try await WebSocket(url: url,
headers: headers.asDictionary,
timeoutInterval: TimeInterval(timeoutMs) / 1000)
return WsConnectResult(connection: WsConnectionAdapter(socket))
} catch {
throw error.asTransportError
}
}
}

private final class WsConnectionAdapter: WsConnection {
private let socket: WebSocket

init(_ socket: WebSocket) {
self.socket = socket
}

func send(frame: Data) async throws {
do {
try await socket.send(data: frame)
} catch {
throw error.asTransportError
}
}

func recv() async throws -> Data? {
do {
// Signalling is binary-only, so `for case let` skips any other frame kind and
// keeps reading. Iterating afresh per call is the same stream: the iterator
// buffers nothing, it reads straight off the task.
for try await case let .data(frame) in socket {
return frame
}
return nil
} catch {
let transportError = error.asTransportError
// livekit-net treats an abrupt teardown as end-of-stream, not an error:
// a reset without a close handshake, or TLS closed without close_notify.
if case .Closed = transportError { return nil }
throw transportError
}
}

func close() async {
socket.close()
}
}

// MARK: - HTTP

final class HttpClientAdapter: HttpClient {
/// Translate one `livekit-net` HTTP call into a `URLRequest`. Cache policy and timeout
/// belong to ``HTTP/request(_:cachePolicy:timeoutInterval:)``, not here.
static func makeRequest(method: HttpMethod, url: URL, headers: [Header], body: Data?) -> URLRequest {
var request = URLRequest(url: url)
request.httpMethod = switch method {
case .get: "GET"
case .post: "POST"
Comment thread
pblazej marked this conversation as resolved.
@unknown default: "GET"

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 Unknown HTTP verbs become GET requests

When livekit-net adds another HttpMethod, makeRequest silently sends it as GET. The server receives a different operation instead of an unsupported-method failure.

Learn more

HttpMethod is non-frozen, so a newer LiveKitUniFFI binary can introduce a case that this SDK did not compile against. The adapter cannot preserve an unknown verb, but substituting GET creates a valid request with different semantics. The throwable request boundary can reject such a case without crashing consumer code.

Example: A future .put request enters the unknown branch and reaches the server as GET. The intended update does not occur, and the caller receives the GET response rather than an unsupported-method error.

Recommended fix: Move the method conversion into a throwing helper or switch inside request. Throw TransportError.Other for unknown cases, then add coverage for the rejection path when the generated API permits constructing an unknown case.

Devin Review


Was this helpful? React with 👍 or 👎 to provide feedback.

}
request.httpBody = body
for header in headers {
request.addValue(header.value, forHTTPHeaderField: header.name)
}
return request
}

func request(method: HttpMethod, url: String, headers: [Header], body: Data?) async throws -> HttpResponse {
guard let url = URL(string: url) else {
throw TransportError.Other("invalid url: \(url)")
}
do {
// A 4xx/5xx is a response, not a transport failure — only the transport throws.
let request = Self.makeRequest(method: method, url: url, headers: headers, body: body)
let (data, response) = try await HTTP.request(request)
return HttpResponse(status: UInt16(clamping: response.statusCode),
headers: response.headers,
body: data)
} catch {
throw error.asTransportError
}
}
}

// MARK: - Conversions

private extension [Header] {
var asDictionary: [String: String] {
reduce(into: [:]) { $0[$1.name] = $1.value }
}
}

private extension HTTPURLResponse {
/// - Note: `livekit-net` documents receipt order, which `HTTPURLResponse` does not keep.
var headers: [Header] {
allHeaderFields.compactMap { name, value in
guard let name = name as? String else { return nil }
return Header(name: name, value: String(describing: value))
}
}
}

extension Error {
/// Map onto `livekit-net`'s taxonomy, which the Rust signal client branches on:
/// `Http` is a rejected upgrade (fail fast), `Timeout`/`Closed`/`Connection` all
/// drive a reconnect.
var asTransportError: TransportError {
if let transportError = self as? TransportError { return transportError }

let underlying = (self as? LiveKitError)?.internalError ?? self

if let upgrade = underlying as? WebSocketUpgradeFailure {
return .Http(status: UInt16(clamping: upgrade.statusCode))
}
if let urlError = underlying as? URLError {
switch urlError.code {
case .timedOut:
return .Timeout
// Our own `close()` surfaces as `.cancelled`; a peer reset arrives as
// `.networkConnectionLost`. Both are end-of-stream to livekit-net.
case .cancelled, .networkConnectionLost:
return .Closed
default:
// The numeric code, not `localizedDescription`: this string is what Rust
// logs and what gets grepped and grouped, and it must not change with the
// device's language.
return .Connection("URLError \(urlError.errorCode)")
}
}
let nsError = underlying as NSError
if nsError.domain == NSPOSIXErrorDomain,
nsError.code == Int(ECONNRESET) || nsError.code == Int(ENOTCONN)
{
return .Closed
}
if let type = (self as? LiveKitError)?.type {
switch type {
case .timedOut: return .Timeout
case .cancelled: return .Closed
default: break
}
}
return .Connection(String(describing: underlying))
}
}
30 changes: 23 additions & 7 deletions Sources/LiveKit/Support/Network/WebSocket.swift
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,12 @@
import Foundation
import Network

/// The server answered the upgrade request with an HTTP response instead of switching
/// protocols. Carried as the `internalError` of the thrown ``LiveKitError``.
struct WebSocketUpgradeFailure: Error, Sendable {
let statusCode: Int
}

actor WebSocket: Loggable, AsyncSequence {
typealias Element = URLSessionWebSocketTask.Message

Expand All @@ -28,7 +34,7 @@
let config = URLSessionConfiguration.default
config.timeoutIntervalForRequest = 60
config.timeoutIntervalForResource = 604_800
config.shouldUseExtendedBackgroundIdleMode = true

Check warning on line 37 in Sources/LiveKit/Support/Network/WebSocket.swift

View workflow job for this annotation

GitHub Actions / Build & Test (xcode-27, latest, visionOS Simulator,name=Apple Vision Pro,OS=27.0)

'shouldUseExtendedBackgroundIdleMode' was deprecated in visionOS 2.4: Not supported [#DeprecatedDeclaration]
config.networkServiceType = .callSignaling
#if os(iOS) || os(visionOS)
// https://developer.apple.com/documentation/foundation/urlsessionconfiguration/improving_network_reliability_using_multipath_tcp
Expand All @@ -37,11 +43,13 @@
return config
}

init(url: URL, token: String, connectOptions: ConnectOptions?) async throws {
init(url: URL, headers: [String: String], timeoutInterval: TimeInterval) async throws {
var request = URLRequest(url: url,
cachePolicy: .useProtocolCachePolicy,
timeoutInterval: connectOptions?.socketConnectTimeoutInterval ?? .defaultSocketConnect)
request.addValue("Bearer \(token)", forHTTPHeaderField: "Authorization")
timeoutInterval: timeoutInterval)
for (name, value) in headers {
request.addValue(value, forHTTPHeaderField: name)
}

#if targetEnvironment(simulator)
if #available(iOS 26.0, *) {
Expand Down Expand Up @@ -136,13 +144,21 @@
}
}

func urlSession(_: URLSession, task _: URLSessionTask, didCompleteWithError error: Error?) {
func urlSession(_: URLSession, task: URLSessionTask, didCompleteWithError error: Error?) {
log("didCompleteWithError: \(String(describing: error))", error != nil ? .error : .debug)

// A rejected upgrade (401 on a bad token, 404 on a wrong path) arrives as a
// plain URLError with the HTTP response still attached. Keep the status: it is
// the difference between failing fast and reconnecting forever.
let upgradeStatus = (task.response as? HTTPURLResponse)
.map(\.statusCode)
.flatMap { $0 == 101 ? nil : WebSocketUpgradeFailure(statusCode: $0) }

_continuation.mutate {
if let error {
let lkError = LiveKitError.from(error: error) ?? LiveKitError(.unknown)
$0?.resume(throwing: lkError)
if let upgradeStatus {
$0?.resume(throwing: LiveKitError(.network, internalError: upgradeStatus))
} else if let error {
$0?.resume(throwing: LiveKitError.from(error: error) ?? LiveKitError(.unknown))
} else {
$0?.resume()
}
Expand Down
Loading
Loading