Skip to content
Draft
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
14 changes: 13 additions & 1 deletion .github/workflows/ios-production-source.yml
Original file line number Diff line number Diff line change
Expand Up @@ -3,7 +3,7 @@ name: iOS production dependency source
on:
pull_request:
push:
branches: ['codex/ios-production-dependency-source', 'feature/neon-fix-fearless']
branches: ['codex/ios-production-dependency-source', 'codex/ios-authorized-rpc-transport', 'feature/neon-fix-fearless']

permissions:
contents: read
Expand All @@ -22,3 +22,15 @@ jobs:
run: python3 scripts/test-scrypt-architectures.py
- name: Reject tracked source mutations during validation
run: git diff --exit-code

authorized-rpc:
runs-on: macos-15
timeout-minutes: 20
steps:
- uses: actions/checkout@11d5960a326750d5838078e36cf38b85af677262 # v4
with:
persist-credentials: false
- name: Test production RPC sources against the pinned transport
run: python3 scripts/test-authorized-rpc.py
- name: Reject tracked source mutations during validation
run: git diff --exit-code
2 changes: 1 addition & 1 deletion Package.swift
Original file line number Diff line number Diff line change
Expand Up @@ -43,7 +43,7 @@ let package = Package(
.package(url: "https://github.com/Boilertalk/secp256k1.swift.git", from: "0.1.7"),
.package(url: "https://github.com/bitmark-inc/tweetnacl-swiftwrap", from: "1.1.0"),
.package(url: "https://github.com/ashleymills/Reachability.swift", from: "5.0.0"),
.package(url: "https://github.com/soramitsu/fearless-starscream", from: "4.0.12"),
.package(url: "https://github.com/soramitsu/fearless-starscream", .revision("b6ef58590241babdb4fe52e916a02c9e2b749e3d")),
.package(url: "https://github.com/google/GoogleSignIn-iOS", from: "7.0.0"),
.package(url: "https://github.com/google/google-api-objectivec-client-for-rest.git", from: "3.3.0"),
.package(url: "https://github.com/attaswift/BigInt.git", from: "5.3.0"),
Expand Down
24 changes: 24 additions & 0 deletions Sources/SSFUtils/SSFUtils/Classes/Network/JSONRPCEngine.swift
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,8 @@ public enum JSONRPCEngineError: Error {
case clientCancelled
case unknownError
case timeout
case requestNotSent
case submissionOutcomeUnknown
}

public protocol JSONRPCResponseHandling {
Expand Down Expand Up @@ -42,11 +44,26 @@ struct JSONRPCResponseHandler<T: Decodable>: JSONRPCResponseHandling {
}
}

/// Application-owned final authorization. Called after SDK framing and blocking
/// writer lock waits. Synchronously invoke handoff once while holding fresh
/// authority; never reenter this engine, await, or retain the nonescaping action.
public protocol JSONRPCWriteAuthorizing: AnyObject {
func authorize(_ handoff: () throws -> Void) throws
}

public struct JSONRPCOptions {
public let resendOnReconnect: Bool
public let writeAuthorization: JSONRPCWriteAuthorizing?

public init(resendOnReconnect: Bool = true) {
self.resendOnReconnect = resendOnReconnect
writeAuthorization = nil
}

/// Guarded mutations can never opt into reconnection replay.
public init(writeAuthorization: JSONRPCWriteAuthorizing) {
resendOnReconnect = false
self.writeAuthorization = writeAuthorization
}
}

Expand Down Expand Up @@ -117,6 +134,11 @@ public protocol JSONRPCEngine: AnyObject {

func cancelForIdentifier(_ identifier: UInt16)

/// Cancel only this authorization's request, even if a UInt16 ID has been
/// reused. Engines without identity-aware removal still cannot send after
/// the operation-owned final authorizer has been cancelled.
func cancelForIdentifier(_ identifier: UInt16, writeAuthorization: JSONRPCWriteAuthorizing)

func generateRequestId() -> UInt16
func addSubscription(_ subscription: JSONRPCSubscribing)
func reconnect(url: URL)
Expand All @@ -127,6 +149,8 @@ public protocol JSONRPCEngine: AnyObject {
}

public extension JSONRPCEngine {
func cancelForIdentifier(_ identifier: UInt16, writeAuthorization: JSONRPCWriteAuthorizing) {}

func callMethod<P: Codable, T: Decodable>(
_ method: String,
params: P?,
Expand Down
137 changes: 103 additions & 34 deletions Sources/SSFUtils/SSFUtils/Classes/Network/JSONRPCOperation.swift
Original file line number Diff line number Diff line change
Expand Up @@ -5,23 +5,78 @@ enum JSONRPCOperationError: Error {
case timeout
}

/// Cancellation is independent of request-ID publication. The final check is
/// inside the application's fresh-authority scope and never waits for a lock.
private final class OperationWriteAuthorization: JSONRPCWriteAuthorizing {
private let lock = NSLock()
private let applicationAuthorization: JSONRPCWriteAuthorizing
private var cancelled = false
private var handoffAttempted = false

init(_ applicationAuthorization: JSONRPCWriteAuthorizing) {
self.applicationAuthorization = applicationAuthorization
}

func authorize(_ handoff: () throws -> Void) throws {
lock.lock()
let denied = cancelled || handoffAttempted
lock.unlock()
guard !denied else { throw JSONRPCEngineError.requestNotSent }
try applicationAuthorization.authorize {
guard lock.try() else { throw JSONRPCEngineError.requestNotSent }
defer { lock.unlock() }
guard !cancelled, !handoffAttempted else { throw JSONRPCEngineError.requestNotSent }
// Once entered, an error can no longer prove non-submission. Keep
// the uncertainty even when cancellation precedes ID publication.
handoffAttempted = true
try handoff()
}
}

func cancel() -> JSONRPCEngineError {
lock.lock()
defer { lock.unlock() }
cancelled = true
return handoffAttempted ? .submissionOutcomeUnknown : .requestNotSent
}
}

public class JSONRPCOperation<P: Codable, T: Decodable>: BaseOperation<T> {
public let engine: JSONRPCEngine
private(set) var requestId: UInt16?
private let requestLock = NSLock()
private var currentRequestId: UInt16?
private(set) var requestId: UInt16? {
get { requestLock.lock(); defer { requestLock.unlock() }; return currentRequestId }
set { requestLock.lock(); currentRequestId = newValue; requestLock.unlock() }
}
public let requestOptions: JSONRPCOptions
private let writeAuthorization: OperationWriteAuthorization?
private let completionSignal = DispatchSemaphore(value: 0)
private let resultLock = NSLock()
private var storedResult: Result<T, Error>?
override public var result: Result<T, Error>? {
get { resultLock.lock(); defer { resultLock.unlock() }; return storedResult }
set { resultLock.lock(); storedResult = newValue; resultLock.unlock() }
}
public let method: String
public var parameters: P?
public let timeout: Int

public init(engine: JSONRPCEngine, method: String, parameters: P? = nil, timeout: Int = 10) {
public init(engine: JSONRPCEngine, method: String, parameters: P? = nil, timeout: Int = 10,
requestOptions: JSONRPCOptions = JSONRPCOptions()) {
self.engine = engine
self.method = method
self.parameters = parameters
self.timeout = timeout
let authorization = requestOptions.writeAuthorization.map(OperationWriteAuthorization.init)
writeAuthorization = authorization
self.requestOptions = authorization.map { JSONRPCOptions(writeAuthorization: $0) } ?? requestOptions

super.init()
}

override public func main() {
defer { requestId = nil }
super.main()

if isCancelled {
Expand All @@ -33,55 +88,69 @@ public class JSONRPCOperation<P: Codable, T: Decodable>: BaseOperation<T> {
}

do {
let semaphore = DispatchSemaphore(value: 0)

var optionalCallResult: Result<T, Error>?

requestId = try engine.callMethod(method, params: parameters) { (result: Result<
T,
Error
>) in
optionalCallResult = result

semaphore.signal()
requestId = try engine.callMethod(method, params: parameters, options: requestOptions) { [weak self] (result: Result<T, Error>) in
guard let self = self else { return }
if self.writeAuthorization == nil, self.isCancelled {
self.completionSignal.signal()
return
}
if case .failure(let error) = result, error as? JSONRPCEngineError == .clientCancelled {
if let authorization = self.writeAuthorization {
self.finish(.failure(authorization.cancel()))
}
} else {
self.finish(result)
}
self.completionSignal.signal()
}

let status = semaphore.wait(timeout: .now() + .seconds(timeout))

if status == .timedOut {
result = .failure(JSONRPCOperationError.timeout)
// Cancellation may race with the synchronous request enqueue.
// Publish the ID first, then cancel again if that race occurred.
if isCancelled, let identifier = requestId {
cancelRequest(identifier)
return
}

guard let callResult = optionalCallResult else {
return
}
let status = completionSignal.wait(timeout: .now() + .seconds(timeout))

if case let .failure(error) = callResult,
let jsonRPCEngineError = error as? JSONRPCEngineError,
jsonRPCEngineError == .clientCancelled
{
if status == .timedOut {
if let authorization = writeAuthorization {
finish(.failure(authorization.cancel()))
if let identifier = requestId { cancelRequest(identifier) }
} else {
finish(.failure(JSONRPCOperationError.timeout))
}
return
}

switch callResult {
case let .success(response):
result = .success(response)
case let .failure(error):
result = .failure(error)
}

} catch {
result = .failure(error)
finish(.failure(error))
}
}

override public func cancel() {
if let authorization = writeAuthorization {
finish(.failure(authorization.cancel()))
completionSignal.signal()
}
super.cancel()
if let requestId = requestId {
engine.cancelForIdentifier(requestId)
cancelRequest(requestId)
}
}

super.cancel()
private func cancelRequest(_ identifier: UInt16) {
if let authorization = writeAuthorization {
engine.cancelForIdentifier(identifier, writeAuthorization: authorization)
} else {
engine.cancelForIdentifier(identifier)
}
}

private func finish(_ value: Result<T, Error>) {
resultLock.lock()
if storedResult == nil { storedResult = value }
resultLock.unlock()
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -2,8 +2,12 @@ import Foundation
import Starscream

extension WebSocketEngine: WebSocketDelegate {
public func didReceive(event: WebSocketEvent, client _: WebSocketClient) {
public func didReceive(event: WebSocketEvent, client: WebSocketClient) {
mutex.lock()
defer { mutex.unlock() }
// A late callback from a replaced endpoint cannot settle or cancel a
// request on the current connection, even if a UInt16 ID is reused.
guard client === connection else { return }

switch event {
case let .binary(data):
Expand All @@ -26,7 +30,6 @@ extension WebSocketEngine: WebSocketDelegate {
logger?.warning("Unhandled event \(event)")
}

mutex.unlock()
}

private func handleTimeout() {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -59,7 +59,7 @@ extension WebSocketEngine: JSONRPCEngine {
failureClosure: failureClosure
)

addSubscription(subscription)
guard addSubscriptionLocked(subscription) else { throw JSONRPCEngineError.unknownError }

updateConnectionForRequest(request)

Expand All @@ -74,17 +74,49 @@ extension WebSocketEngine: JSONRPCEngine {
mutex.unlock()
}

public func reconnect(url: URL) {
self.connection.delegate = nil
public func cancelForIdentifier(_ identifier: UInt16, writeAuthorization: JSONRPCWriteAuthorizing) {
mutex.lock()
defer { mutex.unlock() }
let request = pendingRequests.first { $0.requestId == identifier } ?? inProgressRequests[identifier]
if request?.options.writeAuthorization === writeAuthorization {
cancelRequestForLocalId(identifier)
} else if request == nil,
subscriptions[identifier]?.requestOptions.writeAuthorization === writeAuthorization {
processSubscriptionError(identifier, error: JSONRPCEngineError.submissionOutcomeUnknown, shouldUnsubscribe: true)
}
}

public func reconnect(url: URL) {
mutex.lock()
let shouldResume: Bool
switch state {
case .notConnected: shouldResume = false
case .connecting, .connected, .waitingReconnection, .notReachable: shouldResume = true
}
cancelPendingGuardedRequests()
let cancelled = resetInProgress()
notify(requests: cancelled, error: JSONRPCEngineError.remoteCancelled)
let previous = connection
previous.delegate = nil
reconnectionScheduler.cancel()
pingScheduler.cancel()
self.url = url
let request = URLRequest(url: url, timeoutInterval: 10)
let engine = self.connection.engine

let connection = WebSocket(request: request, engine: engine)
self.connection = connection

connection.callbackQueue = Self.sharedProcessingQueue
connection.delegate = self
// Reusing an already-connected engine would send a new URL's requests
// to the previous socket. A new endpoint requires a new handshake.
let next = replacementConnectionFactory?(request) ?? WebSocket(
request: request, engine: WSEngine(transport: TCPTransport(), certPinner: FoundationSecurity())
)
next.callbackQueue = completionQueue
next.delegate = self
connection = next
changeState(.notConnected)
mutex.unlock()
// No library writer wait occurs while holding the RPC request mutex.
previous.forceDisconnect()
// The previous retry timer belongs to the retired endpoint and was
// cancelled above. Resume explicitly so pending reads/subscriptions do
// not depend on an unrelated future request to connect the new URL.
if shouldResume { connectIfNeeded() }
}
}
Loading
Loading