Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
19 commits
Select commit Hold shift + click to select a range
25b87b3
test: preserve Iroh sessions through application stalls
azooz2003-bit Sep 12, 2026
9d62c24
fix: let Iroh own established connection lifetime
azooz2003-bit Sep 12, 2026
bcbece4
fix: clear cloud panel failures through the workspace API
azooz2003-bit Sep 12, 2026
5f0830f
test: cover coalesced tail of maximum-size mobile frame
azooz2003-bit Sep 12, 2026
73b7806
fix: apply mobile size limits per frame rather than per read
azooz2003-bit Sep 12, 2026
499f8f4
test: derive mobile lane quota expectation
azooz2003-bit Sep 12, 2026
50baaaa
test: include simulator streams in mobile lane quota
azooz2003-bit Sep 13, 2026
ebe4c47
test: keep mobile host connections alive while idle
azooz2003-bit Sep 14, 2026
c0b792e
fix: let native transport own idle session lifetime
azooz2003-bit Sep 14, 2026
a83a9c0
Merge remote-tracking branch 'origin/main' into feat-ios-iroh-connect…
azooz2003-bit Sep 14, 2026
39c35ac
test: distinguish retired transport from recovery presentation
azooz2003-bit Sep 14, 2026
4edfea4
fix: retain recovery presentation while retiring dead transports
azooz2003-bit Sep 14, 2026
e7ca3cb
test: cover live transport foreground recovery
azooz2003-bit Sep 15, 2026
7205608
fix: retain live Iroh sessions during foreground recovery
azooz2003-bit Sep 15, 2026
6f64d4b
test: avoid AppKit keyWindow collision in Cloud drag fixture
austinywang Sep 14, 2026
a18ca14
test: convert Cloud hover fixture coordinates as a point
austinywang Sep 14, 2026
844a43a
fix: disable first-frame timeout for admitted Iroh
azooz2003-bit Sep 15, 2026
6ff2525
Merge origin/main IROH v2 while preserving native liveness authority
azooz2003-bit Sep 15, 2026
20535b6
fix: keep closed control transports bound to their session
azooz2003-bit Sep 15, 2026
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
Original file line number Diff line number Diff line change
Expand Up @@ -201,6 +201,10 @@ public struct MobileSyncFrameCodec {
return frame
}

/// Decodes at most one bounded batch, leaving later frames in `buffer`.
/// Callers drain another batch before awaiting more network bytes whenever
/// the returned count reaches `maximumDecodedFrameCount`. A transport read
/// may coalesce any number of valid frames; its boundary is not a message limit.
public static func decodeFrames(
from buffer: inout Data,
maximumFrameByteCount: Int = defaultMaximumFrameByteCount,
Expand All @@ -222,7 +226,8 @@ public struct MobileSyncFrameCodec {
}
}

