Merge nucleic/olive-dewy-urchin-megz into dev
This commit is contained in:
@@ -154,6 +154,15 @@ final class HostConnection {
|
|||||||
private var reconnectAttempts = 0
|
private var reconnectAttempts = 0
|
||||||
private var seenSeq: Set<UInt64> = []
|
private var seenSeq: Set<UInt64> = []
|
||||||
|
|
||||||
|
// Covalence direct (SYNC §3.3): when the current attempt is over the relay, its channel is
|
||||||
|
// wrapped in a `MultipathFrameChannel` so a background hole punch can migrate the session
|
||||||
|
// off the relay. `directBox` routes the host's `directAnswer`/`directGo`/`directDecline`
|
||||||
|
// replies to the driver; `directTask` is the background upgrade attempt.
|
||||||
|
private var directMux: MultipathFrameChannel?
|
||||||
|
private var directBox: DirectReplyBox?
|
||||||
|
private var directTask: Task<Void, Never>?
|
||||||
|
private var directEnabled: Bool { DirectPathSetting.enabled() }
|
||||||
|
|
||||||
/// The relay's short-lived connection token, cached for reuse across reconnects. The token is a
|
/// The relay's short-lived connection token, cached for reuse across reconnects. The token is a
|
||||||
/// stateless HMAC bearer that isn't enforced single-use (`handleRelay` re-verifies it every
|
/// stateless HMAC bearer that isn't enforced single-use (`handleRelay` re-verifies it every
|
||||||
/// upgrade with no per-token state), so within its TTL a reconnect reuses it and skips the
|
/// upgrade with no per-token state), so within its TTL a reconnect reuses it and skips the
|
||||||
@@ -167,7 +176,7 @@ final class HostConnection {
|
|||||||
private enum TransportAttempt {
|
private enum TransportAttempt {
|
||||||
case lan(NWEndpoint)
|
case lan(NWEndpoint)
|
||||||
case tailnet(host: String, port: UInt16)
|
case tailnet(host: String, port: UInt16)
|
||||||
/// Nucleic Private Relay (mesh P2) — always the last candidate: works from anywhere,
|
/// Covalence Relay (mesh P2) — always the last candidate: works from anywhere,
|
||||||
/// but a direct path beats a brokered one when both exist.
|
/// but a direct path beats a brokered one when both exist.
|
||||||
case relay(base: URL, membershipToken: String)
|
case relay(base: URL, membershipToken: String)
|
||||||
}
|
}
|
||||||
@@ -381,16 +390,27 @@ final class HostConnection {
|
|||||||
connectTask = Task { [weak self] in
|
connectTask = Task { [weak self] in
|
||||||
guard let self else { return }
|
guard let self else { return }
|
||||||
do {
|
do {
|
||||||
let channel = try await self.dialRelay(
|
let relayChannel = try await self.dialRelay(
|
||||||
base: base, membershipToken: membershipToken)
|
base: base, membershipToken: membershipToken)
|
||||||
guard !Task.isCancelled else { channel.close(); return }
|
guard !Task.isCancelled else { relayChannel.close(); return }
|
||||||
|
// Wrap the relay leg in the multipath mux (passthrough until armed at
|
||||||
|
// `.ready`) so the session can migrate to a hole-punched direct path.
|
||||||
|
let channel: any FrameChannel
|
||||||
|
if self.directEnabled {
|
||||||
|
let mux = MultipathFrameChannel(primary: relayChannel)
|
||||||
|
self.directMux = mux
|
||||||
|
channel = mux
|
||||||
|
} else {
|
||||||
|
self.directMux = nil
|
||||||
|
channel = relayChannel
|
||||||
|
}
|
||||||
// A room with no live host swallows frames silently until presence says
|
// A room with no live host swallows frames silently until presence says
|
||||||
// otherwise — bound the handshake like the LAN path does.
|
// otherwise — bound the handshake like the LAN path does.
|
||||||
self.lanConnectTimeout?.cancel()
|
self.lanConnectTimeout?.cancel()
|
||||||
self.lanConnectTimeout = Task { [weak self, weak channel] in
|
self.lanConnectTimeout = Task { [weak self, weak relayChannel] in
|
||||||
try? await Task.sleep(for: .seconds(10))
|
try? await Task.sleep(for: .seconds(10))
|
||||||
guard !Task.isCancelled, let self, !self.connectivity.isLive else { return }
|
guard !Task.isCancelled, let self, !self.connectivity.isLive else { return }
|
||||||
channel?.close()
|
relayChannel?.close()
|
||||||
}
|
}
|
||||||
self.startClient(
|
self.startClient(
|
||||||
channel: channel, hostStaticKey: plan.hostStaticKey,
|
channel: channel, hostStaticKey: plan.hostStaticKey,
|
||||||
@@ -466,7 +486,9 @@ final class HostConnection {
|
|||||||
canProvideIntelligence: PhoneIntelligenceExecutor.isSupported,
|
canProvideIntelligence: PhoneIntelligenceExecutor.isSupported,
|
||||||
intelligenceProfile: PhoneIntelligenceExecutor.profile,
|
intelligenceProfile: PhoneIntelligenceExecutor.profile,
|
||||||
// Mesh casting: receive the activity feed / messages / runner presence.
|
// Mesh casting: receive the activity feed / messages / runner presence.
|
||||||
canCast: true))
|
canCast: true,
|
||||||
|
// Covalence direct (SYNC §3.3): this phone runs the relay→direct upgrade.
|
||||||
|
canDirectConnect: directEnabled))
|
||||||
self.client = client
|
self.client = client
|
||||||
consume(client, pairingPayload: pairingPayload)
|
consume(client, pairingPayload: pairingPayload)
|
||||||
}
|
}
|
||||||
@@ -569,6 +591,7 @@ final class HostConnection {
|
|||||||
relayURL: payload.relayURL))
|
relayURL: payload.relayURL))
|
||||||
}
|
}
|
||||||
callbacks.didUpdate()
|
callbacks.didUpdate()
|
||||||
|
await startDirectUpgradeIfPossible(welcome: welcome)
|
||||||
send(.listSessions)
|
send(.listSessions)
|
||||||
send(.listDashboard)
|
send(.listDashboard)
|
||||||
// Mesh casting: subscribe with the merged ledger's cursors (every origin we hold) —
|
// Mesh casting: subscribe with the merged ledger's cursors (every origin we hold) —
|
||||||
@@ -759,6 +782,12 @@ final class HostConnection {
|
|||||||
case .settings(let settings):
|
case .settings(let settings):
|
||||||
// The account-level synced settings changed on the host — adopt the new truth.
|
// The account-level synced settings changed on the host — adopt the new truth.
|
||||||
callbacks.settingsChanged(settings)
|
callbacks.settingsChanged(settings)
|
||||||
|
case .directAnswer(let offer):
|
||||||
|
directBox?.deliver(.answer(offer))
|
||||||
|
case .directGo:
|
||||||
|
directBox?.deliver(.go)
|
||||||
|
case .directDecline(let reason):
|
||||||
|
directBox?.deliver(.decline(reason))
|
||||||
case .credentialNeeded, .credentialUpdate,
|
case .credentialNeeded, .credentialUpdate,
|
||||||
// The owner's runner-pool credential (item 4) — inert until the phone grows a
|
// The owner's runner-pool credential (item 4) — inert until the phone grows a
|
||||||
// pool-management surface; Macs are the managers today.
|
// pool-management surface; Macs are the managers today.
|
||||||
@@ -805,6 +834,46 @@ final class HostConnection {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// MARK: - Covalence direct upgrade (SYNC §3.3)
|
||||||
|
|
||||||
|
/// Once a relay session is live, arm its mux and start the background hole-punch upgrade.
|
||||||
|
/// The relay stays the session's spine; a successful punch just moves the frames off it.
|
||||||
|
/// `onPathChanged` flips the "Connected · Direct" chip (and back on failback). No-op unless
|
||||||
|
/// this is a relay session, both ends advertise the capability, and the kill switch is off.
|
||||||
|
private func startDirectUpgradeIfPossible(welcome: Welcome) async {
|
||||||
|
guard directEnabled, activeTransport == .relay, directTask == nil,
|
||||||
|
let mux = directMux, welcome.capabilities.canDirectConnect,
|
||||||
|
let hash = await client?.handshakeHash()
|
||||||
|
else { return }
|
||||||
|
mux.arm(bindingKey: MultipathFrameChannel.bindingKey(handshakeHash: hash))
|
||||||
|
mux.onPathChanged = { [weak self] hint in
|
||||||
|
Task { @MainActor in self?.directPathChanged(hint) }
|
||||||
|
}
|
||||||
|
let box = DirectReplyBox()
|
||||||
|
directBox = box
|
||||||
|
let client = self.client
|
||||||
|
directTask = Task { [box, mux, client] in
|
||||||
|
_ = await DirectUpgradeDriver.run(mux: mux, replies: box) { msg in
|
||||||
|
await client?.send(msg)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// The live data path migrated (direct upgrade / failback). Reflect it in the status chip
|
||||||
|
/// without disturbing the relay-backed session underneath.
|
||||||
|
private func directPathChanged(_ hint: SyncTransportHint) {
|
||||||
|
guard connectivity.isLive else { return }
|
||||||
|
activeTransport = hint
|
||||||
|
connectivity = .connected(hint)
|
||||||
|
callbacks.didUpdate()
|
||||||
|
}
|
||||||
|
|
||||||
|
private func teardownDirect() {
|
||||||
|
directTask?.cancel(); directTask = nil
|
||||||
|
directBox?.finish(); directBox = nil
|
||||||
|
directMux = nil
|
||||||
|
}
|
||||||
|
|
||||||
// MARK: - Send / retry / teardown
|
// MARK: - Send / retry / teardown
|
||||||
|
|
||||||
func send(_ msg: ClientMsg) {
|
func send(_ msg: ClientMsg) {
|
||||||
@@ -1005,6 +1074,7 @@ final class HostConnection {
|
|||||||
private func teardownClient() {
|
private func teardownClient() {
|
||||||
lanConnectTimeout?.cancel(); lanConnectTimeout = nil
|
lanConnectTimeout?.cancel(); lanConnectTimeout = nil
|
||||||
eventTask?.cancel(); eventTask = nil
|
eventTask?.cancel(); eventTask = nil
|
||||||
|
teardownDirect()
|
||||||
if let client { Task { await client.disconnect() } }
|
if let client { Task { await client.disconnect() } }
|
||||||
client = nil
|
client = nil
|
||||||
// A dropped connection loses the in-flight prefetch's remaining chunks — settle it empty
|
// A dropped connection loses the in-flight prefetch's remaining chunks — settle it empty
|
||||||
|
|||||||
@@ -21,7 +21,7 @@ struct PairedHost: Codable, Equatable {
|
|||||||
var transport: String?
|
var transport: String?
|
||||||
var tailnetHost: String?
|
var tailnetHost: String?
|
||||||
var tailnetPort: UInt16?
|
var tailnetPort: UInt16?
|
||||||
/// Nucleic Private Relay bootstrap (mesh P2): the host's room, this device's membership
|
/// Covalence Relay bootstrap (mesh P2): the host's room, this device's membership
|
||||||
/// token, and an optional base-URL override. Set from the pairing QR and refreshed by
|
/// token, and an optional base-URL override. Set from the pairing QR and refreshed by
|
||||||
/// the host's `relayMembership` push on every connect; nil while the host has the relay
|
/// the host's `relayMembership` push on every connect; nil while the host has the relay
|
||||||
/// method off (records predating the relay decode them as nil).
|
/// method off (records predating the relay decode them as nil).
|
||||||
|
|||||||
@@ -2274,7 +2274,10 @@ final class RemoteStore: ObservableObject {
|
|||||||
.composerTyping,
|
.composerTyping,
|
||||||
// Account-level settings go straight to every eligible host from
|
// Account-level settings go straight to every eligible host from
|
||||||
// `updateSyncedSettings`, not through this owner-routing switch.
|
// `updateSyncedSettings`, not through this owner-routing switch.
|
||||||
.updateSettings:
|
.updateSettings,
|
||||||
|
// Covalence direct handshake (SYNC §3.3) goes straight to the owning
|
||||||
|
// `HostConnection`'s `SyncClient` from the upgrade driver — never routed here.
|
||||||
|
.directOffer, .directSelect:
|
||||||
break
|
break
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -2473,7 +2476,9 @@ final class RemoteStore: ObservableObject {
|
|||||||
.composerTyping,
|
.composerTyping,
|
||||||
// Account-level settings sync — demo has no host to persist/enforce them; the local
|
// Account-level settings sync — demo has no host to persist/enforce them; the local
|
||||||
// `syncedSettings` mirror is already updated optimistically by `updateSyncedSettings`.
|
// `syncedSettings` mirror is already updated optimistically by `updateSyncedSettings`.
|
||||||
.updateSettings:
|
.updateSettings,
|
||||||
|
// Covalence direct handshake (SYNC §3.3) — demo has no real relay session to upgrade.
|
||||||
|
.directOffer, .directSelect:
|
||||||
break // passive / already handled by the seeded fixtures (demo has no mesh peers)
|
break // passive / already handled by the seeded fixtures (demo has no mesh peers)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user