mirror of
https://github.com/permissionlesstech/bitchat.git
synced 2026-07-26 10:05:19 +00:00
Harden Nostr envelope migration compatibility
This commit is contained in:
@@ -17,9 +17,10 @@ struct NostrProtocol {
|
||||
enum EventKind: Int {
|
||||
case metadata = 0
|
||||
case textNote = 1
|
||||
// Bounded compatibility for BitChat releases that incorrectly emitted
|
||||
// the proprietary payload under standard NIP kinds. Only kind 1059 is
|
||||
// temporarily published during migration; all three remain readable.
|
||||
// Compatibility for BitChat releases that incorrectly emitted the
|
||||
// proprietary payload under standard NIP kinds. Kind 1059 continues
|
||||
// to be published and read until a coordinated cross-platform release
|
||||
// explicitly removes it; all three legacy layers remain readable.
|
||||
case legacyNIP59Seal = 13
|
||||
case legacyNIP17DirectMessage = 14
|
||||
case legacyNIP59GiftWrap = 1059
|
||||
@@ -57,16 +58,10 @@ struct NostrProtocol {
|
||||
/// the layer-specific cap below the public ciphertext ceiling.
|
||||
private static let maximumPrivateEnvelopeSealPlaintextBytes = 48 * 1024
|
||||
|
||||
/// Compatibility-only publication stops at this instant. New-format kind
|
||||
/// 1402 remains first/primary throughout the window; the legacy kind-1059
|
||||
/// copy exists solely so pre-migration BitChat clients can receive it.
|
||||
static let legacyPrivateEnvelopePublicationDeadline = Date(
|
||||
timeIntervalSince1970: 1_792_022_400 // 2026-10-15T00:00:00Z
|
||||
)
|
||||
|
||||
/// New clients subscribe to the provisional BitChat-specific kind and the
|
||||
/// compatibility-only legacy kind so both sides of a rolling rollout can
|
||||
/// recover stored messages.
|
||||
/// compatibility legacy kind so both sides of a rolling rollout can
|
||||
/// recover stored messages. Do not remove kind 1059 here until all
|
||||
/// supported iOS and Android releases have migrated.
|
||||
static let acceptedPrivateEnvelopeKinds = [
|
||||
EventKind.privateEnvelope.rawValue,
|
||||
EventKind.legacyNIP59GiftWrap.rawValue
|
||||
@@ -145,24 +140,22 @@ struct NostrProtocol {
|
||||
}
|
||||
|
||||
/// Events to publish for one logical private payload. The primary
|
||||
/// BitChat-specific format is always first. Until the explicit migration
|
||||
/// deadline, a legacy copy follows for clients that still subscribe only
|
||||
/// to kind 1059. Both encrypt the exact same embedded BitChat payload, so
|
||||
/// receive-side logical-payload dedup collapses the pair.
|
||||
/// BitChat-specific format is always first and a legacy copy follows for
|
||||
/// clients that still subscribe only to kind 1059. There is deliberately
|
||||
/// no date-based cutoff: removal requires a coordinated iOS/Android
|
||||
/// release after supported old clients have migrated. Both encrypt the
|
||||
/// exact same embedded BitChat payload, so receive-side logical-payload
|
||||
/// dedup collapses the pair.
|
||||
static func createPrivateEnvelopePublicationBatch(
|
||||
content: String,
|
||||
recipientPubkey: String,
|
||||
senderIdentity: NostrIdentity,
|
||||
now: Date = Date()
|
||||
senderIdentity: NostrIdentity
|
||||
) throws -> [NostrEvent] {
|
||||
let primary = try createPrivateEnvelope(
|
||||
content: content,
|
||||
recipientPubkey: recipientPubkey,
|
||||
senderIdentity: senderIdentity
|
||||
)
|
||||
guard now < legacyPrivateEnvelopePublicationDeadline else {
|
||||
return [primary]
|
||||
}
|
||||
let compatibilityCopy = try createPrivateEnvelope(
|
||||
content: content,
|
||||
recipientPubkey: recipientPubkey,
|
||||
@@ -176,14 +169,15 @@ struct NostrProtocol {
|
||||
content: String,
|
||||
recipientPubkey: String,
|
||||
senderIdentity: NostrIdentity,
|
||||
format: PrivateEnvelopeWireFormat
|
||||
format: PrivateEnvelopeWireFormat,
|
||||
messageTags: [[String]] = []
|
||||
) throws -> NostrEvent {
|
||||
// 1. Create the unsigned inner BitChat message.
|
||||
let message = NostrEvent(
|
||||
pubkey: senderIdentity.publicKeyHex,
|
||||
createdAt: Date(),
|
||||
kind: format.messageKind,
|
||||
tags: [],
|
||||
tags: messageTags,
|
||||
content: content
|
||||
)
|
||||
|
||||
@@ -279,16 +273,37 @@ struct NostrProtocol {
|
||||
)
|
||||
}
|
||||
|
||||
static func createLegacyPrivateEnvelopeForTesting(
|
||||
static func createPrivateEnvelopeWithInnerTagsForTesting(
|
||||
content: String,
|
||||
recipientPubkey: String,
|
||||
senderIdentity: NostrIdentity
|
||||
senderIdentity: NostrIdentity,
|
||||
innerMessageTags: [[String]]
|
||||
) throws -> NostrEvent {
|
||||
try createPrivateEnvelope(
|
||||
content: content,
|
||||
recipientPubkey: recipientPubkey,
|
||||
senderIdentity: senderIdentity,
|
||||
format: .legacyMislabelledV2
|
||||
format: .bitchatV1,
|
||||
messageTags: innerMessageTags
|
||||
)
|
||||
}
|
||||
|
||||
static func createLegacyPrivateEnvelopeForTesting(
|
||||
content: String,
|
||||
recipientPubkey: String,
|
||||
senderIdentity: NostrIdentity,
|
||||
innerMessageTags: [[String]] = []
|
||||
) throws -> NostrEvent {
|
||||
// Current Android legacy envelopes use exactly one recipient `p` tag
|
||||
// on the unsigned inner kind-14 event; released iOS envelopes use no
|
||||
// inner tags. Tests pass the Android shape explicitly so this helper
|
||||
// cannot silently make the production encoder depend on that quirk.
|
||||
try createPrivateEnvelope(
|
||||
content: content,
|
||||
recipientPubkey: recipientPubkey,
|
||||
senderIdentity: senderIdentity,
|
||||
format: .legacyMislabelledV2,
|
||||
messageTags: innerMessageTags
|
||||
)
|
||||
}
|
||||
|
||||
@@ -637,7 +652,11 @@ struct NostrProtocol {
|
||||
// comes from the seal. Bind its claimed sender and custom kind to that
|
||||
// authenticated layer before exposing content.
|
||||
guard message.kind == format.messageKind.rawValue,
|
||||
message.tags.isEmpty,
|
||||
validInnerMessageTags(
|
||||
message.tags,
|
||||
format: format,
|
||||
recipientPubkey: recipientIdentity.publicKeyHex
|
||||
),
|
||||
message.sig == nil,
|
||||
seal.pubkey == message.pubkey else {
|
||||
throw NostrError.invalidEvent
|
||||
@@ -646,6 +665,23 @@ struct NostrProtocol {
|
||||
return (seal, message)
|
||||
}
|
||||
|
||||
/// Released iOS legacy envelopes used no inner tags, while current
|
||||
/// Android legacy envelopes use exactly the authenticated recipient tag.
|
||||
/// Accept only those two historical shapes for kind 1059. The new kind
|
||||
/// 1402 format remains strict and rejects every inner tag.
|
||||
private static func validInnerMessageTags(
|
||||
_ tags: [[String]],
|
||||
format: PrivateEnvelopeWireFormat,
|
||||
recipientPubkey: String
|
||||
) -> Bool {
|
||||
switch format {
|
||||
case .bitchatV1:
|
||||
return tags.isEmpty
|
||||
case .legacyMislabelledV2:
|
||||
return tags.isEmpty || tags == [["p", recipientPubkey]]
|
||||
}
|
||||
}
|
||||
|
||||
private static func decodePrivateEnvelopeEventJSON(
|
||||
_ json: String,
|
||||
maximumBytes: Int = maximumPrivateEnvelopePlaintextBytes
|
||||
|
||||
@@ -175,6 +175,7 @@ final class NostrRelayManager: ObservableObject {
|
||||
private var subscribeCoalesce: [String: Date] = [:]
|
||||
private var pendingTorConnectionURLs = Set<String>()
|
||||
private var awaitingTorForConnections = false
|
||||
private var awaitingTorForQueueFlush = false
|
||||
private var torReadyWaitAttempts = 0
|
||||
private var cancellables = Set<AnyCancellable>()
|
||||
|
||||
@@ -209,11 +210,44 @@ final class NostrRelayManager: ObservableObject {
|
||||
private var eoseTrackerEpoch = 0
|
||||
private var pendingEOSECallbacks: [String: () -> Void] = [:]
|
||||
|
||||
// Message queue for reliability
|
||||
// Pending sends held only for relays that are not yet connected.
|
||||
// Message queue for reliability. Pending sends are grouped so a private
|
||||
// envelope's primary and legacy migration copies can never be split by
|
||||
// queue eviction. Private batches are protected; lower-priority traffic
|
||||
// is evicted first and a full protected queue rejects the entire new pair.
|
||||
private enum PendingSendPriority: Int {
|
||||
case ephemeral
|
||||
case regular
|
||||
case privateEnvelope
|
||||
}
|
||||
|
||||
private struct PendingSend {
|
||||
var event: NostrEvent
|
||||
let id: UUID
|
||||
let events: [NostrEvent]
|
||||
var pendingRelays: Set<String>
|
||||
var inFlightConnections: [String: ObjectIdentifier]
|
||||
/// Event IDs still awaiting relay durability for each target. A
|
||||
/// private migration batch completes on a relay only after that same
|
||||
/// relay accepts both the primary and compatibility envelopes.
|
||||
var unconfirmedEventIDsByRelay: [String: Set<String>]
|
||||
/// Invalidates delayed acknowledgement timeouts after a relay settles
|
||||
/// or its connection is replaced.
|
||||
var acknowledgementTokensByRelay: [String: UUID]
|
||||
/// Relay rejections received before every frame in a compatibility
|
||||
/// batch finishes writing. Settlement is deferred until the complete
|
||||
/// pair reaches the socket so a fast rejection of the primary cannot
|
||||
/// suppress the legacy migration copy.
|
||||
var rejectedWhileWritingRelays: Set<String>
|
||||
var hasSuccessfulDelivery: Bool
|
||||
let priority: PendingSendPriority
|
||||
let terminalFailure: (() -> Void)?
|
||||
}
|
||||
|
||||
private struct PendingDelivery {
|
||||
let itemID: UUID
|
||||
let events: [NostrEvent]
|
||||
let priority: PendingSendPriority
|
||||
let relayUrl: String
|
||||
let connection: NostrRelayConnectionProtocol
|
||||
}
|
||||
private var messageQueue: [PendingSend] = []
|
||||
private let messageQueueLock = NSLock()
|
||||
@@ -305,7 +339,8 @@ final class NostrRelayManager: ObservableObject {
|
||||
/// Disconnect from all relays
|
||||
func disconnect() {
|
||||
connectionGeneration &+= 1
|
||||
for (_, task) in connections {
|
||||
for (relayURL, task) in connections {
|
||||
clearPendingSendInFlight(relayUrl: relayURL, connection: task)
|
||||
task.cancel(with: .goingAway, reason: nil)
|
||||
}
|
||||
connections.removeAll()
|
||||
@@ -327,6 +362,7 @@ final class NostrRelayManager: ObservableObject {
|
||||
confirmed.forEach { $0(false) }
|
||||
pendingTorConnectionURLs.removeAll()
|
||||
awaitingTorForConnections = false
|
||||
awaitingTorForQueueFlush = false
|
||||
torReadyWaitAttempts = 0
|
||||
updateConnectionStatus()
|
||||
}
|
||||
@@ -351,6 +387,7 @@ final class NostrRelayManager: ObservableObject {
|
||||
pendingEOSECallbacks.removeAll()
|
||||
pendingTorConnectionURLs.removeAll()
|
||||
awaitingTorForConnections = false
|
||||
awaitingTorForQueueFlush = false
|
||||
torReadyWaitAttempts = 0
|
||||
recentInboundEventKeys.removeAll()
|
||||
recentInboundEventKeyOrder.removeAll()
|
||||
@@ -399,38 +436,128 @@ final class NostrRelayManager: ObservableObject {
|
||||
connectToRelays(targets)
|
||||
}
|
||||
|
||||
/// Send an event to specified relays (or all if none specified)
|
||||
/// Send an event to specified relays (or all if none specified).
|
||||
func sendEvent(_ event: NostrEvent, to relayUrls: [String]? = nil) {
|
||||
// Global network policy gate
|
||||
guard dependencies.activationAllowed() else { return }
|
||||
_ = sendEvents(
|
||||
[event],
|
||||
to: relayUrls,
|
||||
priority: Self.pendingSendPriority(for: event),
|
||||
terminalFailure: nil
|
||||
)
|
||||
}
|
||||
|
||||
/// Atomically admits a complete private-envelope migration pair. The
|
||||
/// queue stores the pair as one protected item, so capacity pressure can
|
||||
/// neither retain only kind 1402 nor only kind 1059. A `false` result means
|
||||
/// the entire pair was rejected before any connected-relay sends began.
|
||||
@discardableResult
|
||||
func sendPrivateEnvelopeBatch(
|
||||
_ events: [NostrEvent],
|
||||
to relayUrls: [String]? = nil,
|
||||
terminalFailure: (() -> Void)? = nil
|
||||
) -> Bool {
|
||||
guard events.map(\.kind) == [
|
||||
NostrProtocol.EventKind.privateEnvelope.rawValue,
|
||||
NostrProtocol.EventKind.legacyNIP59GiftWrap.rawValue
|
||||
] else {
|
||||
SecureLogger.error(
|
||||
"Refusing malformed private-envelope migration batch kinds=\(events.map(\.kind))",
|
||||
category: .session
|
||||
)
|
||||
return false
|
||||
}
|
||||
return sendEvents(
|
||||
events,
|
||||
to: relayUrls,
|
||||
priority: .privateEnvelope,
|
||||
terminalFailure: terminalFailure
|
||||
)
|
||||
}
|
||||
|
||||
@discardableResult
|
||||
private func sendEvents(
|
||||
_ events: [NostrEvent],
|
||||
to relayUrls: [String]?,
|
||||
priority: PendingSendPriority,
|
||||
terminalFailure: (() -> Void)?
|
||||
) -> Bool {
|
||||
guard dependencies.activationAllowed(), !events.isEmpty else { return false }
|
||||
let requestedRelays = relayUrls ?? Self.defaultRelays
|
||||
var targetRelays = allowedRelayList(from: requestedRelays)
|
||||
if priority == .privateEnvelope {
|
||||
// Do not admit a protected batch against a relay already in its
|
||||
// terminal cooldown. Such a target has no connection attempt that
|
||||
// could ever retire the queue record; retry callers can try it
|
||||
// again after the cooldown decays.
|
||||
targetRelays.removeAll(where: isPermanentlyFailed)
|
||||
}
|
||||
guard !targetRelays.isEmpty else { return false }
|
||||
|
||||
if priority == .privateEnvelope {
|
||||
// Protected pairs are admitted for every target before any socket
|
||||
// write, including targets already connected. The queue entry
|
||||
// remains authoritative until each relay explicitly accepts both
|
||||
// events (or that relay reaches a terminal outcome).
|
||||
guard enqueuePendingSendBatch(
|
||||
events,
|
||||
pendingRelays: Set(targetRelays),
|
||||
priority: priority,
|
||||
terminalFailure: terminalFailure
|
||||
) != nil else { return false }
|
||||
ensureConnections(to: targetRelays)
|
||||
// Preserve the same fail-closed Tor gate as ordinary sends. The
|
||||
// queue is already durable; a stale socket that still appears
|
||||
// connected must not be written while enforced Tor is unready.
|
||||
if shouldUseTor && dependencies.torEnforced() && !dependencies.torIsReady() {
|
||||
flushMessageQueueWhenTorIsReady()
|
||||
} else {
|
||||
for relayUrl in targetRelays where connectedConnection(for: relayUrl) != nil {
|
||||
flushMessageQueue(for: relayUrl)
|
||||
}
|
||||
}
|
||||
return true
|
||||
}
|
||||
|
||||
if shouldUseTor && dependencies.torEnforced() && !dependencies.torIsReady() {
|
||||
// Fail-closed: nothing touches the network until Tor is up. Queue the
|
||||
// event locally so it survives a slow bootstrap (queued sends flush
|
||||
// when relays connect), then kick off connection setup, which itself
|
||||
// waits for Tor readiness.
|
||||
let targetRelays = allowedRelayList(from: relayUrls ?? Self.defaultRelays)
|
||||
guard !targetRelays.isEmpty else { return }
|
||||
enqueuePendingSend(event, pendingRelays: Set(targetRelays))
|
||||
// complete item locally so it survives a slow bootstrap, then kick
|
||||
// off connection setup, which itself waits for Tor readiness.
|
||||
guard enqueuePendingSendBatch(
|
||||
events,
|
||||
pendingRelays: Set(targetRelays),
|
||||
priority: priority,
|
||||
terminalFailure: nil
|
||||
) != nil else { return false }
|
||||
ensureConnections(to: targetRelays)
|
||||
return
|
||||
return true
|
||||
}
|
||||
let requestedRelays = relayUrls ?? Self.defaultRelays
|
||||
let targetRelays = allowedRelayList(from: requestedRelays)
|
||||
guard !targetRelays.isEmpty else { return }
|
||||
ensureConnections(to: targetRelays)
|
||||
|
||||
// Attempt immediate send to relays with active connections; queue the rest
|
||||
// Ordinary single events retain their existing immediate-send path.
|
||||
var connectedTargets: [(String, NostrRelayConnectionProtocol)] = []
|
||||
var stillPending = Set<String>()
|
||||
for relayUrl in targetRelays {
|
||||
if let connection = connectedConnection(for: relayUrl) {
|
||||
sendToRelay(event: event, connection: connection, relayUrl: relayUrl)
|
||||
connectedTargets.append((relayUrl, connection))
|
||||
} else {
|
||||
stillPending.insert(relayUrl)
|
||||
}
|
||||
}
|
||||
if !stillPending.isEmpty {
|
||||
enqueuePendingSend(event, pendingRelays: stillPending)
|
||||
guard enqueuePendingSendBatch(
|
||||
events,
|
||||
pendingRelays: stillPending,
|
||||
priority: priority,
|
||||
terminalFailure: nil
|
||||
) != nil else { return false }
|
||||
}
|
||||
|
||||
ensureConnections(to: targetRelays)
|
||||
for (relayUrl, connection) in connectedTargets {
|
||||
for event in events {
|
||||
sendToRelay(event: event, connection: connection, relayUrl: relayUrl)
|
||||
}
|
||||
}
|
||||
return true
|
||||
}
|
||||
|
||||
/// Attempts an event only on currently connected target relays and
|
||||
@@ -524,66 +651,452 @@ final class NostrRelayManager: ObservableObject {
|
||||
state.completion(false)
|
||||
}
|
||||
|
||||
private func enqueuePendingSend(_ event: NostrEvent, pendingRelays: Set<String>) {
|
||||
messageQueueLock.lock()
|
||||
messageQueue.append(PendingSend(event: event, pendingRelays: pendingRelays))
|
||||
let overflow = messageQueue.count - TransportConfig.nostrPendingSendQueueCap
|
||||
if overflow > 0 {
|
||||
messageQueue.removeFirst(overflow)
|
||||
private static func pendingSendPriority(for event: NostrEvent) -> PendingSendPriority {
|
||||
switch event.kind {
|
||||
case NostrProtocol.EventKind.ephemeralEvent.rawValue,
|
||||
NostrProtocol.EventKind.geohashPresence.rawValue:
|
||||
return .ephemeral
|
||||
default:
|
||||
return .regular
|
||||
}
|
||||
}
|
||||
|
||||
/// Atomically append a whole send item while preserving protected private
|
||||
/// batches. Capacity is measured in events, not queue records.
|
||||
private func enqueuePendingSendBatch(
|
||||
_ events: [NostrEvent],
|
||||
pendingRelays: Set<String>,
|
||||
priority: PendingSendPriority,
|
||||
terminalFailure: (() -> Void)?
|
||||
) -> UUID? {
|
||||
guard !events.isEmpty, !pendingRelays.isEmpty,
|
||||
events.count <= TransportConfig.nostrPendingSendQueueCap else {
|
||||
SecureLogger.error("Refusing invalid or oversized relay send batch", category: .session)
|
||||
return nil
|
||||
}
|
||||
|
||||
messageQueueLock.lock()
|
||||
let queuedEventCount = messageQueue.reduce(0) { $0 + $1.events.count }
|
||||
let overflow = queuedEventCount + events.count - TransportConfig.nostrPendingSendQueueCap
|
||||
var removalIndexes: [Int] = []
|
||||
var evictedEventCount = 0
|
||||
if overflow > 0 {
|
||||
var eventsToFree = overflow
|
||||
let candidates = messageQueue.indices
|
||||
.filter {
|
||||
messageQueue[$0].priority != .privateEnvelope &&
|
||||
messageQueue[$0].priority.rawValue <= priority.rawValue
|
||||
}
|
||||
.sorted {
|
||||
let left = messageQueue[$0].priority.rawValue
|
||||
let right = messageQueue[$1].priority.rawValue
|
||||
return left == right ? $0 < $1 : left < right
|
||||
}
|
||||
for index in candidates where eventsToFree > 0 {
|
||||
removalIndexes.append(index)
|
||||
let count = messageQueue[index].events.count
|
||||
evictedEventCount += count
|
||||
eventsToFree -= count
|
||||
}
|
||||
if eventsToFree > 0 {
|
||||
messageQueueLock.unlock()
|
||||
SecureLogger.error(
|
||||
"Relay send queue protected-capacity exhausted; rejected entire \(events.count)-event batch",
|
||||
category: .session
|
||||
)
|
||||
return nil
|
||||
}
|
||||
for index in removalIndexes.sorted(by: >) {
|
||||
messageQueue.remove(at: index)
|
||||
}
|
||||
}
|
||||
let itemID = UUID()
|
||||
messageQueue.append(PendingSend(
|
||||
id: itemID,
|
||||
events: events,
|
||||
pendingRelays: pendingRelays,
|
||||
inFlightConnections: [:],
|
||||
unconfirmedEventIDsByRelay: [:],
|
||||
acknowledgementTokensByRelay: [:],
|
||||
rejectedWhileWritingRelays: [],
|
||||
hasSuccessfulDelivery: false,
|
||||
priority: priority,
|
||||
terminalFailure: terminalFailure
|
||||
))
|
||||
messageQueueLock.unlock()
|
||||
guard overflow > 0 else { return }
|
||||
// Dropped events are ephemeral (presence/geo), so no status surfacing
|
||||
// is needed — but the drops should be visible. Sampled so a sustained
|
||||
// relay stall can't flood the log.
|
||||
pendingSendDropCount += overflow
|
||||
if pendingSendDropCount == 1 ||
|
||||
guard evictedEventCount > 0 else { return itemID }
|
||||
let isFirstEviction = pendingSendDropCount == 0
|
||||
pendingSendDropCount += evictedEventCount
|
||||
if isFirstEviction ||
|
||||
pendingSendDropCount.isMultiple(of: TransportConfig.nostrPendingSendDropLogInterval) {
|
||||
SecureLogger.warning(
|
||||
"📤 Relay send queue full — dropped \(pendingSendDropCount) oldest event(s)",
|
||||
"📤 Relay send queue full — evicted \(pendingSendDropCount) lower-priority event(s)",
|
||||
category: .session
|
||||
)
|
||||
}
|
||||
return itemID
|
||||
}
|
||||
|
||||
/// Try to flush any queued messages for relays that are now connected.
|
||||
private func flushMessageQueue(for relayUrl: String? = nil) {
|
||||
guard !(shouldUseTor && dependencies.torEnforced() && !dependencies.torIsReady()) else {
|
||||
flushMessageQueueWhenTorIsReady()
|
||||
return
|
||||
}
|
||||
var deliveries: [PendingDelivery] = []
|
||||
messageQueueLock.lock()
|
||||
defer { messageQueueLock.unlock() }
|
||||
guard !messageQueue.isEmpty else { return }
|
||||
if let target = relayUrl {
|
||||
// Flush only for a specific relay
|
||||
for i in (0..<messageQueue.count).reversed() {
|
||||
var item = messageQueue[i]
|
||||
if item.pendingRelays.contains(target), let conn = connectedConnection(for: target) {
|
||||
sendToRelay(event: item.event, connection: conn, relayUrl: target)
|
||||
item.pendingRelays.remove(target)
|
||||
if item.pendingRelays.isEmpty {
|
||||
messageQueue.remove(at: i)
|
||||
} else {
|
||||
messageQueue[i] = item
|
||||
for index in (0..<messageQueue.count).reversed() {
|
||||
var item = messageQueue[index]
|
||||
let targets = relayUrl.map { [$0] } ?? Array(item.pendingRelays)
|
||||
for target in targets {
|
||||
guard item.pendingRelays.contains(target),
|
||||
item.inFlightConnections[target] == nil,
|
||||
let connection = connectedConnection(for: target) else { continue }
|
||||
deliveries.append(PendingDelivery(
|
||||
itemID: item.id,
|
||||
events: item.events,
|
||||
priority: item.priority,
|
||||
relayUrl: target,
|
||||
connection: connection
|
||||
))
|
||||
if item.priority == .privateEnvelope {
|
||||
// Keep the relay pending until it explicitly accepts both
|
||||
// events. Socket writes only prove that bytes left this
|
||||
// process. `inFlightConnections` prevents duplicate
|
||||
// flushes while writes or relay acknowledgements are
|
||||
// outstanding.
|
||||
item.inFlightConnections[target] = ObjectIdentifier(connection)
|
||||
if item.unconfirmedEventIDsByRelay[target] == nil {
|
||||
item.unconfirmedEventIDsByRelay[target] = Set(item.events.map(\.id))
|
||||
}
|
||||
item.acknowledgementTokensByRelay.removeValue(forKey: target)
|
||||
} else {
|
||||
item.pendingRelays.remove(target)
|
||||
}
|
||||
}
|
||||
} else {
|
||||
// Flush for any relays that now have connections
|
||||
for i in (0..<messageQueue.count).reversed() {
|
||||
var item = messageQueue[i]
|
||||
for url in item.pendingRelays {
|
||||
if let conn = connectedConnection(for: url) {
|
||||
sendToRelay(event: item.event, connection: conn, relayUrl: url)
|
||||
item.pendingRelays.remove(url)
|
||||
}
|
||||
if item.pendingRelays.isEmpty {
|
||||
messageQueue.remove(at: index)
|
||||
} else {
|
||||
messageQueue[index] = item
|
||||
}
|
||||
}
|
||||
messageQueueLock.unlock()
|
||||
|
||||
// Never invoke WebSocket callbacks while holding messageQueueLock.
|
||||
for delivery in deliveries {
|
||||
if delivery.priority == .privateEnvelope {
|
||||
sendPrivateEnvelopeBatchToRelay(
|
||||
delivery.events,
|
||||
itemID: delivery.itemID,
|
||||
connection: delivery.connection,
|
||||
relayUrl: delivery.relayUrl
|
||||
) { [weak self] succeeded in
|
||||
self?.completePrivateEnvelopeBatchDelivery(
|
||||
itemID: delivery.itemID,
|
||||
relayUrl: delivery.relayUrl,
|
||||
connection: delivery.connection,
|
||||
succeeded: succeeded
|
||||
)
|
||||
}
|
||||
if item.pendingRelays.isEmpty {
|
||||
messageQueue.remove(at: i)
|
||||
} else {
|
||||
messageQueue[i] = item
|
||||
} else {
|
||||
for event in delivery.events {
|
||||
sendToRelay(
|
||||
event: event,
|
||||
connection: delivery.connection,
|
||||
relayUrl: delivery.relayUrl
|
||||
)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private func flushMessageQueueWhenTorIsReady() {
|
||||
guard !awaitingTorForQueueFlush else { return }
|
||||
awaitingTorForQueueFlush = true
|
||||
let generation = connectionGeneration
|
||||
dependencies.awaitTorReady { [weak self] ready in
|
||||
guard let self else { return }
|
||||
guard generation == self.connectionGeneration else { return }
|
||||
self.awaitingTorForQueueFlush = false
|
||||
guard ready else { return }
|
||||
self.flushMessageQueue(for: nil)
|
||||
}
|
||||
}
|
||||
|
||||
/// Send both copies in order and report one result for the complete pair.
|
||||
/// Stop on the first failure; the still-pending queue item retries both
|
||||
/// after reconnection, which is safe because receivers deduplicate twins.
|
||||
private func sendPrivateEnvelopeBatchToRelay(
|
||||
_ events: [NostrEvent],
|
||||
itemID: UUID,
|
||||
connection: NostrRelayConnectionProtocol,
|
||||
relayUrl: String,
|
||||
index: Int = 0,
|
||||
completion: @escaping (Bool) -> Void
|
||||
) {
|
||||
guard isCurrentPrivateEnvelopeDelivery(
|
||||
itemID: itemID,
|
||||
relayUrl: relayUrl,
|
||||
connection: connection
|
||||
) else {
|
||||
completion(false)
|
||||
return
|
||||
}
|
||||
guard index < events.count else {
|
||||
completion(true)
|
||||
return
|
||||
}
|
||||
sendToRelay(
|
||||
event: events[index],
|
||||
connection: connection,
|
||||
relayUrl: relayUrl
|
||||
) { [weak self] succeeded in
|
||||
guard succeeded,
|
||||
let self,
|
||||
self.connections[relayUrl] === connection else {
|
||||
completion(false)
|
||||
return
|
||||
}
|
||||
self.sendPrivateEnvelopeBatchToRelay(
|
||||
events,
|
||||
itemID: itemID,
|
||||
connection: connection,
|
||||
relayUrl: relayUrl,
|
||||
index: index + 1,
|
||||
completion: completion
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
private func completePrivateEnvelopeBatchDelivery(
|
||||
itemID: UUID,
|
||||
relayUrl: String,
|
||||
connection: NostrRelayConnectionProtocol,
|
||||
succeeded: Bool
|
||||
) {
|
||||
var acknowledgementToken: UUID?
|
||||
var wasCurrentDelivery = false
|
||||
var settledEarlyRejection = false
|
||||
var terminalFailure: (() -> Void)?
|
||||
var settledPrivateEnvelopeIDs: [String] = []
|
||||
messageQueueLock.lock()
|
||||
if let index = messageQueue.firstIndex(where: { $0.id == itemID }) {
|
||||
var item = messageQueue[index]
|
||||
guard item.inFlightConnections[relayUrl] == ObjectIdentifier(connection) else {
|
||||
messageQueueLock.unlock()
|
||||
return
|
||||
}
|
||||
wasCurrentDelivery = true
|
||||
if succeeded {
|
||||
if item.rejectedWhileWritingRelays.remove(relayUrl) != nil {
|
||||
// A rejection raced the socket-write callbacks. Both
|
||||
// compatibility frames have now been attempted, so the
|
||||
// relay can settle without suppressing the legacy copy.
|
||||
settledEarlyRejection = true
|
||||
item.pendingRelays.remove(relayUrl)
|
||||
item.inFlightConnections.removeValue(forKey: relayUrl)
|
||||
item.unconfirmedEventIDsByRelay.removeValue(forKey: relayUrl)
|
||||
item.acknowledgementTokensByRelay.removeValue(forKey: relayUrl)
|
||||
if item.pendingRelays.isEmpty {
|
||||
settledPrivateEnvelopeIDs = item.events.map(\.id)
|
||||
if !item.hasSuccessfulDelivery {
|
||||
terminalFailure = item.terminalFailure
|
||||
}
|
||||
messageQueue.remove(at: index)
|
||||
} else {
|
||||
messageQueue[index] = item
|
||||
}
|
||||
} else {
|
||||
// Both frames reached the socket. Retain the queue item and
|
||||
// in-flight ownership until matching NIP-01 OKs arrive.
|
||||
let token = UUID()
|
||||
item.acknowledgementTokensByRelay[relayUrl] = token
|
||||
acknowledgementToken = token
|
||||
messageQueue[index] = item
|
||||
}
|
||||
} else {
|
||||
item.inFlightConnections.removeValue(forKey: relayUrl)
|
||||
item.acknowledgementTokensByRelay.removeValue(forKey: relayUrl)
|
||||
// A partial write must replay the complete pair on the
|
||||
// replacement connection. Do not let a rejection from this
|
||||
// incomplete attempt retire the durable queue item.
|
||||
item.rejectedWhileWritingRelays.remove(relayUrl)
|
||||
messageQueue[index] = item
|
||||
}
|
||||
}
|
||||
messageQueueLock.unlock()
|
||||
|
||||
settledPrivateEnvelopeIDs.forEach {
|
||||
_ = Self.pendingPrivateEnvelopeIDs.remove($0)
|
||||
}
|
||||
terminalFailure?()
|
||||
|
||||
// A relay rejection or timeout may have already settled and removed
|
||||
// this attempt while a stale socket-write callback was outstanding.
|
||||
// It must not tear down an otherwise healthy current connection.
|
||||
guard wasCurrentDelivery else { return }
|
||||
guard !settledEarlyRejection else { return }
|
||||
if let acknowledgementToken {
|
||||
dependencies.scheduleAfter(
|
||||
TransportConfig.nostrConfirmedSendAckTimeoutSeconds
|
||||
) { [weak self] in
|
||||
Task { @MainActor [weak self] in
|
||||
self?.timeoutPrivateEnvelopeBatchDelivery(
|
||||
itemID: itemID,
|
||||
relayUrl: relayUrl,
|
||||
token: acknowledgementToken
|
||||
)
|
||||
}
|
||||
}
|
||||
return
|
||||
}
|
||||
guard !succeeded else { return }
|
||||
// A failed WebSocket write means this connection cannot advance the
|
||||
// durable queue. Keep the original pair pending and reconnect with the
|
||||
// manager's bounded backoff; the next successful ping flushes it.
|
||||
handleDisconnection(
|
||||
relayUrl: relayUrl,
|
||||
error: NSError(
|
||||
domain: "NostrRelayPrivateEnvelopeBatch",
|
||||
code: 1,
|
||||
userInfo: [NSLocalizedDescriptionKey: "private-envelope batch write failed"]
|
||||
),
|
||||
connection: connection
|
||||
)
|
||||
}
|
||||
|
||||
private func isCurrentPrivateEnvelopeDelivery(
|
||||
itemID: UUID,
|
||||
relayUrl: String,
|
||||
connection: NostrRelayConnectionProtocol
|
||||
) -> Bool {
|
||||
messageQueueLock.lock()
|
||||
defer { messageQueueLock.unlock() }
|
||||
guard let item = messageQueue.first(where: { $0.id == itemID }) else {
|
||||
return false
|
||||
}
|
||||
return item.inFlightConnections[relayUrl] == ObjectIdentifier(connection)
|
||||
}
|
||||
|
||||
/// Applies one NIP-01 OK to every matching in-flight private batch. Event
|
||||
/// IDs are immutable, so a standards-compliant `duplicate:` rejection is
|
||||
/// equivalent to acceptance during replay: the relay already has it.
|
||||
@discardableResult
|
||||
private func resolvePrivateEnvelopeAcknowledgement(
|
||||
eventID: String,
|
||||
relayUrl: String,
|
||||
accepted: Bool
|
||||
) -> Bool {
|
||||
guard let connection = connections[relayUrl] else { return false }
|
||||
let connectionID = ObjectIdentifier(connection)
|
||||
var matched = false
|
||||
var terminalFailures: [() -> Void] = []
|
||||
var settledPrivateEnvelopeIDs: [String] = []
|
||||
|
||||
messageQueueLock.lock()
|
||||
for index in (0..<messageQueue.count).reversed() {
|
||||
var item = messageQueue[index]
|
||||
guard item.priority == .privateEnvelope,
|
||||
item.pendingRelays.contains(relayUrl),
|
||||
item.inFlightConnections[relayUrl] == connectionID,
|
||||
item.events.contains(where: { $0.id == eventID }),
|
||||
var unconfirmed = item.unconfirmedEventIDsByRelay[relayUrl],
|
||||
unconfirmed.contains(eventID) else {
|
||||
continue
|
||||
}
|
||||
matched = true
|
||||
|
||||
if accepted {
|
||||
unconfirmed.remove(eventID)
|
||||
if unconfirmed.isEmpty &&
|
||||
!item.rejectedWhileWritingRelays.contains(relayUrl) {
|
||||
item.hasSuccessfulDelivery = true
|
||||
item.pendingRelays.remove(relayUrl)
|
||||
item.inFlightConnections.removeValue(forKey: relayUrl)
|
||||
item.unconfirmedEventIDsByRelay.removeValue(forKey: relayUrl)
|
||||
item.acknowledgementTokensByRelay.removeValue(forKey: relayUrl)
|
||||
} else {
|
||||
item.unconfirmedEventIDsByRelay[relayUrl] = unconfirmed
|
||||
}
|
||||
} else if item.acknowledgementTokensByRelay[relayUrl] == nil {
|
||||
// The primary can be rejected before its send completion
|
||||
// starts the legacy write. Record that outcome, but keep this
|
||||
// delivery authoritative until both compatibility frames
|
||||
// finish writing (or a write failure makes the pair replay).
|
||||
item.rejectedWhileWritingRelays.insert(relayUrl)
|
||||
} else {
|
||||
// One rejected half means this relay never durably received the
|
||||
// complete migration pair. Other target relays may still
|
||||
// succeed; only the final all-relay failure is surfaced.
|
||||
item.pendingRelays.remove(relayUrl)
|
||||
item.inFlightConnections.removeValue(forKey: relayUrl)
|
||||
item.unconfirmedEventIDsByRelay.removeValue(forKey: relayUrl)
|
||||
item.acknowledgementTokensByRelay.removeValue(forKey: relayUrl)
|
||||
item.rejectedWhileWritingRelays.remove(relayUrl)
|
||||
}
|
||||
|
||||
if item.pendingRelays.isEmpty {
|
||||
settledPrivateEnvelopeIDs.append(contentsOf: item.events.map(\.id))
|
||||
if !item.hasSuccessfulDelivery, let terminalFailure = item.terminalFailure {
|
||||
terminalFailures.append(terminalFailure)
|
||||
}
|
||||
messageQueue.remove(at: index)
|
||||
} else {
|
||||
messageQueue[index] = item
|
||||
}
|
||||
}
|
||||
messageQueueLock.unlock()
|
||||
|
||||
settledPrivateEnvelopeIDs.forEach {
|
||||
_ = Self.pendingPrivateEnvelopeIDs.remove($0)
|
||||
}
|
||||
terminalFailures.forEach { $0() }
|
||||
return matched
|
||||
}
|
||||
|
||||
private func timeoutPrivateEnvelopeBatchDelivery(
|
||||
itemID: UUID,
|
||||
relayUrl: String,
|
||||
token: UUID
|
||||
) {
|
||||
var terminalFailure: (() -> Void)?
|
||||
var timedOut = false
|
||||
var settledPrivateEnvelopeIDs: [String] = []
|
||||
|
||||
messageQueueLock.lock()
|
||||
if let index = messageQueue.firstIndex(where: { $0.id == itemID }) {
|
||||
var item = messageQueue[index]
|
||||
if item.acknowledgementTokensByRelay[relayUrl] == token {
|
||||
timedOut = true
|
||||
item.pendingRelays.remove(relayUrl)
|
||||
item.inFlightConnections.removeValue(forKey: relayUrl)
|
||||
item.unconfirmedEventIDsByRelay.removeValue(forKey: relayUrl)
|
||||
item.acknowledgementTokensByRelay.removeValue(forKey: relayUrl)
|
||||
item.rejectedWhileWritingRelays.remove(relayUrl)
|
||||
if item.pendingRelays.isEmpty {
|
||||
settledPrivateEnvelopeIDs = item.events.map(\.id)
|
||||
if !item.hasSuccessfulDelivery {
|
||||
terminalFailure = item.terminalFailure
|
||||
}
|
||||
messageQueue.remove(at: index)
|
||||
} else {
|
||||
messageQueue[index] = item
|
||||
}
|
||||
}
|
||||
}
|
||||
messageQueueLock.unlock()
|
||||
|
||||
guard timedOut else { return }
|
||||
settledPrivateEnvelopeIDs.forEach {
|
||||
_ = Self.pendingPrivateEnvelopeIDs.remove($0)
|
||||
}
|
||||
SecureLogger.warning(
|
||||
"📮 Relay did not acknowledge complete private-envelope pair relay=\(relayUrl)",
|
||||
category: .session
|
||||
)
|
||||
terminalFailure?()
|
||||
}
|
||||
|
||||
private func connectedConnection(for relayUrl: String) -> NostrRelayConnectionProtocol? {
|
||||
guard let connection = connections[relayUrl],
|
||||
relays.first(where: { $0.url == relayUrl })?.isConnected == true else {
|
||||
@@ -599,9 +1112,28 @@ final class NostrRelayManager: ObservableObject {
|
||||
relayUrls: [String]? = nil,
|
||||
handler: @escaping (NostrEvent) -> Void,
|
||||
onEOSE: (() -> Void)? = nil
|
||||
) {
|
||||
subscribe(
|
||||
filters: [filter],
|
||||
id: id,
|
||||
relayUrls: relayUrls,
|
||||
handler: handler,
|
||||
onEOSE: onEOSE
|
||||
)
|
||||
}
|
||||
|
||||
/// Subscribe with independent Nostr filters in one REQ. Each filter keeps
|
||||
/// its own relay-side limit; this is required for mixed-version mailbox
|
||||
/// recovery so kind 1402 cannot starve kind 1059 (or the reverse).
|
||||
func subscribe(
|
||||
filters: [NostrFilter],
|
||||
id: String = UUID().uuidString,
|
||||
relayUrls: [String]? = nil,
|
||||
handler: @escaping (NostrEvent) -> Void,
|
||||
onEOSE: (() -> Void)? = nil
|
||||
) {
|
||||
// Global network policy gate
|
||||
guard dependencies.activationAllowed() else { return }
|
||||
guard dependencies.activationAllowed(), !filters.isEmpty else { return }
|
||||
// Coalesce rapid duplicate subscribe requests even while Tor readiness is pending.
|
||||
let now = dependencies.now()
|
||||
if let last = subscribeCoalesce[id], now.timeIntervalSince(last) < subscribeCoalesceInterval {
|
||||
@@ -610,7 +1142,7 @@ final class NostrRelayManager: ObservableObject {
|
||||
subscribeCoalesce[id] = now
|
||||
messageHandlers[id] = handler
|
||||
|
||||
let req = NostrRequest.subscribe(id: id, filters: [filter])
|
||||
let req = NostrRequest.subscribe(id: id, filters: filters)
|
||||
|
||||
do {
|
||||
let message = try encoder.encode(req)
|
||||
@@ -677,6 +1209,7 @@ final class NostrRelayManager: ObservableObject {
|
||||
ensureConnections(to: Self.defaultRelays)
|
||||
}
|
||||
} else {
|
||||
var terminalFailures: [() -> Void] = []
|
||||
for url in Self.defaultRelays {
|
||||
if let connection = connections[url] {
|
||||
connection.cancel(with: .goingAway, reason: nil)
|
||||
@@ -684,18 +1217,9 @@ final class NostrRelayManager: ObservableObject {
|
||||
connections.removeValue(forKey: url)
|
||||
subscriptions.removeValue(forKey: url)
|
||||
pendingSubscriptions.removeValue(forKey: url)
|
||||
terminalFailures.append(contentsOf: retirePendingSendRelay(url))
|
||||
}
|
||||
messageQueueLock.lock()
|
||||
for index in (0..<messageQueue.count).reversed() {
|
||||
var item = messageQueue[index]
|
||||
item.pendingRelays.subtract(Self.defaultRelaySet)
|
||||
if item.pendingRelays.isEmpty {
|
||||
messageQueue.remove(at: index)
|
||||
} else {
|
||||
messageQueue[index] = item
|
||||
}
|
||||
}
|
||||
messageQueueLock.unlock()
|
||||
terminalFailures.forEach { $0() }
|
||||
relays.removeAll { Self.defaultRelaySet.contains($0.url) }
|
||||
updateConnectionStatus()
|
||||
}
|
||||
@@ -1168,11 +1692,28 @@ final class NostrRelayManager: ObservableObject {
|
||||
}
|
||||
case .ok(let eventId, let success, let reason):
|
||||
resolveConfirmedSend(eventID: eventId, relayURL: relayUrl, accepted: success)
|
||||
if success {
|
||||
let normalizedReason = reason
|
||||
.trimmingCharacters(in: .whitespacesAndNewlines)
|
||||
.lowercased()
|
||||
let isDurableDuplicate = !success && normalizedReason.hasPrefix("duplicate:")
|
||||
let matchedPrivateEnvelope = resolvePrivateEnvelopeAcknowledgement(
|
||||
eventID: eventId,
|
||||
relayUrl: relayUrl,
|
||||
accepted: success || isDurableDuplicate
|
||||
)
|
||||
if success || isDurableDuplicate {
|
||||
_ = Self.pendingPrivateEnvelopeIDs.remove(eventId)
|
||||
SecureLogger.debug("✅ Accepted id=\(eventId.prefix(16))… relay=\(relayUrl)", category: .session)
|
||||
SecureLogger.debug(
|
||||
"✅ Accepted id=\(eventId.prefix(16))… relay=\(relayUrl)\(isDurableDuplicate ? " (already stored)" : "")",
|
||||
category: .session
|
||||
)
|
||||
} else {
|
||||
let isPrivateEnvelope = Self.pendingPrivateEnvelopeIDs.remove(eventId) != nil
|
||||
// Evaluate the removal independently: `||` short-circuiting
|
||||
// must not leak a registered ID when the queue match is true.
|
||||
let wasRegisteredPrivateEnvelope =
|
||||
Self.pendingPrivateEnvelopeIDs.remove(eventId) != nil
|
||||
let isPrivateEnvelope =
|
||||
matchedPrivateEnvelope || wasRegisteredPrivateEnvelope
|
||||
if isPrivateEnvelope {
|
||||
SecureLogger.warning("📮 Rejected id=\(eventId.prefix(16))… relay=\(relayUrl) reason=\(reason)", category: .session)
|
||||
} else {
|
||||
@@ -1285,6 +1826,7 @@ final class NostrRelayManager: ObservableObject {
|
||||
connection: NostrRelayConnectionProtocol? = nil
|
||||
) {
|
||||
if let connection, connections[relayUrl] !== connection { return }
|
||||
clearPendingSendInFlight(relayUrl: relayUrl, connection: connection)
|
||||
connections.removeValue(forKey: relayUrl)
|
||||
subscriptions.removeValue(forKey: relayUrl)
|
||||
let awaitingConfirmation = confirmedSends.compactMap { eventID, state in
|
||||
@@ -1305,7 +1847,10 @@ final class NostrRelayManager: ObservableObject {
|
||||
let ns = error as NSError
|
||||
if errorDescription.contains("hostname could not be found") ||
|
||||
errorDescription.contains("dns") ||
|
||||
(ns.domain == NSURLErrorDomain && ns.code == NSURLErrorBadServerResponse) {
|
||||
(ns.domain == NSURLErrorDomain && (
|
||||
ns.code == NSURLErrorBadServerResponse ||
|
||||
ns.code == NSURLErrorCannotFindHost
|
||||
)) {
|
||||
if relays.first(where: { $0.url == relayUrl })?.lastError == nil {
|
||||
SecureLogger.warning("Nostr relay permanent failure for \(relayUrl) - not retrying (code=\(ns.code))", category: .session)
|
||||
}
|
||||
@@ -1315,6 +1860,8 @@ final class NostrRelayManager: ObservableObject {
|
||||
relays[index].nextReconnectTime = nil
|
||||
}
|
||||
pendingSubscriptions[relayUrl] = nil
|
||||
let terminalFailures = retirePendingSendRelay(relayUrl)
|
||||
terminalFailures.forEach { $0() }
|
||||
return
|
||||
}
|
||||
|
||||
@@ -1326,6 +1873,8 @@ final class NostrRelayManager: ObservableObject {
|
||||
// Stop attempting after max attempts
|
||||
if relays[index].reconnectAttempts >= maxReconnectAttempts {
|
||||
SecureLogger.warning("Max reconnection attempts (\(maxReconnectAttempts)) reached for \(relayUrl)", category: .session)
|
||||
let terminalFailures = retirePendingSendRelay(relayUrl)
|
||||
terminalFailures.forEach { $0() }
|
||||
return
|
||||
}
|
||||
|
||||
@@ -1363,6 +1912,66 @@ final class NostrRelayManager: ObservableObject {
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// A relay in terminal cooldown cannot make progress on queued sends. Drop
|
||||
/// it from every target set so one unavailable default relay cannot pin
|
||||
/// otherwise-delivered protected batches forever. If this was the last
|
||||
/// target and no relay completed the pair, report whole-batch failure to
|
||||
/// the originating transport for visible failure or retry.
|
||||
private func retirePendingSendRelay(_ relayUrl: String) -> [() -> Void] {
|
||||
var terminalFailures: [() -> Void] = []
|
||||
var settledPrivateEnvelopeIDs: [String] = []
|
||||
messageQueueLock.lock()
|
||||
for index in (0..<messageQueue.count).reversed() {
|
||||
var item = messageQueue[index]
|
||||
guard item.pendingRelays.remove(relayUrl) != nil else { continue }
|
||||
item.inFlightConnections.removeValue(forKey: relayUrl)
|
||||
item.unconfirmedEventIDsByRelay.removeValue(forKey: relayUrl)
|
||||
item.acknowledgementTokensByRelay.removeValue(forKey: relayUrl)
|
||||
item.rejectedWhileWritingRelays.remove(relayUrl)
|
||||
if item.pendingRelays.isEmpty {
|
||||
if item.priority == .privateEnvelope {
|
||||
settledPrivateEnvelopeIDs.append(contentsOf: item.events.map(\.id))
|
||||
}
|
||||
if item.priority == .privateEnvelope,
|
||||
!item.hasSuccessfulDelivery,
|
||||
let terminalFailure = item.terminalFailure {
|
||||
terminalFailures.append(terminalFailure)
|
||||
}
|
||||
messageQueue.remove(at: index)
|
||||
} else {
|
||||
messageQueue[index] = item
|
||||
}
|
||||
}
|
||||
messageQueueLock.unlock()
|
||||
settledPrivateEnvelopeIDs.forEach {
|
||||
_ = Self.pendingPrivateEnvelopeIDs.remove($0)
|
||||
}
|
||||
return terminalFailures
|
||||
}
|
||||
|
||||
/// Release the in-flight marker owned by a dead socket without removing
|
||||
/// the relay from the durable pending set. A stale completion from that
|
||||
/// socket is ignored by its ObjectIdentifier check.
|
||||
private func clearPendingSendInFlight(
|
||||
relayUrl: String,
|
||||
connection: NostrRelayConnectionProtocol?
|
||||
) {
|
||||
let expectedIdentifier = connection.map(ObjectIdentifier.init)
|
||||
messageQueueLock.lock()
|
||||
for index in messageQueue.indices {
|
||||
var item = messageQueue[index]
|
||||
guard let inFlightIdentifier = item.inFlightConnections[relayUrl],
|
||||
expectedIdentifier == nil || expectedIdentifier == inFlightIdentifier else {
|
||||
continue
|
||||
}
|
||||
item.inFlightConnections.removeValue(forKey: relayUrl)
|
||||
item.acknowledgementTokensByRelay.removeValue(forKey: relayUrl)
|
||||
item.rejectedWhileWritingRelays.remove(relayUrl)
|
||||
messageQueue[index] = item
|
||||
}
|
||||
messageQueueLock.unlock()
|
||||
}
|
||||
|
||||
// MARK: - Public Utility Methods
|
||||
|
||||
@@ -1378,6 +1987,10 @@ final class NostrRelayManager: ObservableObject {
|
||||
|
||||
// Disconnect if connected
|
||||
if let connection = connections[normalizedRelayUrl] {
|
||||
clearPendingSendInFlight(
|
||||
relayUrl: normalizedRelayUrl,
|
||||
connection: connection
|
||||
)
|
||||
connection.cancel(with: .goingAway, reason: nil)
|
||||
connections.removeValue(forKey: normalizedRelayUrl)
|
||||
}
|
||||
@@ -1399,7 +2012,13 @@ final class NostrRelayManager: ObservableObject {
|
||||
var debugPendingMessageQueueCount: Int {
|
||||
messageQueueLock.lock()
|
||||
defer { messageQueueLock.unlock() }
|
||||
return messageQueue.count
|
||||
return messageQueue.reduce(0) { $0 + $1.events.count }
|
||||
}
|
||||
|
||||
var debugPendingMessageQueueEventIDsByBatch: [[String]] {
|
||||
messageQueueLock.lock()
|
||||
defer { messageQueueLock.unlock() }
|
||||
return messageQueue.map { $0.events.map(\.id) }
|
||||
}
|
||||
|
||||
func debugPendingSubscriptionCount(for relayUrl: String) -> Int {
|
||||
@@ -1612,19 +2231,17 @@ struct NostrFilter: Encodable {
|
||||
}
|
||||
|
||||
// BitChat private envelopes, plus compatibility legacy envelopes emitted
|
||||
// during the bounded migration and stored by older releases as kind 1059.
|
||||
static func privateEnvelopesFor(pubkey: String, since: Date? = nil) -> NostrFilter {
|
||||
var filter = NostrFilter()
|
||||
filter.kinds = NostrProtocol.acceptedPrivateEnvelopeKinds
|
||||
filter.since = since?.timeIntervalSince1970.toInt()
|
||||
filter.tagFilters = ["p": [pubkey]]
|
||||
// Before the migration deadline each logical payload is stored once
|
||||
// per accepted wire kind. Scale the combined filter so compatibility
|
||||
// copies do not halve the number of logical messages/acks recovered
|
||||
// after a reconnect.
|
||||
filter.limit = TransportConfig.nostrRelayDefaultFetchLimit
|
||||
* NostrProtocol.acceptedPrivateEnvelopeKinds.count
|
||||
return filter
|
||||
// throughout the coordinated migration and stored by older releases as
|
||||
// kind 1059.
|
||||
static func privateEnvelopeFiltersFor(pubkey: String, since: Date? = nil) -> [NostrFilter] {
|
||||
NostrProtocol.acceptedPrivateEnvelopeKinds.map { kind in
|
||||
var filter = NostrFilter()
|
||||
filter.kinds = [kind]
|
||||
filter.since = since?.timeIntervalSince1970.toInt()
|
||||
filter.tagFilters = ["p": [pubkey]]
|
||||
filter.limit = TransportConfig.nostrPrivateEnvelopeFetchLimitPerKind
|
||||
return filter
|
||||
}
|
||||
}
|
||||
|
||||
// For location channels: geohash-scoped ephemeral events (kind 20000) and presence (kind 20001)
|
||||
|
||||
@@ -158,6 +158,16 @@ final class MessageRouter {
|
||||
}
|
||||
}
|
||||
|
||||
/// Wire one typed event sink to every route, including the queue-backed
|
||||
/// Nostr transport. Keeping this at the router boundary prevents a relay
|
||||
/// admission failure from disappearing merely because only the primary
|
||||
/// mesh transport was assigned a UI delegate.
|
||||
func setEventDelegate(_ delegate: TransportEventDelegate?) {
|
||||
for transport in transports {
|
||||
transport.eventDelegate = delegate
|
||||
}
|
||||
}
|
||||
|
||||
// MARK: - Transport Selection
|
||||
|
||||
private func reachableTransport(for peerID: PeerID) -> Transport? {
|
||||
|
||||
@@ -12,8 +12,11 @@ final class NostrTransport: Transport, @unchecked Sendable {
|
||||
let favoriteStatusForPeerID: @MainActor (PeerID) -> FavoritesPersistenceService.FavoriteRelationship?
|
||||
let currentIdentity: @MainActor () throws -> NostrIdentity?
|
||||
let registerPendingPrivateEnvelope: @MainActor (String) -> Void
|
||||
let sendEvent: @MainActor (NostrEvent) -> Void
|
||||
let now: @MainActor () -> Date
|
||||
let sendPrivateEnvelopeBatch: @MainActor (
|
||||
[NostrEvent],
|
||||
@escaping @MainActor () -> Void
|
||||
) -> Bool
|
||||
let envelopeRetryQueue: NostrPrivateEnvelopeRetryQueue
|
||||
/// Emits whether a relay that carries private messages is up
|
||||
/// (fail-closed behind Tor). A connected geohash/custom relay alone
|
||||
/// doesn't count: DM sends target the default relay set and would
|
||||
@@ -23,6 +26,7 @@ final class NostrTransport: Transport, @unchecked Sendable {
|
||||
/// serialize behind each other; `live` passes the process-wide one.
|
||||
let ackPacer: AckPacer
|
||||
|
||||
@MainActor
|
||||
init(
|
||||
notificationCenter: NotificationCenter,
|
||||
loadFavorites: @escaping @MainActor () -> [Data: FavoritesPersistenceService.FavoriteRelationship],
|
||||
@@ -30,11 +34,14 @@ final class NostrTransport: Transport, @unchecked Sendable {
|
||||
favoriteStatusForPeerID: @escaping @MainActor (PeerID) -> FavoritesPersistenceService.FavoriteRelationship?,
|
||||
currentIdentity: @escaping @MainActor () throws -> NostrIdentity?,
|
||||
registerPendingPrivateEnvelope: @escaping @MainActor (String) -> Void,
|
||||
sendEvent: @escaping @MainActor (NostrEvent) -> Void,
|
||||
sendPrivateEnvelopeBatch: @escaping @MainActor (
|
||||
[NostrEvent],
|
||||
@escaping @MainActor () -> Void
|
||||
) -> Bool,
|
||||
scheduleAfter: @escaping @Sendable (TimeInterval, @escaping @Sendable () -> Void) -> Void,
|
||||
relayConnectivity: @escaping @MainActor () -> AnyPublisher<Bool, Never>,
|
||||
ackPacer: AckPacer? = nil,
|
||||
now: @escaping @MainActor () -> Date = Date.init
|
||||
envelopeRetryQueue: NostrPrivateEnvelopeRetryQueue? = nil
|
||||
) {
|
||||
self.notificationCenter = notificationCenter
|
||||
self.loadFavorites = loadFavorites
|
||||
@@ -42,8 +49,12 @@ final class NostrTransport: Transport, @unchecked Sendable {
|
||||
self.favoriteStatusForPeerID = favoriteStatusForPeerID
|
||||
self.currentIdentity = currentIdentity
|
||||
self.registerPendingPrivateEnvelope = registerPendingPrivateEnvelope
|
||||
self.sendEvent = sendEvent
|
||||
self.now = now
|
||||
self.sendPrivateEnvelopeBatch = sendPrivateEnvelopeBatch
|
||||
self.envelopeRetryQueue = envelopeRetryQueue ?? NostrPrivateEnvelopeRetryQueue(
|
||||
sendPrivateEnvelopeBatch: sendPrivateEnvelopeBatch,
|
||||
registerPendingPrivateEnvelope: registerPendingPrivateEnvelope,
|
||||
scheduleAfter: scheduleAfter
|
||||
)
|
||||
self.relayConnectivity = relayConnectivity
|
||||
// Default pacer drives its throttle through the same injected
|
||||
// scheduler, so tests that step scheduleAfter manually keep
|
||||
@@ -60,12 +71,18 @@ final class NostrTransport: Transport, @unchecked Sendable {
|
||||
favoriteStatusForPeerID: { FavoritesPersistenceService.shared.getFavoriteStatus(forPeerID: $0) },
|
||||
currentIdentity: { try idBridge.getCurrentNostrIdentity() },
|
||||
registerPendingPrivateEnvelope: { NostrRelayManager.registerPendingPrivateEnvelope(id: $0) },
|
||||
sendEvent: { NostrRelayManager.shared.sendEvent($0) },
|
||||
sendPrivateEnvelopeBatch: { events, terminalFailure in
|
||||
NostrRelayManager.shared.sendPrivateEnvelopeBatch(
|
||||
events,
|
||||
terminalFailure: terminalFailure
|
||||
)
|
||||
},
|
||||
scheduleAfter: { delay, action in
|
||||
DispatchQueue.main.asyncAfter(deadline: .now() + delay, execute: action)
|
||||
},
|
||||
relayConnectivity: { NostrRelayManager.shared.$isDMRelayConnected.eraseToAnyPublisher() },
|
||||
ackPacer: NostrTransport.sharedAckPacer
|
||||
ackPacer: NostrTransport.sharedAckPacer,
|
||||
envelopeRetryQueue: NostrTransport.sharedEnvelopeRetryQueue
|
||||
)
|
||||
}
|
||||
}
|
||||
@@ -131,7 +148,37 @@ final class NostrTransport: Transport, @unchecked Sendable {
|
||||
}
|
||||
}
|
||||
static let sharedAckPacer = AckPacer()
|
||||
// Geohash acknowledgements use short-lived NostrTransport instances, so
|
||||
// the retry owner must be process-wide. A per-transport cap would still be
|
||||
// globally unbounded under outage as throwaway instances accumulated.
|
||||
@MainActor
|
||||
private static let sharedEnvelopeRetryQueue = NostrPrivateEnvelopeRetryQueue(
|
||||
sendPrivateEnvelopeBatch: { events, terminalFailure in
|
||||
NostrRelayManager.shared.sendPrivateEnvelopeBatch(
|
||||
events,
|
||||
terminalFailure: terminalFailure
|
||||
)
|
||||
},
|
||||
registerPendingPrivateEnvelope: {
|
||||
NostrRelayManager.registerPendingPrivateEnvelope(id: $0)
|
||||
},
|
||||
scheduleAfter: { delay, action in
|
||||
DispatchQueue.main.asyncAfter(deadline: .now() + delay, execute: action)
|
||||
}
|
||||
)
|
||||
|
||||
@MainActor
|
||||
static func resetControlRetriesForPanicWipe() {
|
||||
sharedEnvelopeRetryQueue.removeAll()
|
||||
}
|
||||
|
||||
private enum PrivateEnvelopeFailurePolicy {
|
||||
case userMessage(messageID: String)
|
||||
case retry(retryKey: String)
|
||||
}
|
||||
|
||||
private let dependencies: Dependencies
|
||||
private let envelopeRetryQueue: NostrPrivateEnvelopeRetryQueue
|
||||
private var favoriteStatusObserver: NSObjectProtocol?
|
||||
|
||||
// Reachability Cache (thread-safe)
|
||||
@@ -148,7 +195,9 @@ final class NostrTransport: Transport, @unchecked Sendable {
|
||||
idBridge: NostrIdentityBridge,
|
||||
dependencies: Dependencies? = nil
|
||||
) {
|
||||
self.dependencies = dependencies ?? .live(idBridge: idBridge)
|
||||
let resolvedDependencies = dependencies ?? .live(idBridge: idBridge)
|
||||
self.dependencies = resolvedDependencies
|
||||
self.envelopeRetryQueue = resolvedDependencies.envelopeRetryQueue
|
||||
|
||||
setupObservers()
|
||||
|
||||
@@ -175,6 +224,18 @@ final class NostrTransport: Transport, @unchecked Sendable {
|
||||
}
|
||||
}
|
||||
|
||||
#if DEBUG
|
||||
@MainActor
|
||||
func debugEnqueueControlRetry(key: String, events: [NostrEvent]) {
|
||||
envelopeRetryQueue.enqueue(key: key, events: events, registerPending: false)
|
||||
}
|
||||
|
||||
@MainActor
|
||||
var debugControlRetryCount: Int {
|
||||
envelopeRetryQueue.debugPendingCount
|
||||
}
|
||||
#endif
|
||||
|
||||
private func setupObservers() {
|
||||
favoriteStatusObserver = dependencies.notificationCenter.addObserver(
|
||||
forName: .favoriteStatusChanged,
|
||||
@@ -264,7 +325,12 @@ final class NostrTransport: Transport, @unchecked Sendable {
|
||||
SecureLogger.error("NostrTransport: failed to embed PM packet", category: .session)
|
||||
return
|
||||
}
|
||||
sendPrivateEnvelope(content: embedded, recipientHex: recipientHex, senderIdentity: senderIdentity)
|
||||
sendPrivateEnvelope(
|
||||
content: embedded,
|
||||
recipientHex: recipientHex,
|
||||
senderIdentity: senderIdentity,
|
||||
failurePolicy: .userMessage(messageID: messageID)
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -290,7 +356,14 @@ final class NostrTransport: Transport, @unchecked Sendable {
|
||||
SecureLogger.error("NostrTransport: failed to embed favorite notification", category: .session)
|
||||
return
|
||||
}
|
||||
sendPrivateEnvelope(content: embedded, recipientHex: recipientHex, senderIdentity: senderIdentity)
|
||||
sendPrivateEnvelope(
|
||||
content: embedded,
|
||||
recipientHex: recipientHex,
|
||||
senderIdentity: senderIdentity,
|
||||
failurePolicy: .retry(
|
||||
retryKey: privateEnvelopeRetryKey(content: embedded, recipientHex: recipientHex)
|
||||
)
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -322,7 +395,13 @@ extension NostrTransport {
|
||||
SecureLogger.error("NostrTransport: failed to embed geohash PM packet", category: .session)
|
||||
return
|
||||
}
|
||||
sendPrivateEnvelope(content: embedded, recipientHex: recipientHex, senderIdentity: identity, registerPending: true)
|
||||
sendPrivateEnvelope(
|
||||
content: embedded,
|
||||
recipientHex: recipientHex,
|
||||
senderIdentity: identity,
|
||||
registerPending: true,
|
||||
failurePolicy: .userMessage(messageID: messageID)
|
||||
)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -345,21 +424,89 @@ extension NostrTransport {
|
||||
|
||||
/// Creates and sends a BitChat private-envelope event over Nostr.
|
||||
@MainActor
|
||||
private func sendPrivateEnvelope(content: String, recipientHex: String, senderIdentity: NostrIdentity, registerPending: Bool = false) {
|
||||
private func sendPrivateEnvelope(
|
||||
content: String,
|
||||
recipientHex: String,
|
||||
senderIdentity: NostrIdentity,
|
||||
registerPending: Bool = false,
|
||||
failurePolicy: PrivateEnvelopeFailurePolicy
|
||||
) {
|
||||
guard let events = try? NostrProtocol.createPrivateEnvelopePublicationBatch(
|
||||
content: content,
|
||||
recipientPubkey: recipientHex,
|
||||
senderIdentity: senderIdentity,
|
||||
now: dependencies.now()
|
||||
senderIdentity: senderIdentity
|
||||
) else {
|
||||
SecureLogger.error("NostrTransport: failed to build Nostr private-envelope batch", category: .session)
|
||||
return
|
||||
}
|
||||
let accepted = dependencies.sendPrivateEnvelopeBatch(events) { [self] in
|
||||
handlePrivateEnvelopeFailure(
|
||||
events: events,
|
||||
registerPending: registerPending,
|
||||
policy: failurePolicy
|
||||
)
|
||||
}
|
||||
guard accepted else {
|
||||
SecureLogger.error(
|
||||
"NostrTransport: private-envelope migration pair was not accepted for relay delivery",
|
||||
category: .session
|
||||
)
|
||||
handlePrivateEnvelopeFailure(
|
||||
events: events,
|
||||
registerPending: registerPending,
|
||||
policy: failurePolicy
|
||||
)
|
||||
return
|
||||
}
|
||||
registerPendingPrivateEnvelopesIfNeeded(events, registerPending: registerPending)
|
||||
}
|
||||
|
||||
@MainActor
|
||||
private func handlePrivateEnvelopeFailure(
|
||||
events: [NostrEvent],
|
||||
registerPending: Bool,
|
||||
policy: PrivateEnvelopeFailurePolicy
|
||||
) {
|
||||
switch policy {
|
||||
case .userMessage(let messageID):
|
||||
deliverTransportEvent(.messageDeliveryStatusUpdated(
|
||||
messageID: messageID,
|
||||
status: .failed(reason: String(
|
||||
localized: "content.delivery.reason.not_delivered",
|
||||
comment: "Failure reason shown when a private message could not enter the relay delivery queue"
|
||||
))
|
||||
))
|
||||
case .retry(let retryKey):
|
||||
envelopeRetryQueue.enqueue(
|
||||
key: retryKey,
|
||||
events: events,
|
||||
registerPending: registerPending
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
@MainActor
|
||||
private func registerPendingPrivateEnvelopesIfNeeded(
|
||||
_ events: [NostrEvent],
|
||||
registerPending: Bool
|
||||
) {
|
||||
guard registerPending else { return }
|
||||
for event in events {
|
||||
if registerPending {
|
||||
dependencies.registerPendingPrivateEnvelope(event.id)
|
||||
}
|
||||
dependencies.sendEvent(event)
|
||||
dependencies.registerPendingPrivateEnvelope(event.id)
|
||||
}
|
||||
}
|
||||
|
||||
@MainActor
|
||||
private func privateEnvelopeRetryKey(content: String, recipientHex: String) -> String {
|
||||
"\(recipientHex.lowercased()):\(Data(content.utf8).sha256Fingerprint())"
|
||||
}
|
||||
|
||||
@MainActor
|
||||
private func deliverTransportEvent(_ event: TransportEvent) {
|
||||
if let eventDelegate {
|
||||
eventDelegate.didReceiveTransportEvent(event)
|
||||
} else {
|
||||
delegate?.receiveTransportEvent(event)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -377,7 +524,14 @@ extension NostrTransport {
|
||||
SecureLogger.error("NostrTransport: failed to embed READ ack", category: .session)
|
||||
return
|
||||
}
|
||||
sendPrivateEnvelope(content: ack, recipientHex: recipientHex, senderIdentity: senderIdentity)
|
||||
sendPrivateEnvelope(
|
||||
content: ack,
|
||||
recipientHex: recipientHex,
|
||||
senderIdentity: senderIdentity,
|
||||
failurePolicy: .retry(
|
||||
retryKey: privateEnvelopeRetryKey(content: ack, recipientHex: recipientHex)
|
||||
)
|
||||
)
|
||||
|
||||
case .deliveredDirect(let messageID, let peerID):
|
||||
guard let recipientNpub = resolveRecipientNpub(for: peerID),
|
||||
@@ -388,17 +542,40 @@ extension NostrTransport {
|
||||
SecureLogger.error("NostrTransport: failed to embed DELIVERED ack", category: .session)
|
||||
return
|
||||
}
|
||||
sendPrivateEnvelope(content: ack, recipientHex: recipientHex, senderIdentity: senderIdentity)
|
||||
sendPrivateEnvelope(
|
||||
content: ack,
|
||||
recipientHex: recipientHex,
|
||||
senderIdentity: senderIdentity,
|
||||
failurePolicy: .retry(
|
||||
retryKey: privateEnvelopeRetryKey(content: ack, recipientHex: recipientHex)
|
||||
)
|
||||
)
|
||||
|
||||
case .deliveredGeohash(let messageID, let recipientHex, let identity):
|
||||
SecureLogger.debug("GeoDM: send DELIVERED mid=\(messageID.prefix(8))…", category: .session)
|
||||
guard let embedded = NostrEmbeddedBitChat.encodeAckForNostrNoRecipient(type: .delivered, messageID: messageID, senderPeerID: senderPeerID) else { return }
|
||||
sendPrivateEnvelope(content: embedded, recipientHex: recipientHex, senderIdentity: identity, registerPending: true)
|
||||
sendPrivateEnvelope(
|
||||
content: embedded,
|
||||
recipientHex: recipientHex,
|
||||
senderIdentity: identity,
|
||||
registerPending: true,
|
||||
failurePolicy: .retry(
|
||||
retryKey: privateEnvelopeRetryKey(content: embedded, recipientHex: recipientHex)
|
||||
)
|
||||
)
|
||||
|
||||
case .readGeohash(let messageID, let recipientHex, let identity):
|
||||
SecureLogger.debug("GeoDM: send READ mid=\(messageID.prefix(8))…", category: .session)
|
||||
guard let embedded = NostrEmbeddedBitChat.encodeAckForNostrNoRecipient(type: .readReceipt, messageID: messageID, senderPeerID: senderPeerID) else { return }
|
||||
sendPrivateEnvelope(content: embedded, recipientHex: recipientHex, senderIdentity: identity, registerPending: true)
|
||||
sendPrivateEnvelope(
|
||||
content: embedded,
|
||||
recipientHex: recipientHex,
|
||||
senderIdentity: identity,
|
||||
registerPending: true,
|
||||
failurePolicy: .retry(
|
||||
retryKey: privateEnvelopeRetryKey(content: embedded, recipientHex: recipientHex)
|
||||
)
|
||||
)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -418,3 +595,122 @@ extension NostrTransport {
|
||||
return nil
|
||||
}
|
||||
}
|
||||
|
||||
/// Bounded retry owner for non-user private control payloads. It is separate
|
||||
/// from `NostrTransport` so a scheduled retry from a short-lived geohash
|
||||
/// transport remains valid after that transport deinitializes. Scheduler
|
||||
/// callbacks retain this queue, never the transport; an evicted key simply
|
||||
/// becomes a harmless no-op when its already-scheduled callback fires.
|
||||
@MainActor
|
||||
final class NostrPrivateEnvelopeRetryQueue {
|
||||
private struct PendingRetry {
|
||||
let events: [NostrEvent]
|
||||
let registerPending: Bool
|
||||
var attempt: Int
|
||||
var isScheduled: Bool
|
||||
}
|
||||
|
||||
private let sendPrivateEnvelopeBatch: @MainActor (
|
||||
[NostrEvent],
|
||||
@escaping @MainActor () -> Void
|
||||
) -> Bool
|
||||
private let registerPendingPrivateEnvelope: @MainActor (String) -> Void
|
||||
private let scheduleAfter: @Sendable (
|
||||
TimeInterval,
|
||||
@escaping @Sendable () -> Void
|
||||
) -> Void
|
||||
private var pending: [String: PendingRetry] = [:]
|
||||
private var insertionOrder: [String] = []
|
||||
|
||||
init(
|
||||
sendPrivateEnvelopeBatch: @escaping @MainActor (
|
||||
[NostrEvent],
|
||||
@escaping @MainActor () -> Void
|
||||
) -> Bool,
|
||||
registerPendingPrivateEnvelope: @escaping @MainActor (String) -> Void,
|
||||
scheduleAfter: @escaping @Sendable (
|
||||
TimeInterval,
|
||||
@escaping @Sendable () -> Void
|
||||
) -> Void
|
||||
) {
|
||||
self.sendPrivateEnvelopeBatch = sendPrivateEnvelopeBatch
|
||||
self.registerPendingPrivateEnvelope = registerPendingPrivateEnvelope
|
||||
self.scheduleAfter = scheduleAfter
|
||||
}
|
||||
|
||||
func enqueue(key: String, events: [NostrEvent], registerPending: Bool) {
|
||||
guard pending[key] == nil else { return }
|
||||
if pending.count >= TransportConfig.nostrPrivateEnvelopeRetryQueueCap,
|
||||
let evictedKey = insertionOrder.first {
|
||||
insertionOrder.removeFirst()
|
||||
pending.removeValue(forKey: evictedKey)
|
||||
// These are control payloads, never user-authored messages. Keep
|
||||
// the bounded-loss decision explicit rather than silently growing
|
||||
// memory during a prolonged outage.
|
||||
SecureLogger.warning(
|
||||
"📮 Private control retry queue full — evicted oldest whole migration pair",
|
||||
category: .session
|
||||
)
|
||||
}
|
||||
pending[key] = PendingRetry(
|
||||
events: events,
|
||||
registerPending: registerPending,
|
||||
attempt: 0,
|
||||
isScheduled: false
|
||||
)
|
||||
insertionOrder.append(key)
|
||||
schedule(key: key)
|
||||
}
|
||||
|
||||
private func schedule(key: String) {
|
||||
guard var item = pending[key], !item.isScheduled else { return }
|
||||
item.isScheduled = true
|
||||
pending[key] = item
|
||||
let exponent = min(item.attempt, 5)
|
||||
let delay = min(2.0 * pow(2.0, Double(exponent)), 60.0)
|
||||
scheduleAfter(delay) { [self] in
|
||||
Task { @MainActor [self] in
|
||||
self.retry(key: key)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private func retry(key: String) {
|
||||
guard var item = pending[key] else { return }
|
||||
item.isScheduled = false
|
||||
pending[key] = item
|
||||
|
||||
let accepted = sendPrivateEnvelopeBatch(item.events) { [self] in
|
||||
self.enqueue(
|
||||
key: key,
|
||||
events: item.events,
|
||||
registerPending: item.registerPending
|
||||
)
|
||||
}
|
||||
if accepted {
|
||||
remove(key: key)
|
||||
if item.registerPending {
|
||||
for event in item.events {
|
||||
registerPendingPrivateEnvelope(event.id)
|
||||
}
|
||||
}
|
||||
} else {
|
||||
item.attempt += 1
|
||||
pending[key] = item
|
||||
schedule(key: key)
|
||||
}
|
||||
}
|
||||
|
||||
private func remove(key: String) {
|
||||
pending.removeValue(forKey: key)
|
||||
insertionOrder.removeAll { $0 == key }
|
||||
}
|
||||
|
||||
func removeAll() {
|
||||
pending.removeAll()
|
||||
insertionOrder.removeAll()
|
||||
}
|
||||
|
||||
var debugPendingCount: Int { pending.count }
|
||||
func debugContains(key: String) -> Bool { pending[key] != nil }
|
||||
}
|
||||
|
||||
@@ -196,13 +196,22 @@ enum TransportConfig {
|
||||
static let nostrGeoRelayCount: Int = 5
|
||||
static let nostrGeohashSampleLookbackSeconds: TimeInterval = 300
|
||||
static let nostrGeohashSampleLimit: Int = 100
|
||||
/// Public envelope timestamps are deliberately shifted into the past for
|
||||
/// privacy and relay compatibility. Mailbox queries must add the complete
|
||||
/// shift to the 24-hour delivery window or boundary messages disappear
|
||||
/// from `since` filters early.
|
||||
/// New iOS public-envelope timestamps are deliberately shifted into the
|
||||
/// past for privacy and relay compatibility.
|
||||
static let nostrPrivateEnvelopeTimestampFuzzSeconds: TimeInterval = 15 * 60
|
||||
static let nostrDMSubscribeLookbackSeconds: TimeInterval = (24 * 60 * 60)
|
||||
+ nostrPrivateEnvelopeTimestampFuzzSeconds
|
||||
/// Deployed Android clients can shift legacy kind-1059 timestamps by the
|
||||
/// full preceding 48 hours.
|
||||
static let nostrLegacyAndroidTimestampFuzzSeconds: TimeInterval = 48 * 60 * 60
|
||||
/// Private mail remains eligible for delivery for the same 24-hour window
|
||||
/// as the persistent sender outbox. The relay query must add timestamp
|
||||
/// randomization to that window rather than replacing it: an Android event
|
||||
/// sent at t0 can legitimately be stamped t0-48h and fetched at t0+24h.
|
||||
static let nostrPrivateEnvelopeDeliveryWindowSeconds: TimeInterval = 24 * 60 * 60
|
||||
static let nostrDMSubscribeClockSkewSeconds: TimeInterval = 15 * 60
|
||||
static let nostrDMSubscribeLookbackSeconds: TimeInterval =
|
||||
nostrPrivateEnvelopeDeliveryWindowSeconds
|
||||
+ nostrLegacyAndroidTimestampFuzzSeconds
|
||||
+ nostrDMSubscribeClockSkewSeconds
|
||||
// A sampled chat message this recent means "a conversation is happening
|
||||
// there" for the empty-timeline nearby-activity hint.
|
||||
static let uiGeohashChatActivityWindowSeconds: TimeInterval = 900
|
||||
@@ -230,11 +239,20 @@ enum TransportConfig {
|
||||
// Reconnect delays get ±20% random jitter so relays that dropped together
|
||||
// (e.g. a network blip) don't thundering-herd the same reconnect instant.
|
||||
static let nostrRelayBackoffJitterRatio: Double = 0.2
|
||||
static let nostrRelayDefaultFetchLimit: Int = 100
|
||||
/// Migration recovery uses one independent relay filter per wire kind so
|
||||
/// primary traffic cannot consume the legacy result budget (or vice
|
||||
/// versa). Five hundred per kind bounds startup work while preserving a
|
||||
/// materially deeper offline mailbox than the generic feed default.
|
||||
static let nostrPrivateEnvelopeFetchLimitPerKind: Int = 500
|
||||
// How many consecutive Tor-readiness waits (each bounded by TorManager's
|
||||
// bootstrap deadline) to attempt before unblocking pending EOSE callers.
|
||||
static let nostrTorReadyMaxWaitAttempts: Int = 3
|
||||
static let nostrPendingSendQueueCap: Int = 200
|
||||
/// Control-payload pairs (delivery/read acknowledgements and favorite
|
||||
/// notifications) that could not enter the relay queue. User messages do
|
||||
/// not use this queue: they fail visibly, while direct messages also
|
||||
/// remain in the router outbox.
|
||||
static let nostrPrivateEnvelopeRetryQueueCap: Int = 256
|
||||
// Sample interval for the send-queue overflow warning (first + every Nth
|
||||
// dropped event). Drops are ephemeral presence/geo traffic — log-only.
|
||||
static let nostrPendingSendDropLogInterval: Int = 10
|
||||
|
||||
@@ -1654,6 +1654,7 @@ final class ChatViewModel: ObservableObject, BitchatDelegate, SynchronousMessage
|
||||
// Drop relay subscriptions, handlers, pending sends, and replay state.
|
||||
// Geohash DM handlers can capture pre-wipe Nostr identities, so a plain
|
||||
// disconnect is not enough here.
|
||||
NostrTransport.resetControlRetriesForPanicWipe()
|
||||
NostrRelayManager.shared.resetForPanicWipe()
|
||||
nostrRelayManager = nil
|
||||
|
||||
|
||||
@@ -41,10 +41,11 @@ struct ChatViewModelServiceBundle {
|
||||
self.privateChatManager = privateChatManager
|
||||
self.unifiedPeerService = unifiedPeerService
|
||||
self.autocompleteService = AutocompleteService()
|
||||
// Persist processed private-envelope event IDs: BitChat randomizes their
|
||||
// timestamps, so the 24h-lookback DM subscriptions redeliver the same
|
||||
// events on every launch and only a cross-launch record stops the
|
||||
// reprocessing (re-sent DELIVERED bursts, phantom-ack noise).
|
||||
// Persist processed private-envelope event IDs: legacy Android can
|
||||
// randomize timestamps across the full 72h15m mailbox lookback, so DM
|
||||
// subscriptions redeliver the same events on every launch and only a
|
||||
// cross-launch record stops the reprocessing (re-sent DELIVERED
|
||||
// bursts, phantom-ack noise).
|
||||
self.deduplicationService = MessageDeduplicationService(nostrEventStore: NostrProcessedEventStore())
|
||||
self.publicMessagePipeline = PublicMessagePipeline()
|
||||
}
|
||||
@@ -176,7 +177,7 @@ private extension ChatViewModelBootstrapper {
|
||||
|
||||
func configureTransport() {
|
||||
viewModel.meshService.delegate = viewModel
|
||||
viewModel.meshService.eventDelegate = viewModel
|
||||
viewModel.messageRouter.setEventDelegate(viewModel)
|
||||
|
||||
DispatchQueue.main.asyncAfter(deadline: .now() + TransportConfig.uiStartupInitialDelaySeconds) { [weak viewModel] in
|
||||
guard let viewModel else { return }
|
||||
|
||||
@@ -162,11 +162,11 @@ final class GeohashSubscriptionManager {
|
||||
if let identity = try? context.deriveNostrIdentity(forGeohash: channel.geohash) {
|
||||
let dmSub = "geo-dm-\(channel.geohash)"
|
||||
context.setGeoDmSubscriptionID(dmSub)
|
||||
let dmFilter = NostrFilter.privateEnvelopesFor(
|
||||
let dmFilters = NostrFilter.privateEnvelopeFiltersFor(
|
||||
pubkey: identity.publicKeyHex,
|
||||
since: Date().addingTimeInterval(-TransportConfig.nostrDMSubscribeLookbackSeconds)
|
||||
)
|
||||
NostrRelayManager.shared.subscribe(filter: dmFilter, id: dmSub) { [weak self] envelope in
|
||||
NostrRelayManager.shared.subscribe(filters: dmFilters, id: dmSub) { [weak self] envelope in
|
||||
Task { @MainActor [weak self] in
|
||||
self?.inbound.subscribePrivateEnvelope(envelope, id: identity)
|
||||
}
|
||||
@@ -260,11 +260,11 @@ final class GeohashSubscriptionManager {
|
||||
if TorManager.shared.isReady {
|
||||
SecureLogger.debug("GeoDM: subscribing DMs pub=\(identity.publicKeyHex.prefix(8))… sub=\(dmSub)", category: .session)
|
||||
}
|
||||
let dmFilter = NostrFilter.privateEnvelopesFor(
|
||||
let dmFilters = NostrFilter.privateEnvelopeFiltersFor(
|
||||
pubkey: identity.publicKeyHex,
|
||||
since: Date().addingTimeInterval(-TransportConfig.nostrDMSubscribeLookbackSeconds)
|
||||
)
|
||||
NostrRelayManager.shared.subscribe(filter: dmFilter, id: dmSub) { [weak self] envelope in
|
||||
NostrRelayManager.shared.subscribe(filters: dmFilters, id: dmSub) { [weak self] envelope in
|
||||
Task { @MainActor [weak self] in
|
||||
self?.inbound.handlePrivateEnvelope(envelope, id: identity)
|
||||
}
|
||||
@@ -388,12 +388,12 @@ final class GeohashSubscriptionManager {
|
||||
category: .session
|
||||
)
|
||||
|
||||
let filter = NostrFilter.privateEnvelopesFor(
|
||||
let filters = NostrFilter.privateEnvelopeFiltersFor(
|
||||
pubkey: currentIdentity.publicKeyHex,
|
||||
since: Date().addingTimeInterval(-TransportConfig.nostrDMSubscribeLookbackSeconds)
|
||||
)
|
||||
|
||||
context.nostrRelayManager?.subscribe(filter: filter, id: "chat-messages") { [weak self] event in
|
||||
context.nostrRelayManager?.subscribe(filters: filters, id: "chat-messages") { [weak self] event in
|
||||
Task { @MainActor [weak self] in
|
||||
self?.inbound.handleAccountPrivateEnvelope(event)
|
||||
}
|
||||
|
||||
@@ -91,8 +91,8 @@ final class NostrInboundPipeline {
|
||||
private weak var context: (any NostrInboundPipelineContext)?
|
||||
private let presence: GeoPresenceTracker
|
||||
private var geoEventLogCount = 0
|
||||
// During the bounded wire-format migration, one logical private payload
|
||||
// is published under both the primary and compatibility formats. Outer
|
||||
// During the coordinated wire-format migration, one logical private
|
||||
// payload is published under both primary and compatibility formats. Outer
|
||||
// event IDs differ, so collapse the authenticated embedded payload before
|
||||
// invoking message/ack side effects. Keep this bounded like the outer-ID
|
||||
// caches; the recipient and authenticated sender are part of the key.
|
||||
|
||||
Reference in New Issue
Block a user