while buffer.count - consumedByteCount >= headerByteCount {
while frames.count < maximumDecodedFrameCount,
buffer.count - consumedByteCount >= headerByteCount {
let frameStart = buffer.index(
buffer.startIndex,
offsetBy: consumedByteCount
Expand All @@ -241,11 +246,6 @@ public struct MobileSyncFrameCodec {
guard buffer.count - consumedByteCount >= headerByteCount + payloadLength else {
break
}
guard frames.count < maximumDecodedFrameCount else {
throw MobileSyncFrameCodecError.tooManyFrames(
maximumDecodedFrameCount
)
}
let payloadStart = headerEnd
let payloadEnd = buffer.index(
payloadStart,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -199,39 +199,31 @@ import Testing
}
}

@Test func frameCodecRejectsAZeroLengthFrameFloodAtTheCallerLimit() throws {
let emptyFrame = try MobileSyncFrameCodec.encodeFrame(Data())
@Test func frameCodecDrainsLargeBatchesWithoutTreatingPacketSizeAsFailure() throws {
let frame = try MobileSyncFrameCodec.encodeFrame(Data("valid payload".utf8))
var buffer = Data()
for _ in 0..<10_000 {
buffer.append(emptyFrame)
}

do {
_ = try MobileSyncFrameCodec.decodeFrames(
from: &buffer,
maximumDecodedFrameCount: 16
)
Issue.record("Expected the decoded frame count limit to fail closed")
} catch let error as MobileSyncFrameCodecError {
#expect(error == .tooManyFrames(16))
#expect(buffer.count == (10_000 - 16) * emptyFrame.count)
for _ in 0..<1_000 { buffer.append(frame) }
var decoded = 0
while !buffer.isEmpty {
let frames = try MobileSyncFrameCodec.decodeFrames(from: &buffer, maximumDecodedFrameCount: 16)
#expect(!frames.isEmpty)
#expect(frames.count <= 16)
#expect(frames.allSatisfy { $0 == Data("valid payload".utf8) })
decoded += frames.count
}
#expect(decoded == 1_000)
}

@Test func frameCodecDefaultFrameCountLimitIsFinite() throws {
let emptyFrame = try MobileSyncFrameCodec.encodeFrame(Data())
let frame = try MobileSyncFrameCodec.encodeFrame(Data())
var buffer = Data()
for _ in 0...MobileSyncFrameCodec.defaultMaximumDecodedFrameCount {
buffer.append(emptyFrame)
}

#expect(throws: MobileSyncFrameCodecError.tooManyFrames(
MobileSyncFrameCodec.defaultMaximumDecodedFrameCount
)) {
_ = try MobileSyncFrameCodec.decodeFrames(from: &buffer)
}
for _ in 0...MobileSyncFrameCodec.defaultMaximumDecodedFrameCount { buffer.append(frame) }
let frames = try MobileSyncFrameCodec.decodeFrames(from: &buffer)
#expect(frames.count == MobileSyncFrameCodec.defaultMaximumDecodedFrameCount)
#expect(buffer == frame)
}


private func base64URLEncode(_ data: Data) -> String {
data.base64EncodedString()
.replacingOccurrences(of: "+", with: "-")
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -178,7 +178,7 @@ public actor IrxConnection {
}

/// Returns recent diagnostic evidence. A false result does not establish
/// peer death; callers must probe before replacing an open connection.
/// peer death; only native connection closure establishes transport failure.
public func hasRecentKeepalive(within age: Duration) -> Bool {
guard let lastPongAt else { return false }
return ContinuousClock.now - lastPongAt <= age
Expand Down Expand Up @@ -261,42 +261,50 @@ public actor IrxConnection {
return writer
}

/// Accepts the next bidirectional lane (server side, post-admission).
/// Returns nil once the connection is closed.
/// Accepts the next bidirectional lane. A malformed descriptor retires
/// only that stream; native connection termination ends acceptance.
public func acceptLane() async -> IrxLaneStream? {
while true {
while !Task.isCancelled {
do {
let stream = try await connection.acceptBi()
let reader = IrxStreamReader(stream.recv())
let writer = IrxStreamWriter(stream.send())
guard
let descriptor = try await reader.readControlFrame(
IrxLaneDescriptor.self)
else {
await writer.finish()
continue
do {
if let descriptor = try await reader.readControlFrame(IrxLaneDescriptor.self) {
return IrxLaneStream(descriptor: descriptor, writer: writer, reader: reader)
}
} catch {
// Stream framing failure does not change native connection state.
}
return IrxLaneStream(descriptor: descriptor, writer: writer, reader: reader)
await writer.reset(errorCode: 2)
await reader.stop(errorCode: 2)
} catch {
closedFlag = true
return nil
}
}
return nil
}

/// Accepts the next unidirectional lane (client side: the events lane).
/// Accepts the next usable server event lane without treating a malformed
/// descriptor or stream EOF as complete-connection closure.
public func acceptUniLane() async throws -> (IrxLaneDescriptor, IrxStreamReader)? {
do {
let stream = try await connection.acceptUni()
let reader = IrxStreamReader(stream)
guard
let descriptor = try await reader.readControlFrame(IrxLaneDescriptor.self)
else { return nil }
return (descriptor, reader)
} catch {
closedFlag = true
return nil
while !Task.isCancelled {
do {
let stream = try await connection.acceptUni()
let reader = IrxStreamReader(stream)
do {
if let descriptor = try await reader.readControlFrame(IrxLaneDescriptor.self) {
return (descriptor, reader)
}
} catch {
// A later event stream may repair this optional feature.
}
await reader.stop(errorCode: 2)
} catch {
return nil
}
}
return nil
}

/// The selected QUIC path right now, for relay attribution evidence.
Expand All @@ -308,8 +316,9 @@ public actor IrxConnection {
return "\(selected.isRelay ? "relay" : "direct"):\(selected.remoteAddr)"
}

/// Starts one periodic peer probe owner. A failed probe retires its stream;
/// another attempt uses a fresh stream on the same QUIC connection.
/// Samples application round-trip latency on an optional lane.
/// A failed probe retires its stream; another attempt uses a fresh stream
/// on the same QUIC connection. Native closure owns dead-peer detection.
public func startClientKeepalive(
interval: Duration = IrxProtocol.keepaliveInterval,
deadline: Duration = IrxProtocol.keepaliveDeadline,
Expand All @@ -321,7 +330,7 @@ public actor IrxConnection {
}

/// Pauses application probes before suspension without closing QUIC.
/// A resumed loop begins with zero strikes and a fresh probe deadline.
/// A resumed loop starts with a fresh probe stream and deadline.
public func setApplicationActive(_ active: Bool) {
guard applicationActive != active else { return }
applicationActive = active
Expand All @@ -336,7 +345,7 @@ public actor IrxConnection {
/// Tests the existing peer with a bounded ping/pong exchange. Concurrent
/// callers share the current probe; a shorter caller deadline also retires
/// that probe, so no caller retries a stopped receive stream.
/// A false result during suspension is inconclusive and must not cause a close.
/// A false result is inconclusive and never authorizes connection teardown.
public func probeLiveness(deadline: Duration = IrxProtocol.keepaliveDeadline) async -> Bool {
guard applicationActive, !isClosed, !Task.isCancelled else { return false }
if let task = probeTask, let id = probeID {
Expand Down Expand Up @@ -365,25 +374,19 @@ public actor IrxConnection {
keepaliveGeneration &+= 1
let generation = keepaliveGeneration
keepaliveTask = Task {
var strikes = 0
while !Task.isCancelled, self.applicationActive, self.keepaliveGeneration == generation {
if strikes == 0 {
do { try await Task.sleep(for: settings.interval) } catch { return }
}
do { try await Task.sleep(for: settings.interval) } catch { return }
guard !Task.isCancelled, self.applicationActive, self.keepaliveGeneration == generation else { return }
let alive = await self.probeLiveness(deadline: settings.deadline)
guard !Task.isCancelled, self.applicationActive, self.keepaliveGeneration == generation else { return }
if alive { strikes = 0; continue }
strikes += 1
if strikes < IrxProtocol.keepaliveStrikeLimit {
self.journal.record("keepalive", "miss", ["strike": String(strikes),
"path": self.selectedPathDescription()])
continue
if alive { continue }
self.journal.record("keepalive", "miss", ["path": self.selectedPathDescription()])
// Probe silence retires only its diagnostic stream. The native
// watcher continues to observe closure while fresh probes retry.
if self.isClosed {
await settings.onDeath()
return
}
self.journal.record("keepalive", "timeout", ["path": self.selectedPathDescription()])
await self.close(code: .keepaliveTimeout, origin: .transport)
await settings.onDeath()
return
}
}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,9 @@ public import Foundation
/// session, so the transport retires that exact session before closing the
/// complete QUIC connection. Session ownership belongs to
/// ``IrxPeerEngine``.
///
/// Each transport binds to at most one admitted session. Native closure makes
/// that transport terminal; the RPC owner creates a new transport to recover.
public actor IrxControlByteTransport: CmxByteTransport {
/// Factory for an admitted connection and its control lane.
public typealias Establish = @Sendable () async throws -> (IrxConnection, IrxLaneStream)
Expand Down Expand Up @@ -113,7 +116,15 @@ public actor IrxControlByteTransport: CmxByteTransport {

private func establishedPair() async throws -> (IrxConnection, IrxLaneStream) {
guard !isClosed else { throw IrxConnectionError.closed(nil) }
if let pair, await !pair.0.isConnectionClosed() {
if let pair {
let connectionIsClosed = await pair.0.isConnectionClosed()
guard !isClosed else { throw IrxConnectionError.closed(nil) }
if connectionIsClosed {
// Reads and writes may still be unwinding on this pair. Keep
// their eventual close tied to this RPC generation's session.
await close()
throw IrxConnectionError.closed(nil)
}
return pair
}
if let connectInFlight {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -36,7 +36,7 @@ public struct IrxClientSession: Sendable {
}

/// THE single reconnect owner for one Mac peer (iOS side). Every trigger -
/// app recovery, foreground, network change, keepalive death, explicit retry -
/// app recovery, foreground, network change, native closure, explicit retry -
/// is an input; automatic triggers JOIN the in-flight dial, explicit intent
/// replaces it. Transport failures retry on a capped backoff that resets on
/// success; denials and supersession park the engine until an explicit
Expand Down Expand Up @@ -295,9 +295,9 @@ public actor IrxPeerEngine {
if hasConnectionIntent { foregroundKick() }
}

/// Retains an open session only after a fresh bounded liveness check. Each
/// unsuccessful attempt uses a live stream. Known native closure bypasses
/// probing; age of the last pong never forces replacement.
/// Samples the existing session on a fresh probe deadline after suspension.
/// An unanswered probe leaves the native connection in place; application
/// owners repair their streams without replacing the admitted session.
public func foregroundKick() {
guard applicationActive, foregroundTask == nil, hasConnectionIntent, parkedCode == nil else { return }
let generation = activityGeneration
Expand All @@ -310,19 +310,26 @@ public actor IrxPeerEngine {
private func recoverOnForeground(generation: UInt64) async {
let started = clockNow()
if let previous = session {
for _ in 0..<IrxProtocol.keepaliveStrikeLimit {
guard foregroundIsCurrent(generation), session?.admit.session == previous.admit.session else { return }
if await previous.connection.isConnectionClosed() { break }
let alive = await previous.connection.probeLiveness(deadline: config.foregroundProbeDeadline)
guard foregroundIsCurrent(generation), session?.admit.session == previous.admit.session else { return }
if alive {
record("foreground-session-retained", ["session": previous.admit.session,
"elapsed": String(describing: started.duration(to: clockNow()))])
return
}
guard foregroundIsCurrent(generation) else { return }
let probeAnswered: Bool
if await previous.connection.isConnectionClosed() {
probeAnswered = false
} else {
probeAnswered = await previous.connection.probeLiveness(deadline: config.foregroundProbeDeadline)
}
guard foregroundIsCurrent(generation), session?.admit.session == previous.admit.session else { return }
_ = try? await ensureSession(explicit: true, trigger: "foreground-probe-failed")
let nativeClosed = await previous.connection.isConnectionClosed()
guard foregroundIsCurrent(generation), session?.admit.session == previous.admit.session else { return }
if !nativeClosed {
record("foreground-session-retained", ["session": previous.admit.session,
"probe_answered": String(probeAnswered),
"elapsed": String(describing: started.duration(to: clockNow()))])
return
}
// Use the same reason classification as the native watcher. A
// terminal peer close must stay parked even if foreground wins
// the race to observe it.
await sessionDied(previous, viaKeepalive: false)
} else {
guard foregroundIsCurrent(generation) else { return }
// Foreground can bring a local transport retry forward, while a
Expand All @@ -346,8 +353,8 @@ public actor IrxPeerEngine {
}

/// Event-driven relay race: fresh discovery just revealed a different
/// home relay for this peer. An admitted session passing keepalives is
/// never touched. An in-flight dial (aimed at the stale relay, where it
/// home relay for this peer. An admitted session is never touched. An
/// in-flight dial (aimed at the stale relay, where it
/// would sit out a silent black-hole timeout) is cancelled and replaced
/// immediately; a pending backoff redial is pulled forward. Parked
/// denials stay parked: authorization state is not a routing question.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -11,18 +11,12 @@ public enum IrxProtocol {
/// Control frames are small (hello/admit/keepalive/descriptors); anything
/// larger is a protocol error, never buffered.
public static let maximumControlFrameByteCount = 256 * 1024
/// Keepalive cadence: one tiny ping per interval, pong deadline after
/// which the connection is declared dead. Hard closes (the realistic
/// relay-expiry case) are detected instantly by the termination watcher;
/// the ping loop bounds SILENT path blackholes to interval + deadline,
/// keeping worst-case detection-plus-redial inside single-digit seconds.
/// Application latency sampling cadence. Connection lifetime is owned by
/// Iroh's native keepalives and negotiated connection idle timeout.
public static let keepaliveInterval: Duration = .seconds(5)
/// A missed application pong retires only the diagnostic stream.
public static let keepaliveDeadline: Duration = .seconds(2)
/// Consecutive pong misses before the connection is declared dead. One
/// transient stall (relay hiccup, brief peer pause) must never sever a
/// healthy session; a re-ping fires immediately after a miss, so real
/// death still detects in ~strikeLimit x deadline.
public static let keepaliveStrikeLimit = 2

}

/// Machine-readable close/denial codes. The code travels in the QUIC
Expand Down
Loading
Loading