Compare commits

..
Author SHA1 Message Date
Codex 3c06c9e793 Allow manual CI dispatch 2026-06-26 00:31:25 +02:00
Codex 930354194f Use Xcode Swift in CI 2026-06-26 00:28:54 +02:00
Codex 15293f6b4b Salt SwiftPM cache by toolchain 2026-06-26 00:21:10 +02:00
Codex 842d8f2cd3 Align REQUEST_SYNC filters with Android 2026-06-25 23:09:35 +02:00
16 changed files with 514 additions and 699 deletions
+13 -10
View File
@@ -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
+3
View File
@@ -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
+41 -240
View File
@@ -48,14 +48,7 @@ private struct URLSessionAdapter: NostrRelaySessionProtocol {
let base: URLSession let base: URLSession
func webSocketTask(with url: URL) -> NostrRelayConnectionProtocol { func webSocketTask(with url: URL) -> NostrRelayConnectionProtocol {
let task = base.webSocketTask(with: url) URLSessionWebSocketTaskAdapter(base: base.webSocketTask(with: url))
// Byte bound per inbound frame; without it the per-relay buffer cap
// (nostrInboundPerRelayBufferCap) bounds FRAMES but not BYTES, and a
// hostile relay could pile up cap × 1 MiB (URLSession default) per
// connection. See TransportConfig.nostrInboundMaxFrameBytes for the
// sizing rationale. Oversized frames fail the receive with an error.
task.maximumMessageSize = TransportConfig.nostrInboundMaxFrameBytes
return URLSessionWebSocketTaskAdapter(base: task)
} }
} }
@@ -113,10 +106,7 @@ private extension NostrRelayManagerDependencies {
@MainActor @MainActor
final class NostrRelayManager: ObservableObject { final class NostrRelayManager: ObservableObject {
static let shared = NostrRelayManager() static let shared = NostrRelayManager()
// Track gift-wraps (kind 1059) we initiated so we can log OK acks at info. // Track gift-wraps (kind 1059) we initiated so we can log OK acks at info
// Entries are removed only on OK acks (or panic wipe); relays that never
// ack leave entries behind for the process lifetime. Observability-only
// state, bounded in practice by outbound DM volume.
private(set) static var pendingGiftWrapIDs = Set<String>() private(set) static var pendingGiftWrapIDs = Set<String>()
static func registerPendingGiftWrap(id: String) { static func registerPendingGiftWrap(id: String) {
pendingGiftWrapIDs.insert(id) pendingGiftWrapIDs.insert(id)
@@ -222,33 +212,7 @@ final class NostrRelayManager: ObservableObject {
// Bump generation to invalidate scheduled reconnects when we reset/disconnect // Bump generation to invalidate scheduled reconnects when we reset/disconnect
private var connectionGeneration: Int = 0 private var connectionGeneration: Int = 0
// Per-relay off-main inbound pipeline: raw socket frames are parsed and
// Schnorr-verified in arrival order OFF the main actor (this is the single
// signature verification for the whole inbound path downstream handlers
// receive only verified events), then hop back to the main actor for dedup
// recording and handler dispatch.
//
// Each relay connection owns its OWN AsyncStream + consumer task, so N
// relays verify in parallel while every relay's frames stay in arrival
// order (a single subscription's events for a relay all arrive on that
// relay's socket, so per-relay ordering preserves per-subscription
// ordering). A burst of EVENT frames from one busy/malicious relay only
// blocks that relay's own verification backlog DMs, OKs, EOSEs, and
// events from every other relay keep flowing on their own pipelines.
//
// Each stream is bounded (`.bufferingNewest`) so a relay flooding faster
// than its verification drains sheds its own oldest frames instead of
// growing memory without bound; it can never starve other relays.
//
// Continuations live in a lock-guarded, `Sendable` router (see
// `InboundFrameRouter` at file scope) so the raw socket receive callback
// (which is NOT main-actor isolated) can route a frame to the right relay
// stream without a per-frame main hop, while the main actor owns pipeline
// creation/teardown. The expensive work (Schnorr verify) is what runs
// off-main; the yield stays cheap.
private let inboundRouter = InboundFrameRouter()
init() { init() {
self.dependencies = .live() self.dependencies = .live()
hasMutualFavorites = dependencies.hasMutualFavorites() hasMutualFavorites = dependencies.hasMutualFavorites()
@@ -302,70 +266,7 @@ final class NostrRelayManager: ObservableObject {
} }
.store(in: &cancellables) .store(in: &cancellables)
} }
deinit {
inboundRouter.finishAll()
}
/// Ensure a serial off-main consumer pipeline exists for a relay. Called on
/// the main actor when a socket is (re)armed for receiving. Idempotent.
///
/// Ordering within the relay is deliberate and security/performance-critical:
/// 1. `precheckInboundEvent` (main hop): per-relay stats plus a cheap
/// duplicate LOOKUP duplicate fan-in from multiple relays dominates
/// real traffic and must never pay for Schnorr verification.
/// 2. `isValidSignature()` runs here, off the main actor the ONLY
/// signature verification on the inbound path (JSON re-serialization +
/// SHA-256 + secp256k1 Schnorr per event).
/// 3. `deliverVerifiedInboundEvent` (main hop): authoritative
/// check-and-RECORD plus handler dispatch. Recording only after
/// verification means a forged-signature copy can never poison the
/// dedup cache and suppress the genuine event.
private func ensureRelayInboundPipeline(for relayUrl: String) {
let started = inboundRouter.startPipeline(for: relayUrl) { [weak self] stream in
Task.detached(priority: .userInitiated) {
for await frame in stream {
guard let parsed = ParsedInbound(frame.message) else { continue }
guard let self else { return }
switch parsed {
case .event(let subId, let event):
guard await self.precheckInboundEvent(
subscriptionID: subId,
eventID: event.id,
relayUrl: relayUrl
) else {
continue
}
guard event.isValidSignature() else {
SecureLogger.warning(
"⚠️ Dropped invalid Nostr event id=\(event.id.prefix(16))… sub=\(subId) relay=\(relayUrl)",
category: .session
)
continue
}
await self.deliverVerifiedInboundEvent(subscriptionID: subId, event: event, from: relayUrl)
case .eose, .ok, .notice:
await self.handleParsedMessage(parsed, from: relayUrl)
}
}
}
}
if started {
SecureLogger.debug("🧵 Started inbound verify pipeline for \(relayUrl)", category: .session)
}
}
/// Tear down a relay's inbound pipeline (socket gone or state wiped). The
/// consumer drains any already-buffered frames before finishing, so
/// in-flight verified events are still delivered.
private func teardownRelayInboundPipeline(for relayUrl: String) {
inboundRouter.finishPipeline(for: relayUrl)
}
private func teardownAllRelayInboundPipelines() {
inboundRouter.finishAll()
}
/// Connect to all configured relays /// Connect to all configured relays
func connect() { func connect() {
// Global network policy gate // Global network policy gate
@@ -380,8 +281,6 @@ final class NostrRelayManager: ObservableObject {
task.cancel(with: .goingAway, reason: nil) task.cancel(with: .goingAway, reason: nil)
} }
connections.removeAll() connections.removeAll()
// Sockets are gone; drop every relay's inbound verify pipeline.
teardownAllRelayInboundPipelines()
markRelaySocketsClosed(resetState: false) markRelaySocketsClosed(resetState: false)
// Sockets are gone, so per-relay subscription state is cleared but // Sockets are gone, so per-relay subscription state is cleared but
// durable intent (subscriptionRequestState, messageHandlers, parked // durable intent (subscriptionRequestState, messageHandlers, parked
@@ -411,7 +310,6 @@ final class NostrRelayManager: ObservableObject {
task.cancel(with: .goingAway, reason: nil) task.cancel(with: .goingAway, reason: nil)
} }
connections.removeAll() connections.removeAll()
teardownAllRelayInboundPipelines()
markRelaySocketsClosed(resetState: true) markRelaySocketsClosed(resetState: true)
subscriptions.removeAll() subscriptions.removeAll()
pendingSubscriptions.removeAll() pendingSubscriptions.removeAll()
@@ -662,7 +560,6 @@ final class NostrRelayManager: ObservableObject {
connection.cancel(with: .goingAway, reason: nil) connection.cancel(with: .goingAway, reason: nil)
} }
connections.removeValue(forKey: url) connections.removeValue(forKey: url)
teardownRelayInboundPipeline(for: url)
subscriptions.removeValue(forKey: url) subscriptions.removeValue(forKey: url)
pendingSubscriptions.removeValue(forKey: url) pendingSubscriptions.removeValue(forKey: url)
} }
@@ -1002,11 +899,7 @@ final class NostrRelayManager: ObservableObject {
connections[urlString] = task connections[urlString] = task
task.resume() task.resume()
// Bring up this relay's own serial verify pipeline before arming the
// socket, so inbound frames have somewhere to land.
ensureRelayInboundPipeline(for: urlString)
// Start receiving messages // Start receiving messages
receiveMessage(from: task, relayUrl: urlString) receiveMessage(from: task, relayUrl: urlString)
@@ -1065,13 +958,14 @@ final class NostrRelayManager: ObservableObject {
switch result { switch result {
case .success(let message): case .success(let message):
// Hand the raw frame to this relay's serial inbound pipeline: // Parse off-main to reduce UI jank, then hop back for state updates
// parsing and signature verification run off-main, in arrival Task.detached(priority: .utility) {
// order, independently of every other relay's pipeline. Routing guard let parsed = ParsedInbound(message) else { return }
// through the lock-guarded router keeps this off the main actor await MainActor.run {
// (no per-frame main hop). self.handleParsedMessage(parsed, from: relayUrl)
self.inboundRouter.yield(InboundFrame(message: message), to: relayUrl) }
}
// Continue receiving // Continue receiving
Task { @MainActor in Task { @MainActor in
self.receiveMessage(from: task, relayUrl: relayUrl) self.receiveMessage(from: task, relayUrl: relayUrl)
@@ -1089,55 +983,35 @@ final class NostrRelayManager: ObservableObject {
// Note: declared at file scope below to avoid MainActor isolation inside this class // Note: declared at file scope below to avoid MainActor isolation inside this class
// and keep parsing off the main actor. // and keep parsing off the main actor.
/// First main-actor hop for an inbound EVENT: per-relay stats plus a cheap // Handle parsed message on MainActor (state updates and handlers)
/// duplicate LOOKUP (no recording) so duplicate fan-in from multiple
/// relays never pays for Schnorr verification. Recording happens only
/// after the signature verifies (`deliverVerifiedInboundEvent`), so a
/// forged-signature copy can never poison the dedup cache and suppress
/// the genuine event.
private func precheckInboundEvent(subscriptionID: String, eventID: String, relayUrl: String) -> Bool {
if let index = relays.firstIndex(where: { $0.url == relayUrl }) {
relays[index].messagesReceived += 1
}
guard !eventID.isEmpty else { return true }
let key = InboundEventKey(subscriptionID: subscriptionID, eventID: eventID)
if recentInboundEventKeys.contains(key) {
recordDuplicateInboundEventDrop(subscriptionID: subscriptionID)
return false
}
return true
}
/// Second main-actor hop, after off-main signature verification:
/// authoritative check-and-record (the serial pipeline means the same
/// event is never in flight twice, but the record must stay atomic with
/// delivery) and handler dispatch.
private func deliverVerifiedInboundEvent(subscriptionID subId: String, event: NostrEvent, from relayUrl: String) {
guard shouldDeliverInboundEvent(subscriptionID: subId, eventID: event.id) else {
return
}
if event.kind != 1059 {
// Per-event logging floods dev builds in busy geohashes; sample it.
inboundEventLogCount += 1
if inboundEventLogCount == 1 || inboundEventLogCount.isMultiple(of: TransportConfig.nostrInboundEventLogInterval) {
SecureLogger.debug("📥 Event #\(inboundEventLogCount) kind=\(event.kind) id=\(event.id.prefix(16))… relay=\(relayUrl)", category: .session)
}
}
if let handler = self.messageHandlers[subId] {
handler(event)
} else {
SecureLogger.warning("⚠️ No handler for subscription \(subId)", category: .session)
}
}
// Handle parsed non-EVENT messages on MainActor (state updates and handlers)
private func handleParsedMessage(_ parsed: ParsedInbound, from relayUrl: String) { private func handleParsedMessage(_ parsed: ParsedInbound, from relayUrl: String) {
switch parsed { switch parsed {
case .event: case .event(let subId, let event):
// Events flow through the serial inbound pipeline (precheck if let index = self.relays.firstIndex(where: { $0.url == relayUrl }) {
// off-main signature verification deliverVerifiedInboundEvent) self.relays[index].messagesReceived += 1
// and never reach this fallback. }
assertionFailure("inbound EVENT bypassed the verified pipeline") guard event.isValidSignature() else {
SecureLogger.warning(
"⚠️ Dropped invalid Nostr event id=\(event.id.prefix(16))… sub=\(subId) relay=\(relayUrl)",
category: .session
)
return
}
guard shouldDeliverInboundEvent(subscriptionID: subId, eventID: event.id) else {
return
}
if event.kind != 1059 {
// Per-event logging floods dev builds in busy geohashes; sample it.
inboundEventLogCount += 1
if inboundEventLogCount == 1 || inboundEventLogCount.isMultiple(of: TransportConfig.nostrInboundEventLogInterval) {
SecureLogger.debug("📥 Event #\(inboundEventLogCount) kind=\(event.kind) id=\(event.id.prefix(16))… relay=\(relayUrl)", category: .session)
}
}
if let handler = self.messageHandlers[subId] {
handler(event)
} else {
SecureLogger.warning("⚠️ No handler for subscription \(subId)", category: .session)
}
case .eose(let subId): case .eose(let subId):
if var tracker = eoseTrackers[subId] { if var tracker = eoseTrackers[subId] {
tracker.pendingRelays.remove(relayUrl) tracker.pendingRelays.remove(relayUrl)
@@ -1232,7 +1106,6 @@ final class NostrRelayManager: ObservableObject {
private func handleDisconnection(relayUrl: String, error: Error) { private func handleDisconnection(relayUrl: String, error: Error) {
connections.removeValue(forKey: relayUrl) connections.removeValue(forKey: relayUrl)
teardownRelayInboundPipeline(for: relayUrl)
subscriptions.removeValue(forKey: relayUrl) subscriptions.removeValue(forKey: relayUrl)
updateRelayStatus(relayUrl, isConnected: false, error: error) updateRelayStatus(relayUrl, isConnected: false, error: error)
settleEOSETrackers(droppingRelay: relayUrl) settleEOSETrackers(droppingRelay: relayUrl)
@@ -1321,9 +1194,8 @@ final class NostrRelayManager: ObservableObject {
if let connection = connections[normalizedRelayUrl] { if let connection = connections[normalizedRelayUrl] {
connection.cancel(with: .goingAway, reason: nil) connection.cancel(with: .goingAway, reason: nil)
connections.removeValue(forKey: normalizedRelayUrl) connections.removeValue(forKey: normalizedRelayUrl)
teardownRelayInboundPipeline(for: normalizedRelayUrl)
} }
// Attempt immediate reconnection // Attempt immediate reconnection
connectToRelay(normalizedRelayUrl) connectToRelay(normalizedRelayUrl)
} }
@@ -1416,77 +1288,6 @@ final class NostrRelayManager: ObservableObject {
// MARK: - Off-main inbound parsing helpers (file scope, non-isolated) // MARK: - Off-main inbound parsing helpers (file scope, non-isolated)
/// A single raw socket frame awaiting off-main parse + Schnorr verification.
private struct InboundFrame: Sendable {
let message: URLSessionWebSocketTask.Message
}
/// Lock-guarded registry of per-relay inbound streams.
///
/// The raw WebSocket receive callback is not main-actor isolated, so it needs a
/// `Sendable` path to route a frame to the correct relay's stream without a
/// per-frame hop onto the main actor. Pipeline lifecycle (start/finish) is
/// driven from the main actor; frame delivery (`yield`) can come from any
/// thread. All access is serialized by a single lock contention is negligible
/// because the guarded critical section is only a dictionary lookup + yield.
private final class InboundFrameRouter: @unchecked Sendable {
private let lock = NSLock()
private var continuations: [String: AsyncStream<InboundFrame>.Continuation] = [:]
private var tasks: [String: Task<Void, Never>] = [:]
/// Start a relay's stream + consumer if one does not already exist.
/// Returns true when a new pipeline was created. The bounded
/// `.bufferingNewest` policy makes a single relay shed its OWN oldest
/// frames under a flood, never other relays' frames. Buffered memory per
/// relay is bounded (not eliminated) at the frame cap times the per-frame
/// byte cap (`maximumMessageSize`) see TransportConfig.
func startPipeline(
for relayUrl: String,
makeConsumer: (AsyncStream<InboundFrame>) -> Task<Void, Never>
) -> Bool {
lock.lock()
defer { lock.unlock() }
if continuations[relayUrl] != nil { return false }
let (stream, continuation) = AsyncStream<InboundFrame>.makeStream(
bufferingPolicy: .bufferingNewest(TransportConfig.nostrInboundPerRelayBufferCap)
)
continuations[relayUrl] = continuation
tasks[relayUrl] = makeConsumer(stream)
return true
}
/// Route a frame to a relay's stream. No-op if the relay has no live
/// pipeline (socket already torn down) the frame is simply dropped, which
/// is safe for best-effort Nostr inbound.
func yield(_ frame: InboundFrame, to relayUrl: String) {
lock.lock()
let continuation = continuations[relayUrl]
lock.unlock()
continuation?.yield(frame)
}
/// Finish a relay's stream. The consumer drains any already-buffered frames
/// before exiting, so in-flight verified events are still delivered.
func finishPipeline(for relayUrl: String) {
lock.lock()
let continuation = continuations.removeValue(forKey: relayUrl)
tasks.removeValue(forKey: relayUrl)
lock.unlock()
continuation?.finish()
}
func finishAll() {
lock.lock()
let allContinuations = continuations
continuations.removeAll()
tasks.removeAll()
lock.unlock()
for continuation in allContinuations.values {
continuation.finish()
}
}
}
private enum ParsedInbound { private enum ParsedInbound {
case event(subId: String, event: NostrEvent) case event(subId: String, event: NostrEvent)
case ok(eventId: String, success: Bool, reason: String) case ok(eventId: String, success: Bool, reason: String)
-20
View File
@@ -55,26 +55,6 @@ enum TransportConfig {
static let nostrDuplicateEventLogInterval: Int = 50 static let nostrDuplicateEventLogInterval: Int = 50
// Sample interval for per-event debug logs on the inbound hot path. // Sample interval for per-event debug logs on the inbound hot path.
static let nostrInboundEventLogInterval: Int = 100 static let nostrInboundEventLogInterval: Int = 100
// Bounded per-relay inbound frame buffer. Each relay connection owns its
// own serial verify pipeline; if a relay floods faster than its Schnorr
// verification drains, the oldest buffered frames for THAT relay are
// dropped (bufferingNewest) so one relay cannot stall other relays.
// Nostr inbound is already best-effort (relays are redundant and events
// replay), so dropping a flooding relay's backlog is safe. Together with
// nostrInboundMaxFrameBytes this caps buffered inbound bytes at
// cap × maxFrameBytes (128 MiB) per hostile relay bounded, not zero.
static let nostrInboundPerRelayBufferCap: Int = 256
// Hard per-frame byte bound, applied as URLSessionWebSocketTask
// .maximumMessageSize (oversized frames fail the receive instead of
// buffering). BitChat's legitimate Nostr traffic is small: geohash chat /
// presence events (kind 20000/20001), kind-1 notes, and NIP-17
// gift-wrapped DMs carrying text payloads or receipts are all a few KiB,
// and most public relays reject events beyond ~64256 KiB anyway. 512 KiB
// leaves an order-of-magnitude margin over anything we produce or expect
// while halving the URLSession default (1 MiB), so a hostile relay's
// worst-case buffered pile-up per connection is
// nostrInboundPerRelayBufferCap × 512 KiB = 128 MiB instead of 256 MiB.
static let nostrInboundMaxFrameBytes: Int = 512 * 1024
// Conversation store diagnostics (field observability) // Conversation store diagnostics (field observability)
// Sample interval for the periodic store-audit "OK" heartbeat line // Sample interval for the periodic store-audit "OK" heartbeat line
+48 -20
View File
@@ -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
+25 -14
View File
@@ -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
} }
} }
-5
View File
@@ -1177,11 +1177,6 @@ final class ChatViewModel: ObservableObject, BitchatDelegate, TransportEventDele
// Geohash DM handlers can capture pre-wipe Nostr identities, so a plain // Geohash DM handlers can capture pre-wipe Nostr identities, so a plain
// disconnect is not enough here. // disconnect is not enough here.
NostrRelayManager.shared.resetForPanicWipe() NostrRelayManager.shared.resetForPanicWipe()
// Clearing relay handlers stops NEW events, but a detached gift-wrap
// decrypt spawned just before the wipe still holds a pre-wipe key and
// ciphertext; bump the pipeline's wipe generation so its result is
// dropped at the main-actor delivery hop instead of landing here.
nostrCoordinator.inbound.invalidateInFlightDecrypts()
nostrRelayManager = nil nostrRelayManager = nil
// Clear Nostr identity associations // Clear Nostr identity associations
+1 -2
View File
@@ -103,8 +103,7 @@ final class GeoPresenceTracker {
else { else {
return return
} }
// The signature was already verified (exactly once, off the main guard event.isValidSignature() else { return }
// actor) by NostrRelayManager before delivery.
guard shouldProcessGeoSamplingEvent(event.id) else { return } guard shouldProcessGeoSamplingEvent(event.id) else { return }
let existingCount = context.geoParticipantCount(for: gh) let existingCount = context.geoParticipantCount(for: gh)
+89 -133
View File
@@ -75,38 +75,20 @@ extension ChatViewModel: NostrInboundPipelineContext {
} }
} }
/// The inbound Nostr hot path: verified relay events in, chat messages / /// The inbound Nostr hot path: raw relay events in, chat messages / Noise
/// Noise payloads out. Pure transformation plus dedup no relay lifecycle. /// payloads out. Pure transformation plus dedup no relay lifecycle.
/// ///
/// Every event arriving here already had its Schnorr signature verified /// Ordering is deliberate and performance-critical: cheap rejects (kind,
/// exactly once, off the main actor, by `NostrRelayManager`'s serial inbound /// dedup lookup) run BEFORE Schnorr signature verification because duplicates
/// pipeline (which records events into its own dedup cache only AFTER /// dominate real relay traffic; events are recorded only AFTER verification so
/// verification, so forged copies can't suppress genuine events). This /// a forged-signature copy can never poison the dedup set; gift-wrap
/// pipeline therefore never re-verifies; it keeps its own event-ID dedup /// verification for the account mailbox runs off-main with an atomic
/// (cheap main-actor lookups) and moves NIP-17 gift-wrap decryption two /// main-actor check-and-record.
/// ECDH+ChaCha layers off the main actor with an atomic main-actor
/// check-and-record.
final class NostrInboundPipeline { final class NostrInboundPipeline {
private weak var context: (any NostrInboundPipelineContext)? private weak var context: (any NostrInboundPipelineContext)?
private let presence: GeoPresenceTracker private let presence: GeoPresenceTracker
private var geoEventLogCount = 0 private var geoEventLogCount = 0
/// Monotonic panic-wipe generation for this pipeline. A panic wipe clears
/// relay handlers so no NEW events flow, but a detached decrypt task
/// spawned just BEFORE the wipe which strongly captures a pre-wipe Nostr
/// private key and ciphertext survives it. Spawn sites capture this
/// value; the task compares it at its main-actor hops and drops its result
/// (no delivery; the captured identity and plaintext die with the task)
/// if `invalidateInFlightDecrypts()` bumped it in between.
@MainActor private(set) var wipeGeneration: UInt64 = 0
/// Called from `ChatViewModel.panicClearAllData()` so plaintext decrypted
/// with pre-wipe keys can never land in post-wipe state.
@MainActor
func invalidateInFlightDecrypts() {
wipeGeneration &+= 1
}
init(context: any NostrInboundPipelineContext, presence: GeoPresenceTracker) { init(context: any NostrInboundPipelineContext, presence: GeoPresenceTracker) {
self.context = context self.context = context
self.presence = presence self.presence = presence
@@ -115,15 +97,17 @@ final class NostrInboundPipeline {
@MainActor @MainActor
func subscribeNostrEvent(_ event: NostrEvent) { func subscribeNostrEvent(_ event: NostrEvent) {
guard let context else { return } guard let context else { return }
// Cheap rejects (kind, dedup lookup) duplicates dominate real // Cheap rejects (kind, dedup lookup) before Schnorr verification
// traffic. The signature was already verified (exactly once, off the // duplicates dominate real traffic and must not pay for crypto.
// main actor) by NostrRelayManager before delivery. // Only verified events are recorded, so a forged-signature copy can
// never poison the dedup set and suppress the genuine event.
guard (event.kind == NostrProtocol.EventKind.ephemeralEvent.rawValue guard (event.kind == NostrProtocol.EventKind.ephemeralEvent.rawValue
|| event.kind == NostrProtocol.EventKind.geohashPresence.rawValue), || event.kind == NostrProtocol.EventKind.geohashPresence.rawValue),
!context.hasProcessedNostrEvent(event.id) !context.hasProcessedNostrEvent(event.id)
else { else {
return return
} }
guard event.isValidSignature() else { return }
context.recordProcessedNostrEvent(event.id) context.recordProcessedNostrEvent(event.id)
@@ -192,14 +176,15 @@ final class NostrInboundPipeline {
@MainActor @MainActor
func handleNostrEvent(_ event: NostrEvent) { func handleNostrEvent(_ event: NostrEvent) {
guard let context else { return } guard let context else { return }
// Cheap rejects (kind, dedup lookup) the signature was already // Cheap rejects (kind, dedup lookup) before Schnorr verification
// verified (exactly once, off the main actor) by NostrRelayManager. // duplicates dominate real traffic and must not pay for crypto.
guard (event.kind == NostrProtocol.EventKind.ephemeralEvent.rawValue guard (event.kind == NostrProtocol.EventKind.ephemeralEvent.rawValue
|| event.kind == NostrProtocol.EventKind.geohashPresence.rawValue) || event.kind == NostrProtocol.EventKind.geohashPresence.rawValue)
else { else {
return return
} }
if context.hasProcessedNostrEvent(event.id) { return } if context.hasProcessedNostrEvent(event.id) { return }
guard event.isValidSignature() else { return }
context.recordProcessedNostrEvent(event.id) context.recordProcessedNostrEvent(event.id)
// Sampled: fires for every geo event and floods dev logs in busy geohashes. // Sampled: fires for every geo event and floods dev logs in busy geohashes.
@@ -279,126 +264,104 @@ final class NostrInboundPipeline {
@MainActor @MainActor
func subscribeGiftWrap(_ giftWrap: NostrEvent, id: NostrIdentity) { func subscribeGiftWrap(_ giftWrap: NostrEvent, id: NostrIdentity) {
guard let context else { return } guard let context else { return }
// Cheap dedup pre-check only; processGeohashGiftWrap does the // Dedup lookup before Schnorr verification; record only after it passes.
// authoritative main-actor check-and-record before the off-main
// NIP-17 unwrap. The outer signature was already verified (exactly
// once, off the main actor) by NostrRelayManager.
guard !context.hasProcessedNostrEvent(giftWrap.id) else { return } guard !context.hasProcessedNostrEvent(giftWrap.id) else { return }
guard giftWrap.isValidSignature() else { return }
context.recordProcessedNostrEvent(giftWrap.id)
// Capture the wipe generation at spawn, alongside the per-geohash guard let (content, senderPubkey, rumorTs) = try? NostrProtocol.decryptPrivateMessage(
// identity (private key) the detached task strongly captures. A panic giftWrap: giftWrap,
// wipe between spawn and delivery bumps the generation, and the task recipientIdentity: id
// drops its result instead of delivering plaintext post-wipe. ),
let wipeGeneration = self.wipeGeneration let packet = Self.decodeEmbeddedBitChatPacket(from: content),
Task.detached(priority: .userInitiated) { [weak self] in packet.type == MessageType.noiseEncrypted.rawValue,
await self?.processGeohashGiftWrap(giftWrap, id: id, verbose: false, wipeGeneration: wipeGeneration) let noisePayload = NoisePayload.decode(packet.payload)
else {
return
}
let messageTimestamp = Date(timeIntervalSince1970: TimeInterval(rumorTs))
let convKey = PeerID(nostr_: senderPubkey)
context.registerNostrKeyMapping(senderPubkey, for: convKey)
switch noisePayload.type {
case .privateMessage:
context.handlePrivateMessage(
noisePayload,
senderPubkey: senderPubkey,
convKey: convKey,
id: id,
messageTimestamp: messageTimestamp
)
case .delivered:
context.handleDelivered(noisePayload, senderPubkey: senderPubkey, convKey: convKey)
case .readReceipt:
context.handleReadReceipt(noisePayload, senderPubkey: senderPubkey, convKey: convKey)
case .verifyChallenge, .verifyResponse:
break
} }
} }
@MainActor @MainActor
func handleGiftWrap(_ giftWrap: NostrEvent, id: NostrIdentity) { func handleGiftWrap(_ giftWrap: NostrEvent, id: NostrIdentity) {
guard let context else { return } guard let context else { return }
// Cheap dedup pre-check only; see subscribeGiftWrap. // Dedup lookup before Schnorr verification; record only after it passes.
if context.hasProcessedNostrEvent(giftWrap.id) { if context.hasProcessedNostrEvent(giftWrap.id) {
return return
} }
guard giftWrap.isValidSignature() else { return }
// Spawn-time wipe-generation capture; see subscribeGiftWrap. context.recordProcessedNostrEvent(giftWrap.id)
let wipeGeneration = self.wipeGeneration
Task.detached(priority: .userInitiated) { [weak self] in
await self?.processGeohashGiftWrap(giftWrap, id: id, verbose: true, wipeGeneration: wipeGeneration)
}
}
/// Geohash-DM gift wrap ingest. The NIP-17 unwrap (two ECDH+ChaCha
/// layers) runs off the main actor; results hop back for state updates.
/// `verbose` keeps `handleGiftWrap`'s decrypt logging without adding it
/// to the sampling path.
///
/// `wipeGeneration` is this pipeline's generation captured at spawn (the
/// moment the pre-wipe `id` was captured); a mismatch at either main-actor
/// hop means a panic wipe happened in between, so the task bails without
/// decrypting (first hop) or without delivering the plaintext (second
/// hop) the captured identity and any decrypted material are simply
/// dropped with the task.
private func processGeohashGiftWrap(
_ giftWrap: NostrEvent,
id: NostrIdentity,
verbose: Bool,
wipeGeneration: UInt64
) async {
guard let context else { return }
// Authoritative check-and-record, atomic on the main actor so two
// concurrent detached tasks can't both process the same event.
let alreadyProcessed: Bool = await MainActor.run {
guard self.wipeGeneration == wipeGeneration else { return true }
if context.hasProcessedNostrEvent(giftWrap.id) { return true }
context.recordProcessedNostrEvent(giftWrap.id)
return false
}
if alreadyProcessed { return }
guard let (content, senderPubkey, rumorTs) = try? NostrProtocol.decryptPrivateMessage( guard let (content, senderPubkey, rumorTs) = try? NostrProtocol.decryptPrivateMessage(
giftWrap: giftWrap, giftWrap: giftWrap,
recipientIdentity: id recipientIdentity: id
) else { ) else {
if verbose { SecureLogger.warning("GeoDM: failed decrypt giftWrap id=\(giftWrap.id.prefix(8))", category: .session)
SecureLogger.warning("GeoDM: failed decrypt giftWrap id=\(giftWrap.id.prefix(8))", category: .session)
}
return return
} }
if verbose { SecureLogger.debug(
SecureLogger.debug( "GeoDM: decrypted gift-wrap id=\(giftWrap.id.prefix(16))... from=\(senderPubkey.prefix(8))...",
"GeoDM: decrypted gift-wrap id=\(giftWrap.id.prefix(16))... from=\(senderPubkey.prefix(8))...", category: .session
category: .session )
)
guard let packet = Self.decodeEmbeddedBitChatPacket(from: content),
packet.type == MessageType.noiseEncrypted.rawValue,
let payload = NoisePayload.decode(packet.payload)
else {
return
} }
await MainActor.run { let convKey = PeerID(nostr_: senderPubkey)
// A panic wipe during the off-main decrypt must not let the context.registerNostrKeyMapping(senderPubkey, for: convKey)
// pre-wipe plaintext reach post-wipe state; drop it here, atomic
// with the wipe on the main actor.
guard self.wipeGeneration == wipeGeneration else { return }
guard let packet = Self.decodeEmbeddedBitChatPacket(from: content),
packet.type == MessageType.noiseEncrypted.rawValue,
let payload = NoisePayload.decode(packet.payload)
else {
return
}
let convKey = PeerID(nostr_: senderPubkey) switch payload.type {
context.registerNostrKeyMapping(senderPubkey, for: convKey) case .privateMessage:
let messageTimestamp = Date(timeIntervalSince1970: TimeInterval(rumorTs))
switch payload.type { context.handlePrivateMessage(
case .privateMessage: payload,
let messageTimestamp = Date(timeIntervalSince1970: TimeInterval(rumorTs)) senderPubkey: senderPubkey,
context.handlePrivateMessage( convKey: convKey,
payload, id: id,
senderPubkey: senderPubkey, messageTimestamp: messageTimestamp
convKey: convKey, )
id: id, case .delivered:
messageTimestamp: messageTimestamp context.handleDelivered(payload, senderPubkey: senderPubkey, convKey: convKey)
) case .readReceipt:
case .delivered: context.handleReadReceipt(payload, senderPubkey: senderPubkey, convKey: convKey)
context.handleDelivered(payload, senderPubkey: senderPubkey, convKey: convKey) case .verifyChallenge, .verifyResponse:
case .readReceipt: break
context.handleReadReceipt(payload, senderPubkey: senderPubkey, convKey: convKey)
case .verifyChallenge, .verifyResponse:
break
}
} }
} }
@MainActor @MainActor
func handleNostrMessage(_ giftWrap: NostrEvent) { func handleNostrMessage(_ giftWrap: NostrEvent) {
guard let context else { return } guard let context else { return }
// Cheap dedup pre-check only; processNostrMessage does the // Cheap dedup pre-check only; Schnorr verification runs off-main in
// authoritative check-and-record before the off-main NIP-17 unwrap. // processNostrMessage, which then does the authoritative
// The outer signature was already verified (exactly once, off the // check-and-record. Recording stays after verification so a
// main actor) by NostrRelayManager, and only verified events are // forged-signature copy can never poison the dedup set and suppress
// recorded, so a forged-signature copy can never poison the dedup // the genuine event.
// set and suppress the genuine event.
if context.hasProcessedNostrEvent(giftWrap.id) { return } if context.hasProcessedNostrEvent(giftWrap.id) { return }
Task.detached(priority: .userInitiated) { [weak self] in Task.detached(priority: .userInitiated) { [weak self] in
@@ -407,6 +370,7 @@ final class NostrInboundPipeline {
} }
func processNostrMessage(_ giftWrap: NostrEvent) async { func processNostrMessage(_ giftWrap: NostrEvent) async {
guard giftWrap.isValidSignature() else { return }
guard let context else { return } guard let context else { return }
// Authoritative check-and-record, atomic on the main actor so two // Authoritative check-and-record, atomic on the main actor so two
// concurrent detached tasks can't both process the same event. // concurrent detached tasks can't both process the same event.
@@ -416,13 +380,8 @@ final class NostrInboundPipeline {
return false return false
} }
if alreadyProcessed { return } if alreadyProcessed { return }
// Fetch the identity and the wipe generation in ONE main-actor hop: let currentIdentity: NostrIdentity? = await MainActor.run {
// the generation then vouches for exactly this identity. A wipe after context.currentNostrIdentity()
// this point bumps the generation and the delivery hop below drops
// the decrypted result (same guard as processGeohashGiftWrap; this
// account-mailbox path had the identical hazard).
let (currentIdentity, wipeGeneration): (NostrIdentity?, UInt64) = await MainActor.run {
(context.currentNostrIdentity(), self.wipeGeneration)
} }
guard let currentIdentity else { return } guard let currentIdentity else { return }
@@ -454,9 +413,6 @@ final class NostrInboundPipeline {
let payload = NoisePayload.decode(packet.payload) { let payload = NoisePayload.decode(packet.payload) {
let messageTimestamp = Date(timeIntervalSince1970: TimeInterval(rumorTimestamp)) let messageTimestamp = Date(timeIntervalSince1970: TimeInterval(rumorTimestamp))
await MainActor.run { await MainActor.run {
// Drop pre-wipe plaintext if a panic wipe landed
// during the off-main decrypt (see above).
guard self.wipeGeneration == wipeGeneration else { return }
context.registerNostrKeyMapping(senderPubkey, for: targetPeerID) context.registerNostrKeyMapping(senderPubkey, for: targetPeerID)
switch payload.type { switch payload.type {
@@ -382,12 +382,10 @@ struct ChatNostrCoordinatorContextTests {
coordinator.inbound.handleGiftWrap(giftWrap, id: recipient) coordinator.inbound.handleGiftWrap(giftWrap, id: recipient)
// The NIP-17 unwrap runs off the main actor; wait for the hop back.
let convKey = PeerID(nostr_: sender.publicKeyHex) let convKey = PeerID(nostr_: sender.publicKeyHex)
let routed = await TestHelpers.waitUntil({ context.handledPrivateMessages.count == 1 })
#expect(routed)
#expect(context.recordedNostrEventIDs == [giftWrap.id]) #expect(context.recordedNostrEventIDs == [giftWrap.id])
#expect(context.nostrKeyMapping[convKey] == sender.publicKeyHex) #expect(context.nostrKeyMapping[convKey] == sender.publicKeyHex)
#expect(context.handledPrivateMessages.count == 1)
#expect(context.handledPrivateMessages.first?.senderPubkey == sender.publicKeyHex) #expect(context.handledPrivateMessages.first?.senderPubkey == sender.publicKeyHex)
#expect(context.handledPrivateMessages.first?.convKey == convKey) #expect(context.handledPrivateMessages.first?.convKey == convKey)
@@ -400,76 +398,30 @@ struct ChatNostrCoordinatorContextTests {
// The same gift wrap is dropped on replay. // The same gift wrap is dropped on replay.
coordinator.inbound.handleGiftWrap(giftWrap, id: recipient) coordinator.inbound.handleGiftWrap(giftWrap, id: recipient)
await drainMainQueue()
#expect(context.recordedNostrEventIDs == [giftWrap.id]) #expect(context.recordedNostrEventIDs == [giftWrap.id])
#expect(context.handledPrivateMessages.count == 1) #expect(context.handledPrivateMessages.count == 1)
} }
@Test @MainActor @Test @MainActor
func handleGiftWrap_panicWipeAfterSpawnDropsDecryptedResult() async throws { func processNostrMessage_invalidSignatureDoesNotPoisonDedup() async throws {
let context = MockChatNostrContext() let context = MockChatNostrContext()
let coordinator = ChatNostrCoordinator(context: context) let coordinator = ChatNostrCoordinator(context: context)
let recipient = try NostrIdentity.generate() let recipient = try NostrIdentity.generate()
let sender = try NostrIdentity.generate() let sender = try NostrIdentity.generate()
let embedded = try #require(NostrEmbeddedBitChat.encodePMForNostrNoRecipient(
content: "pre-wipe secret",
messageID: "gm-wipe-1",
senderPeerID: PeerID(str: "aabbccddeeff0011")
))
let giftWrap = try NostrProtocol.createPrivateMessage(
content: embedded,
recipientPubkey: recipient.publicKeyHex,
senderIdentity: sender
)
// Spawn the detached decrypt (it strongly captures the pre-wipe
// identity), then panic-wipe in the SAME main-actor turn guaranteed
// to land before the task's first main-actor hop.
coordinator.inbound.handleGiftWrap(giftWrap, id: recipient)
coordinator.inbound.invalidateInFlightDecrypts()
// Give the detached task ample time to have delivered if the wipe
// guard were broken.
try? await Task.sleep(nanoseconds: 200_000_000)
await drainMainQueue()
#expect(context.handledPrivateMessages.isEmpty)
#expect(context.recordedNostrEventIDs.isEmpty)
// The pipeline itself stays usable: a gift wrap spawned AFTER the
// wipe (new generation) still decrypts and delivers.
coordinator.inbound.handleGiftWrap(giftWrap, id: recipient)
let delivered = await TestHelpers.waitUntil({ context.handledPrivateMessages.count == 1 })
#expect(delivered)
}
// NOTE: Inbound Schnorr signature verification (and the forged-copy
// dedup-poisoning invariant) is enforced once, off the main actor, at the
// relay boundary see NostrRelayManagerTests
// `test_receiveEvent_invalidSignatureDoesNotPoisonDuplicateCache` and
// `test_receiveGiftWrap_tamperedSignatureIsDroppedAndDoesNotPoisonDedup`.
// The inbound pipeline only ever sees verified events.
@Test @MainActor
func processNostrMessage_duplicateDeliveryProcessesOnce() async throws {
let context = MockChatNostrContext()
let coordinator = ChatNostrCoordinator(context: context)
let recipient = try NostrIdentity.generate()
let sender = try NostrIdentity.generate()
context.nostrIdentity = recipient
let giftWrap = try NostrProtocol.createPrivateMessage( let giftWrap = try NostrProtocol.createPrivateMessage(
content: "verify:noop", content: "verify:noop",
recipientPubkey: recipient.publicKeyHex, recipientPubkey: recipient.publicKeyHex,
senderIdentity: sender senderIdentity: sender
) )
var invalidGiftWrap = giftWrap
invalidGiftWrap.sig = String(repeating: "0", count: 128)
// Fan-in of the same (already verified) gift wrap from several relays // A forged-signature copy is rejected WITHOUT entering the dedup set...
// records and processes exactly once. await coordinator.inbound.processNostrMessage(invalidGiftWrap)
await coordinator.inbound.processNostrMessage(giftWrap) #expect(context.recordedNostrEventIDs.isEmpty)
#expect(context.recordedNostrEventIDs == [giftWrap.id])
// ...so the genuine event with the same ID still processes and records.
await coordinator.inbound.processNostrMessage(giftWrap) await coordinator.inbound.processNostrMessage(giftWrap)
#expect(context.recordedNostrEventIDs == [giftWrap.id]) #expect(context.recordedNostrEventIDs == [giftWrap.id])
} }
+27 -12
View File
@@ -352,11 +352,31 @@ struct ChatViewModelNostrExtensionTests {
#expect(!viewModel.messages.contains { $0.content == "Blocked" }) #expect(!viewModel.messages.contains { $0.content == "Blocked" })
} }
// NOTE: Tampered-signature rejection is enforced once, off the main @Test @MainActor
// actor, at the relay boundary (events only reach the inbound pipeline func handleNostrEvent_rejectsInvalidSignature() async throws {
// after verification) see NostrRelayManagerTests let (viewModel, _) = makeTestableViewModel()
// `test_receiveEvent_invalidSignatureDoesNotPoisonDuplicateCache` and let geohash = "u4pruydq"
// `test_receiveGiftWrap_tamperedSignatureIsDroppedAndDoesNotPoisonDedup`. let identity = try NostrIdentity.generate()
viewModel.switchLocationChannel(to: .location(GeohashChannel(level: .city, geohash: geohash)))
let event = NostrEvent(
pubkey: identity.publicKeyHex,
createdAt: Date(),
kind: .ephemeralEvent,
tags: [["g", geohash]],
content: "Valid"
)
var signed = try event.sign(with: identity.schnorrSigningKey())
signed.id = "deadbeef"
viewModel.handleNostrEvent(signed)
try? await Task.sleep(nanoseconds: 100_000_000)
viewModel.publicMessagePipeline.flushIfNeeded()
#expect(!viewModel.messages.contains { $0.content == "Tampered" })
}
@Test @MainActor @Test @MainActor
func subscribeGiftWrap_rejectsOversizedEmbeddedPacket() async throws { func subscribeGiftWrap_rejectsOversizedEmbeddedPacket() async throws {
@@ -545,14 +565,9 @@ struct ChatViewModelNostrExtensionTests {
viewModel.handleGiftWrap(giftWrap, id: recipient) viewModel.handleGiftWrap(giftWrap, id: recipient)
// Gift-wrap decryption runs off the main actor; wait for the ack try? await Task.sleep(nanoseconds: 50_000_000)
// (sent even for blocked senders) to know processing finished.
let didAck = await TestHelpers.waitUntil(
{ viewModel.sentGeoDeliveryAcks.contains(messageID) },
timeout: 5.0
)
#expect(didAck)
#expect(viewModel.privateChats[convKey] == nil) #expect(viewModel.privateChats[convKey] == nil)
#expect(viewModel.sentGeoDeliveryAcks.contains(messageID))
} }
@Test @MainActor @Test @MainActor
+20 -5
View File
@@ -326,11 +326,26 @@ struct ChatViewModelPresenceHandlingTests {
#expect(viewModel.geohashParticipantCount(for: activeGeohash) >= 1) #expect(viewModel.geohashParticipantCount(for: activeGeohash) >= 1)
} }
// NOTE: Tampered-signature rejection (and the forged-copy dedup-poisoning @Test func subscribeNostrEvent_samplingInvalidSignatureDoesNotPoisonDedup() async throws {
// invariant) is enforced once, off the main actor, at the relay boundary let (viewModel, _) = makeTestableViewModel()
// the sampling path only ever sees verified events. See let sampleGeohash = "u4pru"
// NostrRelayManagerTests let identity = try NostrIdentity.generate()
// `test_receiveEvent_invalidSignatureDoesNotPoisonDuplicateCache`. let event = NostrEvent(
pubkey: identity.publicKeyHex,
createdAt: Date(),
kind: .geohashPresence,
tags: [["g", sampleGeohash]],
content: ""
)
let signed = try event.sign(with: identity.schnorrSigningKey())
var invalid = signed
invalid.sig = String(repeating: "0", count: 128)
viewModel.subscribeNostrEvent(invalid, gh: sampleGeohash)
viewModel.subscribeNostrEvent(signed, gh: sampleGeohash)
#expect(viewModel.geohashParticipantCount(for: sampleGeohash) == 1)
}
// MARK: - Test Helper // MARK: - Test Helper
+101
View File
@@ -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 {
@@ -77,10 +77,8 @@ final class PerformanceBaselineTests: XCTestCase {
// MARK: - 1a. Nostr inbound event handling (fresh events) // MARK: - 1a. Nostr inbound event handling (fresh events)
/// `NostrInboundPipeline.handleNostrEvent` for never-seen geo events /// `NostrInboundPipeline.handleNostrEvent` for never-seen geo events
/// (kind 20000): dedup record, presence/nickname bookkeeping, and /// (kind 20000): signature verification, dedup record, presence/nickname
/// public-message ingest scheduling. Schnorr signature verification is /// bookkeeping, and public-message ingest scheduling.
/// NOT part of this path anymore it runs exactly once, off the main
/// actor, in `NostrRelayManager` before delivery.
func testNostrInboundEventHandling_freshEvents() throws { func testNostrInboundEventHandling_freshEvents() throws {
let events = try Self.makeSignedGeohashEvents(count: 500) let events = try Self.makeSignedGeohashEvents(count: 500)
// A fresh context per measure pass so every event takes the // A fresh context per measure pass so every event takes the
@@ -108,9 +106,8 @@ final class PerformanceBaselineTests: XCTestCase {
/// The dedup-hit path: identical events replayed. Duplicates dominate /// The dedup-hit path: identical events replayed. Duplicates dominate
/// real relay traffic (the same event arrives from several relays), so /// real relay traffic (the same event arrives from several relays), so
/// this path runs hundreds of times a minute in busy geohashes. It is a /// this path runs hundreds of times a minute in busy geohashes. Note it
/// pure dedup lookup: no crypto (verification happens upstream in /// still pays full Schnorr signature verification before the dedup check.
/// `NostrRelayManager`, and only for the first-seen copy).
func testNostrInboundEventHandling_duplicateEvents() throws { func testNostrInboundEventHandling_duplicateEvents() throws {
let events = try Self.makeSignedGeohashEvents(count: 500) let events = try Self.makeSignedGeohashEvents(count: 500)
let context = PerfNostrContext() let context = PerfNostrContext()
@@ -639,15 +639,11 @@ final class NostrRelayManagerTests: XCTestCase {
try context.sessionFactory.latestConnection(for: firstRelayURL)?.emitEventMessage(subscriptionID: "events", event: event) try context.sessionFactory.latestConnection(for: firstRelayURL)?.emitEventMessage(subscriptionID: "events", event: event)
try context.sessionFactory.latestConnection(for: secondRelayURL)?.emitEventMessage(subscriptionID: "events", event: event) try context.sessionFactory.latestConnection(for: secondRelayURL)?.emitEventMessage(subscriptionID: "events", event: event)
// Wait on the DELIVERY-side state: handler dispatch and the duplicate let countedOnBothRelays = await waitUntil {
// drop both land on the second main hop, after off-main verification.
let settled = await waitUntil {
context.manager.relays.first(where: { $0.url == firstRelayURL })?.messagesReceived == 1 && context.manager.relays.first(where: { $0.url == firstRelayURL })?.messagesReceived == 1 &&
context.manager.relays.first(where: { $0.url == secondRelayURL })?.messagesReceived == 1 && context.manager.relays.first(where: { $0.url == secondRelayURL })?.messagesReceived == 1
receivedIDs == [event.id] &&
context.manager.debugDuplicateInboundEventDropCount == 1
} }
XCTAssertTrue(settled) XCTAssertTrue(countedOnBothRelays)
XCTAssertEqual(receivedIDs, [event.id]) XCTAssertEqual(receivedIDs, [event.id])
XCTAssertEqual(context.manager.debugDuplicateInboundEventDropCount, 1) XCTAssertEqual(context.manager.debugDuplicateInboundEventDropCount, 1)
XCTAssertEqual(context.manager.debugDuplicateInboundEventDropCount(forSubscriptionID: "events"), 1) XCTAssertEqual(context.manager.debugDuplicateInboundEventDropCount(forSubscriptionID: "events"), 1)
@@ -680,17 +676,12 @@ final class NostrRelayManagerTests: XCTestCase {
) )
} }
// Wait on the DELIVERY-side state: the winner's handler dispatch and let countedOnEveryRelay = await waitUntil {
// the losers' duplicate drops both land on the second main hop, after
// off-main verification messagesReceived (first hop) settles sooner.
let settled = await waitUntil {
relayURLs.allSatisfy { relayURL in relayURLs.allSatisfy { relayURL in
context.manager.relays.first(where: { $0.url == relayURL })?.messagesReceived == 1 context.manager.relays.first(where: { $0.url == relayURL })?.messagesReceived == 1
} && }
receivedIDs == [event.id] &&
context.manager.debugDuplicateInboundEventDropCount == relayURLs.count - 1
} }
XCTAssertTrue(settled) XCTAssertTrue(countedOnEveryRelay)
XCTAssertEqual(receivedIDs, [event.id]) XCTAssertEqual(receivedIDs, [event.id])
XCTAssertEqual(context.manager.debugDuplicateInboundEventDropCount, relayURLs.count - 1) XCTAssertEqual(context.manager.debugDuplicateInboundEventDropCount, relayURLs.count - 1)
XCTAssertEqual( XCTAssertEqual(
@@ -723,153 +714,16 @@ final class NostrRelayManagerTests: XCTestCase {
try context.sessionFactory.latestConnection(for: firstRelayURL)?.emitEventMessage(subscriptionID: "events", event: invalidEvent) try context.sessionFactory.latestConnection(for: firstRelayURL)?.emitEventMessage(subscriptionID: "events", event: invalidEvent)
try context.sessionFactory.latestConnection(for: secondRelayURL)?.emitEventMessage(subscriptionID: "events", event: event) try context.sessionFactory.latestConnection(for: secondRelayURL)?.emitEventMessage(subscriptionID: "events", event: event)
// Wait on the DELIVERY-side state (second main hop, after off-main let countedOnBothRelays = await waitUntil {
// verification), not just messagesReceived (first main hop) the
// handler only fires after verify + a second hop.
let genuineDelivered = await waitUntil {
context.manager.relays.first(where: { $0.url == firstRelayURL })?.messagesReceived == 1 && context.manager.relays.first(where: { $0.url == firstRelayURL })?.messagesReceived == 1 &&
context.manager.relays.first(where: { $0.url == secondRelayURL })?.messagesReceived == 1 && context.manager.relays.first(where: { $0.url == secondRelayURL })?.messagesReceived == 1
receivedIDs == [event.id]
} }
XCTAssertTrue(genuineDelivered) XCTAssertTrue(countedOnBothRelays)
XCTAssertEqual(receivedIDs, [event.id]) XCTAssertEqual(receivedIDs, [event.id])
XCTAssertEqual(context.manager.debugDuplicateInboundEventDropCount, 0) XCTAssertEqual(context.manager.debugDuplicateInboundEventDropCount, 0)
XCTAssertEqual(context.manager.debugDuplicateInboundEventDropCount(forSubscriptionID: "events"), 0) XCTAssertEqual(context.manager.debugDuplicateInboundEventDropCount(forSubscriptionID: "events"), 0)
} }
/// The relay boundary is the single signature-verification point for the
/// whole inbound path (downstream pipelines no longer re-verify), so a
/// tampered gift wrap (kind 1059, the DM/mailbox path) must be dropped
/// here and must not poison the dedup cache against the genuine copy.
func test_receiveGiftWrap_tamperedSignatureIsDroppedAndDoesNotPoisonDedup() async throws {
let firstRelayURL = "wss://giftwrap-one.example"
let secondRelayURL = "wss://giftwrap-two.example"
let context = makeContext(permission: .denied)
let sender = try NostrIdentity.generate()
let recipient = try NostrIdentity.generate()
let giftWrap = try NostrProtocol.createPrivateMessage(
content: "psst",
recipientPubkey: recipient.publicKeyHex,
senderIdentity: sender
)
let tampered = invalidSignatureCopy(of: giftWrap)
var receivedIDs: [String] = []
context.manager.subscribe(
filter: makeFilter(),
id: "gift-wraps",
relayUrls: [firstRelayURL, secondRelayURL]
) { event in
receivedIDs.append(event.id)
}
let subscriptionsSent = await waitUntil {
context.sessionFactory.latestConnection(for: firstRelayURL)?.sentStrings.count == 1 &&
context.sessionFactory.latestConnection(for: secondRelayURL)?.sentStrings.count == 1
}
XCTAssertTrue(subscriptionsSent)
try context.sessionFactory.latestConnection(for: firstRelayURL)?.emitEventMessage(subscriptionID: "gift-wraps", event: tampered)
try context.sessionFactory.latestConnection(for: secondRelayURL)?.emitEventMessage(subscriptionID: "gift-wraps", event: giftWrap)
// Wait on the DELIVERY-side state (second main hop, after off-main
// verification), not just messagesReceived (first main hop) the
// handler only fires after verify + a second hop.
let genuineDelivered = await waitUntil {
context.manager.relays.first(where: { $0.url == firstRelayURL })?.messagesReceived == 1 &&
context.manager.relays.first(where: { $0.url == secondRelayURL })?.messagesReceived == 1 &&
receivedIDs == [giftWrap.id]
}
XCTAssertTrue(genuineDelivered)
XCTAssertEqual(receivedIDs, [giftWrap.id])
XCTAssertEqual(context.manager.debugDuplicateInboundEventDropCount, 0)
}
/// Signature verification runs off-main in a per-relay serial consumer;
/// several frames buffered on one socket must still be delivered to the
/// handler in that relay's arrival order.
func test_receiveEvent_deliversBackToBackEventsInArrivalOrder() async throws {
let relayURL = "wss://ordered.example"
let context = makeContext(permission: .denied)
let events = try (0..<12).map { try makeSignedEvent(content: "ordered-\($0)") }
var receivedIDs: [String] = []
context.manager.subscribe(filter: makeFilter(), id: "ordered", relayUrls: [relayURL]) { event in
receivedIDs.append(event.id)
}
let subscriptionSent = await waitUntil {
context.sessionFactory.latestConnection(for: relayURL)?.sentStrings.count == 1
}
XCTAssertTrue(subscriptionSent)
for event in events {
try context.sessionFactory.latestConnection(for: relayURL)?.emitEventMessage(subscriptionID: "ordered", event: event)
}
let allDelivered = await waitUntil(timeout: 5.0) {
receivedIDs.count == events.count
}
XCTAssertTrue(allDelivered)
XCTAssertEqual(receivedIDs, events.map(\.id))
}
/// Each relay owns its own off-main verify pipeline, so a large backlog of
/// EVENT frames on one relay must NOT head-of-line-block a frame that
/// arrives on a different relay. Under the previous single global consumer,
/// relay B's single event (emitted after relay A's whole burst) could only
/// be delivered once every frame in A's backlog had been Schnorr-verified;
/// with per-relay pipelines B verifies concurrently and lands before A's
/// backlog drains.
func test_receiveEvent_busyRelayDoesNotBlockOtherRelayDelivery() async throws {
let busyRelayURL = "wss://busy-relay.example"
let quietRelayURL = "wss://quiet-relay.example"
let context = makeContext(permission: .denied)
// Distinct subscriptions per relay so dedup never coalesces A vs. B.
let busyEvents = try (0..<200).map { try makeSignedEvent(content: "busy-\($0)") }
let quietEvent = try makeSignedEvent(content: "quiet")
var busyDeliveredCount = 0
var quietDeliveredAfterBusyCount = -1 // busy-count observed when B lands
context.manager.subscribe(filter: makeFilter(), id: "busy", relayUrls: [busyRelayURL]) { _ in
busyDeliveredCount += 1
}
context.manager.subscribe(filter: makeFilter(), id: "quiet", relayUrls: [quietRelayURL]) { _ in
if quietDeliveredAfterBusyCount < 0 {
quietDeliveredAfterBusyCount = busyDeliveredCount
}
}
let subscribed = await waitUntil {
context.sessionFactory.latestConnection(for: busyRelayURL)?.sentStrings.count == 1 &&
context.sessionFactory.latestConnection(for: quietRelayURL)?.sentStrings.count == 1
}
XCTAssertTrue(subscribed)
// Flood relay A first, then emit a single frame on relay B.
for event in busyEvents {
try context.sessionFactory.latestConnection(for: busyRelayURL)?.emitEventMessage(subscriptionID: "busy", event: event)
}
try context.sessionFactory.latestConnection(for: quietRelayURL)?.emitEventMessage(subscriptionID: "quiet", event: quietEvent)
let quietDelivered = await waitUntil(timeout: 5.0) { quietDeliveredAfterBusyCount >= 0 }
XCTAssertTrue(quietDelivered, "relay B's event was never delivered")
// The signal: B did not have to wait for A's entire backlog. If the two
// pipelines were globally serialized, B could only land after all 200 of
// A's frames, so busyDeliveredCount would be 200 when B arrived.
XCTAssertLessThan(
quietDeliveredAfterBusyCount,
busyEvents.count,
"relay B was head-of-line blocked behind relay A's backlog"
)
// Both relays still drain fully and in order.
let allDelivered = await waitUntil(timeout: 5.0) {
busyDeliveredCount == busyEvents.count
}
XCTAssertTrue(allDelivered)
}
func test_receiveEvent_withoutHandlerStillTracksReceivedCount() async throws { func test_receiveEvent_withoutHandlerStillTracksReceivedCount() async throws {
let relayURL = "wss://missing-handler.example" let relayURL = "wss://missing-handler.example"
let context = makeContext(permission: .denied) let context = makeContext(permission: .denied)
@@ -1841,11 +1695,7 @@ private final class MockRelayConnection: NostrRelayConnectionProtocol {
} }
func receive(completionHandler: @escaping (Result<URLSessionWebSocketTask.Message, Error>) -> Void) { func receive(completionHandler: @escaping (Result<URLSessionWebSocketTask.Message, Error>) -> Void) {
if !pendingResults.isEmpty { receiveHandler = completionHandler
completionHandler(pendingResults.removeFirst())
} else {
receiveHandler = completionHandler
}
} }
func sendPing(pongReceiveHandler: @escaping (Error?) -> Void) { func sendPing(pongReceiveHandler: @escaping (Error?) -> Void) {
@@ -1878,24 +1728,15 @@ private final class MockRelayConnection: NostrRelayConnectionProtocol {
} }
func emitRawString(_ string: String) throws { func emitRawString(_ string: String) throws {
deliver(.success(.string(string))) let handler = receiveHandler
receiveHandler = nil
handler?(.success(.string(string)))
} }
private func emit(jsonObject: Any) throws { private func emit(jsonObject: Any) throws {
let data = try JSONSerialization.data(withJSONObject: jsonObject) let data = try JSONSerialization.data(withJSONObject: jsonObject)
deliver(.success(.data(data))) let handler = receiveHandler
} receiveHandler = nil
handler?(.success(.data(data)))
// Frames emitted before the manager re-arms `receive` are queued so
// back-to-back emissions model a socket with several buffered frames.
private var pendingResults: [Result<URLSessionWebSocketTask.Message, Error>] = []
private func deliver(_ result: Result<URLSessionWebSocketTask.Message, Error>) {
if let handler = receiveHandler {
receiveHandler = nil
handler(result)
} else {
pendingResults.append(result)
}
} }
} }
@@ -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
}
}