mirror of
https://github.com/permissionlesstech/bitchat.git
synced 2026-07-25 15:45:20 +00:00
Compare commits
4
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
3c06c9e793 | ||
|
|
930354194f | ||
|
|
15293f6b4b | ||
|
|
842d8f2cd3 |
@@ -5,6 +5,7 @@ on:
|
|||||||
branches:
|
branches:
|
||||||
- main
|
- main
|
||||||
pull_request:
|
pull_request:
|
||||||
|
workflow_dispatch:
|
||||||
|
|
||||||
jobs:
|
jobs:
|
||||||
test:
|
test:
|
||||||
@@ -29,22 +30,24 @@ jobs:
|
|||||||
- name: Checkout code
|
- name: Checkout code
|
||||||
uses: actions/checkout@v5
|
uses: actions/checkout@v5
|
||||||
|
|
||||||
# Use the Xcode-bundled Swift toolchain: it always matches the SDK on
|
- name: Compute Swift cache salt
|
||||||
# the runner image. A standalone swift.org toolchain (setup-swift) broke
|
id: swift_cache_salt
|
||||||
# whenever the image's Xcode moved ahead of it ("this SDK is not
|
run: |
|
||||||
# supported by the compiler").
|
{
|
||||||
- name: Note toolchain version (cache key)
|
swift --version
|
||||||
id: swift-version
|
xcrun --show-sdk-platform-path
|
||||||
run: echo "version=$(swift --version 2>/dev/null | head -1 | shasum | cut -c1-12)" >> "$GITHUB_OUTPUT"
|
xcrun --show-sdk-version
|
||||||
|
xcodebuild -version
|
||||||
|
} | shasum -a 256 | awk '{ print "value=" $1 }' >> "$GITHUB_OUTPUT"
|
||||||
|
|
||||||
- name: Cache build artifacts
|
- name: Cache build artifacts
|
||||||
uses: actions/cache@v4
|
uses: actions/cache@v4
|
||||||
with:
|
with:
|
||||||
path: ${{ matrix.path }}/.build
|
path: ${{ matrix.path }}/.build
|
||||||
key: ${{ runner.os }}-${{ steps.swift-version.outputs.version }}-${{ matrix.name }}-${{ hashFiles(format('{0}/**/*.swift', matrix.path), format('{0}/**/Package.resolved', matrix.path)) }}
|
key: ${{ runner.os }}-${{ runner.arch }}-${{ matrix.name }}-${{ steps.swift_cache_salt.outputs.value }}-${{ hashFiles(format('{0}/**/*.swift', matrix.path), format('{0}/**/Package.resolved', matrix.path)) }}
|
||||||
restore-keys: |
|
restore-keys: |
|
||||||
${{ runner.os }}-${{ steps.swift-version.outputs.version }}-${{ matrix.name }}-${{ hashFiles(format('{0}/**/Package.resolved', matrix.path)) }}
|
${{ runner.os }}-${{ runner.arch }}-${{ matrix.name }}-${{ steps.swift_cache_salt.outputs.value }}-${{ hashFiles(format('{0}/**/Package.resolved', matrix.path)) }}
|
||||||
${{ runner.os }}-${{ steps.swift-version.outputs.version }}-${{ matrix.name }}-
|
${{ runner.os }}-${{ runner.arch }}-${{ matrix.name }}-${{ steps.swift_cache_salt.outputs.value }}-
|
||||||
|
|
||||||
- name: Build tests
|
- name: Build tests
|
||||||
# Built separately so the hang watchdog below times only test
|
# Built separately so the hang watchdog below times only test
|
||||||
|
|||||||
@@ -4,6 +4,9 @@ import Foundation
|
|||||||
// - 0x01: P (uint8) — Golomb-Rice parameter
|
// - 0x01: P (uint8) — Golomb-Rice parameter
|
||||||
// - 0x02: M (uint32, big-endian) — hash range (N * 2^P)
|
// - 0x02: M (uint32, big-endian) — hash range (N * 2^P)
|
||||||
// - 0x03: data (opaque) — GR bitstream bytes (MSB-first)
|
// - 0x03: data (opaque) — GR bitstream bytes (MSB-first)
|
||||||
|
// - 0x04: wanted types — raw MessageType bytes
|
||||||
|
// - 0x05: minimum timestamp (uint64, big-endian) — epoch milliseconds
|
||||||
|
// - 0x06: fragment id filter (utf8)
|
||||||
struct RequestSyncPacket {
|
struct RequestSyncPacket {
|
||||||
let p: Int
|
let p: Int
|
||||||
let m: UInt32
|
let m: UInt32
|
||||||
|
|||||||
@@ -54,7 +54,7 @@ private extension GeoRelayDirectoryDependencies {
|
|||||||
refreshCheckInterval: TransportConfig.geoRelayRefreshCheckIntervalSeconds,
|
refreshCheckInterval: TransportConfig.geoRelayRefreshCheckIntervalSeconds,
|
||||||
retryInitialSeconds: TransportConfig.geoRelayRetryInitialSeconds,
|
retryInitialSeconds: TransportConfig.geoRelayRetryInitialSeconds,
|
||||||
retryMaxSeconds: TransportConfig.geoRelayRetryMaxSeconds,
|
retryMaxSeconds: TransportConfig.geoRelayRetryMaxSeconds,
|
||||||
awaitTorReady: { await TorManager.shared.awaitEgressReady() },
|
awaitTorReady: { await TorManager.shared.awaitReady() },
|
||||||
makeFetchData: {
|
makeFetchData: {
|
||||||
let session = TorURLSession.shared.session
|
let session = TorURLSession.shared.session
|
||||||
return { request in
|
return { request in
|
||||||
|
|||||||
@@ -61,11 +61,6 @@ struct NostrRelayManagerDependencies {
|
|||||||
var locationPermissionPublisher: AnyPublisher<LocationChannelManager.PermissionState, Never>
|
var locationPermissionPublisher: AnyPublisher<LocationChannelManager.PermissionState, Never>
|
||||||
var torEnforced: () -> Bool
|
var torEnforced: () -> Bool
|
||||||
var torIsReady: () -> Bool
|
var torIsReady: () -> Bool
|
||||||
/// Synchronous cached egress-gate check: `true` only while a positive Tor
|
|
||||||
/// egress verification is within its TTL (or Tor is not enforced). When
|
|
||||||
/// `false`, connections must be queued behind `awaitTorReady`, which runs
|
|
||||||
/// the async egress self-check.
|
|
||||||
var torEgressVerified: () -> Bool
|
|
||||||
var torIsForeground: () -> Bool
|
var torIsForeground: () -> Bool
|
||||||
var awaitTorReady: (@escaping (Bool) -> Void) -> Void
|
var awaitTorReady: (@escaping (Bool) -> Void) -> Void
|
||||||
var makeSession: () -> NostrRelaySessionProtocol
|
var makeSession: () -> NostrRelaySessionProtocol
|
||||||
@@ -88,14 +83,10 @@ private extension NostrRelayManagerDependencies {
|
|||||||
locationPermissionPublisher: LocationChannelManager.shared.$permissionState.eraseToAnyPublisher(),
|
locationPermissionPublisher: LocationChannelManager.shared.$permissionState.eraseToAnyPublisher(),
|
||||||
torEnforced: { TorManager.shared.torEnforced },
|
torEnforced: { TorManager.shared.torEnforced },
|
||||||
torIsReady: { TorManager.shared.isReady },
|
torIsReady: { TorManager.shared.isReady },
|
||||||
torEgressVerified: { TorManager.shared.isEgressVerified },
|
|
||||||
torIsForeground: { TorManager.shared.isForeground() },
|
torIsForeground: { TorManager.shared.isForeground() },
|
||||||
awaitTorReady: { completion in
|
awaitTorReady: { completion in
|
||||||
Task.detached {
|
Task.detached {
|
||||||
// Require both Tor bootstrap AND a positive egress self-check
|
let ready = await TorManager.shared.awaitReady()
|
||||||
// so relay sockets never open unless traffic is proven to
|
|
||||||
// route through Tor (fail-closed).
|
|
||||||
let ready = await TorManager.shared.awaitEgressReady()
|
|
||||||
await MainActor.run {
|
await MainActor.run {
|
||||||
completion(ready)
|
completion(ready)
|
||||||
}
|
}
|
||||||
@@ -634,14 +625,8 @@ final class NostrRelayManager: ObservableObject {
|
|||||||
|
|
||||||
// MARK: - Private Methods
|
// MARK: - Private Methods
|
||||||
|
|
||||||
/// Every path that opens a relay socket funnels through this check (initial
|
|
||||||
/// connect, reconnect backoff timers, subscription-triggered connects,
|
|
||||||
/// manual retry). It must hold connections back when Tor isn't bootstrapped
|
|
||||||
/// OR when the runtime egress self-check has no fresh positive verdict —
|
|
||||||
/// otherwise the already-bootstrapped path would open sockets without ever
|
|
||||||
/// running the egress canary.
|
|
||||||
private var shouldWaitForTorBeforeConnecting: Bool {
|
private var shouldWaitForTorBeforeConnecting: Bool {
|
||||||
shouldUseTor && (!dependencies.torIsReady() || !dependencies.torEgressVerified())
|
shouldUseTor && !dependencies.torIsReady()
|
||||||
}
|
}
|
||||||
|
|
||||||
private func connectToRelays(_ relayUrls: [String], shouldLog: Bool = false) {
|
private func connectToRelays(_ relayUrls: [String], shouldLog: Bool = false) {
|
||||||
@@ -689,20 +674,16 @@ final class NostrRelayManager: ObservableObject {
|
|||||||
guard ready else {
|
guard ready else {
|
||||||
self.torReadyWaitAttempts += 1
|
self.torReadyWaitAttempts += 1
|
||||||
if self.torReadyWaitAttempts < TransportConfig.nostrTorReadyMaxWaitAttempts {
|
if self.torReadyWaitAttempts < TransportConfig.nostrTorReadyMaxWaitAttempts {
|
||||||
SecureLogger.warning("Tor not ready or egress unverified; re-queueing \(pending.count) relay connection(s) (attempt \(self.torReadyWaitAttempts))", category: .session)
|
SecureLogger.warning("Tor not ready; re-queueing \(pending.count) relay connection(s) (attempt \(self.torReadyWaitAttempts))", category: .session)
|
||||||
self.queueConnectionsUntilTorReady(pending)
|
self.queueConnectionsUntilTorReady(pending)
|
||||||
} else {
|
} else {
|
||||||
// Still fail-closed (no network), but unblock any callers
|
// Still fail-closed (no network), but unblock any callers
|
||||||
// waiting on EOSE so the UI doesn't hang indefinitely.
|
// waiting on EOSE so the UI doesn't hang indefinitely.
|
||||||
// Queued subscriptions/sends are kept; a bounded-cadence
|
// Queued subscriptions/sends are kept and flush if a later
|
||||||
// retry (below) re-enters the gate so a transient failure
|
// trigger (e.g. app foreground) brings Tor up.
|
||||||
// (Tor stall, canary outage keeping the egress unverified)
|
SecureLogger.error("❌ Tor not ready after \(self.torReadyWaitAttempts) wait(s); aborting relay connections (fail-closed)", category: .session)
|
||||||
// recovers automatically, and any later trigger (e.g. app
|
|
||||||
// foreground) also re-enters it.
|
|
||||||
SecureLogger.error("❌ Tor not ready or egress unverified after \(self.torReadyWaitAttempts) wait(s); relays stay closed (fail-closed), retrying in \(Int(TransportConfig.nostrTorGateRetrySeconds))s", category: .session)
|
|
||||||
self.torReadyWaitAttempts = 0
|
self.torReadyWaitAttempts = 0
|
||||||
self.unblockPendingEOSECallbacks(reason: "tor-unavailable")
|
self.unblockPendingEOSECallbacks(reason: "tor-unavailable")
|
||||||
self.scheduleTorGateRetry(pending)
|
|
||||||
}
|
}
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
@@ -712,24 +693,6 @@ final class NostrRelayManager: ObservableObject {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// After the Tor-gate wait attempts are exhausted, keep a low-frequency
|
|
||||||
/// retry alive so the gate re-opens without an external trigger once the
|
|
||||||
/// transient failure clears. Bounded cadence: one wait cycle per
|
|
||||||
/// `nostrTorGateRetrySeconds`; the egress verifier additionally throttles
|
|
||||||
/// actual canary probes to one per its `minRetryInterval`.
|
|
||||||
private func scheduleTorGateRetry(_ relayUrls: [String]) {
|
|
||||||
guard !relayUrls.isEmpty else { return }
|
|
||||||
let generation = connectionGeneration
|
|
||||||
dependencies.scheduleAfter(TransportConfig.nostrTorGateRetrySeconds) { [weak self] in
|
|
||||||
Task { @MainActor [weak self] in
|
|
||||||
guard let self else { return }
|
|
||||||
// Void after disconnect/reset: those paths bump the generation.
|
|
||||||
guard generation == self.connectionGeneration else { return }
|
|
||||||
self.queueConnectionsUntilTorReady(relayUrls)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Park an EOSE callback while Tor is not yet ready, and schedule the same
|
/// Park an EOSE callback while Tor is not yet ready, and schedule the same
|
||||||
/// fallback timeout `startEOSETracking` uses. Without it, a parked callback
|
/// fallback timeout `startEOSETracking` uses. Without it, a parked callback
|
||||||
/// would only be unblocked by Tor-readiness retry exhaustion (several
|
/// would only be unblocked by Tor-readiness retry exhaustion (several
|
||||||
|
|||||||
@@ -171,10 +171,6 @@ enum TransportConfig {
|
|||||||
// How many consecutive Tor-readiness waits (each bounded by TorManager's
|
// How many consecutive Tor-readiness waits (each bounded by TorManager's
|
||||||
// bootstrap deadline) to attempt before unblocking pending EOSE callers.
|
// bootstrap deadline) to attempt before unblocking pending EOSE callers.
|
||||||
static let nostrTorReadyMaxWaitAttempts: Int = 3
|
static let nostrTorReadyMaxWaitAttempts: Int = 3
|
||||||
// After Tor-gate wait attempts are exhausted (Tor never ready, or egress
|
|
||||||
// self-check unverified), retry the whole gate at this bounded cadence so
|
|
||||||
// a transient failure (e.g. canary outage) recovers without user action.
|
|
||||||
static let nostrTorGateRetrySeconds: TimeInterval = 30.0
|
|
||||||
static let nostrPendingSendQueueCap: Int = 200
|
static let nostrPendingSendQueueCap: Int = 200
|
||||||
// Sample interval for the send-queue overflow warning (first + every Nth
|
// Sample interval for the send-queue overflow warning (first + every Nth
|
||||||
// dropped event). Drops are ephemeral presence/geo traffic — log-only.
|
// dropped event). Drops are ephemeral presence/geo traffic — log-only.
|
||||||
|
|||||||
@@ -167,6 +167,16 @@ final class GossipSyncManager {
|
|||||||
return packet.timestamp >= cutoffMs
|
return packet.timestamp >= cutoffMs
|
||||||
}
|
}
|
||||||
|
|
||||||
|
private func normalizedSyncTypes(_ types: SyncTypeFlags?) -> SyncTypeFlags {
|
||||||
|
guard let types, !types.isEmpty else { return .publicMessages }
|
||||||
|
return types
|
||||||
|
}
|
||||||
|
|
||||||
|
private func isAtOrAfterMinimumTimestamp(_ packet: BitchatPacket, sinceTimestamp: UInt64?) -> Bool {
|
||||||
|
guard let sinceTimestamp else { return true }
|
||||||
|
return packet.timestamp >= sinceTimestamp
|
||||||
|
}
|
||||||
|
|
||||||
private func _onPublicPacketSeen(_ packet: BitchatPacket) {
|
private func _onPublicPacketSeen(_ packet: BitchatPacket) {
|
||||||
guard let messageType = MessageType(rawValue: packet.type) else { return }
|
guard let messageType = MessageType(rawValue: packet.type) else { return }
|
||||||
let isBroadcastRecipient: Bool = {
|
let isBroadcastRecipient: Bool = {
|
||||||
@@ -218,8 +228,8 @@ final class GossipSyncManager {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
private func sendRequestSync(for types: SyncTypeFlags) {
|
func sendRequestSync(for types: SyncTypeFlags? = nil, sinceTimestamp: UInt64? = nil) {
|
||||||
let payload = buildGcsPayload(for: types)
|
let payload = buildGcsPayload(for: types, sinceTimestamp: sinceTimestamp)
|
||||||
let pkt = BitchatPacket(
|
let pkt = BitchatPacket(
|
||||||
type: MessageType.requestSync.rawValue,
|
type: MessageType.requestSync.rawValue,
|
||||||
senderID: Data(hexString: myPeerID.id) ?? Data(),
|
senderID: Data(hexString: myPeerID.id) ?? Data(),
|
||||||
@@ -233,11 +243,11 @@ final class GossipSyncManager {
|
|||||||
delegate?.sendPacket(signed)
|
delegate?.sendPacket(signed)
|
||||||
}
|
}
|
||||||
|
|
||||||
private func sendRequestSync(to peerID: PeerID, types: SyncTypeFlags) {
|
func sendRequestSync(to peerID: PeerID, types: SyncTypeFlags? = nil, sinceTimestamp: UInt64? = nil) {
|
||||||
// Register the request for RSR validation
|
// Register the request for RSR validation
|
||||||
requestSyncManager.registerRequest(to: peerID)
|
requestSyncManager.registerRequest(to: peerID)
|
||||||
|
|
||||||
let payload = buildGcsPayload(for: types)
|
let payload = buildGcsPayload(for: types, sinceTimestamp: sinceTimestamp)
|
||||||
var recipient = Data()
|
var recipient = Data()
|
||||||
var temp = peerID.id
|
var temp = peerID.id
|
||||||
while temp.count >= 2 && recipient.count < 8 {
|
while temp.count >= 2 && recipient.count < 8 {
|
||||||
@@ -265,7 +275,8 @@ final class GossipSyncManager {
|
|||||||
}
|
}
|
||||||
|
|
||||||
private func _handleRequestSync(from peerID: PeerID, request: RequestSyncPacket) {
|
private func _handleRequestSync(from peerID: PeerID, request: RequestSyncPacket) {
|
||||||
let requestedTypes = (request.types ?? .publicMessages)
|
let requestedTypes = normalizedSyncTypes(request.types)
|
||||||
|
let sinceTimestamp = request.sinceTimestamp
|
||||||
// Decode GCS into sorted set and prepare membership checker
|
// Decode GCS into sorted set and prepare membership checker
|
||||||
let sorted = GCSFilter.decodeToSortedSet(p: request.p, m: request.m, data: request.data)
|
let sorted = GCSFilter.decodeToSortedSet(p: request.p, m: request.m, data: request.data)
|
||||||
func mightContain(_ id: Data) -> Bool {
|
func mightContain(_ id: Data) -> Bool {
|
||||||
@@ -276,7 +287,7 @@ final class GossipSyncManager {
|
|||||||
if requestedTypes.contains(.announce) {
|
if requestedTypes.contains(.announce) {
|
||||||
for (_, pair) in latestAnnouncementByPeer {
|
for (_, pair) in latestAnnouncementByPeer {
|
||||||
let (idHex, pkt) = pair
|
let (idHex, pkt) = pair
|
||||||
guard isPacketFresh(pkt) else { continue }
|
guard isPacketFresh(pkt), isAtOrAfterMinimumTimestamp(pkt, sinceTimestamp: sinceTimestamp) else { continue }
|
||||||
let idBytes = Data(hexString: idHex) ?? Data()
|
let idBytes = Data(hexString: idHex) ?? Data()
|
||||||
if !mightContain(idBytes) {
|
if !mightContain(idBytes) {
|
||||||
var toSend = pkt
|
var toSend = pkt
|
||||||
@@ -290,6 +301,7 @@ final class GossipSyncManager {
|
|||||||
if requestedTypes.contains(.message) {
|
if requestedTypes.contains(.message) {
|
||||||
let toSendMsgs = messages.allPackets(isFresh: isPacketFresh)
|
let toSendMsgs = messages.allPackets(isFresh: isPacketFresh)
|
||||||
for pkt in toSendMsgs {
|
for pkt in toSendMsgs {
|
||||||
|
guard isAtOrAfterMinimumTimestamp(pkt, sinceTimestamp: sinceTimestamp) else { continue }
|
||||||
let idBytes = PacketIdUtil.computeId(pkt)
|
let idBytes = PacketIdUtil.computeId(pkt)
|
||||||
if !mightContain(idBytes) {
|
if !mightContain(idBytes) {
|
||||||
var toSend = pkt
|
var toSend = pkt
|
||||||
@@ -303,6 +315,7 @@ final class GossipSyncManager {
|
|||||||
if requestedTypes.contains(.fragment) {
|
if requestedTypes.contains(.fragment) {
|
||||||
let frags = fragments.allPackets(isFresh: isPacketFresh)
|
let frags = fragments.allPackets(isFresh: isPacketFresh)
|
||||||
for pkt in frags {
|
for pkt in frags {
|
||||||
|
guard isAtOrAfterMinimumTimestamp(pkt, sinceTimestamp: sinceTimestamp) else { continue }
|
||||||
let idBytes = PacketIdUtil.computeId(pkt)
|
let idBytes = PacketIdUtil.computeId(pkt)
|
||||||
if !mightContain(idBytes) {
|
if !mightContain(idBytes) {
|
||||||
var toSend = pkt
|
var toSend = pkt
|
||||||
@@ -316,6 +329,7 @@ final class GossipSyncManager {
|
|||||||
if requestedTypes.contains(.fileTransfer) {
|
if requestedTypes.contains(.fileTransfer) {
|
||||||
let files = fileTransfers.allPackets(isFresh: isPacketFresh)
|
let files = fileTransfers.allPackets(isFresh: isPacketFresh)
|
||||||
for pkt in files {
|
for pkt in files {
|
||||||
|
guard isAtOrAfterMinimumTimestamp(pkt, sinceTimestamp: sinceTimestamp) else { continue }
|
||||||
let idBytes = PacketIdUtil.computeId(pkt)
|
let idBytes = PacketIdUtil.computeId(pkt)
|
||||||
if !mightContain(idBytes) {
|
if !mightContain(idBytes) {
|
||||||
var toSend = pkt
|
var toSend = pkt
|
||||||
@@ -328,25 +342,33 @@ final class GossipSyncManager {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// Build REQUEST_SYNC payload using current candidates and GCS params
|
// Build REQUEST_SYNC payload using current candidates and GCS params
|
||||||
private func buildGcsPayload(for types: SyncTypeFlags) -> Data {
|
private func buildGcsPayload(for types: SyncTypeFlags?, sinceTimestamp: UInt64? = nil) -> Data {
|
||||||
|
let requestedTypes = normalizedSyncTypes(types)
|
||||||
|
let encodedTypes = types.flatMap { $0.isEmpty ? nil : $0 }
|
||||||
var candidates: [BitchatPacket] = []
|
var candidates: [BitchatPacket] = []
|
||||||
if types.contains(.announce) {
|
if requestedTypes.contains(.announce) {
|
||||||
for (_, pair) in latestAnnouncementByPeer where isPacketFresh(pair.packet) {
|
for (_, pair) in latestAnnouncementByPeer where isPacketFresh(pair.packet) && isAtOrAfterMinimumTimestamp(pair.packet, sinceTimestamp: sinceTimestamp) {
|
||||||
candidates.append(pair.packet)
|
candidates.append(pair.packet)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
if types.contains(.message) {
|
if requestedTypes.contains(.message) {
|
||||||
candidates.append(contentsOf: messages.allPackets(isFresh: isPacketFresh))
|
candidates.append(contentsOf: messages.allPackets(isFresh: isPacketFresh).filter {
|
||||||
|
isAtOrAfterMinimumTimestamp($0, sinceTimestamp: sinceTimestamp)
|
||||||
|
})
|
||||||
}
|
}
|
||||||
if types.contains(.fragment) {
|
if requestedTypes.contains(.fragment) {
|
||||||
candidates.append(contentsOf: fragments.allPackets(isFresh: isPacketFresh))
|
candidates.append(contentsOf: fragments.allPackets(isFresh: isPacketFresh).filter {
|
||||||
|
isAtOrAfterMinimumTimestamp($0, sinceTimestamp: sinceTimestamp)
|
||||||
|
})
|
||||||
}
|
}
|
||||||
if types.contains(.fileTransfer) {
|
if requestedTypes.contains(.fileTransfer) {
|
||||||
candidates.append(contentsOf: fileTransfers.allPackets(isFresh: isPacketFresh))
|
candidates.append(contentsOf: fileTransfers.allPackets(isFresh: isPacketFresh).filter {
|
||||||
|
isAtOrAfterMinimumTimestamp($0, sinceTimestamp: sinceTimestamp)
|
||||||
|
})
|
||||||
}
|
}
|
||||||
if candidates.isEmpty {
|
if candidates.isEmpty {
|
||||||
let p = GCSFilter.deriveP(targetFpr: config.gcsTargetFpr)
|
let p = GCSFilter.deriveP(targetFpr: config.gcsTargetFpr)
|
||||||
let req = RequestSyncPacket(p: p, m: 1, data: Data(), types: types)
|
let req = RequestSyncPacket(p: p, m: 1, data: Data(), types: encodedTypes, sinceTimestamp: sinceTimestamp)
|
||||||
return req.encode()
|
return req.encode()
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -356,21 +378,21 @@ final class GossipSyncManager {
|
|||||||
let p = GCSFilter.deriveP(targetFpr: config.gcsTargetFpr)
|
let p = GCSFilter.deriveP(targetFpr: config.gcsTargetFpr)
|
||||||
let nMax = GCSFilter.estimateMaxElements(sizeBytes: config.gcsMaxBytes, p: p)
|
let nMax = GCSFilter.estimateMaxElements(sizeBytes: config.gcsMaxBytes, p: p)
|
||||||
let cap: Int
|
let cap: Int
|
||||||
if types == .fragment {
|
if requestedTypes == .fragment {
|
||||||
cap = max(1, config.fragmentCapacity)
|
cap = max(1, config.fragmentCapacity)
|
||||||
} else if types == .fileTransfer {
|
} else if requestedTypes == .fileTransfer {
|
||||||
cap = max(1, config.fileTransferCapacity)
|
cap = max(1, config.fileTransferCapacity)
|
||||||
} else {
|
} else {
|
||||||
cap = max(1, config.seenCapacity)
|
cap = max(1, config.seenCapacity)
|
||||||
}
|
}
|
||||||
let takeN = min(candidates.count, min(nMax, cap))
|
let takeN = min(candidates.count, min(nMax, cap))
|
||||||
if takeN <= 0 {
|
if takeN <= 0 {
|
||||||
let req = RequestSyncPacket(p: p, m: 1, data: Data(), types: types)
|
let req = RequestSyncPacket(p: p, m: 1, data: Data(), types: encodedTypes, sinceTimestamp: sinceTimestamp)
|
||||||
return req.encode()
|
return req.encode()
|
||||||
}
|
}
|
||||||
let ids: [Data] = candidates.prefix(takeN).map { PacketIdUtil.computeId($0) }
|
let ids: [Data] = candidates.prefix(takeN).map { PacketIdUtil.computeId($0) }
|
||||||
let params = GCSFilter.buildFilter(ids: ids, maxBytes: config.gcsMaxBytes, targetFpr: config.gcsTargetFpr)
|
let params = GCSFilter.buildFilter(ids: ids, maxBytes: config.gcsMaxBytes, targetFpr: config.gcsTargetFpr)
|
||||||
let req = RequestSyncPacket(p: params.p, m: params.m, data: params.data, types: types)
|
let req = RequestSyncPacket(p: params.p, m: params.m, data: params.data, types: encodedTypes, sinceTimestamp: sinceTimestamp)
|
||||||
return req.encode()
|
return req.encode()
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -461,5 +483,11 @@ extension GossipSyncManager {
|
|||||||
messages.allPackets { _ in true }.filter { PeerID(hexData: $0.senderID) == peerID }.count
|
messages.allPackets { _ in true }.filter { PeerID(hexData: $0.senderID) == peerID }.count
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func _buildGcsPayloadSynchronously(for types: SyncTypeFlags? = nil, sinceTimestamp: UInt64? = nil) -> Data {
|
||||||
|
queue.sync {
|
||||||
|
buildGcsPayload(for: types, sinceTimestamp: sinceTimestamp)
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
#endif
|
#endif
|
||||||
|
|||||||
@@ -1,8 +1,8 @@
|
|||||||
import BitFoundation
|
import BitFoundation
|
||||||
import Foundation
|
import Foundation
|
||||||
|
|
||||||
/// Bitfield describing which message types are covered by a REQUEST_SYNC round.
|
/// Internal bitfield describing which message types are covered by a REQUEST_SYNC round.
|
||||||
/// Matches the Android mapping (bit index -> message type).
|
/// The wire TLV uses raw `MessageType` bytes, matching Android's wantedTypes format.
|
||||||
struct SyncTypeFlags: OptionSet {
|
struct SyncTypeFlags: OptionSet {
|
||||||
let rawValue: UInt64
|
let rawValue: UInt64
|
||||||
|
|
||||||
@@ -80,26 +80,37 @@ struct SyncTypeFlags: OptionSet {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func toData() -> Data? {
|
func toData() -> Data? {
|
||||||
guard rawValue != 0 else { return nil }
|
let bytes = toMessageTypes().map(\.rawValue)
|
||||||
var value = rawValue
|
guard !bytes.isEmpty else { return nil }
|
||||||
var bytes: [UInt8] = []
|
|
||||||
while value > 0 && bytes.count < 8 {
|
|
||||||
bytes.append(UInt8(value & 0xFF))
|
|
||||||
value >>= 8
|
|
||||||
}
|
|
||||||
while let last = bytes.last, last == 0 {
|
|
||||||
bytes.removeLast()
|
|
||||||
}
|
|
||||||
guard !bytes.isEmpty, bytes.count <= 8 else { return nil }
|
|
||||||
return Data(bytes)
|
return Data(bytes)
|
||||||
}
|
}
|
||||||
|
|
||||||
static func decode(_ data: Data) -> SyncTypeFlags? {
|
static func decode(_ data: Data) -> SyncTypeFlags? {
|
||||||
|
guard !data.isEmpty else { return nil }
|
||||||
|
|
||||||
|
// Prior experimental iOS builds encoded announce+message as the
|
||||||
|
// bitfield byte 0x03. Android's upgraded wire format uses the same
|
||||||
|
// value for LEAVE, which this sync manager does not store, so prefer
|
||||||
|
// the legacy interpretation for the single-byte ambiguous case.
|
||||||
|
if data.count == 1 && data[0] == 0x03 {
|
||||||
|
return .publicMessages
|
||||||
|
}
|
||||||
|
|
||||||
|
let types = data.compactMap { MessageType(rawValue: $0) }
|
||||||
|
if !types.isEmpty {
|
||||||
|
return SyncTypeFlags(messageTypes: types)
|
||||||
|
}
|
||||||
|
|
||||||
|
return decodeLegacyBitfield(data)
|
||||||
|
}
|
||||||
|
|
||||||
|
private static func decodeLegacyBitfield(_ data: Data) -> SyncTypeFlags? {
|
||||||
guard (1...8).contains(data.count) else { return nil }
|
guard (1...8).contains(data.count) else { return nil }
|
||||||
var raw: UInt64 = 0
|
var raw: UInt64 = 0
|
||||||
for (index, byte) in data.enumerated() {
|
for (index, byte) in data.enumerated() {
|
||||||
raw |= UInt64(byte) << UInt64(index * 8)
|
raw |= UInt64(byte) << UInt64(index * 8)
|
||||||
}
|
}
|
||||||
return SyncTypeFlags(rawValue: raw)
|
let flags = SyncTypeFlags(rawValue: raw)
|
||||||
|
return flags.isEmpty ? nil : flags
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -279,6 +279,107 @@ struct GossipSyncManagerTests {
|
|||||||
#expect(sentPackets.count == 1)
|
#expect(sentPackets.count == 1)
|
||||||
#expect(sentPackets[0].type == MessageType.fragment.rawValue)
|
#expect(sentPackets[0].type == MessageType.fragment.rawValue)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Test func handleRequestSyncHonorsSinceTimestamp() async throws {
|
||||||
|
var config = GossipSyncManager.Config()
|
||||||
|
config.seenCapacity = 5
|
||||||
|
config.fragmentCapacity = 0
|
||||||
|
config.fileTransferCapacity = 0
|
||||||
|
config.messageSyncIntervalSeconds = 0
|
||||||
|
config.fragmentSyncIntervalSeconds = 0
|
||||||
|
config.fileTransferSyncIntervalSeconds = 0
|
||||||
|
|
||||||
|
let requestSyncManager = RequestSyncManager()
|
||||||
|
let manager = GossipSyncManager(myPeerID: myPeerID, config: config, requestSyncManager: requestSyncManager)
|
||||||
|
let delegate = RecordingDelegate()
|
||||||
|
manager.delegate = delegate
|
||||||
|
|
||||||
|
let sender = try #require(Data(hexString: "aabbccddeeff0011"))
|
||||||
|
let threshold = UInt64(Date().timeIntervalSince1970 * 1000)
|
||||||
|
|
||||||
|
let oldMessage = BitchatPacket(
|
||||||
|
type: MessageType.message.rawValue,
|
||||||
|
senderID: sender,
|
||||||
|
recipientID: nil,
|
||||||
|
timestamp: threshold - 1,
|
||||||
|
payload: Data([0x10]),
|
||||||
|
signature: nil,
|
||||||
|
ttl: 1
|
||||||
|
)
|
||||||
|
|
||||||
|
let freshMessage = BitchatPacket(
|
||||||
|
type: MessageType.message.rawValue,
|
||||||
|
senderID: sender,
|
||||||
|
recipientID: nil,
|
||||||
|
timestamp: threshold,
|
||||||
|
payload: Data([0x20]),
|
||||||
|
signature: nil,
|
||||||
|
ttl: 1
|
||||||
|
)
|
||||||
|
|
||||||
|
manager.onPublicPacketSeen(oldMessage)
|
||||||
|
manager.onPublicPacketSeen(freshMessage)
|
||||||
|
|
||||||
|
let peer = PeerID(str: "FFFFFFFFFFFFFFFF")
|
||||||
|
let request = RequestSyncPacket(p: 4, m: 1, data: Data(), types: .message, sinceTimestamp: threshold)
|
||||||
|
manager.handleRequestSync(from: peer, request: request)
|
||||||
|
|
||||||
|
try await TestHelpers.waitFor({ delegate.packets.count == 1 }, timeout: TestConstants.shortTimeout)
|
||||||
|
let sentPacket = try #require(delegate.packets.first)
|
||||||
|
#expect(sentPacket.payload == Data([0x20]))
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test func buildGcsPayloadHonorsSinceTimestamp() throws {
|
||||||
|
var config = GossipSyncManager.Config()
|
||||||
|
config.seenCapacity = 5
|
||||||
|
config.fragmentCapacity = 0
|
||||||
|
config.fileTransferCapacity = 0
|
||||||
|
config.messageSyncIntervalSeconds = 0
|
||||||
|
config.fragmentSyncIntervalSeconds = 0
|
||||||
|
config.fileTransferSyncIntervalSeconds = 0
|
||||||
|
config.gcsTargetFpr = 0.000001
|
||||||
|
|
||||||
|
let requestSyncManager = RequestSyncManager()
|
||||||
|
let manager = GossipSyncManager(myPeerID: myPeerID, config: config, requestSyncManager: requestSyncManager)
|
||||||
|
|
||||||
|
let sender = try #require(Data(hexString: "1122334455667788"))
|
||||||
|
let threshold = UInt64(Date().timeIntervalSince1970 * 1000)
|
||||||
|
|
||||||
|
let oldMessage = BitchatPacket(
|
||||||
|
type: MessageType.message.rawValue,
|
||||||
|
senderID: sender,
|
||||||
|
recipientID: nil,
|
||||||
|
timestamp: threshold - 1,
|
||||||
|
payload: Data([0xA0]),
|
||||||
|
signature: nil,
|
||||||
|
ttl: 1
|
||||||
|
)
|
||||||
|
|
||||||
|
let freshMessage = BitchatPacket(
|
||||||
|
type: MessageType.message.rawValue,
|
||||||
|
senderID: sender,
|
||||||
|
recipientID: nil,
|
||||||
|
timestamp: threshold + 1,
|
||||||
|
payload: Data([0xB0]),
|
||||||
|
signature: nil,
|
||||||
|
ttl: 1
|
||||||
|
)
|
||||||
|
|
||||||
|
manager.onPublicPacketSeen(oldMessage)
|
||||||
|
manager.onPublicPacketSeen(freshMessage)
|
||||||
|
manager._performMaintenanceSynchronously()
|
||||||
|
|
||||||
|
let payload = manager._buildGcsPayloadSynchronously(for: .message, sinceTimestamp: threshold)
|
||||||
|
let request = try #require(RequestSyncPacket.decode(from: payload))
|
||||||
|
let sorted = GCSFilter.decodeToSortedSet(p: request.p, m: request.m, data: request.data)
|
||||||
|
let oldBucket = GCSFilter.bucket(for: PacketIdUtil.computeId(oldMessage), modulus: request.m)
|
||||||
|
let freshBucket = GCSFilter.bucket(for: PacketIdUtil.computeId(freshMessage), modulus: request.m)
|
||||||
|
|
||||||
|
#expect(request.types?.contains(.message) == true)
|
||||||
|
#expect(request.sinceTimestamp == threshold)
|
||||||
|
#expect(GCSFilter.contains(sortedValues: sorted, candidate: freshBucket))
|
||||||
|
#expect(GCSFilter.contains(sortedValues: sorted, candidate: oldBucket) == false)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
private final class RecordingDelegate: GossipSyncManager.Delegate {
|
private final class RecordingDelegate: GossipSyncManager.Delegate {
|
||||||
|
|||||||
@@ -1,403 +0,0 @@
|
|||||||
//
|
|
||||||
// TorEgressVerifierTests.swift
|
|
||||||
// bitchatTests
|
|
||||||
//
|
|
||||||
// Unit tests for the runtime Tor-egress self-check policy/caching. The network
|
|
||||||
// probe is injected, so these tests are deterministic and offline.
|
|
||||||
//
|
|
||||||
|
|
||||||
import Testing
|
|
||||||
import Foundation
|
|
||||||
@testable import bitchat
|
|
||||||
import Tor
|
|
||||||
|
|
||||||
@Suite(.serialized)
|
|
||||||
struct TorEgressVerifierTests {
|
|
||||||
|
|
||||||
/// Deterministic, controllable clock + probe.
|
|
||||||
private final class Harness: @unchecked Sendable {
|
|
||||||
private let lock = NSLock()
|
|
||||||
private var _now = Date(timeIntervalSince1970: 1_000_000)
|
|
||||||
private var _result: TorEgressVerifier.ProbeResult = .verifiedTor
|
|
||||||
private var _hanging = false
|
|
||||||
private var _gated = false
|
|
||||||
private var _released = false
|
|
||||||
private var _probeCount = 0
|
|
||||||
private var _cancelledCount = 0
|
|
||||||
|
|
||||||
var now: Date {
|
|
||||||
lock.lock(); defer { lock.unlock() }; return _now
|
|
||||||
}
|
|
||||||
var probeCount: Int {
|
|
||||||
lock.lock(); defer { lock.unlock() }; return _probeCount
|
|
||||||
}
|
|
||||||
/// Number of hung probes that observed cooperative cancellation.
|
|
||||||
var cancelledCount: Int {
|
|
||||||
lock.lock(); defer { lock.unlock() }; return _cancelledCount
|
|
||||||
}
|
|
||||||
func advance(_ seconds: TimeInterval) {
|
|
||||||
lock.lock(); _now = _now.addingTimeInterval(seconds); lock.unlock()
|
|
||||||
}
|
|
||||||
func setResult(_ r: TorEgressVerifier.ProbeResult) {
|
|
||||||
lock.lock(); _result = r; lock.unlock()
|
|
||||||
}
|
|
||||||
/// When `true`, probes park forever and only exit via cooperative
|
|
||||||
/// cancellation — models a canary request wedged by
|
|
||||||
/// `waitsForConnectivity` deferring the request timer.
|
|
||||||
func setHanging(_ hanging: Bool) {
|
|
||||||
lock.lock(); _hanging = hanging; lock.unlock()
|
|
||||||
}
|
|
||||||
/// When `true`, probes wait for `release()` before returning — models
|
|
||||||
/// a slow-but-completing canary for join-semantics tests.
|
|
||||||
func setGated(_ gated: Bool) {
|
|
||||||
lock.lock(); _gated = gated; lock.unlock()
|
|
||||||
}
|
|
||||||
func release() {
|
|
||||||
lock.lock(); _released = true; lock.unlock()
|
|
||||||
}
|
|
||||||
/// Suspends until at least `n` probes have started.
|
|
||||||
func waitUntilProbeCount(atLeast n: Int) async {
|
|
||||||
while probeCount < n {
|
|
||||||
try? await Task.sleep(nanoseconds: 1_000_000)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
func makeProbe() -> @Sendable () async -> TorEgressVerifier.ProbeResult {
|
|
||||||
return { [self] in
|
|
||||||
lock.lock()
|
|
||||||
_probeCount += 1
|
|
||||||
let r = _result
|
|
||||||
let hang = _hanging
|
|
||||||
let gated = _gated
|
|
||||||
lock.unlock()
|
|
||||||
if hang {
|
|
||||||
// Park until cancelled (verifier watchdog or invalidate());
|
|
||||||
// cancellation-responsive so no task outlives the test.
|
|
||||||
while !Task.isCancelled {
|
|
||||||
try? await Task.sleep(nanoseconds: 2_000_000)
|
|
||||||
}
|
|
||||||
lock.lock(); _cancelledCount += 1; lock.unlock()
|
|
||||||
return .unreachable("hung probe cancelled")
|
|
||||||
}
|
|
||||||
if gated {
|
|
||||||
while !Task.isCancelled {
|
|
||||||
lock.lock(); let released = _released; lock.unlock()
|
|
||||||
if released { break }
|
|
||||||
try? await Task.sleep(nanoseconds: 1_000_000)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
return r
|
|
||||||
}
|
|
||||||
}
|
|
||||||
func nowProvider() -> @Sendable () -> Date {
|
|
||||||
return { [self] in self.now }
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
private func makeVerifier(
|
|
||||||
_ h: Harness,
|
|
||||||
ttl: TimeInterval = 300,
|
|
||||||
minRetry: TimeInterval = 5,
|
|
||||||
probeTimeout: TimeInterval = TorEgressVerifier.defaultProbeTimeout
|
|
||||||
) -> TorEgressVerifier {
|
|
||||||
TorEgressVerifier(
|
|
||||||
ttl: ttl,
|
|
||||||
minRetryInterval: minRetry,
|
|
||||||
probeTimeout: probeTimeout,
|
|
||||||
now: h.nowProvider(),
|
|
||||||
probe: h.makeProbe()
|
|
||||||
)
|
|
||||||
}
|
|
||||||
|
|
||||||
@Test("verifiedTor allows and is cached within TTL (single probe)")
|
|
||||||
func verifiedIsCached() async {
|
|
||||||
let h = Harness()
|
|
||||||
h.setResult(.verifiedTor)
|
|
||||||
let v = makeVerifier(h, ttl: 300)
|
|
||||||
|
|
||||||
#expect(await v.verify() == true)
|
|
||||||
// Second call within TTL must not re-probe.
|
|
||||||
#expect(await v.verify() == true)
|
|
||||||
#expect(h.probeCount == 1)
|
|
||||||
}
|
|
||||||
|
|
||||||
@Test("cache expires after TTL and re-probes")
|
|
||||||
func cacheExpires() async {
|
|
||||||
let h = Harness()
|
|
||||||
h.setResult(.verifiedTor)
|
|
||||||
let v = makeVerifier(h, ttl: 300)
|
|
||||||
|
|
||||||
#expect(await v.verify() == true)
|
|
||||||
#expect(h.probeCount == 1)
|
|
||||||
|
|
||||||
h.advance(301)
|
|
||||||
#expect(await v.verify() == true)
|
|
||||||
#expect(h.probeCount == 2)
|
|
||||||
}
|
|
||||||
|
|
||||||
@Test("notTor refuses (leak detected) and is never cached as allowed")
|
|
||||||
func notTorRefuses() async {
|
|
||||||
let h = Harness()
|
|
||||||
h.setResult(.notTor)
|
|
||||||
let v = makeVerifier(h, ttl: 300, minRetry: 0)
|
|
||||||
|
|
||||||
#expect(await v.verify() == false)
|
|
||||||
// A subsequent success recovers.
|
|
||||||
h.setResult(.verifiedTor)
|
|
||||||
#expect(await v.verify() == true)
|
|
||||||
}
|
|
||||||
|
|
||||||
@Test("unreachable refuses: an unverified egress must not proceed (fail-closed)")
|
|
||||||
func unreachableRefuses() async {
|
|
||||||
let h = Harness()
|
|
||||||
h.setResult(.unreachable("down"))
|
|
||||||
let v = makeVerifier(h, ttl: 300, minRetry: 0)
|
|
||||||
|
|
||||||
#expect(await v.verify() == false)
|
|
||||||
#expect(v.hasFreshVerification == false)
|
|
||||||
}
|
|
||||||
|
|
||||||
@Test("unreachable-then-reachable recovers via retry")
|
|
||||||
func unreachableRecoversWhenCanaryReturns() async {
|
|
||||||
let h = Harness()
|
|
||||||
h.setResult(.unreachable("down"))
|
|
||||||
let v = makeVerifier(h, ttl: 300, minRetry: 5)
|
|
||||||
|
|
||||||
#expect(await v.verify() == false)
|
|
||||||
#expect(h.probeCount == 1)
|
|
||||||
|
|
||||||
// Canary comes back; the next probe (after the retry throttle window)
|
|
||||||
// verifies and allows again.
|
|
||||||
h.setResult(.verifiedTor)
|
|
||||||
h.advance(5)
|
|
||||||
#expect(await v.verify() == true)
|
|
||||||
#expect(h.probeCount == 2)
|
|
||||||
#expect(v.hasFreshVerification == true)
|
|
||||||
}
|
|
||||||
|
|
||||||
@Test("cached verifiedTor within TTL allows during a canary blip without re-probing")
|
|
||||||
func cachedVerifiedAllowsDuringCanaryBlip() async {
|
|
||||||
let h = Harness()
|
|
||||||
h.setResult(.verifiedTor)
|
|
||||||
let v = makeVerifier(h, ttl: 300, minRetry: 0)
|
|
||||||
|
|
||||||
#expect(await v.verify() == true)
|
|
||||||
#expect(h.probeCount == 1)
|
|
||||||
|
|
||||||
// The canary goes down inside the TTL window: the cached positive
|
|
||||||
// verdict is authoritative, no probe runs, traffic stays allowed.
|
|
||||||
h.setResult(.unreachable("down"))
|
|
||||||
h.advance(100)
|
|
||||||
#expect(await v.verify() == true)
|
|
||||||
#expect(h.probeCount == 1)
|
|
||||||
#expect(v.hasFreshVerification == true)
|
|
||||||
}
|
|
||||||
|
|
||||||
@Test("TTL expiry + unreachable refuses until a probe succeeds again")
|
|
||||||
func expiredCacheWithUnreachableRefuses() async {
|
|
||||||
let h = Harness()
|
|
||||||
h.setResult(.verifiedTor)
|
|
||||||
let v = makeVerifier(h, ttl: 300, minRetry: 0)
|
|
||||||
|
|
||||||
#expect(await v.verify() == true)
|
|
||||||
|
|
||||||
// Past the TTL the old verdict no longer stands: an unreachable canary
|
|
||||||
// means unverified egress, so connection opens are refused.
|
|
||||||
h.setResult(.unreachable("down"))
|
|
||||||
h.advance(301)
|
|
||||||
#expect(v.hasFreshVerification == false)
|
|
||||||
#expect(await v.verify() == false)
|
|
||||||
#expect(h.probeCount == 2)
|
|
||||||
|
|
||||||
// A subsequent successful probe restores service.
|
|
||||||
h.setResult(.verifiedTor)
|
|
||||||
#expect(await v.verify() == true)
|
|
||||||
#expect(v.hasFreshVerification == true)
|
|
||||||
}
|
|
||||||
|
|
||||||
@Test("minRetryInterval bounds re-probing while unverified (no canary hammering)")
|
|
||||||
func throttleReprobe() async {
|
|
||||||
let h = Harness()
|
|
||||||
h.setResult(.unreachable("down"))
|
|
||||||
let v = makeVerifier(h, ttl: 300, minRetry: 5)
|
|
||||||
|
|
||||||
#expect(await v.verify() == false)
|
|
||||||
#expect(h.probeCount == 1)
|
|
||||||
// Within minRetry window: reuse last (refusing) decision, no new probe.
|
|
||||||
h.advance(1)
|
|
||||||
#expect(await v.verify() == false)
|
|
||||||
#expect(h.probeCount == 1)
|
|
||||||
// After the window: re-probe.
|
|
||||||
h.advance(5)
|
|
||||||
#expect(await v.verify() == false)
|
|
||||||
#expect(h.probeCount == 2)
|
|
||||||
}
|
|
||||||
|
|
||||||
@Test("invalidate clears the synchronous cache snapshot")
|
|
||||||
func invalidateClearsSnapshot() async {
|
|
||||||
let h = Harness()
|
|
||||||
h.setResult(.verifiedTor)
|
|
||||||
let v = makeVerifier(h, ttl: 300)
|
|
||||||
|
|
||||||
#expect(await v.verify() == true)
|
|
||||||
#expect(v.hasFreshVerification == true)
|
|
||||||
await v.invalidate()
|
|
||||||
#expect(v.hasFreshVerification == false)
|
|
||||||
}
|
|
||||||
|
|
||||||
@Test("notTor drops any cached verification snapshot")
|
|
||||||
func notTorClearsSnapshot() async {
|
|
||||||
let h = Harness()
|
|
||||||
h.setResult(.verifiedTor)
|
|
||||||
let v = makeVerifier(h, ttl: 300, minRetry: 0)
|
|
||||||
|
|
||||||
#expect(await v.verify() == true)
|
|
||||||
#expect(v.hasFreshVerification == true)
|
|
||||||
|
|
||||||
h.setResult(.notTor)
|
|
||||||
h.advance(301)
|
|
||||||
#expect(await v.verify() == false)
|
|
||||||
#expect(v.hasFreshVerification == false)
|
|
||||||
}
|
|
||||||
|
|
||||||
@Test("invalidate forces a fresh probe")
|
|
||||||
func invalidateForcesReprobe() async {
|
|
||||||
let h = Harness()
|
|
||||||
h.setResult(.verifiedTor)
|
|
||||||
let v = makeVerifier(h, ttl: 300)
|
|
||||||
|
|
||||||
#expect(await v.verify() == true)
|
|
||||||
#expect(h.probeCount == 1)
|
|
||||||
await v.invalidate()
|
|
||||||
#expect(await v.verify() == true)
|
|
||||||
#expect(h.probeCount == 2)
|
|
||||||
}
|
|
||||||
|
|
||||||
@Test("lastProbeResult reflects the most recent outcome")
|
|
||||||
func lastResultTracked() async {
|
|
||||||
let h = Harness()
|
|
||||||
h.setResult(.notTor)
|
|
||||||
let v = makeVerifier(h, ttl: 300, minRetry: 0)
|
|
||||||
_ = await v.verify()
|
|
||||||
#expect(await v.lastProbeResult() == .notTor)
|
|
||||||
}
|
|
||||||
|
|
||||||
// MARK: - Liveness (probe timeout watchdog + invalidate cancellation)
|
|
||||||
|
|
||||||
@Test("a probe that never completes is bounded by probeTimeout and fails closed")
|
|
||||||
func hungProbeIsBoundedByTimeout() async {
|
|
||||||
let h = Harness()
|
|
||||||
h.setHanging(true)
|
|
||||||
let v = makeVerifier(h, ttl: 300, minRetry: 0, probeTimeout: 0.05)
|
|
||||||
|
|
||||||
// Without the watchdog this would wedge: waitsForConnectivity can
|
|
||||||
// defer the request timer, leaving the canary bounded only by the
|
|
||||||
// 7-day resource timeout.
|
|
||||||
#expect(await v.verify() == false)
|
|
||||||
let last = await v.lastProbeResult()
|
|
||||||
switch last {
|
|
||||||
case .unreachable:
|
|
||||||
break // fail-closed timeout verdict recorded
|
|
||||||
default:
|
|
||||||
Issue.record("expected .unreachable after probe timeout, got \(String(describing: last))")
|
|
||||||
}
|
|
||||||
#expect(v.hasFreshVerification == false)
|
|
||||||
|
|
||||||
// The hung probe task itself was cancelled (URLSession task would be
|
|
||||||
// torn down), not abandoned.
|
|
||||||
while h.cancelledCount < 1 {
|
|
||||||
try? await Task.sleep(nanoseconds: 1_000_000)
|
|
||||||
}
|
|
||||||
#expect(h.cancelledCount == 1)
|
|
||||||
|
|
||||||
// The in-flight slot was cleared: the next verify() starts a fresh
|
|
||||||
// probe (does not join the hung one) and recovers.
|
|
||||||
h.setHanging(false)
|
|
||||||
h.setResult(.verifiedTor)
|
|
||||||
#expect(await v.verify() == true)
|
|
||||||
#expect(h.probeCount == 2)
|
|
||||||
#expect(v.hasFreshVerification == true)
|
|
||||||
}
|
|
||||||
|
|
||||||
@Test("invalidate() cancels the in-flight probe and the awaiting caller fails closed")
|
|
||||||
func invalidateCancelsInFlightProbe() async {
|
|
||||||
let h = Harness()
|
|
||||||
h.setHanging(true)
|
|
||||||
// Long (real-time) timeout and throttle: only invalidate() can
|
|
||||||
// unblock the caller, and only invalidate() clearing the throttle
|
|
||||||
// lets the follow-up probe run without advancing the clock.
|
|
||||||
let v = makeVerifier(h, ttl: 300, minRetry: 600, probeTimeout: 600)
|
|
||||||
|
|
||||||
let first = Task { await v.verify() }
|
|
||||||
await h.waitUntilProbeCount(atLeast: 1)
|
|
||||||
await v.invalidate()
|
|
||||||
|
|
||||||
// The awaiting caller resolves promptly (no 600s wait) and refuses.
|
|
||||||
#expect(await first.value == false)
|
|
||||||
#expect(v.hasFreshVerification == false)
|
|
||||||
|
|
||||||
// The hung probe observed cancellation (a Tor restart genuinely
|
|
||||||
// shakes the wedged canary request).
|
|
||||||
while h.cancelledCount < 1 {
|
|
||||||
try? await Task.sleep(nanoseconds: 1_000_000)
|
|
||||||
}
|
|
||||||
|
|
||||||
// Recovery: a fresh verify() runs a NEW probe — it neither joins the
|
|
||||||
// cancelled one nor inherits its throttle/last-result state.
|
|
||||||
h.setHanging(false)
|
|
||||||
h.setResult(.verifiedTor)
|
|
||||||
#expect(await v.verify() == true)
|
|
||||||
#expect(h.probeCount == 2)
|
|
||||||
#expect(v.hasFreshVerification == true)
|
|
||||||
}
|
|
||||||
|
|
||||||
@Test("recovery after a hung probe survives invalidate + Tor restart cycle")
|
|
||||||
func hungThenInvalidatedThenRecovers() async {
|
|
||||||
let h = Harness()
|
|
||||||
h.setHanging(true)
|
|
||||||
let v = makeVerifier(h, ttl: 300, minRetry: 5, probeTimeout: 600)
|
|
||||||
|
|
||||||
// Wedge one probe, then simulate a Tor restart mid-flight.
|
|
||||||
let wedged = Task { await v.verify() }
|
|
||||||
await h.waitUntilProbeCount(atLeast: 1)
|
|
||||||
await v.invalidate()
|
|
||||||
#expect(await wedged.value == false)
|
|
||||||
|
|
||||||
// Canary still down right after restart: fresh probe, fail closed —
|
|
||||||
// and the throttle applies to the FRESH result (bounded retry intact).
|
|
||||||
h.setHanging(false)
|
|
||||||
h.setResult(.unreachable("circuit not built"))
|
|
||||||
#expect(await v.verify() == false)
|
|
||||||
#expect(h.probeCount == 2)
|
|
||||||
h.advance(1)
|
|
||||||
#expect(await v.verify() == false)
|
|
||||||
#expect(h.probeCount == 2) // throttled, no hammering
|
|
||||||
|
|
||||||
// Canary returns after the retry window: verification recovers.
|
|
||||||
h.setResult(.verifiedTor)
|
|
||||||
h.advance(5)
|
|
||||||
#expect(await v.verify() == true)
|
|
||||||
#expect(v.hasFreshVerification == true)
|
|
||||||
}
|
|
||||||
|
|
||||||
@Test("concurrent verify() callers still share a single in-flight probe")
|
|
||||||
func concurrentCallersShareOneProbe() async {
|
|
||||||
let h = Harness()
|
|
||||||
h.setResult(.verifiedTor)
|
|
||||||
h.setGated(true)
|
|
||||||
let v = makeVerifier(h, ttl: 300, minRetry: 0)
|
|
||||||
|
|
||||||
let t1 = Task { await v.verify() }
|
|
||||||
await h.waitUntilProbeCount(atLeast: 1)
|
|
||||||
let t2 = Task { await v.verify() }
|
|
||||||
// Give t2 a chance to join the in-flight probe before releasing it.
|
|
||||||
// Either way the invariant holds: t2 joins the shared probe, or (if
|
|
||||||
// scheduled after completion) hits the fresh TTL cache — exactly one
|
|
||||||
// probe runs.
|
|
||||||
try? await Task.sleep(nanoseconds: 20_000_000)
|
|
||||||
h.release()
|
|
||||||
|
|
||||||
#expect(await t1.value == true)
|
|
||||||
#expect(await t2.value == true)
|
|
||||||
#expect(h.probeCount == 1)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
@@ -111,116 +111,6 @@ final class NostrRelayManagerTests: XCTestCase {
|
|||||||
XCTAssertTrue(connected)
|
XCTAssertTrue(connected)
|
||||||
}
|
}
|
||||||
|
|
||||||
func test_connect_whenTorAlreadyReady_waitsForEgressVerificationBeforeCreatingSessions() async {
|
|
||||||
// Tor is bootstrapped, but the egress self-check has no fresh verdict:
|
|
||||||
// the ready path must still queue behind the async egress gate instead
|
|
||||||
// of opening sockets directly.
|
|
||||||
let context = makeContext(
|
|
||||||
permission: .authorized,
|
|
||||||
userTorEnabled: true,
|
|
||||||
torEnforced: true,
|
|
||||||
torIsReady: true,
|
|
||||||
torEgressVerified: false
|
|
||||||
)
|
|
||||||
|
|
||||||
context.manager.connect()
|
|
||||||
|
|
||||||
XCTAssertTrue(context.sessionFactory.requestedURLs.isEmpty)
|
|
||||||
XCTAssertEqual(context.torWaiter.awaitCallCount, 1)
|
|
||||||
|
|
||||||
// Egress verification succeeds (awaitEgressReady returned true, which
|
|
||||||
// implies the verifier now holds a fresh cached verdict).
|
|
||||||
context.torEgressVerified.value = true
|
|
||||||
context.torWaiter.resolve(true)
|
|
||||||
|
|
||||||
let connectedAfterVerification = await waitUntil {
|
|
||||||
context.sessionFactory.requestedURLs.count == self.expectedDefaultRelayCount &&
|
|
||||||
context.manager.relays.allSatisfy(\.isConnected)
|
|
||||||
}
|
|
||||||
XCTAssertTrue(connectedAfterVerification)
|
|
||||||
}
|
|
||||||
|
|
||||||
func test_reconnect_requeuesBehindEgressGateWhenVerificationLapses() async {
|
|
||||||
// Reconnect backoff timers call connectToRelay directly; when the
|
|
||||||
// cached egress verification has lapsed by then, the reconnect must go
|
|
||||||
// back through the gate rather than opening a socket.
|
|
||||||
let relayURL = "wss://egress-reconnect.example"
|
|
||||||
let context = makeContext(
|
|
||||||
permission: .denied,
|
|
||||||
userTorEnabled: true,
|
|
||||||
torEnforced: true,
|
|
||||||
torIsReady: true,
|
|
||||||
torEgressVerified: true
|
|
||||||
)
|
|
||||||
|
|
||||||
context.manager.ensureConnections(to: [relayURL])
|
|
||||||
let connected = await waitUntil {
|
|
||||||
context.manager.relays.first(where: { $0.url == relayURL })?.isConnected == true
|
|
||||||
}
|
|
||||||
XCTAssertTrue(connected)
|
|
||||||
XCTAssertEqual(context.sessionFactory.requestedURLs.count, 1)
|
|
||||||
|
|
||||||
// The socket drops and the cached verification expires meanwhile.
|
|
||||||
context.torEgressVerified.value = false
|
|
||||||
context.sessionFactory.latestConnection(for: relayURL)?
|
|
||||||
.fail(error: NSError(domain: NSURLErrorDomain, code: NSURLErrorTimedOut))
|
|
||||||
let reconnectScheduled = await waitUntil { context.scheduler.scheduled.count == 1 }
|
|
||||||
XCTAssertTrue(reconnectScheduled)
|
|
||||||
|
|
||||||
context.scheduler.runNext()
|
|
||||||
let queuedBehindGate = await waitUntil { context.torWaiter.awaitCallCount == 1 }
|
|
||||||
XCTAssertTrue(queuedBehindGate)
|
|
||||||
// No new socket until the egress gate passes.
|
|
||||||
XCTAssertEqual(context.sessionFactory.requestedURLs.count, 1)
|
|
||||||
|
|
||||||
context.torEgressVerified.value = true
|
|
||||||
context.torWaiter.resolve(true)
|
|
||||||
|
|
||||||
let reconnected = await waitUntil {
|
|
||||||
context.sessionFactory.requestedURLs.count == 2 &&
|
|
||||||
context.manager.relays.first(where: { $0.url == relayURL })?.isConnected == true
|
|
||||||
}
|
|
||||||
XCTAssertTrue(reconnected)
|
|
||||||
}
|
|
||||||
|
|
||||||
func test_connect_egressGateExhaustionSchedulesBoundedRetryAndRecovers() async {
|
|
||||||
// A persistent unverified egress exhausts the wait attempts; a bounded
|
|
||||||
// low-frequency retry must then recover automatically once
|
|
||||||
// verification succeeds (e.g. transient canary outage ends).
|
|
||||||
let relayURL = "wss://egress-gate-retry.example"
|
|
||||||
let context = makeContext(
|
|
||||||
permission: .denied,
|
|
||||||
userTorEnabled: true,
|
|
||||||
torEnforced: true,
|
|
||||||
torIsReady: true,
|
|
||||||
torEgressVerified: false
|
|
||||||
)
|
|
||||||
|
|
||||||
context.manager.ensureConnections(to: [relayURL])
|
|
||||||
for _ in 0..<TransportConfig.nostrTorReadyMaxWaitAttempts {
|
|
||||||
context.torWaiter.resolve(false)
|
|
||||||
}
|
|
||||||
|
|
||||||
// Fail-closed, with a single bounded-cadence retry scheduled.
|
|
||||||
XCTAssertTrue(context.sessionFactory.requestedURLs.isEmpty)
|
|
||||||
XCTAssertEqual(context.scheduler.scheduled.count, 1)
|
|
||||||
XCTAssertEqual(context.scheduler.scheduled.first?.delay, TransportConfig.nostrTorGateRetrySeconds)
|
|
||||||
|
|
||||||
// The outage ends before the retry fires.
|
|
||||||
let attemptsBefore = context.torWaiter.awaitCallCount
|
|
||||||
context.scheduler.runNext()
|
|
||||||
let regated = await waitUntil { context.torWaiter.awaitCallCount == attemptsBefore + 1 }
|
|
||||||
XCTAssertTrue(regated)
|
|
||||||
|
|
||||||
context.torEgressVerified.value = true
|
|
||||||
context.torWaiter.resolve(true)
|
|
||||||
|
|
||||||
let recovered = await waitUntil {
|
|
||||||
context.manager.relays.first(where: { $0.url == relayURL })?.isConnected == true
|
|
||||||
}
|
|
||||||
XCTAssertTrue(recovered)
|
|
||||||
}
|
|
||||||
|
|
||||||
func test_subscribe_unblocksDeferredEOSEWhenTorWaitAttemptsExhausted() async {
|
func test_subscribe_unblocksDeferredEOSEWhenTorWaitAttemptsExhausted() async {
|
||||||
let relayURL = "wss://tor-eose-unblock.example"
|
let relayURL = "wss://tor-eose-unblock.example"
|
||||||
let context = makeContext(permission: .denied, userTorEnabled: true, torEnforced: true, torIsReady: false)
|
let context = makeContext(permission: .denied, userTorEnabled: true, torEnforced: true, torIsReady: false)
|
||||||
@@ -1559,7 +1449,6 @@ final class NostrRelayManagerTests: XCTestCase {
|
|||||||
userTorEnabled: Bool = false,
|
userTorEnabled: Bool = false,
|
||||||
torEnforced: Bool = false,
|
torEnforced: Bool = false,
|
||||||
torIsReady: Bool = true,
|
torIsReady: Bool = true,
|
||||||
torEgressVerified: Bool = true,
|
|
||||||
torIsForeground: Bool = true,
|
torIsForeground: Bool = true,
|
||||||
jitterUnit: @escaping () -> Double = { 0.5 } // 0.5 -> jitter factor 1.0 (no jitter)
|
jitterUnit: @escaping () -> Double = { 0.5 } // 0.5 -> jitter factor 1.0 (no jitter)
|
||||||
) -> RelayManagerTestContext {
|
) -> RelayManagerTestContext {
|
||||||
@@ -1569,7 +1458,6 @@ final class NostrRelayManagerTests: XCTestCase {
|
|||||||
let scheduler = MockRelayScheduler()
|
let scheduler = MockRelayScheduler()
|
||||||
let clock = MutableClock(now: Date(timeIntervalSince1970: 1_700_000_000))
|
let clock = MutableClock(now: Date(timeIntervalSince1970: 1_700_000_000))
|
||||||
let torWaiter = MockTorWaiter(isReady: torIsReady)
|
let torWaiter = MockTorWaiter(isReady: torIsReady)
|
||||||
let torEgressVerifiedFlag = MutableBool(value: torEgressVerified)
|
|
||||||
let torForeground = MutableBool(value: torIsForeground)
|
let torForeground = MutableBool(value: torIsForeground)
|
||||||
let activationFlag = MutableBool(value: activationAllowed)
|
let activationFlag = MutableBool(value: activationAllowed)
|
||||||
let manager = NostrRelayManager(
|
let manager = NostrRelayManager(
|
||||||
@@ -1582,7 +1470,6 @@ final class NostrRelayManagerTests: XCTestCase {
|
|||||||
locationPermissionPublisher: permissionSubject.eraseToAnyPublisher(),
|
locationPermissionPublisher: permissionSubject.eraseToAnyPublisher(),
|
||||||
torEnforced: { torEnforced },
|
torEnforced: { torEnforced },
|
||||||
torIsReady: { torWaiter.isReady },
|
torIsReady: { torWaiter.isReady },
|
||||||
torEgressVerified: { torEgressVerifiedFlag.value },
|
|
||||||
torIsForeground: { torForeground.value },
|
torIsForeground: { torForeground.value },
|
||||||
awaitTorReady: torWaiter.await(completion:),
|
awaitTorReady: torWaiter.await(completion:),
|
||||||
makeSession: { sessionFactory },
|
makeSession: { sessionFactory },
|
||||||
@@ -1602,7 +1489,6 @@ final class NostrRelayManagerTests: XCTestCase {
|
|||||||
clock: clock,
|
clock: clock,
|
||||||
activationAllowed: activationFlag,
|
activationAllowed: activationFlag,
|
||||||
torWaiter: torWaiter,
|
torWaiter: torWaiter,
|
||||||
torEgressVerified: torEgressVerifiedFlag,
|
|
||||||
torForeground: torForeground
|
torForeground: torForeground
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
@@ -1657,7 +1543,6 @@ private struct RelayManagerTestContext {
|
|||||||
let clock: MutableClock
|
let clock: MutableClock
|
||||||
let activationAllowed: MutableBool
|
let activationAllowed: MutableBool
|
||||||
let torWaiter: MockTorWaiter
|
let torWaiter: MockTorWaiter
|
||||||
let torEgressVerified: MutableBool
|
|
||||||
let torForeground: MutableBool
|
let torForeground: MutableBool
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -0,0 +1,118 @@
|
|||||||
|
//
|
||||||
|
// RequestSyncPacketTests.swift
|
||||||
|
// bitchat
|
||||||
|
//
|
||||||
|
// This is free and unencumbered software released into the public domain.
|
||||||
|
// For more information, see <https://unlicense.org>
|
||||||
|
//
|
||||||
|
|
||||||
|
import Foundation
|
||||||
|
import Testing
|
||||||
|
import BitFoundation
|
||||||
|
@testable import bitchat
|
||||||
|
|
||||||
|
struct RequestSyncPacketTests {
|
||||||
|
@Test func baseFieldsRoundTrip() throws {
|
||||||
|
let original = RequestSyncPacket(
|
||||||
|
p: 7,
|
||||||
|
m: 12_800,
|
||||||
|
data: Data([1, 2, 3, 4, 5])
|
||||||
|
)
|
||||||
|
|
||||||
|
let decoded = try #require(RequestSyncPacket.decode(from: original.encode()))
|
||||||
|
|
||||||
|
#expect(decoded.p == 7)
|
||||||
|
#expect(decoded.m == 12_800)
|
||||||
|
#expect(decoded.data == Data([1, 2, 3, 4, 5]))
|
||||||
|
#expect(decoded.types == nil)
|
||||||
|
#expect(decoded.sinceTimestamp == nil)
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test func upgradedFieldsRoundTripAsAndroidWantedTypes() throws {
|
||||||
|
let original = RequestSyncPacket(
|
||||||
|
p: 8,
|
||||||
|
m: 25_600,
|
||||||
|
data: Data([10, 20, 30]),
|
||||||
|
types: .publicMessages,
|
||||||
|
sinceTimestamp: 1_700_000_000_000
|
||||||
|
)
|
||||||
|
|
||||||
|
let encoded = original.encode()
|
||||||
|
let wantedTypes = try #require(tlvValue(type: 0x04, in: encoded))
|
||||||
|
let decoded = try #require(RequestSyncPacket.decode(from: encoded))
|
||||||
|
|
||||||
|
#expect(wantedTypes == Data([MessageType.announce.rawValue, MessageType.message.rawValue]))
|
||||||
|
#expect(decoded.p == 8)
|
||||||
|
#expect(decoded.m == 25_600)
|
||||||
|
#expect(decoded.data == Data([10, 20, 30]))
|
||||||
|
#expect(decoded.types?.contains(.announce) == true)
|
||||||
|
#expect(decoded.types?.contains(.message) == true)
|
||||||
|
#expect(decoded.sinceTimestamp == 1_700_000_000_000)
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test func decodesLegacyPayloadWithoutUpgradeFields() throws {
|
||||||
|
let payload = Data([
|
||||||
|
0x01, 0x00, 0x01, 0x07,
|
||||||
|
0x02, 0x00, 0x04, 0x00, 0x00, 0x32, 0x00,
|
||||||
|
0x03, 0x00, 0x03, 0x01, 0x02, 0x03
|
||||||
|
])
|
||||||
|
|
||||||
|
let decoded = try #require(RequestSyncPacket.decode(from: payload))
|
||||||
|
|
||||||
|
#expect(decoded.p == 7)
|
||||||
|
#expect(decoded.m == 12_800)
|
||||||
|
#expect(decoded.data == Data([1, 2, 3]))
|
||||||
|
#expect(decoded.types == nil)
|
||||||
|
#expect(decoded.sinceTimestamp == nil)
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test func decodesAndroidWantedTypesAndMinTimestamp() throws {
|
||||||
|
let payload = Data([
|
||||||
|
0x01, 0x00, 0x01, 0x08,
|
||||||
|
0x02, 0x00, 0x04, 0x00, 0x00, 0x64, 0x00,
|
||||||
|
0x03, 0x00, 0x02, 0xAA, 0xBB,
|
||||||
|
0x04, 0x00, 0x02, MessageType.announce.rawValue, MessageType.message.rawValue,
|
||||||
|
0x05, 0x00, 0x08, 0x00, 0x00, 0x01, 0x8B, 0xCF, 0xE5, 0x68, 0x00
|
||||||
|
])
|
||||||
|
|
||||||
|
let decoded = try #require(RequestSyncPacket.decode(from: payload))
|
||||||
|
|
||||||
|
#expect(decoded.p == 8)
|
||||||
|
#expect(decoded.m == 25_600)
|
||||||
|
#expect(decoded.data == Data([0xAA, 0xBB]))
|
||||||
|
#expect(decoded.types?.contains(.announce) == true)
|
||||||
|
#expect(decoded.types?.contains(.message) == true)
|
||||||
|
#expect(decoded.sinceTimestamp == 1_700_000_000_000)
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test func decodesLegacyIOSPublicMessageBitfield() throws {
|
||||||
|
let payload = Data([
|
||||||
|
0x01, 0x00, 0x01, 0x08,
|
||||||
|
0x02, 0x00, 0x04, 0x00, 0x00, 0x64, 0x00,
|
||||||
|
0x03, 0x00, 0x00,
|
||||||
|
0x04, 0x00, 0x01, 0x03
|
||||||
|
])
|
||||||
|
|
||||||
|
let decoded = try #require(RequestSyncPacket.decode(from: payload))
|
||||||
|
|
||||||
|
#expect(decoded.types?.contains(.announce) == true)
|
||||||
|
#expect(decoded.types?.contains(.message) == true)
|
||||||
|
}
|
||||||
|
|
||||||
|
private func tlvValue(type: UInt8, in data: Data) -> Data? {
|
||||||
|
var offset = 0
|
||||||
|
while offset + 3 <= data.count {
|
||||||
|
let currentType = data[offset]
|
||||||
|
offset += 1
|
||||||
|
let length = (Int(data[offset]) << 8) | Int(data[offset + 1])
|
||||||
|
offset += 2
|
||||||
|
guard offset + length <= data.count else { return nil }
|
||||||
|
let value = data.subdata(in: offset..<(offset + length))
|
||||||
|
offset += length
|
||||||
|
if currentType == type {
|
||||||
|
return value
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -30,7 +30,6 @@ let package = Package(
|
|||||||
"TorManager.swift",
|
"TorManager.swift",
|
||||||
"TorURLSession.swift",
|
"TorURLSession.swift",
|
||||||
"TorNotifications.swift",
|
"TorNotifications.swift",
|
||||||
"TorEgressVerifier.swift",
|
|
||||||
],
|
],
|
||||||
linkerSettings: [
|
linkerSettings: [
|
||||||
.linkedLibrary("resolv"),
|
.linkedLibrary("resolv"),
|
||||||
|
|||||||
@@ -1,365 +0,0 @@
|
|||||||
import BitLogger
|
|
||||||
import Foundation
|
|
||||||
|
|
||||||
/// Runtime self-check that the proxied `URLSession` egress is *actually* routed
|
|
||||||
/// through Tor — defense-in-depth for the case where a platform silently ignores
|
|
||||||
/// `URLSessionConfiguration.connectionProxyDictionary` SOCKS settings and lets
|
|
||||||
/// traffic egress directly (leaking the real IP while Tor appears enabled).
|
|
||||||
///
|
|
||||||
/// Runtime verification (see `scripts/tor-egress-verification/`) showed that
|
|
||||||
/// macOS and the iOS simulator DO honor the SOCKS proxy for both plain HTTPS and
|
|
||||||
/// `URLSessionWebSocketTask`, and that the proxied session is fail-closed (every
|
|
||||||
/// request errors when the SOCKS proxy is down). Apple does not officially
|
|
||||||
/// support SOCKS for URLSession on iOS, so on a physical device the behavior is
|
|
||||||
/// not contractually guaranteed. This verifier closes that gap: before relay
|
|
||||||
/// connections are opened under enforced Tor, it performs a canary request whose
|
|
||||||
/// response positively reports whether the egress hit the network via Tor.
|
|
||||||
///
|
|
||||||
/// Policy (`verify()` return value) — fail-closed on unverified egress:
|
|
||||||
/// - `.verifiedTor` → allow, and cache the positive result for `ttl`.
|
|
||||||
/// - `.notTor` → REFUSE, and drop any cached verification. The canary
|
|
||||||
/// reached the internet but the exit is NOT a Tor node:
|
|
||||||
/// a real leak. Never allow relays.
|
|
||||||
/// - `.unreachable` → REFUSE (egress unverified). The canary itself failed
|
|
||||||
/// (endpoint down / circuit not built), so we cannot tell
|
|
||||||
/// whether the platform honored the SOCKS proxy — the
|
|
||||||
/// exact ambiguity this verifier exists to resolve.
|
|
||||||
/// Unverified traffic must not proceed on the
|
|
||||||
/// enforced-Tor path.
|
|
||||||
///
|
|
||||||
/// TTL / retry semantics:
|
|
||||||
/// - A `verifiedTor` verdict allows connection *opens* for `ttl` without
|
|
||||||
/// re-probing, so a brief canary blip inside the TTL window does not take
|
|
||||||
/// relays offline (a fresh positive verdict is authoritative for the
|
|
||||||
/// window). Already-open sockets are never torn down by verification —
|
|
||||||
/// they were opened under a verified egress and the proxied session is
|
|
||||||
/// fail-closed by construction.
|
|
||||||
/// - After TTL expiry (or `invalidate()` on Tor restart/dormant/shutdown),
|
|
||||||
/// the next `verify()` re-probes; while the canary stays `.unreachable`,
|
|
||||||
/// new connection opens are refused until a probe succeeds again.
|
|
||||||
/// - Probe cadence is bounded: at most one probe per `minRetryInterval`
|
|
||||||
/// (callers within the window reuse the last decision), and concurrent
|
|
||||||
/// `verify()` calls share one in-flight probe. Recovery from a transient
|
|
||||||
/// canary outage is automatic: callers that keep retrying (relay connect
|
|
||||||
/// gate, GeoRelayDirectory backoff) re-probe and succeed once the canary
|
|
||||||
/// is reachable again.
|
|
||||||
///
|
|
||||||
/// Liveness:
|
|
||||||
/// - Every probe is hard-bounded by `probeTimeout` via an independent async
|
|
||||||
/// watchdog. The proxied session sets `waitsForConnectivity = true`, which
|
|
||||||
/// can defer the per-request timer indefinitely, leaving the request
|
|
||||||
/// bounded only by the default 7-day resource timeout — without the
|
|
||||||
/// watchdog a hung canary would wedge every subsequent `verify()` caller.
|
|
||||||
/// On timeout the probe task is cancelled (cooperatively cancelling the
|
|
||||||
/// underlying `URLSessionTask`), the verdict is the fail-closed
|
|
||||||
/// `.unreachable`, and the in-flight slot is cleared so the next
|
|
||||||
/// `verify()` (after the retry throttle) starts a fresh probe.
|
|
||||||
/// - `invalidate()` cancels any in-flight probe and clears all cached state
|
|
||||||
/// (including the retry throttle), so a Tor restart/dormant/shutdown
|
|
||||||
/// genuinely resets the verifier: a hung probe cannot survive it.
|
|
||||||
///
|
|
||||||
/// The probe is injectable so the policy/caching logic is unit-tested without a
|
|
||||||
/// live network (see `TorEgressVerifierTests`).
|
|
||||||
public actor TorEgressVerifier {
|
|
||||||
public enum ProbeResult: Equatable, Sendable {
|
|
||||||
/// Canary succeeded and the exit is a Tor node.
|
|
||||||
case verifiedTor
|
|
||||||
/// Canary succeeded but the exit is NOT Tor — a direct-egress leak.
|
|
||||||
case notTor
|
|
||||||
/// Canary could not complete (endpoint down, no circuit, parse error).
|
|
||||||
case unreachable(String)
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Hard upper bound for a single canary probe, enforced independently of
|
|
||||||
/// URLSession timers (see the "Liveness" section of the type doc). Matches
|
|
||||||
/// the live probe's per-request timeout.
|
|
||||||
public static let defaultProbeTimeout: TimeInterval = 20
|
|
||||||
|
|
||||||
private let probe: @Sendable () async -> ProbeResult
|
|
||||||
private let now: @Sendable () -> Date
|
|
||||||
private let ttl: TimeInterval
|
|
||||||
/// Minimum spacing between probes when not currently verified, so a
|
|
||||||
/// persistent `.unreachable` cannot hammer the canary endpoint on every
|
|
||||||
/// reconnect burst.
|
|
||||||
private let minRetryInterval: TimeInterval
|
|
||||||
/// Outer wall-clock bound on a single probe (watchdog; fail-closed).
|
|
||||||
private let probeTimeout: TimeInterval
|
|
||||||
|
|
||||||
private var lastVerifiedAt: Date?
|
|
||||||
private var lastProbeAt: Date?
|
|
||||||
private var lastResult: ProbeResult?
|
|
||||||
private var inFlight: Task<Bool, Never>?
|
|
||||||
/// Bumped whenever a new probe starts or `invalidate()` runs. A completing
|
|
||||||
/// probe only records its outcome (cache/throttle) and clears `inFlight`
|
|
||||||
/// if its generation is still current, so a cancelled/superseded probe
|
|
||||||
/// cannot clobber state owned by a fresh one (actor-reentrancy safety).
|
|
||||||
private var probeGeneration = 0
|
|
||||||
|
|
||||||
/// Lock-protected mirror of "verified within TTL" so synchronous gates
|
|
||||||
/// (e.g. `NostrRelayManager`'s connect path) can consult the cache without
|
|
||||||
/// awaiting the actor.
|
|
||||||
private let verifiedSnapshot = VerifiedSnapshot()
|
|
||||||
|
|
||||||
private final class VerifiedSnapshot: @unchecked Sendable {
|
|
||||||
private let lock = NSLock()
|
|
||||||
private var verifiedUntil: Date?
|
|
||||||
|
|
||||||
func update(_ until: Date?) {
|
|
||||||
lock.lock()
|
|
||||||
verifiedUntil = until
|
|
||||||
lock.unlock()
|
|
||||||
}
|
|
||||||
|
|
||||||
func isFresh(at date: Date) -> Bool {
|
|
||||||
lock.lock()
|
|
||||||
defer { lock.unlock() }
|
|
||||||
guard let verifiedUntil else { return false }
|
|
||||||
return date < verifiedUntil
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
public init(
|
|
||||||
ttl: TimeInterval,
|
|
||||||
minRetryInterval: TimeInterval = 5.0,
|
|
||||||
probeTimeout: TimeInterval = TorEgressVerifier.defaultProbeTimeout,
|
|
||||||
now: @escaping @Sendable () -> Date = Date.init,
|
|
||||||
probe: @escaping @Sendable () async -> ProbeResult
|
|
||||||
) {
|
|
||||||
self.ttl = ttl
|
|
||||||
self.minRetryInterval = minRetryInterval
|
|
||||||
self.probeTimeout = probeTimeout
|
|
||||||
self.now = now
|
|
||||||
self.probe = probe
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Drop any cached verification (e.g. after a Tor restart or when the
|
|
||||||
/// network path changes) AND cancel any in-flight probe. The next
|
|
||||||
/// `verify()` starts a fresh probe — it neither joins the cancelled one
|
|
||||||
/// nor is throttled by its outcome, so a probe hung from before a Tor
|
|
||||||
/// restart cannot wedge callers after it.
|
|
||||||
public func invalidate() {
|
|
||||||
probeGeneration += 1
|
|
||||||
inFlight?.cancel()
|
|
||||||
inFlight = nil
|
|
||||||
lastVerifiedAt = nil
|
|
||||||
lastProbeAt = nil
|
|
||||||
lastResult = nil
|
|
||||||
verifiedSnapshot.update(nil)
|
|
||||||
}
|
|
||||||
|
|
||||||
/// The most recent probe outcome, for diagnostics/tests.
|
|
||||||
public func lastProbeResult() -> ProbeResult? { lastResult }
|
|
||||||
|
|
||||||
/// Synchronous view of the cache: `true` while a `verifiedTor` verdict is
|
|
||||||
/// within its TTL. Callers that get `false` must route through the async
|
|
||||||
/// `verify()` gate (which probes) before opening connections.
|
|
||||||
public nonisolated var hasFreshVerification: Bool {
|
|
||||||
verifiedSnapshot.isFresh(at: now())
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Returns `true` only when the proxied egress is verified to exit via Tor
|
|
||||||
/// (a fresh probe or a cached `verifiedTor` verdict within TTL). Returns
|
|
||||||
/// `false` when a non-Tor egress was positively detected *or* when the
|
|
||||||
/// egress could not be verified. See the type doc for the full policy.
|
|
||||||
public func verify() async -> Bool {
|
|
||||||
if isFreshlyVerified() { return true }
|
|
||||||
// Throttle re-probes when the last attempt did not verify.
|
|
||||||
if let last = lastProbeAt,
|
|
||||||
let result = lastResult,
|
|
||||||
now().timeIntervalSince(last) < minRetryInterval {
|
|
||||||
return decision(for: result)
|
|
||||||
}
|
|
||||||
if let inFlight { return await inFlight.value }
|
|
||||||
|
|
||||||
probeGeneration += 1
|
|
||||||
let generation = probeGeneration
|
|
||||||
let task = Task<Bool, Never> { await self.runProbe(generation: generation) }
|
|
||||||
inFlight = task
|
|
||||||
let allowed = await task.value
|
|
||||||
// Only clear the slot if this probe is still the current one: an
|
|
||||||
// `invalidate()` while we were suspended has already cleared it and a
|
|
||||||
// newer probe may occupy it (do not clobber the fresh task).
|
|
||||||
if probeGeneration == generation { inFlight = nil }
|
|
||||||
return allowed
|
|
||||||
}
|
|
||||||
|
|
||||||
private func isFreshlyVerified() -> Bool {
|
|
||||||
guard let last = lastVerifiedAt else { return false }
|
|
||||||
return now().timeIntervalSince(last) < ttl
|
|
||||||
}
|
|
||||||
|
|
||||||
private func decision(for result: ProbeResult) -> Bool {
|
|
||||||
switch result {
|
|
||||||
case .verifiedTor: return true
|
|
||||||
// Fail closed: both a positively detected leak and an unverifiable
|
|
||||||
// egress refuse connections. Only a fresh `verifiedTor` allows.
|
|
||||||
case .unreachable, .notTor: return false
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
private func runProbe(generation: Int) async -> Bool {
|
|
||||||
let result = await boundedProbe()
|
|
||||||
// Superseded by `invalidate()` (Tor restart/dormant/shutdown) while the
|
|
||||||
// probe ran: its verdict predates the reset, so discard it — recording
|
|
||||||
// it would re-seed the throttle/cache that invalidate() just cleared.
|
|
||||||
// Fail closed for the callers that were awaiting this probe.
|
|
||||||
guard generation == probeGeneration else { return false }
|
|
||||||
lastProbeAt = now()
|
|
||||||
lastResult = result
|
|
||||||
switch result {
|
|
||||||
case .verifiedTor:
|
|
||||||
lastVerifiedAt = now()
|
|
||||||
verifiedSnapshot.update(now().addingTimeInterval(ttl))
|
|
||||||
return true
|
|
||||||
case .notTor:
|
|
||||||
lastVerifiedAt = nil
|
|
||||||
verifiedSnapshot.update(nil)
|
|
||||||
SecureLogger.error(
|
|
||||||
"🧅 Tor egress self-check FAILED: request exited via a NON-Tor address — refusing relay connections (possible IP leak)",
|
|
||||||
category: .session
|
|
||||||
)
|
|
||||||
return false
|
|
||||||
case .unreachable(let why):
|
|
||||||
// Note: a probe only runs when no fresh cached verdict exists, so
|
|
||||||
// there is no still-valid cache to preserve or drop here.
|
|
||||||
SecureLogger.warning(
|
|
||||||
"🧅 Tor egress self-check could not complete (\(why)) — egress UNVERIFIED; refusing relay connections until the canary succeeds (bounded retry)",
|
|
||||||
category: .session
|
|
||||||
)
|
|
||||||
return false
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Runs the injected probe raced against `probeTimeout`, guaranteeing a
|
|
||||||
/// result in bounded time regardless of URLSession timer behavior (the
|
|
||||||
/// proxied session's `waitsForConnectivity` can defer the per-request
|
|
||||||
/// timeout indefinitely). Whichever side loses the race is cancelled:
|
|
||||||
/// - on timeout, the probe task is cancelled (URLSession's async APIs
|
|
||||||
/// cancel the underlying `URLSessionTask` cooperatively) and the result
|
|
||||||
/// is the fail-closed `.unreachable`;
|
|
||||||
/// - on completion, the watchdog's sleep is cancelled so no timer lingers.
|
|
||||||
/// Cancelling the enclosing task (`invalidate()`) resolves immediately as
|
|
||||||
/// `.unreachable` and cancels both sides.
|
|
||||||
///
|
|
||||||
/// `nonisolated` so the race body never re-enters the actor; it touches
|
|
||||||
/// only immutable `Sendable` state.
|
|
||||||
nonisolated private func boundedProbe() async -> ProbeResult {
|
|
||||||
let probe = self.probe
|
|
||||||
let timeout = self.probeTimeout
|
|
||||||
let race = ProbeRace()
|
|
||||||
return await withTaskCancellationHandler {
|
|
||||||
await withCheckedContinuation { (continuation: CheckedContinuation<ProbeResult, Never>) in
|
|
||||||
race.install(continuation)
|
|
||||||
let probeTask = Task { race.finish(await probe()) }
|
|
||||||
let watchdog = Task {
|
|
||||||
try? await Task.sleep(nanoseconds: UInt64(max(0, timeout) * 1_000_000_000))
|
|
||||||
guard !Task.isCancelled else { return }
|
|
||||||
race.finish(.unreachable("probe timed out after \(Int(timeout))s"))
|
|
||||||
}
|
|
||||||
race.register(probeTask: probeTask, watchdog: watchdog)
|
|
||||||
}
|
|
||||||
} onCancel: {
|
|
||||||
race.finish(.unreachable("probe cancelled"))
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Resolve-once rendezvous for the probe/watchdog race. Lock-protected
|
|
||||||
/// (never held across an await); the first `finish()` wins, resumes the
|
|
||||||
/// continuation exactly once, and cancels both tasks.
|
|
||||||
private final class ProbeRace: @unchecked Sendable {
|
|
||||||
private let lock = NSLock()
|
|
||||||
private var continuation: CheckedContinuation<ProbeResult, Never>?
|
|
||||||
private var pendingResult: ProbeResult?
|
|
||||||
private var resolved = false
|
|
||||||
private var probeTask: Task<Void, Never>?
|
|
||||||
private var watchdog: Task<Void, Never>?
|
|
||||||
|
|
||||||
func install(_ continuation: CheckedContinuation<ProbeResult, Never>) {
|
|
||||||
lock.lock()
|
|
||||||
if let result = pendingResult {
|
|
||||||
// finish() ran before the continuation existed (e.g. the
|
|
||||||
// enclosing task was already cancelled): resolve immediately.
|
|
||||||
pendingResult = nil
|
|
||||||
lock.unlock()
|
|
||||||
continuation.resume(returning: result)
|
|
||||||
return
|
|
||||||
}
|
|
||||||
self.continuation = continuation
|
|
||||||
lock.unlock()
|
|
||||||
}
|
|
||||||
|
|
||||||
func register(probeTask: Task<Void, Never>, watchdog: Task<Void, Never>) {
|
|
||||||
lock.lock()
|
|
||||||
if resolved {
|
|
||||||
lock.unlock()
|
|
||||||
probeTask.cancel()
|
|
||||||
watchdog.cancel()
|
|
||||||
return
|
|
||||||
}
|
|
||||||
self.probeTask = probeTask
|
|
||||||
self.watchdog = watchdog
|
|
||||||
lock.unlock()
|
|
||||||
}
|
|
||||||
|
|
||||||
func finish(_ result: ProbeResult) {
|
|
||||||
lock.lock()
|
|
||||||
guard !resolved else { lock.unlock(); return }
|
|
||||||
resolved = true
|
|
||||||
let continuation = self.continuation
|
|
||||||
self.continuation = nil
|
|
||||||
if continuation == nil { pendingResult = result }
|
|
||||||
let probeTask = self.probeTask
|
|
||||||
let watchdog = self.watchdog
|
|
||||||
self.probeTask = nil
|
|
||||||
self.watchdog = nil
|
|
||||||
lock.unlock()
|
|
||||||
probeTask?.cancel()
|
|
||||||
watchdog?.cancel()
|
|
||||||
continuation?.resume(returning: result)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// MARK: - Live probe
|
|
||||||
|
|
||||||
public extension TorEgressVerifier {
|
|
||||||
/// Default canary: fetch Tor Project's connectivity check API through the
|
|
||||||
/// shared proxied session and assert `IsTor == true`. Because the response
|
|
||||||
/// is served from the *exit's* vantage point, a silent direct egress is
|
|
||||||
/// caught here as `.notTor`. `check.torproject.org` is clearnet, so this
|
|
||||||
/// works without onion-service support.
|
|
||||||
///
|
|
||||||
/// Follow-up (see PR): make the canary endpoint configurable and add an
|
|
||||||
/// onion-service canary so verification does not depend on a single host.
|
|
||||||
/// Note: the per-request `timeoutInterval` below is best-effort only — the
|
|
||||||
/// proxied session's `waitsForConnectivity` can defer it. The authoritative
|
|
||||||
/// bound is the verifier's `probeTimeout` watchdog, whose cancellation
|
|
||||||
/// propagates into `session.data(for:)` and cancels the URLSessionTask.
|
|
||||||
static func liveProbe(
|
|
||||||
endpoint: URL = URL(string: "https://check.torproject.org/api/ip")!,
|
|
||||||
timeout: TimeInterval = TorEgressVerifier.defaultProbeTimeout
|
|
||||||
) -> @Sendable () async -> ProbeResult {
|
|
||||||
return {
|
|
||||||
var request = URLRequest(url: endpoint)
|
|
||||||
request.timeoutInterval = timeout
|
|
||||||
request.cachePolicy = .reloadIgnoringLocalAndRemoteCacheData
|
|
||||||
let session = TorURLSession.shared.session
|
|
||||||
do {
|
|
||||||
let (data, response) = try await session.data(for: request)
|
|
||||||
guard let http = response as? HTTPURLResponse,
|
|
||||||
(200..<300).contains(http.statusCode) else {
|
|
||||||
return .unreachable("http status \((response as? HTTPURLResponse)?.statusCode ?? -1)")
|
|
||||||
}
|
|
||||||
guard let json = try? JSONSerialization.jsonObject(with: data) as? [String: Any] else {
|
|
||||||
return .unreachable("unparseable canary response")
|
|
||||||
}
|
|
||||||
if let isTor = json["IsTor"] as? Bool {
|
|
||||||
return isTor ? .verifiedTor : .notTor
|
|
||||||
}
|
|
||||||
return .unreachable("canary response missing IsTor")
|
|
||||||
} catch {
|
|
||||||
return .unreachable(error.localizedDescription)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
@@ -59,14 +59,6 @@ public final class TorManager: ObservableObject {
|
|||||||
private var socksReady: Bool = false { didSet { recomputeReady() } }
|
private var socksReady: Bool = false { didSet { recomputeReady() } }
|
||||||
private var restarting: Bool = false
|
private var restarting: Bool = false
|
||||||
|
|
||||||
/// Runtime egress self-check: proves the proxied session actually exits via
|
|
||||||
/// Tor before relay connections are opened (defense-in-depth against a
|
|
||||||
/// platform silently ignoring the SOCKS proxy). Cached for a few minutes.
|
|
||||||
public let egressVerifier = TorEgressVerifier(
|
|
||||||
ttl: 300,
|
|
||||||
probe: TorEgressVerifier.liveProbe()
|
|
||||||
)
|
|
||||||
|
|
||||||
// Whether the app must enforce Tor for all connections (fail-closed).
|
// Whether the app must enforce Tor for all connections (fail-closed).
|
||||||
public var torEnforced: Bool {
|
public var torEnforced: Bool {
|
||||||
#if BITCHAT_DEV_ALLOW_CLEARNET
|
#if BITCHAT_DEV_ALLOW_CLEARNET
|
||||||
@@ -133,32 +125,6 @@ public final class TorManager: ObservableObject {
|
|||||||
return await MainActor.run(body: { self.networkPermitted })
|
return await MainActor.run(body: { self.networkPermitted })
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Synchronous, cached view of the egress gate: `true` while a positive
|
|
||||||
/// egress verification is within its TTL (or when Tor is not enforced).
|
|
||||||
/// When this is `false`, callers must route through `awaitEgressReady()`
|
|
||||||
/// (which probes) before opening any connection — never connect directly.
|
|
||||||
public var isEgressVerified: Bool {
|
|
||||||
guard torEnforced else { return true }
|
|
||||||
return egressVerifier.hasFreshVerification
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Like `awaitReady`, but additionally requires that a canary request
|
|
||||||
/// through the proxied session positively verifies Tor egress. Returns
|
|
||||||
/// `false` if Tor never became ready, or if the egress self-check could
|
|
||||||
/// not positively verify a Tor exit (non-Tor egress detected, or canary
|
|
||||||
/// unreachable → unverified). Callers must fail closed on `false` — never
|
|
||||||
/// fall back to a direct connection.
|
|
||||||
nonisolated
|
|
||||||
public func awaitEgressReady(timeout: TimeInterval = 75.0) async -> Bool {
|
|
||||||
let ready = await awaitReady(timeout: timeout)
|
|
||||||
guard ready else { return false }
|
|
||||||
// Clearnet dev builds don't route through Tor, so the canary would
|
|
||||||
// (correctly) report non-Tor; skip it there.
|
|
||||||
let enforced = await MainActor.run { self.torEnforced }
|
|
||||||
guard enforced else { return true }
|
|
||||||
return await egressVerifier.verify()
|
|
||||||
}
|
|
||||||
|
|
||||||
// MARK: - Filesystem
|
// MARK: - Filesystem
|
||||||
|
|
||||||
func dataDirectoryURL() -> URL? {
|
func dataDirectoryURL() -> URL? {
|
||||||
@@ -359,8 +325,6 @@ public final class TorManager: ObservableObject {
|
|||||||
self.socksReady = false
|
self.socksReady = false
|
||||||
self.isStarting = false
|
self.isStarting = false
|
||||||
}
|
}
|
||||||
// Force a fresh egress self-check once Tor comes back.
|
|
||||||
Task { await egressVerifier.invalidate() }
|
|
||||||
}
|
}
|
||||||
|
|
||||||
public func shutdownCompletely() {
|
public func shutdownCompletely() {
|
||||||
@@ -389,7 +353,6 @@ public final class TorManager: ObservableObject {
|
|||||||
// Note: Don't clear startedAt here - it will be set fresh on next startIfNeeded()
|
// Note: Don't clear startedAt here - it will be set fresh on next startIfNeeded()
|
||||||
// Clearing it here races with startup and defeats the grace period
|
// Clearing it here races with startup and defeats the grace period
|
||||||
}
|
}
|
||||||
await self.egressVerifier.invalidate()
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -405,8 +368,6 @@ public final class TorManager: ObservableObject {
|
|||||||
self.isDormant = false
|
self.isDormant = false
|
||||||
self.lastRestartAt = Date()
|
self.lastRestartAt = Date()
|
||||||
}
|
}
|
||||||
// New Arti instance means new circuits; re-verify egress after restart.
|
|
||||||
await egressVerifier.invalidate()
|
|
||||||
|
|
||||||
_ = arti_stop()
|
_ = arti_stop()
|
||||||
|
|
||||||
|
|||||||
@@ -1,90 +0,0 @@
|
|||||||
// Standalone harness: does Apple's URLSession honor connectionProxyDictionary
|
|
||||||
// SOCKS settings for a plain HTTPS GET and for URLSessionWebSocketTask?
|
|
||||||
//
|
|
||||||
// Build: swiftc -O proxy_probe.swift -o proxy_probe
|
|
||||||
// Usage: proxy_probe <http|ws> <cf|raw> <proxyPort> [targetURL]
|
|
||||||
//
|
|
||||||
// Prints a single RESULT line: RESULT <mode> <keyStyle> <outcome> <detail>
|
|
||||||
// The caller correlates this with the SOCKS proxy's connection log to decide
|
|
||||||
// whether the request was proxied.
|
|
||||||
#if canImport(CFNetwork)
|
|
||||||
import CFNetwork
|
|
||||||
#endif
|
|
||||||
import Foundation
|
|
||||||
|
|
||||||
let args = CommandLine.arguments
|
|
||||||
guard args.count >= 4 else {
|
|
||||||
FileHandle.standardError.write(Data("usage: proxy_probe <http|ws> <cf|raw> <proxyPort> [targetURL]\n".utf8))
|
|
||||||
exit(2)
|
|
||||||
}
|
|
||||||
let mode = args[1]
|
|
||||||
let keyStyle = args[2]
|
|
||||||
let proxyPort = Int(args[3]) ?? 19999
|
|
||||||
let host = "127.0.0.1"
|
|
||||||
|
|
||||||
func makeProxyDict() -> [AnyHashable: Any] {
|
|
||||||
switch keyStyle {
|
|
||||||
#if os(macOS)
|
|
||||||
case "cf":
|
|
||||||
// The exact constants the app uses on macOS.
|
|
||||||
return [
|
|
||||||
kCFNetworkProxiesSOCKSEnable as String: 1,
|
|
||||||
kCFNetworkProxiesSOCKSProxy as String: host,
|
|
||||||
kCFNetworkProxiesSOCKSPort as String: proxyPort
|
|
||||||
]
|
|
||||||
#endif
|
|
||||||
default:
|
|
||||||
// The exact raw string keys the app uses on iOS.
|
|
||||||
return [
|
|
||||||
"SOCKSEnable": 1,
|
|
||||||
"SOCKSProxy": host,
|
|
||||||
"SOCKSPort": proxyPort
|
|
||||||
]
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
let cfg = URLSessionConfiguration.ephemeral
|
|
||||||
cfg.waitsForConnectivity = false
|
|
||||||
cfg.timeoutIntervalForRequest = 20
|
|
||||||
cfg.connectionProxyDictionary = makeProxyDict()
|
|
||||||
let session = URLSession(configuration: cfg)
|
|
||||||
|
|
||||||
func emit(_ outcome: String, _ detail: String) {
|
|
||||||
print("RESULT \(mode) \(keyStyle) \(outcome) \(detail)")
|
|
||||||
exit(outcome == "ERROR" ? 1 : 0)
|
|
||||||
}
|
|
||||||
|
|
||||||
let sem = DispatchSemaphore(value: 0)
|
|
||||||
|
|
||||||
if mode == "http" {
|
|
||||||
let target = URL(string: args.count >= 5 ? args[4] : "https://raw.githubusercontent.com/permissionlesstech/georelays/refs/heads/main/nostr_relays.csv")!
|
|
||||||
let task = session.dataTask(with: target) { data, resp, err in
|
|
||||||
if let err = err {
|
|
||||||
emit("ERROR", "\(err.localizedDescription)")
|
|
||||||
} else if let http = resp as? HTTPURLResponse {
|
|
||||||
emit("OK", "status=\(http.statusCode) bytes=\(data?.count ?? 0)")
|
|
||||||
} else {
|
|
||||||
emit("OK", "bytes=\(data?.count ?? 0)")
|
|
||||||
}
|
|
||||||
}
|
|
||||||
task.resume()
|
|
||||||
} else {
|
|
||||||
// WebSocket
|
|
||||||
let target = URL(string: args.count >= 5 ? args[4] : "wss://relay.damus.io")!
|
|
||||||
let ws = session.webSocketTask(with: target)
|
|
||||||
ws.resume()
|
|
||||||
// A successful ping proves the TLS+WS handshake completed end-to-end.
|
|
||||||
ws.sendPing { err in
|
|
||||||
if let err = err {
|
|
||||||
emit("ERROR", "\(err.localizedDescription)")
|
|
||||||
} else {
|
|
||||||
emit("OK", "ws-ping-ok")
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// Global watchdog so we never hang.
|
|
||||||
DispatchQueue.global().asyncAfter(deadline: .now() + 25) {
|
|
||||||
emit("ERROR", "timeout")
|
|
||||||
}
|
|
||||||
sem.wait()
|
|
||||||
@@ -1,73 +0,0 @@
|
|||||||
#!/usr/bin/env bash
|
|
||||||
# Orchestrates the Tor-egress proxy-honoring verification on macOS.
|
|
||||||
#
|
|
||||||
# For each (request-type x key-style) it runs two experiments:
|
|
||||||
# A) proxy UP — did a connection arrive at the SOCKS proxy? (log grows)
|
|
||||||
# B) proxy DOWN — pointed at a dead port; does the request still SUCCEED?
|
|
||||||
# If it succeeds with no proxy, egress went DIRECT (proxy ignored).
|
|
||||||
# If it fails, the proxy setting is being enforced (fail-closed).
|
|
||||||
#
|
|
||||||
# Discriminator: PROXIED = connection observed at proxy AND fails when proxy down
|
|
||||||
# DIRECT = no connection at proxy OR succeeds when proxy down
|
|
||||||
set -u
|
|
||||||
DIR="$(cd "$(dirname "$0")" && pwd)"
|
|
||||||
PORT=19999
|
|
||||||
DEADPORT=19998 # nothing listens here
|
|
||||||
LOG="$(mktemp -t sockslog)"
|
|
||||||
BIN="$(mktemp -t proxyprobe)"
|
|
||||||
|
|
||||||
echo "== building swift probe =="
|
|
||||||
swiftc -O "$DIR/proxy_probe.swift" -o "$BIN" || { echo "swiftc failed"; exit 1; }
|
|
||||||
|
|
||||||
echo "== starting SOCKS proxy on $PORT =="
|
|
||||||
: > "$LOG"
|
|
||||||
python3 "$DIR/socks5_probe_proxy.py" "$PORT" "$LOG" >/tmp/socksproxy.out 2>&1 &
|
|
||||||
PROXY_PID=$!
|
|
||||||
trap 'kill $PROXY_PID 2>/dev/null' EXIT
|
|
||||||
# wait for READY
|
|
||||||
for _ in $(seq 1 50); do
|
|
||||||
grep -q READY /tmp/socksproxy.out 2>/dev/null && break
|
|
||||||
sleep 0.1
|
|
||||||
done
|
|
||||||
|
|
||||||
run_case() {
|
|
||||||
local mode="$1" key="$2"
|
|
||||||
# Experiment A: proxy up, watch log
|
|
||||||
local before after target
|
|
||||||
before=$(wc -l < "$LOG" | tr -d ' ')
|
|
||||||
local outA
|
|
||||||
outA=$("$BIN" "$mode" "$key" "$PORT" 2>/dev/null | grep '^RESULT' || echo "RESULT $mode $key ERROR no-output")
|
|
||||||
sleep 0.3
|
|
||||||
after=$(wc -l < "$LOG" | tr -d ' ')
|
|
||||||
local proxied="NO"
|
|
||||||
if [ "$after" -gt "$before" ]; then proxied="YES"; fi
|
|
||||||
local newlines
|
|
||||||
newlines=$(tail -n +"$((before+1))" "$LOG" | tr '\t' ' ' | tr '\n' '|')
|
|
||||||
|
|
||||||
# Experiment B: proxy down (dead port), same request
|
|
||||||
local outB
|
|
||||||
outB=$("$BIN" "$mode" "$key" "$DEADPORT" 2>/dev/null | grep '^RESULT' || echo "RESULT $mode $key ERROR no-output")
|
|
||||||
|
|
||||||
echo "----------------------------------------"
|
|
||||||
echo "CASE mode=$mode key=$key"
|
|
||||||
echo " A(proxy up): $outA | connection_at_proxy=$proxied [$newlines]"
|
|
||||||
echo " B(proxy down): $outB"
|
|
||||||
# verdict
|
|
||||||
local a_ok b_ok
|
|
||||||
a_ok=$(echo "$outA" | awk '{print $4}')
|
|
||||||
b_ok=$(echo "$outB" | awk '{print $4}')
|
|
||||||
local verdict="UNKNOWN"
|
|
||||||
if [ "$proxied" = "YES" ] && [ "$b_ok" = "ERROR" ]; then verdict="PROXIED (enforced)"; fi
|
|
||||||
if [ "$proxied" = "NO" ] && [ "$b_ok" = "OK" ]; then verdict="DIRECT (proxy ignored)"; fi
|
|
||||||
if [ "$proxied" = "YES" ] && [ "$b_ok" = "OK" ]; then verdict="AMBIGUOUS (uses proxy if up, but egresses direct if down)"; fi
|
|
||||||
if [ "$proxied" = "NO" ] && [ "$b_ok" = "ERROR" ]; then verdict="BLOCKED both (network/target issue?)"; fi
|
|
||||||
echo " VERDICT: $verdict"
|
|
||||||
}
|
|
||||||
|
|
||||||
for mode in http ws; do
|
|
||||||
for key in cf raw; do
|
|
||||||
run_case "$mode" "$key"
|
|
||||||
done
|
|
||||||
done
|
|
||||||
echo "========================================"
|
|
||||||
echo "raw proxy log:"; cat "$LOG" | tr '\t' ' '
|
|
||||||
@@ -1,166 +0,0 @@
|
|||||||
#!/usr/bin/env python3
|
|
||||||
"""
|
|
||||||
Minimal threaded SOCKS5 CONNECT proxy used to verify whether Apple's
|
|
||||||
URLSession actually honors `connectionProxyDictionary` SOCKS settings for
|
|
||||||
different request types (plain HTTPS vs URLSessionWebSocketTask).
|
|
||||||
|
|
||||||
Behavior:
|
|
||||||
- Speaks enough SOCKS5 (no-auth) to complete a CONNECT and then relays
|
|
||||||
bytes bidirectionally to the real destination.
|
|
||||||
- Every accepted CONNECT is appended to a log file as one line:
|
|
||||||
<iso8601>\tCONNECT\t<host>:<port>
|
|
||||||
- Any raw connection that is NOT valid SOCKS5 is logged as:
|
|
||||||
<iso8601>\tNON_SOCKS\t<first-bytes-hex>
|
|
||||||
(this catches the feared case where URLSession sends a raw TLS/HTTP
|
|
||||||
ClientHello straight at the proxy port instead of a SOCKS greeting).
|
|
||||||
|
|
||||||
If a request egresses DIRECTLY (proxy ignored), nothing is logged at all.
|
|
||||||
|
|
||||||
Usage: socks5_probe_proxy.py <listen_port> <log_file>
|
|
||||||
"""
|
|
||||||
import selectors
|
|
||||||
import socket
|
|
||||||
import sys
|
|
||||||
import threading
|
|
||||||
from datetime import datetime, timezone
|
|
||||||
|
|
||||||
LOG_LOCK = threading.Lock()
|
|
||||||
|
|
||||||
|
|
||||||
def log(logfile, kind, detail):
|
|
||||||
line = f"{datetime.now(timezone.utc).isoformat()}\t{kind}\t{detail}\n"
|
|
||||||
with LOG_LOCK:
|
|
||||||
with open(logfile, "a") as f:
|
|
||||||
f.write(line)
|
|
||||||
sys.stderr.write("[proxy] " + line)
|
|
||||||
sys.stderr.flush()
|
|
||||||
|
|
||||||
|
|
||||||
def recv_exact(sock, n):
|
|
||||||
buf = b""
|
|
||||||
while len(buf) < n:
|
|
||||||
chunk = sock.recv(n - len(buf))
|
|
||||||
if not chunk:
|
|
||||||
return None
|
|
||||||
buf += chunk
|
|
||||||
return buf
|
|
||||||
|
|
||||||
|
|
||||||
def handle(client, logfile):
|
|
||||||
client.settimeout(15)
|
|
||||||
try:
|
|
||||||
# SOCKS5 greeting: VER=0x05, NMETHODS, METHODS...
|
|
||||||
head = recv_exact(client, 2)
|
|
||||||
if not head:
|
|
||||||
return
|
|
||||||
if head[0] != 0x05:
|
|
||||||
# Not SOCKS5 at all — this is the smoking gun for a direct egress
|
|
||||||
# that mistakenly hit the proxy port. Log the first bytes.
|
|
||||||
rest = b""
|
|
||||||
try:
|
|
||||||
client.setblocking(False)
|
|
||||||
rest = client.recv(64)
|
|
||||||
except Exception:
|
|
||||||
pass
|
|
||||||
log(logfile, "NON_SOCKS", (head + rest).hex())
|
|
||||||
return
|
|
||||||
nmethods = head[1]
|
|
||||||
if nmethods:
|
|
||||||
recv_exact(client, nmethods)
|
|
||||||
# Reply: no authentication required
|
|
||||||
client.sendall(b"\x05\x00")
|
|
||||||
|
|
||||||
# Request: VER, CMD, RSV, ATYP, ADDR, PORT
|
|
||||||
req = recv_exact(client, 4)
|
|
||||||
if not req or req[1] != 0x01: # only CONNECT
|
|
||||||
client.sendall(b"\x05\x07\x00\x01\x00\x00\x00\x00\x00\x00")
|
|
||||||
return
|
|
||||||
atyp = req[3]
|
|
||||||
if atyp == 0x01: # IPv4
|
|
||||||
addr = socket.inet_ntoa(recv_exact(client, 4))
|
|
||||||
elif atyp == 0x03: # domain
|
|
||||||
ln = recv_exact(client, 1)[0]
|
|
||||||
addr = recv_exact(client, ln).decode("ascii", errors="replace")
|
|
||||||
elif atyp == 0x04: # IPv6
|
|
||||||
addr = socket.inet_ntop(socket.AF_INET6, recv_exact(client, 16))
|
|
||||||
else:
|
|
||||||
client.sendall(b"\x05\x08\x00\x01\x00\x00\x00\x00\x00\x00")
|
|
||||||
return
|
|
||||||
port = int.from_bytes(recv_exact(client, 2), "big")
|
|
||||||
|
|
||||||
log(logfile, "CONNECT", f"{addr}:{port}")
|
|
||||||
|
|
||||||
# Connect to the real destination and reply success.
|
|
||||||
try:
|
|
||||||
remote = socket.create_connection((addr, port), timeout=15)
|
|
||||||
except Exception as e:
|
|
||||||
log(logfile, "CONNECT_FAIL", f"{addr}:{port} {e}")
|
|
||||||
client.sendall(b"\x05\x01\x00\x01\x00\x00\x00\x00\x00\x00")
|
|
||||||
return
|
|
||||||
client.sendall(b"\x05\x00\x00\x01\x00\x00\x00\x00\x00\x00")
|
|
||||||
|
|
||||||
relay(client, remote)
|
|
||||||
except Exception:
|
|
||||||
pass
|
|
||||||
finally:
|
|
||||||
try:
|
|
||||||
client.close()
|
|
||||||
except Exception:
|
|
||||||
pass
|
|
||||||
|
|
||||||
|
|
||||||
def relay(a, b):
|
|
||||||
a.setblocking(False)
|
|
||||||
b.setblocking(False)
|
|
||||||
sel = selectors.DefaultSelector()
|
|
||||||
sel.register(a, selectors.EVENT_READ, b)
|
|
||||||
sel.register(b, selectors.EVENT_READ, a)
|
|
||||||
try:
|
|
||||||
while True:
|
|
||||||
events = sel.select(timeout=30)
|
|
||||||
if not events:
|
|
||||||
break
|
|
||||||
for key, _ in events:
|
|
||||||
src = key.fileobj
|
|
||||||
dst = key.data
|
|
||||||
try:
|
|
||||||
data = src.recv(65536)
|
|
||||||
except (BlockingIOError, InterruptedError):
|
|
||||||
continue
|
|
||||||
except Exception:
|
|
||||||
return
|
|
||||||
if not data:
|
|
||||||
return
|
|
||||||
try:
|
|
||||||
dst.sendall(data)
|
|
||||||
except Exception:
|
|
||||||
return
|
|
||||||
finally:
|
|
||||||
sel.close()
|
|
||||||
for s in (a, b):
|
|
||||||
try:
|
|
||||||
s.close()
|
|
||||||
except Exception:
|
|
||||||
pass
|
|
||||||
|
|
||||||
|
|
||||||
def main():
|
|
||||||
if len(sys.argv) != 3:
|
|
||||||
print("usage: socks5_probe_proxy.py <port> <logfile>", file=sys.stderr)
|
|
||||||
sys.exit(2)
|
|
||||||
port = int(sys.argv[1])
|
|
||||||
logfile = sys.argv[2]
|
|
||||||
srv = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
|
|
||||||
srv.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
|
|
||||||
srv.bind(("127.0.0.1", port))
|
|
||||||
srv.listen(64)
|
|
||||||
sys.stderr.write(f"[proxy] listening on 127.0.0.1:{port}, log={logfile}\n")
|
|
||||||
sys.stderr.flush()
|
|
||||||
print("READY", flush=True)
|
|
||||||
while True:
|
|
||||||
client, _ = srv.accept()
|
|
||||||
threading.Thread(target=handle, args=(client, logfile), daemon=True).start()
|
|
||||||
|
|
||||||
|
|
||||||
if __name__ == "__main__":
|
|
||||||
main()
|
|
||||||
Reference in New Issue
Block a user