mirror of
https://github.com/permissionlesstech/bitchat.git
synced 2026-07-25 02:05:19 +00:00
Mesh and geohash timelines are now store conversations. All public mutation sites flow through store intents; PublicMessagePipeline keeps its 80ms UI batching but commits batches via store appends with each buffered entry carrying its destination conversation (a mid-batch channel switch now flushes instead of dropping the buffer). ChatViewModel.messages becomes a cached get-only view of the active conversation, invalidated through the change subject. The mesh late-insert threshold is consciously removed: it only ever ordered the non-rendered messages copy, so strict timestamp insertion makes the working set agree with rendered order. PublicTimelineStore and the per-message full-array legacy sync are deleted; the coalescing bridge mirrors public conversations for the remaining legacy readers. pipeline.publicIngest: 6.6k -> 9.5k msg/s (+45%); private steady; store.append 237k/s. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
198 lines
7.9 KiB
Swift
198 lines
7.9 KiB
Swift
import BitFoundation
|
|
import BitLogger
|
|
import Foundation
|
|
import SwiftUI
|
|
|
|
/// The narrow surface `GeoPresenceTracker` needs from its owner.
|
|
///
|
|
/// Split out of `ChatNostrContext`: member names are shared with the sibling
|
|
/// component contexts so `ChatViewModel` provides a single witness for each.
|
|
@MainActor
|
|
protocol GeoPresenceContext: AnyObject {
|
|
var activeChannel: ChannelID { get }
|
|
/// Per-geohash notification cooldown: geohash -> last notify time.
|
|
var lastGeoNotificationAt: [String: Date] { get set }
|
|
var geoNicknames: [String: String] { get }
|
|
var teleportedGeoCount: Int { get }
|
|
|
|
func deriveNostrIdentity(forGeohash geohash: String) throws -> NostrIdentity
|
|
func isNostrBlocked(pubkeyHexLowercased: String) -> Bool
|
|
func parseMentions(from content: String) -> [String]
|
|
|
|
func recordGeoParticipant(pubkeyHex: String, geohash: String)
|
|
func geoParticipantCount(for geohash: String) -> Int
|
|
func markGeoTeleported(_ pubkeyHexLowercased: String)
|
|
|
|
/// Appends a geohash message if absent (single-writer store intent).
|
|
/// Returns `true` when stored.
|
|
@discardableResult
|
|
func appendGeohashMessageIfAbsent(_ message: BitchatMessage, toGeohash geohash: String) -> Bool
|
|
|
|
/// Posts the sampled-geohash-activity local notification.
|
|
func notifyGeohashActivity(geohash: String, bodyPreview: String)
|
|
}
|
|
|
|
extension ChatViewModel: GeoPresenceContext {
|
|
// `activeChannel`, `lastGeoNotificationAt`, `geoNicknames`, the Nostr
|
|
// identity/blocking members, and the
|
|
// `appendGeohashMessageIfAbsent(_:toGeohash:)` store intent already have
|
|
// witnesses on `ChatViewModel`. The members below flatten nested service
|
|
// accesses into intent-named calls.
|
|
|
|
var teleportedGeoCount: Int {
|
|
locationPresenceStore.teleportedGeo.count
|
|
}
|
|
|
|
func recordGeoParticipant(pubkeyHex: String, geohash: String) {
|
|
participantTracker.recordParticipant(pubkeyHex: pubkeyHex, geohash: geohash)
|
|
}
|
|
|
|
func geoParticipantCount(for geohash: String) -> Int {
|
|
participantTracker.participantCount(for: geohash)
|
|
}
|
|
|
|
func markGeoTeleported(_ pubkeyHexLowercased: String) {
|
|
locationPresenceStore.markTeleported(pubkeyHexLowercased)
|
|
}
|
|
|
|
func notifyGeohashActivity(geohash: String, bodyPreview: String) {
|
|
NotificationService.shared.sendGeohashActivityNotification(geohash: geohash, bodyPreview: bodyPreview)
|
|
}
|
|
}
|
|
|
|
/// Geohash presence bookkeeping that is independent of relay subscriptions:
|
|
/// teleport-tag detection and marking, the sampling-event LRU dedup, and the
|
|
/// per-geohash notification cooldown for sampled activity.
|
|
final class GeoPresenceTracker {
|
|
private weak var context: (any GeoPresenceContext)?
|
|
private var recentGeoSamplingEventIDs = Set<String>()
|
|
private var recentGeoSamplingEventIDOrder: [String] = []
|
|
|
|
init(context: any GeoPresenceContext) {
|
|
self.context = context
|
|
}
|
|
|
|
/// True when the event carries a `["t", "teleport"]` tag.
|
|
static func hasTeleportTag(_ event: NostrEvent) -> Bool {
|
|
event.tags.contains { tag in
|
|
tag.count >= 2 && tag[0].lowercased() == "t" && tag[1].lowercased() == "teleport"
|
|
}
|
|
}
|
|
|
|
/// Marks a peer teleported on a follow-up main-actor hop (keeps the
|
|
/// inbound hot path free of presence-store writes).
|
|
@MainActor
|
|
func scheduleMarkPeerTeleported(_ key: String, logged: Bool) {
|
|
Task { @MainActor [weak context] in
|
|
guard let context else { return }
|
|
context.markGeoTeleported(key)
|
|
if logged {
|
|
SecureLogger.info(
|
|
"GeoTeleport: mark peer teleported key=\(key.prefix(8))… total=\(context.teleportedGeoCount)",
|
|
category: .session
|
|
)
|
|
}
|
|
}
|
|
}
|
|
|
|
@MainActor
|
|
func subscribeNostrEvent(_ event: NostrEvent, gh: String) {
|
|
guard let context else { return }
|
|
guard (event.kind == NostrProtocol.EventKind.ephemeralEvent.rawValue
|
|
|| event.kind == NostrProtocol.EventKind.geohashPresence.rawValue)
|
|
else {
|
|
return
|
|
}
|
|
guard event.isValidSignature() else { return }
|
|
guard shouldProcessGeoSamplingEvent(event.id) else { return }
|
|
|
|
let existingCount = context.geoParticipantCount(for: gh)
|
|
context.recordGeoParticipant(pubkeyHex: event.pubkey, geohash: gh)
|
|
|
|
guard let content = event.content.trimmedOrNilIfEmpty else { return }
|
|
if context.isNostrBlocked(pubkeyHexLowercased: event.pubkey.lowercased()) { return }
|
|
if let my = try? context.deriveNostrIdentity(forGeohash: gh),
|
|
my.publicKeyHex.lowercased() == event.pubkey.lowercased() {
|
|
return
|
|
}
|
|
guard existingCount == 0 else { return }
|
|
|
|
let eventTime = Date(timeIntervalSince1970: TimeInterval(event.created_at))
|
|
if Date().timeIntervalSince(eventTime) > 30 { return }
|
|
|
|
#if os(iOS)
|
|
guard UIApplication.shared.applicationState == .active else { return }
|
|
if case .location(let channel) = context.activeChannel, channel.geohash == gh { return }
|
|
#elseif os(macOS)
|
|
guard NSApplication.shared.isActive else { return }
|
|
if case .location(let channel) = context.activeChannel, channel.geohash == gh { return }
|
|
#endif
|
|
|
|
cooldownPerGeohash(gh, content: content, event: event)
|
|
}
|
|
|
|
@MainActor
|
|
func cooldownPerGeohash(_ gh: String, content: String, event: NostrEvent) {
|
|
guard let context else { return }
|
|
let now = Date()
|
|
let last = context.lastGeoNotificationAt[gh] ?? .distantPast
|
|
if now.timeIntervalSince(last) < TransportConfig.uiGeoNotifyCooldownSeconds { return }
|
|
|
|
let preview: String = {
|
|
let maxLen = TransportConfig.uiGeoNotifySnippetMaxLen
|
|
if content.count <= maxLen { return content }
|
|
let idx = content.index(content.startIndex, offsetBy: maxLen)
|
|
return String(content[..<idx]) + "…"
|
|
}()
|
|
|
|
Task { @MainActor [weak context] in
|
|
guard let context else { return }
|
|
context.lastGeoNotificationAt[gh] = now
|
|
let senderSuffix = String(event.pubkey.suffix(4))
|
|
let nick = context.geoNicknames[event.pubkey.lowercased()]
|
|
let senderName = (nick?.isEmpty == false ? nick! : "anon") + "#" + senderSuffix
|
|
|
|
let rawTs = Date(timeIntervalSince1970: TimeInterval(event.created_at))
|
|
let ts = min(rawTs, Date())
|
|
let mentions = context.parseMentions(from: content)
|
|
let message = BitchatMessage(
|
|
id: event.id,
|
|
sender: senderName,
|
|
content: content,
|
|
timestamp: ts,
|
|
isRelay: false,
|
|
senderPeerID: PeerID(nostr: event.pubkey),
|
|
mentions: mentions.isEmpty ? nil : mentions
|
|
)
|
|
if context.appendGeohashMessageIfAbsent(message, toGeohash: gh) {
|
|
context.notifyGeohashActivity(geohash: gh, bodyPreview: preview)
|
|
}
|
|
}
|
|
}
|
|
|
|
/// First-seen check for sampled geohash events with LRU eviction so the
|
|
/// dedup set stays bounded across long sampling sessions.
|
|
func shouldProcessGeoSamplingEvent(_ eventID: String) -> Bool {
|
|
guard !eventID.isEmpty else { return true }
|
|
guard recentGeoSamplingEventIDs.insert(eventID).inserted else {
|
|
return false
|
|
}
|
|
recentGeoSamplingEventIDOrder.append(eventID)
|
|
|
|
let cap = TransportConfig.geoSamplingEventLRUCap
|
|
if recentGeoSamplingEventIDOrder.count > cap {
|
|
let removeCount = recentGeoSamplingEventIDOrder.count - cap
|
|
for staleID in recentGeoSamplingEventIDOrder.prefix(removeCount) {
|
|
recentGeoSamplingEventIDs.remove(staleID)
|
|
}
|
|
recentGeoSamplingEventIDOrder.removeFirst(removeCount)
|
|
}
|
|
return true
|
|
}
|
|
|
|
func clearGeoSamplingEventDedup() {
|
|
recentGeoSamplingEventIDs.removeAll()
|
|
recentGeoSamplingEventIDOrder.removeAll()
|
|
}
|
|
}
|