diff --git a/app/src/main/java/com/bitchat/android/BitchatApplication.kt b/app/src/main/java/com/bitchat/android/BitchatApplication.kt index 282f3295..ecd32327 100644 --- a/app/src/main/java/com/bitchat/android/BitchatApplication.kt +++ b/app/src/main/java/com/bitchat/android/BitchatApplication.kt @@ -13,6 +13,9 @@ class BitchatApplication : Application() { override fun onCreate() { super.onCreate() + // Start the single process-wide power policy before transport components are constructed. + com.bitchat.android.mesh.PowerManager.getInstance(this).start() + // Initialize Tor first so any early network goes over Tor try { val torProvider = ArtiTorManager.getInstance() @@ -53,6 +56,10 @@ class BitchatApplication : Application() { com.bitchat.android.nostr.GeohashConversationRegistry.initialize(this) } catch (_: Exception) { } + // Own relay connectivity, selected-channel subscriptions, and presence scheduling at the + // process level so closing the Activity does not disconnect Nostr. + try { com.bitchat.android.nostr.NostrBackgroundRuntime.initialize(this) } catch (_: Exception) { } + // Initialize mesh service preferences try { com.bitchat.android.service.MeshServicePreferences.init(this) } catch (_: Exception) { } diff --git a/app/src/main/java/com/bitchat/android/mesh/BluetoothConnectionManager.kt b/app/src/main/java/com/bitchat/android/mesh/BluetoothConnectionManager.kt index faf67378..f2bb54c6 100644 --- a/app/src/main/java/com/bitchat/android/mesh/BluetoothConnectionManager.kt +++ b/app/src/main/java/com/bitchat/android/mesh/BluetoothConnectionManager.kt @@ -18,7 +18,7 @@ class BluetoothConnectionManager( private val context: Context, private val myPeerID: String, private val fragmentManager: FragmentManager? = null -) : PowerManagerDelegate { +) { companion object { private const val TAG = "BluetoothConnectionManager" @@ -30,7 +30,7 @@ class BluetoothConnectionManager( private val bluetoothAdapter: BluetoothAdapter? = bluetoothManager.adapter // Power management - private val powerManager = PowerManager(context.applicationContext) + private val powerManager = PowerManager.getInstance(context.applicationContext) // Coroutines private val connectionScope = CoroutineScope(Dispatchers.IO + SupervisorJob()) @@ -117,7 +117,19 @@ class BluetoothConnectionManager( } init { - powerManager.delegate = this + connectionScope.launch { + var previousMode: PowerManager.PowerMode? = null + powerManager.profile.collect { profile -> + val modeChanged = previousMode != null && previousMode != profile.mode + previousMode = profile.mode + if (!isActive || !isBleTransportEnabled()) return@collect + + if (modeChanged && isGattServerEnabled()) { + serverManager.restartAdvertising() + } + clientManager.applyPowerProfile(profile) + } + } // Observe debug settings to enforce role state while active try { val dbg = com.bitchat.android.ui.debug.DebugSettingsManager.getInstance() @@ -294,9 +306,6 @@ class BluetoothConnectionManager( clientManager.stop() serverManager.stop() - // Stop power manager - powerManager.stop() - // Stop connection tracker connectionTracker.stop() @@ -461,54 +470,6 @@ class BluetoothConnectionManager( } } - // MARK: - PowerManagerDelegate Implementation - - override fun onPowerModeChanged(newMode: PowerManager.PowerMode) { - Log.i(TAG, "Power mode changed to: $newMode") - - connectionScope.launch { - if (!isActive || !isBleTransportEnabled()) { - serverManager.stop() - clientManager.stop() - return@launch - } - - // Avoid rapid scan restarts by checking if we need to change scan behavior - val wasUsingDutyCycle = powerManager.shouldUseDutyCycle() - - // Update advertising with new power settings if server enabled - val serverEnabled = isGattServerEnabled() - if (serverEnabled) { - serverManager.restartAdvertising() - } else { - serverManager.stop() - } - - // Only restart scanning if the duty cycle behavior changed - val nowUsingDutyCycle = powerManager.shouldUseDutyCycle() - if (wasUsingDutyCycle != nowUsingDutyCycle) { - val clientEnabled = isGattClientEnabled() - if (clientEnabled) { - clientManager.restartScanning() - } else { - clientManager.stop() - } - } - - // Enforce connection limits - enforceStrictLimits() - } - } - - override fun onScanStateChanged(shouldScan: Boolean) { - if (!isActive || !isBleTransportEnabled()) { - clientManager.onScanStateChanged(false) - return - } - clientManager.onScanStateChanged(shouldScan) - } - - // MARK: - Private Implementation - All moved to component managers } /** diff --git a/app/src/main/java/com/bitchat/android/mesh/BluetoothGattClientManager.kt b/app/src/main/java/com/bitchat/android/mesh/BluetoothGattClientManager.kt index cc13d5df..498339ce 100644 --- a/app/src/main/java/com/bitchat/android/mesh/BluetoothGattClientManager.kt +++ b/app/src/main/java/com/bitchat/android/mesh/BluetoothGattClientManager.kt @@ -93,6 +93,7 @@ class BluetoothGattClientManager( @Volatile private var lastScanResultTime = 0L private var scanRetryCount = 0 private var scanWatchdogJob: Job? = null + private var scanDutyCycleJob: Job? = null // RSSI monitoring state private var rssiMonitoringJob: Job? = null @@ -131,18 +132,9 @@ class BluetoothGattClientManager( isActive = true connectionScope.launch { - if (powerManager.shouldUseDutyCycle()) { - Log.i(TAG, "Using power-aware duty cycling") - // Duty cycle drives onScanStateChanged(true/false); scanningDesired follows that. - } else { - scanningDesired = true - startScanning() - } - + applyPowerProfile(powerManager.profile.value) // Start RSSI monitoring startRSSIMonitoring() - // Start the scan watchdog so a silently-dead or wedged scanner self-heals. - startScanWatchdog() } return true @@ -153,6 +145,8 @@ class BluetoothGattClientManager( */ fun stop() { scanningDesired = false + scanDutyCycleJob?.cancel() + scanDutyCycleJob = null stopScanWatchdog() if (!isActive) { // Idempotent stop @@ -208,10 +202,10 @@ class BluetoothGattClientManager( Log.d(TAG, "Failed to request RSSI from ${deviceConn.device.address}: ${e.message}") } } - delay(AppConstants.Mesh.RSSI_UPDATE_INTERVAL_MS) + delay(powerManager.profile.value.ble.rssiPollIntervalMs) } catch (e: Exception) { Log.w(TAG, "Error in RSSI monitoring: ${e.message}") - delay(AppConstants.Mesh.RSSI_UPDATE_INTERVAL_MS) + delay(powerManager.profile.value.ble.rssiPollIntervalMs) } } } @@ -661,13 +655,37 @@ class BluetoothGattClientManager( connectionScope.launch { stopScanning() delay(1000) // Extra delay to avoid rate limiting - - if (powerManager.shouldUseDutyCycle()) { - Log.i(TAG, "Switching to duty cycle scanning mode") - // Duty cycle will handle scanning - } else { - Log.i(TAG, "Switching to continuous scanning mode") - startScanning() + applyPowerProfile(powerManager.profile.value) + } + } + + /** + * Apply the current process-wide profile without ever disabling background discovery. + */ + fun applyPowerProfile(profile: PowerManager.RuntimePerformanceProfile) { + scanDutyCycleJob?.cancel() + scanDutyCycleJob = null + if (!isActive || !isClientRoleEnabled()) { + onScanStateChanged(false) + return + } + + if (profile.ble.continuousScan) { + startScanWatchdog() + onScanStateChanged(true) + return + } + + // Duty-cycled scans are re-armed every window, so the continuous-scan watchdog would only + // create background wakeups during intentional OFF periods. + stopScanWatchdog() + scanDutyCycleJob = connectionScope.launch { + while (isActive && isClientRoleEnabled()) { + onScanStateChanged(true) + delay(profile.ble.scanOnMs) + if (!isActive || !isClientRoleEnabled()) break + onScanStateChanged(false) + delay(profile.ble.scanOffMs) } } } diff --git a/app/src/main/java/com/bitchat/android/mesh/BluetoothMeshService.kt b/app/src/main/java/com/bitchat/android/mesh/BluetoothMeshService.kt index 61924738..f828af9e 100644 --- a/app/src/main/java/com/bitchat/android/mesh/BluetoothMeshService.kt +++ b/app/src/main/java/com/bitchat/android/mesh/BluetoothMeshService.kt @@ -126,7 +126,6 @@ class BluetoothMeshService(private val context: Context) : TransportBridgeServic var delegate: BluetoothMeshDelegate? = null // Coroutines - private var announceJob: Job? = null // Tracks whether this instance has been terminated via stopServices() private var terminated = false @@ -240,23 +239,6 @@ class BluetoothMeshService(private val context: Context) : TransportBridgeServic } } - /** - * Send broadcast announcement every 30 seconds - */ - private fun sendPeriodicBroadcastAnnounce() { - announceJob?.cancel() - announceJob = serviceScope.launch { - while (isActive) { - try { - delay(30000) // 30 seconds - sendBroadcastAnnounce() - } catch (e: Exception) { - Log.e(TAG, "Error in periodic broadcast announce: ${e.message}") - } - } - } - } - /** * Setup delegate connections between components */ @@ -794,8 +776,6 @@ class BluetoothMeshService(private val context: Context) : TransportBridgeServic isActive = true TransportBridgeService.register("BLE", this) - // Start periodic announcements for peer discovery and connectivity - sendPeriodicBroadcastAnnounce() // Start periodic syncs com.bitchat.android.service.MeshServiceHolder.startSharedGossip("BLE") } else { @@ -817,8 +797,6 @@ class BluetoothMeshService(private val context: Context) : TransportBridgeServic private fun pauseServicesForTransportDisable() { Log.i(TAG, "Disabling BLE mesh transport") isActive = false - announceJob?.cancel() - announceJob = null com.bitchat.android.service.MeshServiceHolder.stopSharedGossip("BLE") TransportBridgeService.unregister("BLE") try { com.bitchat.android.services.AppStateStore.clearTransportPeers("BLE") } catch (_: Exception) { } @@ -838,8 +816,6 @@ class BluetoothMeshService(private val context: Context) : TransportBridgeServic Log.i(TAG, "Stopping Bluetooth mesh service") isActive = false - announceJob?.cancel() - announceJob = null TransportBridgeService.unregister("BLE") try { com.bitchat.android.services.AppStateStore.clearTransportPeers("BLE") } catch (_: Exception) { } try { com.bitchat.android.services.AppStateStore.clearTransportDirectPeers("BLE") } catch (_: Exception) { } diff --git a/app/src/main/java/com/bitchat/android/mesh/MeshCore.kt b/app/src/main/java/com/bitchat/android/mesh/MeshCore.kt index bd7e08aa..ae1f535a 100644 --- a/app/src/main/java/com/bitchat/android/mesh/MeshCore.kt +++ b/app/src/main/java/com/bitchat/android/mesh/MeshCore.kt @@ -20,7 +20,6 @@ import com.bitchat.android.service.TransportBridgeService import com.bitchat.android.sync.GossipSyncManager import com.bitchat.android.util.toHexString import kotlinx.coroutines.CoroutineScope -import kotlinx.coroutines.Job import kotlinx.coroutines.delay import kotlinx.coroutines.launch import kotlinx.coroutines.runBlocking @@ -116,7 +115,6 @@ class MeshCore( var delegate: MeshDelegate? = null - private var announceJob: Job? = null private var isActive = false init { @@ -145,7 +143,6 @@ class MeshCore( fun startCore() { if (isActive) return isActive = true - startPeriodicBroadcastAnnounce() if (ownsGossipManager) { gossipSyncManager.start() } @@ -154,8 +151,6 @@ class MeshCore( fun stopCore() { if (!isActive) return isActive = false - announceJob?.cancel() - announceJob = null directPeers.clear() if (ownsGossipManager) { gossipSyncManager.stop() @@ -208,18 +203,6 @@ class MeshCore( return acceptedByLocalTransport || acceptedByBridgedTransport } - private fun startPeriodicBroadcastAnnounce() { - announceJob?.cancel() - announceJob = scope.launch { - while (isActive) { - try { - delay(30_000) - sendBroadcastAnnounce() - } catch (_: Exception) { } - } - } - } - private fun setupDelegates() { peerManager.delegate = object : PeerManagerDelegate { override fun onPeerListUpdated(peerIDs: List) { diff --git a/app/src/main/java/com/bitchat/android/mesh/PowerManager.kt b/app/src/main/java/com/bitchat/android/mesh/PowerManager.kt index 3096828a..9714c63a 100644 --- a/app/src/main/java/com/bitchat/android/mesh/PowerManager.kt +++ b/app/src/main/java/com/bitchat/android/mesh/PowerManager.kt @@ -14,327 +14,261 @@ import androidx.lifecycle.Lifecycle import androidx.lifecycle.LifecycleEventObserver import androidx.lifecycle.LifecycleOwner import androidx.lifecycle.ProcessLifecycleOwner -import kotlinx.coroutines.* -import kotlin.math.max +import com.bitchat.android.services.AppStateStore +import com.bitchat.android.util.AppConstants +import kotlinx.coroutines.CoroutineScope +import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.SupervisorJob +import kotlinx.coroutines.cancel +import kotlinx.coroutines.flow.MutableStateFlow +import kotlinx.coroutines.flow.StateFlow +import kotlinx.coroutines.flow.asStateFlow +import kotlinx.coroutines.launch /** - * Power-aware Bluetooth management for bitchat - * Adjusts scanning, advertising, and connection behavior based on battery state + * Process-wide source of truth for background and battery-aware scheduling. + * + * The manager deliberately exposes policy only. Consumers retain ownership of their hardware and + * network jobs, so changing profile never implies disconnecting a transport or suspending Tor. */ -class PowerManager(private val context: Context) : LifecycleEventObserver { - +class PowerManager private constructor(context: Context) : LifecycleEventObserver { + + enum class PowerMode { + PERFORMANCE, + BALANCED, + POWER_SAVER, + ULTRA_LOW_POWER + } + + enum class BatteryBand { NORMAL, LOW, CRITICAL } + + data class BleSchedule( + val scanOnMs: Long, + val scanOffMs: Long, + val continuousScan: Boolean, + val rssiPollIntervalMs: Long, + val maxConnections: Int = 8 + ) + + data class WifiAwareSchedule( + val tcpKeepAliveMs: Long, + val discoveryKeepAliveMs: Long, + val connectionMaintenanceMs: Long, + val discoveryIdleRefreshMs: Long, + val discoverySessionRefreshMinMs: Long, + val discoveryStaleMs: Long + ) + + data class NostrSchedule( + val presenceHeartbeatMinMs: Long, + val presenceHeartbeatMaxMs: Long, + val subscriptionValidationMs: Long + ) + + data class RuntimePerformanceProfile( + val mode: PowerMode, + val batteryBand: BatteryBand, + val isBackground: Boolean, + val isCharging: Boolean, + val hasDirectPeers: Boolean, + val ble: BleSchedule, + val meshAnnouncementIntervalMs: Long, + val wifiAware: WifiAwareSchedule, + val nostr: NostrSchedule + ) + companion object { private const val TAG = "PowerManager" - - // Battery thresholds - private const val CRITICAL_BATTERY = com.bitchat.android.util.AppConstants.Power.CRITICAL_BATTERY_PERCENT - private const val LOW_BATTERY = com.bitchat.android.util.AppConstants.Power.LOW_BATTERY_PERCENT - private const val MEDIUM_BATTERY = com.bitchat.android.util.AppConstants.Power.MEDIUM_BATTERY_PERCENT - - // Scan duty cycle periods (ms) - private const val SCAN_ON_DURATION_NORMAL = com.bitchat.android.util.AppConstants.Power.SCAN_ON_DURATION_NORMAL_MS // 8 seconds on - private const val SCAN_OFF_DURATION_NORMAL = com.bitchat.android.util.AppConstants.Power.SCAN_OFF_DURATION_NORMAL_MS // 2 seconds off - private const val SCAN_ON_DURATION_POWER_SAVE = com.bitchat.android.util.AppConstants.Power.SCAN_ON_DURATION_POWER_SAVE_MS // 2 seconds on - private const val SCAN_OFF_DURATION_POWER_SAVE = com.bitchat.android.util.AppConstants.Power.SCAN_OFF_DURATION_POWER_SAVE_MS // 8 seconds off - private const val SCAN_ON_DURATION_ULTRA_LOW = com.bitchat.android.util.AppConstants.Power.SCAN_ON_DURATION_ULTRA_LOW_MS // 1 second on - private const val SCAN_OFF_DURATION_ULTRA_LOW = com.bitchat.android.util.AppConstants.Power.SCAN_OFF_DURATION_ULTRA_LOW_MS // 10 seconds off - - // Connection limits - private const val MAX_CONNECTIONS_NORMAL = com.bitchat.android.util.AppConstants.Power.MAX_CONNECTIONS_NORMAL - private const val MAX_CONNECTIONS_POWER_SAVE = com.bitchat.android.util.AppConstants.Power.MAX_CONNECTIONS_POWER_SAVE - private const val MAX_CONNECTIONS_ULTRA_LOW = com.bitchat.android.util.AppConstants.Power.MAX_CONNECTIONS_ULTRA_LOW + + @Volatile + private var INSTANCE: PowerManager? = null + + fun getInstance(context: Context): PowerManager = + INSTANCE ?: synchronized(this) { + INSTANCE ?: PowerManager(context.applicationContext).also { INSTANCE = it } + } } - - enum class PowerMode { - PERFORMANCE, // Full power, no restrictions - BALANCED, // Moderate power saving - POWER_SAVER, // Aggressive power saving - ULTRA_LOW_POWER // Minimal operations only - } - - private var currentMode = PowerMode.BALANCED - private var isCharging = false - private var batteryLevel = 100 - private var isAppInBackground = true - - private val powerScope = CoroutineScope(Dispatchers.IO + SupervisorJob()) - private var dutyCycleJob: Job? = null - - var delegate: PowerManagerDelegate? = null - - // Battery monitoring + + private val appContext = context.applicationContext + private val scope = CoroutineScope(Dispatchers.Default + SupervisorJob()) + + @Volatile private var batteryLevel = 100 + @Volatile private var isCharging = false + @Volatile private var isAppInBackground = true + @Volatile private var hasDirectPeers = false + @Volatile private var shutdown = false + + private val _profile = MutableStateFlow( + PowerProfileResolver.resolve( + batteryLevel = batteryLevel, + isCharging = isCharging, + isBackground = isAppInBackground, + hasDirectPeers = hasDirectPeers + ) + ) + val profile: StateFlow = _profile.asStateFlow() + private val batteryReceiver = object : BroadcastReceiver() { override fun onReceive(context: Context?, intent: Intent?) { when (intent?.action) { Intent.ACTION_BATTERY_CHANGED -> { val level = intent.getIntExtra(BatteryManager.EXTRA_LEVEL, -1) val scale = intent.getIntExtra(BatteryManager.EXTRA_SCALE, -1) - if (level != -1 && scale != -1) { - batteryLevel = (level * 100) / scale - } - + if (level >= 0 && scale > 0) batteryLevel = (level * 100) / scale + val status = intent.getIntExtra(BatteryManager.EXTRA_STATUS, -1) isCharging = status == BatteryManager.BATTERY_STATUS_CHARGING || - status == BatteryManager.BATTERY_STATUS_FULL - - updatePowerMode() + status == BatteryManager.BATTERY_STATUS_FULL + refreshProfile() } Intent.ACTION_POWER_CONNECTED -> { isCharging = true - updatePowerMode() + refreshProfile() } Intent.ACTION_POWER_DISCONNECTED -> { isCharging = false - updatePowerMode() + refreshProfile() } } } } - + init { registerBatteryReceiver() - - // Register for process lifecycle events on the main thread Handler(Looper.getMainLooper()).post { try { ProcessLifecycleOwner.get().lifecycle.addObserver(this) + isAppInBackground = !ProcessLifecycleOwner.get().lifecycle.currentState + .isAtLeast(Lifecycle.State.STARTED) + refreshProfile() } catch (e: Exception) { - Log.e(TAG, "Failed to register lifecycle observer: ${e.message}") + Log.e(TAG, "Failed to register process lifecycle observer: ${e.message}") } } - - updatePowerMode() + scope.launch { + AppStateStore.directPeers + .collect { peers -> + hasDirectPeers = peers.isNotEmpty() + refreshProfile() + } + } } - - fun start() { - Log.i(TAG, "Starting power management") - startDutyCycle() - } - - fun stop() { - Log.i(TAG, "Stopping power management") - powerScope.cancel() - unregisterBatteryReceiver() - - // Unregister lifecycle observer + + /** Kept for existing callers; initialization is eager and idempotent. */ + fun start() = Unit + + /** + * Process-level teardown for explicit full application shutdown only. + * Transport stop/restart paths must not call this. + */ + fun shutdown() { + if (shutdown) return + shutdown = true + scope.cancel() + try { appContext.unregisterReceiver(batteryReceiver) } catch (_: Exception) { } Handler(Looper.getMainLooper()).post { - try { - ProcessLifecycleOwner.get().lifecycle.removeObserver(this) - } catch (e: Exception) { - Log.e(TAG, "Failed to remove lifecycle observer: ${e.message}") - } + try { ProcessLifecycleOwner.get().lifecycle.removeObserver(this) } catch (_: Exception) { } } } - + override fun onStateChanged(source: LifecycleOwner, event: Lifecycle.Event) { when (event) { Lifecycle.Event.ON_START -> { - Log.d(TAG, "Process lifecycle: ON_START (App coming to foreground)") isAppInBackground = false - updatePowerMode() + refreshProfile() } Lifecycle.Event.ON_STOP -> { - Log.d(TAG, "Process lifecycle: ON_STOP (App going to background)") isAppInBackground = true - updatePowerMode() + refreshProfile() } - else -> {} + else -> Unit } } - - /** - * Get scan settings optimized for current power mode - */ + fun getScanSettings(): ScanSettings { + val current = profile.value val builder = ScanSettings.Builder() .setCallbackType(ScanSettings.CALLBACK_TYPE_ALL_MATCHES) - when (currentMode) { + when (current.mode) { PowerMode.PERFORMANCE -> builder .setScanMode(ScanSettings.SCAN_MODE_LOW_LATENCY) .setMatchMode(ScanSettings.MATCH_MODE_AGGRESSIVE) .setNumOfMatches(ScanSettings.MATCH_NUM_MAX_ADVERTISEMENT) - PowerMode.BALANCED -> builder .setScanMode(ScanSettings.SCAN_MODE_BALANCED) .setMatchMode(ScanSettings.MATCH_MODE_AGGRESSIVE) .setNumOfMatches(ScanSettings.MATCH_NUM_ONE_ADVERTISEMENT) - - PowerMode.POWER_SAVER -> builder - .setScanMode(ScanSettings.SCAN_MODE_LOW_POWER) - .setMatchMode(ScanSettings.MATCH_MODE_STICKY) - .setNumOfMatches(ScanSettings.MATCH_NUM_ONE_ADVERTISEMENT) - + PowerMode.POWER_SAVER, PowerMode.ULTRA_LOW_POWER -> builder .setScanMode(ScanSettings.SCAN_MODE_LOW_POWER) .setMatchMode(ScanSettings.MATCH_MODE_STICKY) .setNumOfMatches(ScanSettings.MATCH_NUM_ONE_ADVERTISEMENT) } - return builder.setReportDelay(0).build() } - - /** - * Get advertising settings optimized for current power mode - */ - fun getAdvertiseSettings(): AdvertiseSettings { - return when (currentMode) { - PowerMode.PERFORMANCE -> AdvertiseSettings.Builder() - .setAdvertiseMode(AdvertiseSettings.ADVERTISE_MODE_LOW_LATENCY) - .setTxPowerLevel(AdvertiseSettings.ADVERTISE_TX_POWER_HIGH) - .setConnectable(true) - .setTimeout(0) - .build() - - PowerMode.BALANCED -> AdvertiseSettings.Builder() - .setAdvertiseMode(AdvertiseSettings.ADVERTISE_MODE_BALANCED) - .setTxPowerLevel(AdvertiseSettings.ADVERTISE_TX_POWER_MEDIUM) - .setConnectable(true) - .setTimeout(0) - .build() - - PowerMode.POWER_SAVER -> AdvertiseSettings.Builder() - .setAdvertiseMode(AdvertiseSettings.ADVERTISE_MODE_LOW_POWER) - .setTxPowerLevel(AdvertiseSettings.ADVERTISE_TX_POWER_LOW) - .setConnectable(true) - .setTimeout(0) - .build() - - PowerMode.ULTRA_LOW_POWER -> AdvertiseSettings.Builder() - .setAdvertiseMode(AdvertiseSettings.ADVERTISE_MODE_LOW_POWER) - .setTxPowerLevel(AdvertiseSettings.ADVERTISE_TX_POWER_ULTRA_LOW) - .setConnectable(true) - .setTimeout(0) - .build() - } + + fun getAdvertiseSettings(): AdvertiseSettings = when (profile.value.mode) { + PowerMode.PERFORMANCE -> AdvertiseSettings.Builder() + .setAdvertiseMode(AdvertiseSettings.ADVERTISE_MODE_LOW_LATENCY) + .setTxPowerLevel(AdvertiseSettings.ADVERTISE_TX_POWER_HIGH) + .setConnectable(true).setTimeout(0).build() + PowerMode.BALANCED -> AdvertiseSettings.Builder() + .setAdvertiseMode(AdvertiseSettings.ADVERTISE_MODE_BALANCED) + .setTxPowerLevel(AdvertiseSettings.ADVERTISE_TX_POWER_MEDIUM) + .setConnectable(true).setTimeout(0).build() + PowerMode.POWER_SAVER -> AdvertiseSettings.Builder() + .setAdvertiseMode(AdvertiseSettings.ADVERTISE_MODE_LOW_POWER) + .setTxPowerLevel(AdvertiseSettings.ADVERTISE_TX_POWER_LOW) + .setConnectable(true).setTimeout(0).build() + PowerMode.ULTRA_LOW_POWER -> AdvertiseSettings.Builder() + .setAdvertiseMode(AdvertiseSettings.ADVERTISE_MODE_LOW_POWER) + .setTxPowerLevel(AdvertiseSettings.ADVERTISE_TX_POWER_ULTRA_LOW) + .setConnectable(true).setTimeout(0).build() } - - /** - * Get maximum allowed connections for current power mode - */ - fun getMaxConnections(): Int { - return when (currentMode) { - PowerMode.PERFORMANCE -> MAX_CONNECTIONS_NORMAL - PowerMode.BALANCED -> MAX_CONNECTIONS_NORMAL - PowerMode.POWER_SAVER -> MAX_CONNECTIONS_POWER_SAVE - PowerMode.ULTRA_LOW_POWER -> MAX_CONNECTIONS_ULTRA_LOW - } + + fun getMaxConnections(): Int = profile.value.ble.maxConnections + + fun getRSSIThreshold(): Int = when (profile.value.mode) { + PowerMode.PERFORMANCE -> -95 + PowerMode.BALANCED -> -85 + PowerMode.POWER_SAVER -> -75 + PowerMode.ULTRA_LOW_POWER -> -65 } - - /** - * Get RSSI filter threshold for current power mode - */ - fun getRSSIThreshold(): Int { - return when (currentMode) { - PowerMode.PERFORMANCE -> -95 - PowerMode.BALANCED -> -85 - PowerMode.POWER_SAVER -> -75 - PowerMode.ULTRA_LOW_POWER -> -65 - } - } - - /** - * Check if duty cycling should be used - */ - fun shouldUseDutyCycle(): Boolean { - return currentMode != PowerMode.PERFORMANCE - } - - /** - * Get current power mode information - */ - fun getPowerInfo(): String { - return buildString { + + fun shouldUseDutyCycle(): Boolean = !profile.value.ble.continuousScan + + fun getPowerInfo(): String = profile.value.let { current -> + buildString { appendLine("=== Power Manager Status ===") - appendLine("Current Mode: $currentMode") - appendLine("Battery Level: $batteryLevel%") - appendLine("Is Charging: $isCharging") - appendLine("App In Background: $isAppInBackground") - appendLine("Max Connections: ${getMaxConnections()}") - appendLine("RSSI Threshold: ${getRSSIThreshold()} dBm") - appendLine("Use Duty Cycle: ${shouldUseDutyCycle()}") + appendLine("Current Mode: ${current.mode}") + appendLine("Battery Band: ${current.batteryBand} ($batteryLevel%)") + appendLine("Is Charging: ${current.isCharging}") + appendLine("App In Background: ${current.isBackground}") + appendLine("Has Direct Peers: ${current.hasDirectPeers}") + appendLine("BLE Scan: ${current.ble.scanOnMs}ms ON / ${current.ble.scanOffMs}ms OFF") + appendLine("Max Connections: ${current.ble.maxConnections}") } } - - private fun updatePowerMode() { - // Determine the base mode from battery/charging state only - val baseMode = when { - // Charging in foreground may use performance - isCharging && !isAppInBackground -> PowerMode.PERFORMANCE - // Critical battery - force ultra low power regardless of foreground/background - batteryLevel <= CRITICAL_BATTERY -> PowerMode.ULTRA_LOW_POWER - - // Low battery - prefer power saver - batteryLevel <= LOW_BATTERY -> PowerMode.POWER_SAVER - - // Otherwise balanced - else -> PowerMode.BALANCED - } - - // If app is in background (including when running as a foreground service), - // cap the power mode to at least POWER_SAVER. Preserve ULTRA_LOW_POWER. - val newMode = if (isAppInBackground) { - if (baseMode == PowerMode.ULTRA_LOW_POWER) PowerMode.ULTRA_LOW_POWER else PowerMode.POWER_SAVER - } else { - baseMode - } - - if (newMode != currentMode) { - val oldMode = currentMode - currentMode = newMode - Log.i(TAG, "Power mode changed: $oldMode → $newMode (battery: $batteryLevel%, charging: $isCharging, background: $isAppInBackground)") - - delegate?.onPowerModeChanged(currentMode) - - // Restart duty cycle with new parameters - if (shouldUseDutyCycle()) { - startDutyCycle() - } else { - stopDutyCycle() - } + private fun refreshProfile() { + if (shutdown) return + val next = PowerProfileResolver.resolve( + batteryLevel = batteryLevel, + isCharging = isCharging, + isBackground = isAppInBackground, + hasDirectPeers = hasDirectPeers + ) + if (_profile.value != next) { + Log.i( + TAG, + "Profile changed: mode=${next.mode}, battery=${next.batteryBand}, " + + "background=${next.isBackground}, directPeers=${next.hasDirectPeers}" + ) + _profile.value = next } } - - private fun startDutyCycle() { - stopDutyCycle() - - if (!shouldUseDutyCycle()) { - delegate?.onScanStateChanged(true) // Always scan in performance mode - return - } - - val (onDuration, offDuration) = when (currentMode) { - PowerMode.BALANCED -> SCAN_ON_DURATION_NORMAL to SCAN_OFF_DURATION_NORMAL - PowerMode.POWER_SAVER -> SCAN_ON_DURATION_POWER_SAVE to SCAN_OFF_DURATION_POWER_SAVE - PowerMode.ULTRA_LOW_POWER -> SCAN_ON_DURATION_ULTRA_LOW to SCAN_OFF_DURATION_ULTRA_LOW - PowerMode.PERFORMANCE -> return // No duty cycle - } - - dutyCycleJob = powerScope.launch { - while (isActive && shouldUseDutyCycle()) { - // Scan ON period - Log.d(TAG, "Duty cycle: Scan ON for ${onDuration}ms") - delegate?.onScanStateChanged(true) - delay(onDuration) - - // Scan OFF period (keep advertising active) - if (isActive && shouldUseDutyCycle()) { - Log.d(TAG, "Duty cycle: Scan OFF for ${offDuration}ms") - delegate?.onScanStateChanged(false) - delay(offDuration) - } - } - } - - Log.i(TAG, "Started duty cycle: ${onDuration}ms ON, ${offDuration}ms OFF") - } - - private fun stopDutyCycle() { - dutyCycleJob?.cancel() - dutyCycleJob = null - } - + private fun registerBatteryReceiver() { try { val filter = IntentFilter().apply { @@ -342,25 +276,141 @@ class PowerManager(private val context: Context) : LifecycleEventObserver { addAction(Intent.ACTION_POWER_CONNECTED) addAction(Intent.ACTION_POWER_DISCONNECTED) } - context.registerReceiver(batteryReceiver, filter) + appContext.registerReceiver(batteryReceiver, filter) } catch (e: Exception) { Log.w(TAG, "Failed to register battery receiver: ${e.message}") } } - - private fun unregisterBatteryReceiver() { - try { - context.unregisterReceiver(batteryReceiver) - } catch (e: Exception) { - Log.w(TAG, "Failed to unregister battery receiver: ${e.message}") - } - } } /** - * Delegate interface for power management callbacks + * Pure policy resolver so cadence and battery-boundary behavior can be tested without Android. */ -interface PowerManagerDelegate { - fun onPowerModeChanged(newMode: PowerManager.PowerMode) - fun onScanStateChanged(shouldScan: Boolean) +internal object PowerProfileResolver { + fun resolve( + batteryLevel: Int, + isCharging: Boolean, + isBackground: Boolean, + hasDirectPeers: Boolean + ): PowerManager.RuntimePerformanceProfile { + val batteryBand = when { + batteryLevel <= AppConstants.Power.CRITICAL_BATTERY_PERCENT -> + PowerManager.BatteryBand.CRITICAL + batteryLevel <= AppConstants.Power.LOW_BATTERY_PERCENT -> + PowerManager.BatteryBand.LOW + else -> PowerManager.BatteryBand.NORMAL + } + + val mode = when { + isBackground && batteryBand == PowerManager.BatteryBand.CRITICAL -> + PowerManager.PowerMode.ULTRA_LOW_POWER + isBackground -> PowerManager.PowerMode.POWER_SAVER + isCharging -> PowerManager.PowerMode.PERFORMANCE + batteryBand == PowerManager.BatteryBand.CRITICAL -> + PowerManager.PowerMode.ULTRA_LOW_POWER + batteryBand == PowerManager.BatteryBand.LOW -> + PowerManager.PowerMode.POWER_SAVER + else -> PowerManager.PowerMode.BALANCED + } + + val ble = when { + isBackground && hasDirectPeers -> + PowerManager.BleSchedule( + scanOnMs = 1_000L, + scanOffMs = 29_000L, + continuousScan = false, + rssiPollIntervalMs = backgroundRssiInterval(batteryBand) + ) + isBackground -> + PowerManager.BleSchedule( + scanOnMs = 1_000L, + scanOffMs = 59_000L, + continuousScan = false, + rssiPollIntervalMs = backgroundRssiInterval(batteryBand) + ) + mode == PowerManager.PowerMode.PERFORMANCE -> + PowerManager.BleSchedule(Long.MAX_VALUE, 0L, true, 5_000L) + mode == PowerManager.PowerMode.BALANCED -> + PowerManager.BleSchedule(8_000L, 2_000L, false, 10_000L) + mode == PowerManager.PowerMode.POWER_SAVER -> + PowerManager.BleSchedule(2_000L, 28_000L, false, 30_000L) + else -> + PowerManager.BleSchedule(1_000L, 29_000L, false, 60_000L) + } + + val announcementInterval = when { + isBackground && batteryBand == PowerManager.BatteryBand.NORMAL -> 60_000L + isBackground && batteryBand == PowerManager.BatteryBand.LOW -> 120_000L + isBackground -> 300_000L + mode == PowerManager.PowerMode.POWER_SAVER -> 60_000L + mode == PowerManager.PowerMode.ULTRA_LOW_POWER -> 120_000L + else -> 30_000L + } + + val wifi = when { + isBackground && batteryBand == PowerManager.BatteryBand.NORMAL -> + wifi(10, 60, 60, 5, 5, 10) + isBackground && batteryBand == PowerManager.BatteryBand.LOW -> + wifi(20, 120, 120, 10, 10, 15) + isBackground -> + wifi(30, 180, 180, 15, 15, 20) + mode == PowerManager.PowerMode.PERFORMANCE -> + PowerManager.WifiAwareSchedule(2_000L, 20_000L, 15_000L, 120_000L, 90_000L, 300_000L) + mode == PowerManager.PowerMode.BALANCED -> + PowerManager.WifiAwareSchedule(5_000L, 30_000L, 30_000L, 180_000L, 120_000L, 300_000L) + mode == PowerManager.PowerMode.POWER_SAVER -> + PowerManager.WifiAwareSchedule(10_000L, 60_000L, 60_000L, 300_000L, 180_000L, 600_000L) + else -> + PowerManager.WifiAwareSchedule(15_000L, 90_000L, 90_000L, 600_000L, 300_000L, 900_000L) + } + + val nostr = when { + isBackground && batteryBand == PowerManager.BatteryBand.NORMAL -> + PowerManager.NostrSchedule(300_000L, 330_000L, 300_000L) + isBackground && batteryBand == PowerManager.BatteryBand.LOW -> + PowerManager.NostrSchedule(600_000L, 660_000L, 600_000L) + isBackground -> + PowerManager.NostrSchedule(900_000L, 990_000L, 900_000L) + mode == PowerManager.PowerMode.POWER_SAVER -> + PowerManager.NostrSchedule(120_000L, 132_000L, 60_000L) + mode == PowerManager.PowerMode.ULTRA_LOW_POWER -> + PowerManager.NostrSchedule(300_000L, 330_000L, 120_000L) + else -> + PowerManager.NostrSchedule(40_000L, 80_000L, 30_000L) + } + + return PowerManager.RuntimePerformanceProfile( + mode = mode, + batteryBand = batteryBand, + isBackground = isBackground, + isCharging = isCharging, + hasDirectPeers = hasDirectPeers, + ble = ble, + meshAnnouncementIntervalMs = announcementInterval, + wifiAware = wifi, + nostr = nostr + ) + } + + private fun wifi( + tcpSeconds: Long, + discoverySeconds: Long, + maintenanceSeconds: Long, + idleMinutes: Long, + minRefreshMinutes: Long, + staleMinutes: Long + ) = PowerManager.WifiAwareSchedule( + tcpKeepAliveMs = tcpSeconds * 1_000L, + discoveryKeepAliveMs = discoverySeconds * 1_000L, + connectionMaintenanceMs = maintenanceSeconds * 1_000L, + discoveryIdleRefreshMs = idleMinutes * 60_000L, + discoverySessionRefreshMinMs = minRefreshMinutes * 60_000L, + discoveryStaleMs = staleMinutes * 60_000L + ) + + private fun backgroundRssiInterval(band: PowerManager.BatteryBand): Long = when (band) { + PowerManager.BatteryBand.NORMAL -> 60_000L + PowerManager.BatteryBand.LOW -> 120_000L + PowerManager.BatteryBand.CRITICAL -> 300_000L + } } diff --git a/app/src/main/java/com/bitchat/android/mesh/UnifiedMeshService.kt b/app/src/main/java/com/bitchat/android/mesh/UnifiedMeshService.kt index ca49def9..5a0cc697 100644 --- a/app/src/main/java/com/bitchat/android/mesh/UnifiedMeshService.kt +++ b/app/src/main/java/com/bitchat/android/mesh/UnifiedMeshService.kt @@ -7,6 +7,16 @@ import com.bitchat.android.model.BitchatFilePacket import com.bitchat.android.model.BitchatMessage import com.bitchat.android.noise.NoiseSession import com.bitchat.android.wifiaware.WifiAwareController +import kotlinx.coroutines.CoroutineScope +import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.Job +import kotlinx.coroutines.SupervisorJob +import kotlinx.coroutines.delay +import kotlinx.coroutines.flow.collectLatest +import kotlinx.coroutines.flow.distinctUntilChanged +import kotlinx.coroutines.flow.map +import kotlinx.coroutines.isActive +import kotlinx.coroutines.launch /** * Feature-facing mesh service that hides local transport selection from the rest of the app. @@ -24,6 +34,10 @@ class UnifiedMeshService( private const val TAG = "UnifiedMeshService" } + private val serviceScope = CoroutineScope(Dispatchers.Default + SupervisorJob()) + private val powerManager = PowerManager.getInstance(context.applicationContext) + private var announcementJob: Job? = null + override val myPeerID: String get() = bluetooth.myPeerID @@ -49,14 +63,37 @@ class UnifiedMeshService( try { WifiAwareController.startIfPossible() } catch (e: Exception) { Log.w(TAG, "Failed to start Wi-Fi Aware transport: ${e.message}") } + startAnnouncementScheduler() refreshDelegates() } override fun stopServices() { + announcementJob?.cancel() + announcementJob = null try { bluetooth.stopServices() } catch (_: Exception) { } try { WifiAwareController.stop() } catch (_: Exception) { } } + private fun startAnnouncementScheduler() { + if (announcementJob?.isActive == true) return + announcementJob = serviceScope.launch { + powerManager.profile + .map { profile -> + profile.meshAnnouncementIntervalMs to profile.hasDirectPeers + } + .distinctUntilChanged() + .collectLatest { (intervalMs, hasRecipients) -> + if (!hasRecipients) return@collectLatest + // Connection-specific paths already send an immediate announce. Begin the + // periodic cadence after the configured interval to avoid a transition burst. + while (isActive) { + delay(intervalMs) + if (powerManager.profile.value.hasDirectPeers) sendBroadcastAnnounce() + } + } + } + } + override fun sendMessage(content: String, mentions: List, channel: String?) { when { isBleEnabled() -> bluetooth.sendMessage(content, mentions, channel) diff --git a/app/src/main/java/com/bitchat/android/net/ArtiTorManager.kt b/app/src/main/java/com/bitchat/android/net/ArtiTorManager.kt index d560c648..b1eed65e 100644 --- a/app/src/main/java/com/bitchat/android/net/ArtiTorManager.kt +++ b/app/src/main/java/com/bitchat/android/net/ArtiTorManager.kt @@ -133,6 +133,9 @@ class ArtiTorManager private constructor() { lastLogTime.set(System.currentTimeMillis()) _statusFlow.update { it.copy(lastLogLine = s) } handleArtiLogLine(s) + if (_statusFlow.value.bootstrapPercent < 100) { + armBootstrapInactivityWatchdog() + } } artiProxy = ArtiProxy.Builder(application) @@ -374,29 +377,30 @@ class ArtiTorManager private constructor() { } private fun startInactivityMonitoring() { - inactivityJob?.cancel() - inactivityJob = appScope.launch { - while (true) { - delay(INACTIVITY_TIMEOUT_MS) - val currentTime = System.currentTimeMillis() - val lastActivity = lastLogTime.get() - val timeSinceLastActivity = currentTime - lastActivity + armBootstrapInactivityWatchdog() + } - if (timeSinceLastActivity > INACTIVITY_TIMEOUT_MS) { - val currentMode = _statusFlow.value.mode - if (currentMode == TorMode.ON) { - val bootstrapPercent = _statusFlow.value.bootstrapPercent - if (bootstrapPercent < 100) { - Log.w(TAG, "Inactivity detected (${timeSinceLastActivity}ms), restarting Arti") - currentApplication?.let { app -> - appScope.launch { - restartArti(app) - } - } - break - } - } - } + /** + * One-shot startup watchdog. Arti log activity re-arms it and successful bootstrap cancels it, + * so a healthy background Tor process has no permanent five-second polling coroutine. + */ + private fun armBootstrapInactivityWatchdog() { + inactivityJob?.cancel() + if (lifecycleState != LifecycleState.RUNNING || _statusFlow.value.bootstrapPercent >= 100) { + inactivityJob = null + return + } + inactivityJob = appScope.launch { + delay(INACTIVITY_TIMEOUT_MS) + val timeSinceLastActivity = System.currentTimeMillis() - lastLogTime.get() + if ( + timeSinceLastActivity >= INACTIVITY_TIMEOUT_MS && + _statusFlow.value.mode == TorMode.ON && + _statusFlow.value.bootstrapPercent < 100 && + lifecycleState == LifecycleState.RUNNING + ) { + Log.w(TAG, "Bootstrap inactivity detected (${timeSinceLastActivity}ms), restarting Arti") + currentApplication?.let { restartArti(it) } } } } @@ -494,6 +498,7 @@ class ArtiTorManager private constructor() { running = true ) } + stopInactivityMonitoring() completeWaitersIf(TorState.RUNNING) } diff --git a/app/src/main/java/com/bitchat/android/nostr/GeohashMessageHandler.kt b/app/src/main/java/com/bitchat/android/nostr/GeohashMessageHandler.kt index 1fc802d9..88d41466 100644 --- a/app/src/main/java/com/bitchat/android/nostr/GeohashMessageHandler.kt +++ b/app/src/main/java/com/bitchat/android/nostr/GeohashMessageHandler.kt @@ -3,8 +3,6 @@ package com.bitchat.android.nostr import android.app.Application import android.util.Log import com.bitchat.android.model.BitchatMessage -import com.bitchat.android.ui.ChatState -import com.bitchat.android.ui.MessageManager import kotlinx.coroutines.CoroutineScope import kotlinx.coroutines.launch import java.util.Date @@ -17,11 +15,10 @@ import java.util.Date */ class GeohashMessageHandler( private val application: Application, - private val state: ChatState, - private val messageManager: MessageManager, private val repo: GeohashRepository, private val scope: CoroutineScope, - private val dataManager: com.bitchat.android.ui.DataManager + private val dataManager: com.bitchat.android.ui.DataManager, + private val addChannelMessage: (String, BitchatMessage) -> Unit ) { companion object { private const val TAG = "GeohashMessageHandler" } @@ -103,7 +100,7 @@ class GeohashMessageHandler( if (hasNonce) NostrProofOfWork.calculateDifficulty(event.id).takeIf { it > 0 } else null } catch (_: Exception) { null } ) - messageManager.addChannelMessage("geo:$subscribedGeohash", msg) + addChannelMessage("geo:$subscribedGeohash", msg) } catch (e: Exception) { Log.e(TAG, "onEvent error: ${e.message}") } diff --git a/app/src/main/java/com/bitchat/android/nostr/GeohashRepository.kt b/app/src/main/java/com/bitchat/android/nostr/GeohashRepository.kt index 3a7ced32..50c72cd4 100644 --- a/app/src/main/java/com/bitchat/android/nostr/GeohashRepository.kt +++ b/app/src/main/java/com/bitchat/android/nostr/GeohashRepository.kt @@ -28,14 +28,17 @@ class GeohashRepository( // conversation key (e.g., "nostr_") -> source geohash it belongs to private val conversationGeohash: MutableMap = mutableMapOf() + @Synchronized fun setConversationGeohash(convKey: String, geohash: String) { if (geohash.isNotEmpty()) { conversationGeohash[convKey] = geohash } } + @Synchronized fun getConversationGeohash(convKey: String): String? = conversationGeohash[convKey] + @Synchronized fun findPubkeyByNickname(targetNickname: String): String? { return geoNicknames.entries.firstOrNull { (_, nickname) -> val base = nickname.split("#").firstOrNull() ?: nickname @@ -43,6 +46,7 @@ class GeohashRepository( }?.key } + @Synchronized fun findPubkeyByShortId(shortId: String): String? { // First check cached nicknames (fastest) var found = geoNicknames.keys.firstOrNull { it.startsWith(shortId, ignoreCase = true) } @@ -66,6 +70,7 @@ class GeohashRepository( fun setCurrentGeohash(geo: String?) { currentGeohash = geo } fun getCurrentGeohash(): String? = currentGeohash + @Synchronized fun clearAll() { geohashParticipants.clear() geoNicknames.clear() @@ -76,6 +81,7 @@ class GeohashRepository( currentGeohash = null } + @Synchronized fun cacheNickname(pubkeyHex: String, nickname: String) { val lower = pubkeyHex.lowercase() val previous = geoNicknames[lower] @@ -85,8 +91,10 @@ class GeohashRepository( } } + @Synchronized fun getCachedNickname(pubkeyHex: String): String? = geoNicknames[pubkeyHex.lowercase()] + @Synchronized fun markTeleported(pubkeyHex: String) { val set = state.getTeleportedGeoValue().toMutableSet() val key = pubkeyHex.lowercase() @@ -97,10 +105,12 @@ class GeohashRepository( } } + @Synchronized fun isPersonTeleported(pubkeyHex: String): Boolean { return state.getTeleportedGeoValue().contains(pubkeyHex.lowercase()) } + @Synchronized fun updateParticipant(geohash: String, participantId: String, lastSeen: Date) { val participants = geohashParticipants.getOrPut(geohash) { mutableMapOf() } // Cap to now: prevents future-timestamped events (clock skew / malicious created_at) @@ -117,6 +127,7 @@ class GeohashRepository( updateReactiveParticipantCounts() } + @Synchronized fun geohashParticipantCount(geohash: String): Int { val cutoff = Date(System.currentTimeMillis() - 5 * 60 * 1000) val participants = geohashParticipants[geohash] ?: return 0 @@ -130,6 +141,7 @@ class GeohashRepository( return participants.keys.count { !dataManager.isGeohashUserBlocked(it) } } + @Synchronized fun refreshGeohashPeople() { val geohash = currentGeohash if (geohash == null) { @@ -168,6 +180,7 @@ class GeohashRepository( state.setGeohashPeople(people) } + @Synchronized fun updateReactiveParticipantCounts() { val cutoff = Date(System.currentTimeMillis() - 5 * 60 * 1000) val counts = mutableMapOf() @@ -180,12 +193,15 @@ class GeohashRepository( state.setGeohashParticipantCounts(counts) } + @Synchronized fun putNostrKeyMapping(tempKeyOrPeer: String, pubkeyHex: String) { nostrKeyMapping[tempKeyOrPeer] = pubkeyHex } + @Synchronized fun getNostrKeyMapping(): Map = nostrKeyMapping.toMap() + @Synchronized fun displayNameForNostrPubkey(pubkeyHex: String): String { val suffix = pubkeyHex.takeLast(4) val lower = pubkeyHex.lowercase() @@ -203,6 +219,7 @@ class GeohashRepository( return "$nick#$suffix" } + @Synchronized fun displayNameForNostrPubkeyUI(pubkeyHex: String): String { val lower = pubkeyHex.lowercase() val suffix = pubkeyHex.takeLast(4) @@ -234,6 +251,7 @@ class GeohashRepository( /** * Get display name for any geohash (not just current one) for header titles */ + @Synchronized fun displayNameForGeohashConversation(pubkeyHex: String, sourceGeohash: String): String { val lower = pubkeyHex.lowercase() val suffix = pubkeyHex.takeLast(4) diff --git a/app/src/main/java/com/bitchat/android/nostr/NostrBackgroundEventProcessor.kt b/app/src/main/java/com/bitchat/android/nostr/NostrBackgroundEventProcessor.kt new file mode 100644 index 00000000..442840eb --- /dev/null +++ b/app/src/main/java/com/bitchat/android/nostr/NostrBackgroundEventProcessor.kt @@ -0,0 +1,115 @@ +package com.bitchat.android.nostr + +import android.app.Application +import com.bitchat.android.model.DeliveryStatus +import com.bitchat.android.services.AppStateStore +import com.bitchat.android.ui.ChatState +import com.bitchat.android.ui.DataManager +import com.bitchat.android.ui.MessageManager +import com.bitchat.android.ui.NoiseSessionDelegate +import com.bitchat.android.ui.PrivateChatManager +import kotlinx.coroutines.CoroutineScope +import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.launch + +/** + * Process-owned Nostr event processing. + * + * Relay subscriptions must remain useful when no Activity exists, so their handlers cannot be + * borrowed from a ViewModel. This processor owns only application-scoped collaborators and writes + * messages to [AppStateStore], which the next UI instance hydrates from. + */ +internal class NostrBackgroundEventProcessor( + application: Application, + parentScope: CoroutineScope +) { + private val scope = CoroutineScope( + parentScope.coroutineContext + Dispatchers.IO.limitedParallelism(1) + ) + private val state = ChatState(scope) + private val dataManager = DataManager(application.applicationContext).apply { + state.setNickname(loadNickname()) + loadBlockedUsers() + loadGeohashBlockedUsers() + } + private val messageManager = MessageManager(state) + private val geohashRepository = GeohashRepository(application, state, dataManager) + private val privateChatManager = PrivateChatManager( + state = state, + messageManager = messageManager, + dataManager = dataManager, + noiseSessionDelegate = object : NoiseSessionDelegate { + override fun hasEstablishedSession(peerID: String): Boolean = false + override fun initiateHandshake(peerID: String) = Unit + override fun getMyPeerID(): String = "" + }, + trackUnreadMessages = false + ) + private val geohashMessageHandler = GeohashMessageHandler( + application = application, + repo = geohashRepository, + scope = scope, + dataManager = dataManager, + addChannelMessage = AppStateStore::addChannelMessage + ) + private val directMessageHandler = NostrDirectMessageHandler( + application = application, + state = state, + privateChatManager = privateChatManager, + updateDeliveryStatus = ::updateDeliveryStatus, + scope = scope, + repo = geohashRepository, + dataManager = dataManager + ) + + init { + // Keep the headless state aligned with messages sent or received through other transports. + // This preserves duplicate detection and focused-conversation behavior without retaining UI. + scope.launch { + AppStateStore.privateMessages.collect(state::setPrivateChats) + } + scope.launch { + AppStateStore.nickname.collect(state::setNickname) + } + scope.launch { + AppStateStore.selectedPrivateChatPeer.collect(state::setSelectedPrivateChatPeer) + } + } + + fun onAccountDm(event: NostrEvent, identity: NostrIdentity) { + refreshBlockLists() + directMessageHandler.onGiftWrap(event, "", identity) + } + + fun onGeohashMessage(event: NostrEvent, geohash: String) { + refreshBlockLists() + geohashMessageHandler.onEvent(event, geohash) + } + + fun onGeohashDm(event: NostrEvent, geohash: String, identity: NostrIdentity) { + refreshBlockLists() + directMessageHandler.onGiftWrap(event, geohash, identity) + } + + fun conversationGeohash(conversationKey: String): String? = + geohashRepository.getConversationGeohash(conversationKey) + ?: GeohashConversationRegistry.get(conversationKey) + + fun displayNameForNostrPubkey(pubkeyHex: String): String = + geohashRepository.displayNameForNostrPubkeyUI(pubkeyHex) + + fun displayNameForGeohashConversation(pubkeyHex: String, sourceGeohash: String): String = + geohashRepository.displayNameForGeohashConversation(pubkeyHex, sourceGeohash) + + private fun updateDeliveryStatus(messageId: String, status: DeliveryStatus) { + messageManager.updateMessageDeliveryStatus(messageId, status) + // The headless state may not yet contain a just-sent UI message. Update the process store + // unconditionally so a delivery/read receipt can never be lost during Activity handoff. + AppStateStore.updatePrivateMessageStatus(messageId, status) + } + + private fun refreshBlockLists() { + dataManager.loadBlockedUsers() + dataManager.loadGeohashBlockedUsers() + } +} diff --git a/app/src/main/java/com/bitchat/android/nostr/NostrBackgroundRuntime.kt b/app/src/main/java/com/bitchat/android/nostr/NostrBackgroundRuntime.kt new file mode 100644 index 00000000..a895f612 --- /dev/null +++ b/app/src/main/java/com/bitchat/android/nostr/NostrBackgroundRuntime.kt @@ -0,0 +1,272 @@ +package com.bitchat.android.nostr + +import android.app.Application +import android.util.Log +import com.bitchat.android.geohash.ChannelID +import com.bitchat.android.geohash.GeohashNostrPrivacyPolicy +import com.bitchat.android.geohash.LiveLocationPrivacyGate +import com.bitchat.android.geohash.LocationChannelManager +import com.bitchat.android.mesh.PowerManager +import kotlinx.coroutines.CancellationException +import kotlinx.coroutines.CoroutineScope +import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.SupervisorJob +import kotlinx.coroutines.delay +import kotlinx.coroutines.flow.collectLatest +import kotlinx.coroutines.flow.combine +import kotlinx.coroutines.flow.distinctUntilChanged +import kotlinx.coroutines.flow.map +import kotlinx.coroutines.launch +import java.security.SecureRandom +import kotlin.random.asKotlinRandom + +/** + * Process-owned Nostr connectivity and low-volume background subscriptions. + * + * Stable subscriptions dispatch directly to a process-owned event processor. The UI hydrates from + * the process state store, so relay events remain useful without retaining a cleared ViewModel. + */ +object NostrBackgroundRuntime { + private const val TAG = "NostrBackground" + + private val scope = CoroutineScope(Dispatchers.IO + SupervisorJob()) + private val random = SecureRandom().asKotlinRandom() + private val lock = Any() + + @Volatile private var initialized = false + @Volatile private var activeGeohash: String? = null + @Volatile private var activeGeohashLiveToken: Long? = null + @Volatile private var conversationGeohash: String? = null + private lateinit var application: Application + private lateinit var subscriptions: NostrSubscriptionManager + private lateinit var locationChannels: LocationChannelManager + private lateinit var eventProcessor: NostrBackgroundEventProcessor + + fun initialize(app: Application) { + synchronized(lock) { + if (initialized) return + application = app + locationChannels = LocationChannelManager.getInstance(app) + eventProcessor = NostrBackgroundEventProcessor(app, scope) + subscriptions = NostrSubscriptionManager( + app, + owner = NostrRelayManager.OWNER_BACKGROUND + ) + // Publish readiness only after every process-owned dependency is available. + initialized = true + } + + subscriptions.connect() + subscribeAccountDm() + observeSelectedChannel() + startPresenceScheduler() + } + + fun resetSubscriptions() { + if (!initialized) return + subscriptions.unsubscribeAllOwned() + scope.launch { + // Let CLOSE frames be queued before replacing the deterministic IDs. + delay(100) + subscribeAccountDm() + activeGeohash?.let { geohash -> + subscribeSelectedGeohash(geohash, activeGeohashLiveToken) + } + } + } + + fun ensureConversationDm(geohash: String) { + if (!initialized || geohash == conversationGeohash) return + val selectedLiveToken = activeGeohashLiveToken + val selectedChannelSubscriptionIsUsable = + geohash == activeGeohash && + (selectedLiveToken == null || + LiveLocationPrivacyGate.accepts(selectedLiveToken)) + if (selectedChannelSubscriptionIsUsable) return + conversationGeohash?.let { subscriptions.unsubscribe("geo-dm-conversation-$it") } + conversationGeohash = geohash + subscribeGeohashDm( + geohash = geohash, + subscriptionId = "geo-dm-conversation-$geohash", + liveLocationToken = null + ) + } + + private fun subscribeAccountDm() { + val identity = NostrIdentityBridge.getCurrentNostrIdentity(application) ?: return + subscriptions.subscribeGiftWraps( + pubkey = identity.publicKeyHex, + sinceMs = System.currentTimeMillis() - 172_800_000L, + id = "chat-messages", + handler = { event -> + eventProcessor.onAccountDm(event, identity) + } + ) + } + + private fun observeSelectedChannel() { + scope.launch { + locationChannels.selectedChannel.collectLatest { channel -> + val locationChannel = channel as? ChannelID.Location + val next = locationChannel?.channel?.geohash + val nextToken = locationChannel?.let { + locationChannels.liveLocationTokenForSelectedChannel(it.channel) + } + val previous = activeGeohash + if (previous == next && activeGeohashLiveToken == nextToken) { + return@collectLatest + } + + previous?.let { + subscriptions.unsubscribe("geohash-$it") + subscriptions.unsubscribe("geo-dm-$it") + } + activeGeohash = next + activeGeohashLiveToken = nextToken + if (conversationGeohash == next) { + subscriptions.unsubscribe("geo-dm-conversation-$next") + conversationGeohash = null + } + next?.let { geohash -> + val isLiveDerived = + locationChannels.isSelectedChannelLiveDerived(locationChannel.channel) + if (!isLiveDerived || nextToken != null) { + subscribeSelectedGeohash(geohash, nextToken) + } + } + } + } + } + + private fun subscribeSelectedGeohash( + geohash: String, + liveLocationToken: Long? + ) { + subscriptions.subscribeGeohashMessages( + geohash = geohash, + sinceMs = System.currentTimeMillis() - 3_600_000L, + limit = 200, + id = "geohash-$geohash", + handler = { event -> eventProcessor.onGeohashMessage(event, geohash) }, + liveLocationToken = liveLocationToken + ) + subscribeGeohashDm(geohash, "geo-dm-$geohash", liveLocationToken) + } + + private fun subscribeGeohashDm( + geohash: String, + subscriptionId: String, + liveLocationToken: Long? + ) { + scope.launch { + val subscribe = { + val identity = NostrIdentityBridge.deriveIdentity(geohash, application) + subscriptions.subscribeGiftWraps( + pubkey = identity.publicKeyHex, + sinceMs = System.currentTimeMillis() - 172_800_000L, + id = subscriptionId, + handler = { event -> + eventProcessor.onGeohashDm(event, geohash, identity) + }, + liveLocationToken = liveLocationToken + ) + GeohashAliasRegistry.put( + "nostr_${identity.publicKeyHex.take(16)}", + identity.publicKeyHex + ) + } + + if (liveLocationToken == null) { + subscribe() + } else { + LiveLocationPrivacyGate.runIfAllowed(liveLocationToken, subscribe) + } + } + } + + private fun startPresenceScheduler() { + val powerManager = PowerManager.getInstance(application) + scope.launch { + val nostrSchedule = powerManager.profile + .map { it.nostr } + .distinctUntilChanged() + combine( + locationChannels.availableChannels, + LiveLocationPrivacyGate.enabled, + nostrSchedule + ) { channels, liveLocationEnabled, schedule -> + val targets = GeohashNostrPrivacyPolicy.livePresenceTargets( + availableChannels = channels, + liveLocationEnabled = liveLocationEnabled + ) + targets to schedule + }.collectLatest { (targets, schedule) -> + if (targets.isEmpty()) return@collectLatest + while (true) { + val waitMs = if (schedule.presenceHeartbeatMaxMs > schedule.presenceHeartbeatMinMs) { + random.nextLong( + schedule.presenceHeartbeatMinMs, + schedule.presenceHeartbeatMaxMs + 1 + ) + } else { + schedule.presenceHeartbeatMinMs + } + delay(waitMs) + + // Send every target in one wake window; do not spread the batch over seconds. + targets.forEach { geohash -> + try { + val token = LiveLocationPrivacyGate.captureToken() + ?: return@forEach + if (geohash !in GeohashNostrPrivacyPolicy.livePresenceTargets( + availableChannels = locationChannels.availableChannels.value, + liveLocationEnabled = true + ) + ) return@forEach + + var identity: NostrIdentity? = null + LiveLocationPrivacyGate.runIfAllowed(token) { + identity = NostrIdentityBridge.deriveIdentity(geohash, application) + } + val preparedIdentity = identity ?: return@forEach + if (!LiveLocationPrivacyGate.accepts(token)) return@forEach + val event = NostrProtocol.createGeohashPresenceEvent( + geohash, + preparedIdentity + ) + LiveLocationPrivacyGate.runIfAllowed(token) { + NostrRelayManager.getInstance(application).sendEventToGeohash( + event = event, + geohash = geohash, + includeDefaults = false, + nRelays = 5, + liveLocationToken = token + ) + } + } catch (e: CancellationException) { + throw e + } catch (e: Exception) { + Log.w(TAG, "Presence heartbeat failed for $geohash: ${e.message}") + } + } + } + } + } + } + + fun conversationGeohash(conversationKey: String): String? = + if (initialized) eventProcessor.conversationGeohash(conversationKey) + else GeohashConversationRegistry.get(conversationKey) + + fun displayNameForNostrPubkey(pubkeyHex: String): String? = + if (initialized) eventProcessor.displayNameForNostrPubkey(pubkeyHex) else null + + fun displayNameForGeohashConversation( + pubkeyHex: String, + sourceGeohash: String + ): String? = if (initialized) { + eventProcessor.displayNameForGeohashConversation(pubkeyHex, sourceGeohash) + } else { + null + } +} diff --git a/app/src/main/java/com/bitchat/android/nostr/NostrDirectMessageHandler.kt b/app/src/main/java/com/bitchat/android/nostr/NostrDirectMessageHandler.kt index fce95c2d..2fdc3180 100644 --- a/app/src/main/java/com/bitchat/android/nostr/NostrDirectMessageHandler.kt +++ b/app/src/main/java/com/bitchat/android/nostr/NostrDirectMessageHandler.kt @@ -15,7 +15,6 @@ import com.bitchat.android.services.ContactDirectory import com.bitchat.android.services.ContactIdentityResolver import com.bitchat.android.services.SeenMessageStore import com.bitchat.android.ui.ChatState -import com.bitchat.android.ui.MeshDelegateHandler import com.bitchat.android.ui.PrivateChatManager import com.bitchat.android.ui.PrivateMessageOrigin import kotlinx.coroutines.CoroutineScope @@ -28,7 +27,7 @@ class NostrDirectMessageHandler( private val application: Application, private val state: ChatState, private val privateChatManager: PrivateChatManager, - private val meshDelegateHandler: MeshDelegateHandler, + private val updateDeliveryStatus: (String, DeliveryStatus) -> Unit, private val scope: CoroutineScope, private val repo: GeohashRepository, private val dataManager: com.bitchat.android.ui.DataManager, @@ -57,7 +56,7 @@ class NostrDirectMessageHandler( } fun onGiftWrap(giftWrap: NostrEvent, geohash: String, identity: NostrIdentity) { - scope.launch(Dispatchers.Default) { + scope.launch { try { if (dedupe(giftWrap.id)) return@launch @@ -182,13 +181,19 @@ class NostrDirectMessageHandler( NoisePayloadType.DELIVERED -> { val messageId = String(payload.data, Charsets.UTF_8) withContext(Dispatchers.Main) { - meshDelegateHandler.didReceiveDeliveryAck(messageId, conversationID) + updateDeliveryStatus( + messageId, + DeliveryStatus.Delivered(conversationID, Date()) + ) } } NoisePayloadType.READ_RECEIPT -> { val messageId = String(payload.data, Charsets.UTF_8) withContext(Dispatchers.Main) { - meshDelegateHandler.didReceiveReadReceipt(messageId, conversationID) + updateDeliveryStatus( + messageId, + DeliveryStatus.Read(conversationID, Date()) + ) } } NoisePayloadType.FILE_TRANSFER -> { diff --git a/app/src/main/java/com/bitchat/android/nostr/NostrPendingEventQueue.kt b/app/src/main/java/com/bitchat/android/nostr/NostrPendingEventQueue.kt new file mode 100644 index 00000000..dd9e16d1 --- /dev/null +++ b/app/src/main/java/com/bitchat/android/nostr/NostrPendingEventQueue.kt @@ -0,0 +1,90 @@ +package com.bitchat.android.nostr + +/** + * Thread-safe bounded queue of relay deliveries awaiting a usable WebSocket. + * + * Queue entries have a local ID rather than using the Nostr event ID: the same signed event may be + * intentionally published more than once with different relay sets or privacy provenance. + */ +internal class NostrPendingEventQueue( + private val capacity: Int +) { + init { + require(capacity > 0) + } + + data class Delivery( + val queueId: Long, + val event: NostrEvent, + val liveLocationToken: Long? + ) + + private data class Entry( + val queueId: Long, + val event: NostrEvent, + val pendingRelayUrls: MutableSet, + val liveLocationToken: Long? + ) + + private val lock = Any() + private val entries = ArrayDeque() + private var nextQueueId = 1L + + fun enqueue( + event: NostrEvent, + relayUrls: Collection, + liveLocationToken: Long? + ): Long? { + val pendingRelays = relayUrls.filterTo(linkedSetOf()) { it.isNotBlank() } + if (pendingRelays.isEmpty()) return null + + return synchronized(lock) { + if (entries.size >= capacity) entries.removeFirst() + val queueId = nextQueueId++ + entries.addLast( + Entry( + queueId = queueId, + event = event, + pendingRelayUrls = pendingRelays, + liveLocationToken = liveLocationToken + ) + ) + queueId + } + } + + fun pendingForRelay(relayUrl: String): List = synchronized(lock) { + entries + .asSequence() + .filter { relayUrl in it.pendingRelayUrls } + .map { Delivery(it.queueId, it.event, it.liveLocationToken) } + .toList() + } + + fun markDelivered(queueId: Long, relayUrl: String) { + synchronized(lock) { + val iterator = entries.iterator() + while (iterator.hasNext()) { + val entry = iterator.next() + if (entry.queueId != queueId) continue + entry.pendingRelayUrls.remove(relayUrl) + if (entry.pendingRelayUrls.isEmpty()) iterator.remove() + return + } + } + } + + fun removeLiveLocationEvents() { + synchronized(lock) { + entries.removeAll { it.liveLocationToken != null } + } + } + + fun clear() { + synchronized(lock) { + entries.clear() + } + } + + internal fun size(): Int = synchronized(lock) { entries.size } +} diff --git a/app/src/main/java/com/bitchat/android/nostr/NostrRelayManager.kt b/app/src/main/java/com/bitchat/android/nostr/NostrRelayManager.kt index 5e82612f..3386a014 100644 --- a/app/src/main/java/com/bitchat/android/nostr/NostrRelayManager.kt +++ b/app/src/main/java/com/bitchat/android/nostr/NostrRelayManager.kt @@ -2,17 +2,18 @@ package com.bitchat.android.nostr import android.util.Log import com.bitchat.android.geohash.LiveLocationPrivacyGate -import com.google.gson.Gson import kotlinx.coroutines.flow.MutableStateFlow import kotlinx.coroutines.flow.StateFlow import kotlinx.coroutines.flow.asStateFlow -import com.google.gson.JsonArray +import kotlinx.coroutines.flow.collectLatest +import kotlinx.coroutines.flow.distinctUntilChanged +import kotlinx.coroutines.flow.map import com.google.gson.JsonParser import kotlinx.coroutines.* import okhttp3.* import java.util.UUID import java.util.concurrent.ConcurrentHashMap -import java.util.concurrent.TimeUnit +import java.util.concurrent.atomic.AtomicBoolean import kotlin.math.min import kotlin.math.pow @@ -27,11 +28,15 @@ class NostrRelayManager private constructor() { val shared = NostrRelayManager() private const val TAG = "NostrRelayManager" + private const val MAX_QUEUED_EVENTS = 500 + const val OWNER_LEGACY = "legacy" + const val OWNER_BACKGROUND = "background" /** * Get instance for Android compatibility (context-aware calls) */ fun getInstance(context: android.content.Context): NostrRelayManager { + shared.appContext = context.applicationContext return shared } @@ -84,6 +89,8 @@ class NostrRelayManager private constructor() { // Internal state private val relaysList = mutableListOf() private val connections = ConcurrentHashMap() + private val reconnectJobs = ConcurrentHashMap() + private val desiredConnected = AtomicBoolean(false) private val subscriptions = ConcurrentHashMap>() // relay URL -> subscription IDs private val messageHandlers = ConcurrentHashMap Unit>() @@ -99,28 +106,25 @@ class NostrRelayManager private constructor() { val handler: (NostrEvent) -> Unit, val targetRelayUrls: Set? = null, // null means all relays val createdAt: Long = System.currentTimeMillis(), + val originGeohash: String? = null, + val owner: String = OWNER_LEGACY, val liveLocationToken: Long? = null ) // Event deduplication system private val eventDeduplicator = NostrEventDeduplicator.getInstance() - // Message queue for reliability - private data class QueuedEvent( - val event: NostrEvent, - val targetRelays: List, - val liveLocationToken: Long? = null - ) - - private val messageQueue = mutableListOf() - private val messageQueueLock = Any() + // Bounded per-relay delivery queue for reconnect reliability. + private val messageQueue = NostrPendingEventQueue(MAX_QUEUED_EVENTS) // Coroutine scope for background operations private val scope = CoroutineScope(Dispatchers.IO + SupervisorJob()) // Subscription validation timer private var subscriptionValidationJob: Job? = null - private val SUBSCRIPTION_VALIDATION_INTERVAL = com.bitchat.android.util.AppConstants.Nostr.SUBSCRIPTION_VALIDATION_INTERVAL_MS // 30 seconds + private val powerManager: com.bitchat.android.mesh.PowerManager? + get() = appContext?.let { com.bitchat.android.mesh.PowerManager.getInstance(it) } + @Volatile private var appContext: android.content.Context? = null // OkHttp client for WebSocket connections (via provider to honor Tor) private val httpClient: OkHttpClient @@ -193,6 +197,7 @@ class NostrRelayManager private constructor() { handler: (NostrEvent) -> Unit, includeDefaults: Boolean = false, nRelays: Int = 5, + owner: String = OWNER_LEGACY, liveLocationToken: Long? = null ): String { if (!isNetworkActionAllowed(liveLocationToken)) return id @@ -209,8 +214,14 @@ class NostrRelayManager private constructor() { id = id, handler = handler, targetRelayUrls = relayUrls, + owner = owner, liveLocationToken = liveLocationToken - ) + ).also { subscriptionId -> + activeSubscriptions[subscriptionId]?.let { subscription -> + activeSubscriptions[subscriptionId] = + subscription.copy(originGeohash = geohash) + } + } } /** @@ -269,14 +280,24 @@ class NostrRelayManager private constructor() { ) closeTargets.forEach { (relayUrl, relaySubscriptionIds) -> val webSocket = connections[relayUrl] ?: return@forEach - relaySubscriptionIds.forEach { subscriptionId -> + for (subscriptionId in relaySubscriptionIds) { val request = NostrRequest.Close(subscriptionId) val message = gson.toJson(request, NostrRequest::class.java) val closeQueued = runCatching { webSocket.send(message) } .getOrDefault(false) if (!closeQueued) { connections.remove(relayUrl, webSocket) + subscriptions.remove(relayUrl) webSocket.cancel() + updateRelayStatus( + relayUrl, + isConnected = false, + error = IllegalStateException("Failed to close revoked subscription") + ) + if (desiredConnected.get() && relayUrl in nonLiveRelayUrls) { + scope.launch { connectToRelay(relayUrl, liveLocationToken = null) } + } + break } } } @@ -296,9 +317,7 @@ class NostrRelayManager private constructor() { } subscriptions.replaceAll { _, ids -> ids - liveSubscriptionIds } - synchronized(messageQueueLock) { - messageQueue.removeAll { it.liveLocationToken != null } - } + messageQueue.removeLiveLocationEvents() liveGeohashTokens.keys.forEach(geohashToRelays::remove) liveGeohashTokens.clear() @@ -307,6 +326,8 @@ class NostrRelayManager private constructor() { .filterNotTo(mutableSetOf()) { it in nonLiveRelayUrls } liveOnlyRelayUrls.forEach { relayUrl -> connections.remove(relayUrl)?.cancel() + subscriptions.remove(relayUrl) + reconnectJobs.remove(relayUrl)?.cancel() } synchronized(relaysList) { relaysList.removeAll { it.url in liveOnlyRelayUrls } @@ -329,11 +350,15 @@ class NostrRelayManager private constructor() { } updateRelaysList() + if (!desiredConnected.get()) return val job = scope.launch { - if (!isNetworkActionAllowed(liveLocationToken)) return@launch + if (!desiredConnected.get() || + !isNetworkActionAllowed(liveLocationToken) + ) return@launch relayUrls.forEach { relayUrl -> launch { - if (!connections.containsKey(relayUrl) && + if (desiredConnected.get() && + !connections.containsKey(relayUrl) && isNetworkActionAllowed(liveLocationToken) ) { connectToRelay(relayUrl, liveLocationToken) @@ -373,6 +398,8 @@ class NostrRelayManager private constructor() { * Connect to all configured relays */ fun connect() { + desiredConnected.set(true) + Log.i(TAG, "Connecting to ${relaysList.size} Nostr relays") scope.launch { relaysList.forEach { relay -> launch { @@ -393,17 +420,29 @@ class NostrRelayManager private constructor() { * Disconnect from all relays */ fun disconnect() { + Log.i(TAG, "Disconnecting from all Nostr relays") + desiredConnected.set(false) + // Stop subscription validation stopSubscriptionValidation() - - connections.values.forEach { webSocket -> + reconnectJobs.values.forEach(Job::cancel) + reconnectJobs.clear() + liveLocationConnectionJobs.forEach(Job::cancel) + liveLocationConnectionJobs.clear() + + val sockets = connections.values.toList() + connections.clear() + sockets.forEach { webSocket -> webSocket.close(1000, "Manual disconnect") } - connections.clear() - - // Clear subscriptions + + // Preserve logical subscriptions for controlled resets, but forget per-socket state. subscriptions.clear() - + relaysList.forEach { + it.isConnected = false + it.nextReconnectTime = null + } + updateRelaysList() updateConnectionStatus() } @@ -415,18 +454,25 @@ class NostrRelayManager private constructor() { relayUrls: List? = null, liveLocationToken: Long? = null ) { - val targetRelays = relayUrls ?: relaysList.map { it.url } + val targetRelays = (relayUrls ?: relaysList.map { it.url }) + .filter { it.isNotBlank() } + .distinct() + if (targetRelays.isEmpty()) return val queued = runNetworkAction(liveLocationToken) { - synchronized(messageQueueLock) { - messageQueue.add(QueuedEvent(event, targetRelays, liveLocationToken)) - } + val queueId = messageQueue.enqueue( + event = event, + relayUrls = targetRelays, + liveLocationToken = liveLocationToken + ) ?: return@runNetworkAction scope.launch { if (!isNetworkActionAllowed(liveLocationToken)) return@launch targetRelays.forEach { relayUrl -> val webSocket = connections[relayUrl] if (webSocket != null) { - sendToRelay(event, webSocket, relayUrl, liveLocationToken) + if (sendToRelay(event, webSocket, relayUrl, liveLocationToken)) { + messageQueue.markDelivered(queueId, relayUrl) + } } } } @@ -443,6 +489,7 @@ class NostrRelayManager private constructor() { id: String = generateSubscriptionId(), handler: (NostrEvent) -> Unit, targetRelayUrls: List? = null, + owner: String = OWNER_LEGACY, liveLocationToken: Long? = null ): String { val subscriptionInfo = SubscriptionInfo( @@ -450,6 +497,7 @@ class NostrRelayManager private constructor() { filter = filter, handler = handler, targetRelayUrls = targetRelayUrls?.toSet(), + owner = owner, liveLocationToken = liveLocationToken ) @@ -546,12 +594,20 @@ class NostrRelayManager private constructor() { } } } + + fun unsubscribeOwner(owner: String) { + activeSubscriptions.values + .filter { it.owner == owner } + .map { it.id } + .forEach(::unsubscribe) + } /** * Manually retry connection to a specific relay */ fun retryConnection(relayUrl: String) { val relay = relaysList.find { it.url == relayUrl } ?: return + desiredConnected.set(true) val liveToken = liveLocationRelayTokens[relayUrl] ?.takeIf { relayUrl !in nonLiveRelayUrls } if (!isNetworkActionAllowed(liveToken)) return @@ -561,8 +617,8 @@ class NostrRelayManager private constructor() { relay.nextReconnectTime = null // Disconnect if connected - connections[relayUrl]?.close(1000, "Manual retry") - connections.remove(relayUrl) + reconnectJobs.remove(relayUrl)?.cancel() + connections.remove(relayUrl)?.close(1000, "Manual retry") // Attempt immediate reconnection scope.launch { @@ -575,6 +631,7 @@ class NostrRelayManager private constructor() { * This will automatically restore all subscriptions when reconnected */ fun resetAllConnections() { + val shouldReconnect = desiredConnected.get() disconnect() // Reset all relay states @@ -584,8 +641,8 @@ class NostrRelayManager private constructor() { relay.lastError = null } - // Reconnect - subscriptions will be automatically restored in onOpen - connect() + // Reconnect only when connectivity was desired before the controlled reset. + if (shouldReconnect) connect() } /** @@ -615,9 +672,7 @@ class NostrRelayManager private constructor() { geohashToRelays.clear() // Clear any queued messages waiting to be sent - synchronized(messageQueueLock) { - messageQueue.clear() - } + messageQueue.clear() Log.i(TAG, "Cleared all Nostr subscriptions and routing caches") } catch (e: Exception) { @@ -709,35 +764,52 @@ class NostrRelayManager private constructor() { stopSubscriptionValidation() // Stop any existing validation subscriptionValidationJob = scope.launch { - while (isActive) { - delay(SUBSCRIPTION_VALIDATION_INTERVAL) - - try { - val report = validateSubscriptionConsistency() - if (!report.isConsistent && report.connectedRelayCount > 0) { - Log.w(TAG, "Nostr subscription inconsistencies detected") - - // Auto-repair: re-establish subscriptions for relays with missing ones - connections.forEach { (relayUrl, webSocket) -> - val currentSubs = subscriptions[relayUrl] ?: emptySet() - val expectedSubs = activeSubscriptions.keys.filter { subId -> - val subInfo = activeSubscriptions[subId] - subInfo?.targetRelayUrls == null || subInfo.targetRelayUrls.contains(relayUrl) - }.toSet() - - val missingSubs = expectedSubs - currentSubs - if (missingSubs.isNotEmpty()) { - Log.i(TAG, "Auto-repairing ${missingSubs.size} missing subscriptions") - restoreSubscriptionsForRelay(relayUrl, webSocket) - } - } - } - } catch (e: Exception) { - Log.e(TAG, "Error during subscription validation: ${e.message}") + val manager = powerManager + if (manager == null) { + runSubscriptionValidationLoop( + com.bitchat.android.util.AppConstants.Nostr + .SUBSCRIPTION_VALIDATION_INTERVAL_MS + ) + return@launch + } + + manager.profile + .map { it.nostr.subscriptionValidationMs } + .distinctUntilChanged() + .collectLatest(::runSubscriptionValidationLoop) + } + } + + private suspend fun runSubscriptionValidationLoop(intervalMs: Long) { + while (currentCoroutineContext().isActive && desiredConnected.get()) { + delay(intervalMs) + if (!desiredConnected.get()) break + validateAndRepairSubscriptions() + } + } + + private fun validateAndRepairSubscriptions() { + try { + val report = validateSubscriptionConsistency() + if (report.isConsistent || report.connectedRelayCount == 0) return + + Log.w(TAG, "Nostr subscription inconsistencies detected") + connections.forEach { (relayUrl, webSocket) -> + val currentSubs = subscriptions[relayUrl] ?: emptySet() + val expectedSubs = activeSubscriptions.keys.filter { subId -> + val subInfo = activeSubscriptions[subId] + subInfo?.targetRelayUrls == null || + subInfo.targetRelayUrls.contains(relayUrl) + }.toSet() + + if ((expectedSubs - currentSubs).isNotEmpty()) { + Log.i(TAG, "Auto-repairing missing subscriptions") + restoreSubscriptionsForRelay(relayUrl, webSocket) } } + } catch (e: Exception) { + Log.e(TAG, "Error during subscription validation: ${e.message}") } - } /** @@ -754,6 +826,7 @@ class NostrRelayManager private constructor() { urlString: String, liveLocationToken: Long? = null ) { + if (!desiredConnected.get()) return val connectionToken = liveLocationToken ?.takeIf { urlString !in nonLiveRelayUrls } if (!isNetworkActionAllowed(connectionToken)) return @@ -767,17 +840,25 @@ class NostrRelayManager private constructor() { .url(urlString) .build() - runNetworkAction(connectionToken) { + val started = runNetworkAction(connectionToken) { val webSocket = httpClient.newWebSocket( request, RelayWebSocketListener(urlString, connectionToken) ) - connections[urlString] = webSocket + val existing = connections.putIfAbsent(urlString, webSocket) + when { + existing != null -> webSocket.close(1000, "Duplicate connection") + !desiredConnected.get() -> { + connections.remove(urlString, webSocket) + webSocket.close(1000, "Connection no longer desired") + } + } } + if (!started) return } catch (e: Exception) { Log.e(TAG, "Failed to create WebSocket connection") - handleDisconnection(urlString, e, liveLocationToken) + handleConnectionCreationFailure(urlString, e, connectionToken) } } @@ -786,9 +867,9 @@ class NostrRelayManager private constructor() { webSocket: WebSocket, relayUrl: String, liveLocationToken: Long? = null - ) { - if (!isNetworkActionAllowed(liveLocationToken)) return - try { + ): Boolean { + if (!isNetworkActionAllowed(liveLocationToken)) return false + return try { val request = NostrRequest.Event(event) val message = gson.toJson(request, NostrRequest::class.java) @@ -798,17 +879,21 @@ class NostrRelayManager private constructor() { } if (success) { // Update relay stats - val relay = relaysList.find { it.url == relayUrl } - relay?.messagesSent = (relay?.messagesSent ?: 0) + 1 + relaysList.find { it.url == relayUrl }?.let { relay -> + relay.messagesSent += 1 + } updateRelaysList() + true } else { Log.e(TAG, "Failed to send event: WebSocket send failed") + false } } catch (e: Exception) { Log.e(TAG, "Failed to send event") + false } } - + private fun handleMessage(message: String, relayUrl: String) { try { val jsonElement = JsonParser.parseString(message) @@ -822,8 +907,9 @@ class NostrRelayManager private constructor() { when (response) { is NostrResponse.Event -> { // Update relay stats - val relay = relaysList.find { it.url == relayUrl } - relay?.messagesReceived = (relay?.messagesReceived ?: 0) + 1 + relaysList.find { it.url == relayUrl }?.let { relay -> + relay.messagesReceived += 1 + } updateRelaysList() // CLIENT-SIDE FILTER ENFORCEMENT: Ensure this event matches the subscription's filter @@ -882,17 +968,38 @@ class NostrRelayManager private constructor() { private fun handleDisconnection( relayUrl: String, + webSocket: WebSocket, error: Throwable, liveLocationToken: Long? = null + ) { + // Ignore callbacks from intentionally closed or replaced sockets. They must not remove a + // newer socket or schedule a reconnect after a controlled disconnect/privacy revocation. + if (!connections.remove(relayUrl, webSocket)) return + subscriptions.remove(relayUrl) + handleCurrentDisconnection(relayUrl, error, liveLocationToken) + } + + private fun handleConnectionCreationFailure( + relayUrl: String, + error: Throwable, + liveLocationToken: Long? + ) { + if (!desiredConnected.get()) return + handleCurrentDisconnection(relayUrl, error, liveLocationToken) + } + + private fun handleCurrentDisconnection( + relayUrl: String, + error: Throwable, + liveLocationToken: Long? ) { val connectionToken = liveLocationToken ?.takeIf { relayUrl !in nonLiveRelayUrls } - connections.remove(relayUrl) - // NOTE: Don't remove subscriptions here - keep them for restoration on reconnection - // subscriptions.remove(relayUrl) // REMOVED - this was causing subscription loss - + updateRelayStatus(relayUrl, false, error) - if (!isNetworkActionAllowed(connectionToken)) return + if (!desiredConnected.get() || + !isNetworkActionAllowed(connectionToken) + ) return // Check if this is a DNS error val errorMessage = error.message?.lowercase() ?: "" @@ -927,13 +1034,19 @@ class NostrRelayManager private constructor() { Log.d(TAG, "Scheduling Nostr relay reconnection") - // Schedule reconnection - scope.launch { + reconnectJobs.remove(relayUrl)?.cancel() + val reconnectJob = scope.launch { delay(backoffInterval) - if (isNetworkActionAllowed(connectionToken)) { + if (desiredConnected.get() && + isNetworkActionAllowed(connectionToken) + ) { connectToRelay(relayUrl, connectionToken) } } + reconnectJobs[relayUrl] = reconnectJob + reconnectJob.invokeOnCompletion { + reconnectJobs.remove(relayUrl, reconnectJob) + } } private fun updateRelayStatus(url: String, isConnected: Boolean, error: Throwable? = null) { @@ -1014,36 +1127,38 @@ class NostrRelayManager private constructor() { ) : WebSocketListener() { override fun onOpen(webSocket: WebSocket, response: Response) { - if (!isNetworkActionAllowed(liveLocationToken)) { - connections.remove(relayUrl) - webSocket.cancel() + if (!desiredConnected.get() || + connections[relayUrl] !== webSocket || + !isNetworkActionAllowed(liveLocationToken) + ) { + connections.remove(relayUrl, webSocket) + webSocket.close(1000, "Stale connection") return } + reconnectJobs.remove(relayUrl)?.cancel() updateRelayStatus(relayUrl, true) // Restore all active subscriptions for this relay restoreSubscriptionsForRelay(relayUrl, webSocket) - // Process any queued messages for this relay - synchronized(messageQueueLock) { - val iterator = messageQueue.iterator() - while (iterator.hasNext()) { - val queued = iterator.next() - if (relayUrl in queued.targetRelays && - isNetworkActionAllowed(queued.liveLocationToken) - ) { - sendToRelay( - queued.event, - webSocket, - relayUrl, - queued.liveLocationToken - ) - } + // Process only events still pending for this relay, outside the queue lock. + val queuedForRelay = messageQueue.pendingForRelay(relayUrl) + .filter { isNetworkActionAllowed(it.liveLocationToken) } + queuedForRelay.forEach { delivery -> + if (sendToRelay( + delivery.event, + webSocket, + relayUrl, + delivery.liveLocationToken + ) + ) { + messageQueue.markDelivered(delivery.queueId, relayUrl) } } } override fun onMessage(webSocket: WebSocket, text: String) { + if (connections[relayUrl] !== webSocket) return handleMessage(text, relayUrl) } @@ -1053,12 +1168,12 @@ class NostrRelayManager private constructor() { override fun onClosed(webSocket: WebSocket, code: Int, reason: String) { val error = Exception("WebSocket closed: $code $reason") - handleDisconnection(relayUrl, error, liveLocationToken) + handleDisconnection(relayUrl, webSocket, error, liveLocationToken) } override fun onFailure(webSocket: WebSocket, t: Throwable, response: Response?) { Log.e(TAG, "Nostr WebSocket failure") - handleDisconnection(relayUrl, t, liveLocationToken) + handleDisconnection(relayUrl, webSocket, t, liveLocationToken) } } } diff --git a/app/src/main/java/com/bitchat/android/nostr/NostrSubscriptionManager.kt b/app/src/main/java/com/bitchat/android/nostr/NostrSubscriptionManager.kt index 4b8365a1..19ebeb9c 100644 --- a/app/src/main/java/com/bitchat/android/nostr/NostrSubscriptionManager.kt +++ b/app/src/main/java/com/bitchat/android/nostr/NostrSubscriptionManager.kt @@ -3,29 +3,49 @@ package com.bitchat.android.nostr import android.app.Application import android.util.Log import com.bitchat.android.geohash.LiveLocationPrivacyGate -import kotlinx.coroutines.CoroutineScope -import kotlinx.coroutines.launch /** * NostrSubscriptionManager - * - Encapsulates subscription lifecycle with NostrRelayManager + * - Encapsulates ordered subscription lifecycle with NostrRelayManager. + * + * Relay-manager operations are already non-blocking and schedule network I/O on their own scope. + * Keeping this facade synchronous prevents lifecycle cancellation from dropping unsubscribe work + * and guarantees that a channel switch closes the old subscription before opening the new one. */ class NostrSubscriptionManager( private val application: Application, - private val scope: CoroutineScope + private val owner: String = NostrRelayManager.OWNER_LEGACY ) { companion object { private const val TAG = "NostrSubscriptionManager" } private val relayManager get() = NostrRelayManager.getInstance(application) - fun connect() = scope.launch { runCatching { relayManager.connect() }.onFailure { Log.e(TAG, "connect failed: ${it.message}") } } - fun disconnect() = scope.launch { runCatching { relayManager.disconnect() }.onFailure { Log.e(TAG, "disconnect failed: ${it.message}") } } + fun connect() { + runCatching { relayManager.connect() } + .onFailure { Log.e(TAG, "connect failed: ${it.message}") } + } - fun subscribeGiftWraps(pubkey: String, sinceMs: Long, id: String, handler: (NostrEvent) -> Unit) { - scope.launch { - val filter = NostrFilter.giftWrapsFor(pubkey, sinceMs) - relayManager.subscribe(filter, id, handler) - } + fun disconnect() { + runCatching { relayManager.disconnect() } + .onFailure { Log.e(TAG, "disconnect failed: ${it.message}") } + } + + fun subscribeGiftWraps( + pubkey: String, + sinceMs: Long, + id: String, + handler: (NostrEvent) -> Unit, + liveLocationToken: Long? = null + ) { + if (!isAllowed(liveLocationToken)) return + val filter = NostrFilter.giftWrapsFor(pubkey, sinceMs) + relayManager.subscribe( + filter = filter, + id = id, + handler = handler, + owner = owner, + liveLocationToken = liveLocationToken + ) } /** Subscribe to geohash chat messages only (kind 20000) — low-volume, kept alive in background. */ @@ -37,21 +57,18 @@ class NostrSubscriptionManager( handler: (NostrEvent) -> Unit, liveLocationToken: Long? = null ) { - scope.launch { - if (liveLocationToken != null && - !LiveLocationPrivacyGate.accepts(liveLocationToken) - ) return@launch - val filter = NostrFilter.geohashMessages(geohash, sinceMs, limit) - relayManager.subscribeForGeohash( - geohash, - filter, - id, - handler, - includeDefaults = false, - nRelays = 5, - liveLocationToken = liveLocationToken - ) - } + if (!isAllowed(liveLocationToken)) return + val filter = NostrFilter.geohashMessages(geohash, sinceMs, limit) + relayManager.subscribeForGeohash( + geohash, + filter, + id, + handler, + includeDefaults = false, + nRelays = 5, + owner = owner, + liveLocationToken = liveLocationToken + ) } /** Subscribe to geohash presence heartbeats only (kind 20001) — high-volume, paused in background. */ @@ -63,22 +80,28 @@ class NostrSubscriptionManager( handler: (NostrEvent) -> Unit, liveLocationToken: Long? = null ) { - scope.launch { - if (liveLocationToken != null && - !LiveLocationPrivacyGate.accepts(liveLocationToken) - ) return@launch - val filter = NostrFilter.geohashPresence(geohash, sinceMs, limit) - relayManager.subscribeForGeohash( - geohash, - filter, - id, - handler, - includeDefaults = false, - nRelays = 5, - liveLocationToken = liveLocationToken - ) - } + if (!isAllowed(liveLocationToken)) return + val filter = NostrFilter.geohashPresence(geohash, sinceMs, limit) + relayManager.subscribeForGeohash( + geohash, + filter, + id, + handler, + includeDefaults = false, + nRelays = 5, + owner = owner, + liveLocationToken = liveLocationToken + ) } - fun unsubscribe(id: String) { scope.launch { runCatching { relayManager.unsubscribe(id) } } } + fun unsubscribe(id: String) { + runCatching { relayManager.unsubscribe(id) } + } + + fun unsubscribeAllOwned() { + runCatching { relayManager.unsubscribeOwner(owner) } + } + + private fun isAllowed(liveLocationToken: Long?): Boolean = + liveLocationToken == null || LiveLocationPrivacyGate.accepts(liveLocationToken) } diff --git a/app/src/main/java/com/bitchat/android/service/AppShutdownCoordinator.kt b/app/src/main/java/com/bitchat/android/service/AppShutdownCoordinator.kt index f4ca9f74..30a641d0 100644 --- a/app/src/main/java/com/bitchat/android/service/AppShutdownCoordinator.kt +++ b/app/src/main/java/com/bitchat/android/service/AppShutdownCoordinator.kt @@ -56,6 +56,8 @@ object AppShutdownCoordinator { // Stop mesh (best-effort) try { mesh?.stopServices() } catch (_: Exception) { } + try { com.bitchat.android.nostr.NostrRelayManager.shared.disconnect() } catch (_: Exception) { } + try { com.bitchat.android.mesh.PowerManager.getInstance(app).shutdown() } catch (_: Exception) { } // Stop Tor temporarily (do not change user setting) val torProvider = ArtiTorManager.getInstance() diff --git a/app/src/main/java/com/bitchat/android/service/MeshForegroundService.kt b/app/src/main/java/com/bitchat/android/service/MeshForegroundService.kt index 2d84b16f..2f70bf84 100644 --- a/app/src/main/java/com/bitchat/android/service/MeshForegroundService.kt +++ b/app/src/main/java/com/bitchat/android/service/MeshForegroundService.kt @@ -19,8 +19,8 @@ import com.bitchat.android.mesh.BluetoothMeshService import kotlinx.coroutines.CoroutineScope import kotlinx.coroutines.Dispatchers import kotlinx.coroutines.Job -import kotlinx.coroutines.delay -import kotlinx.coroutines.isActive +import kotlinx.coroutines.flow.map +import kotlinx.coroutines.flow.distinctUntilChanged import kotlinx.coroutines.launch class MeshForegroundService : Service() { @@ -112,9 +112,11 @@ class MeshForegroundService : Service() { private val unifiedMeshService: com.bitchat.android.mesh.MeshService? get() = MeshServiceHolder.unifiedMeshService private val serviceJob = Job() - private val scope = CoroutineScope(Dispatchers.Default + serviceJob) + // Service lifecycle callbacks and notification state are main-thread confined. + private val scope = CoroutineScope(Dispatchers.Main.immediate + serviceJob) private var isInForeground: Boolean = false private var isShuttingDown: Boolean = false + private var lastNotifiedPeerCount: Int? = null override fun onCreate() { super.onCreate() @@ -131,6 +133,16 @@ class MeshForegroundService : Service() { MeshServiceHolder.attach(created) } MeshServiceHolder.getUnifiedOrCreate(applicationContext) + + // Notification content is driven by peer-state changes, not a permanent timer. + updateJob = scope.launch { + com.bitchat.android.services.AppStateStore.peers + .map { peers -> peers.distinct().size } + .distinctUntilChanged() + .collect { + if (isInForeground) updateNotification(force = false) + } + } } override fun onStartCommand(intent: Intent?, flags: Int, startId: Int): Int { @@ -176,9 +188,11 @@ class MeshForegroundService : Service() { ACTION_UPDATE_NOTIFICATION -> { // If we became eligible and are not in foreground yet, promote once if (MeshServicePreferences.isBackgroundEnabled(true) && hasAllRequiredPermissions() && !isInForeground) { - val n = buildNotification(getUnifiedActivePeerCount()) + val count = getUnifiedActivePeerCount() + val n = buildNotification(count) startForegroundCompat(n) isInForeground = true + lastNotifiedPeerCount = count } else { updateNotification(force = true) } @@ -191,32 +205,11 @@ class MeshForegroundService : Service() { // Promote exactly once when eligible, otherwise stay background (or stop) if (MeshServicePreferences.isBackgroundEnabled(true) && hasAllRequiredPermissions() && !isInForeground) { - val notification = buildNotification(getUnifiedActivePeerCount()) + val count = getUnifiedActivePeerCount() + val notification = buildNotification(count) startForegroundCompat(notification) isInForeground = true - } - - // Periodically refresh the notification with live network size - if (updateJob == null) { - updateJob = scope.launch { - while (isActive) { - // Retry enabling mesh/foreground once permissions become available - ensureMeshStarted() - val eligible = MeshServicePreferences.isBackgroundEnabled(true) && hasAllRequiredPermissions() - if (eligible) { - // Only update the notification; do not re-call startForeground() - updateNotification(force = false) - } else { - // If disabled or perms missing, ensure we are not in foreground and clear notif - if (isInForeground) { - try { stopForeground(false) } catch (_: Exception) { } - isInForeground = false - } - notificationManager.cancel(NOTIFICATION_ID) - } - delay(5000) - } - } + lastNotifiedPeerCount = count } return START_STICKY @@ -230,19 +223,9 @@ class MeshForegroundService : Service() { android.util.Log.e("MeshForegroundService", "Failed to ensure Wi-Fi Aware transport: ${e.message}") } - val bleEnabled = try { - com.bitchat.android.ui.debug.DebugPreferenceManager.getBleEnabled(true) - } catch (_: Exception) { - true - } - if (!bleEnabled) { - try { meshService?.setBleTransportEnabled(false) } catch (_: Exception) { } - return - } - if (!hasBluetoothPermissions()) return try { android.util.Log.d("MeshForegroundService", "Ensuring mesh service is started") - val service = MeshServiceHolder.getOrCreate(applicationContext) + val service = MeshServiceHolder.getUnifiedOrCreate(applicationContext) service.startServices() } catch (e: Exception) { android.util.Log.e("MeshForegroundService", "Failed to start mesh service: ${e.message}") @@ -255,14 +238,17 @@ class MeshForegroundService : Service() { return } val count = getUnifiedActivePeerCount() - val notification = buildNotification(count) if (MeshServicePreferences.isBackgroundEnabled(true) && hasAllRequiredPermissions()) { - notificationManager.notify(NOTIFICATION_ID, notification) + if (lastNotifiedPeerCount != count) { + notificationManager.notify(NOTIFICATION_ID, buildNotification(count)) + lastNotifiedPeerCount = count + } } else if (force) { // If disabled and forced, make sure to remove any prior foreground state try { stopForeground(false) } catch (_: Exception) { } notificationManager.cancel(NOTIFICATION_ID) isInForeground = false + lastNotifiedPeerCount = null } } diff --git a/app/src/main/java/com/bitchat/android/services/AppStateStore.kt b/app/src/main/java/com/bitchat/android/services/AppStateStore.kt index 503f26a6..e484d21a 100644 --- a/app/src/main/java/com/bitchat/android/services/AppStateStore.kt +++ b/app/src/main/java/com/bitchat/android/services/AppStateStore.kt @@ -17,6 +17,8 @@ object AppStateStore { private val peerIdsByTransport = mutableMapOf>() // Direct (single-hop) peer IDs per transport, used to gossip a unified neighbor set. private val directPeerIdsByTransport = mutableMapOf>() + private val _directPeers = MutableStateFlow>(emptySet()) + val directPeers: StateFlow> = _directPeers.asStateFlow() // Connected peer IDs (mesh ephemeral IDs) private val _peers = MutableStateFlow>(emptyList()) val peers: StateFlow> = _peers.asStateFlow() @@ -29,6 +31,12 @@ object AppStateStore { private val _privateMessages = MutableStateFlow>>(emptyMap()) val privateMessages: StateFlow>> = _privateMessages.asStateFlow() + private val _nickname = MutableStateFlow("") + val nickname: StateFlow = _nickname.asStateFlow() + + private val _selectedPrivateChatPeer = MutableStateFlow(null) + val selectedPrivateChatPeer: StateFlow = _selectedPrivateChatPeer.asStateFlow() + // Channel messages by channel name private val _channelMessages = MutableStateFlow>>(emptyMap()) val channelMessages: StateFlow>> = _channelMessages.asStateFlow() @@ -39,6 +47,14 @@ object AppStateStore { } } + fun setNickname(nickname: String) { + _nickname.value = nickname + } + + fun setSelectedPrivateChatPeer(peerID: String?) { + _selectedPrivateChatPeer.value = peerID + } + fun setTransportPeers(transportId: String, ids: List) { synchronized(this) { peerIdsByTransport[transportId] = ids.toSet() @@ -69,20 +85,25 @@ object AppStateStore { fun setTransportDirectPeers(transportId: String, ids: Collection) { synchronized(this) { directPeerIdsByTransport[transportId] = ids.toSet() + publishDirectPeersLocked() } } fun clearTransportDirectPeers(transportId: String) { synchronized(this) { directPeerIdsByTransport.remove(transportId) + publishDirectPeersLocked() } } /** Union of direct peers across all transports. */ - fun getDirectPeers(): Set { - synchronized(this) { - return directPeerIdsByTransport.values.flatten().toSet() - } + fun getDirectPeers(): Set = _directPeers.value + + private fun publishDirectPeersLocked() { + _directPeers.value = directPeerIdsByTransport.values + .asSequence() + .flatten() + .toSet() } fun addPublicMessage(msg: BitchatMessage) { @@ -212,9 +233,12 @@ object AppStateStore { peerIdsByTransport.clear() directPeerIdsByTransport.clear() _peers.value = emptyList() + _directPeers.value = emptySet() _publicMessages.value = emptyList() _privateMessages.value = emptyMap() _channelMessages.value = emptyMap() + _nickname.value = "" + _selectedPrivateChatPeer.value = null } } diff --git a/app/src/main/java/com/bitchat/android/ui/ChatState.kt b/app/src/main/java/com/bitchat/android/ui/ChatState.kt index 13c754a8..f4504cc2 100644 --- a/app/src/main/java/com/bitchat/android/ui/ChatState.kt +++ b/app/src/main/java/com/bitchat/android/ui/ChatState.kt @@ -216,6 +216,7 @@ class ChatState( fun setNickname(nickname: String) { _nickname.value = nickname + com.bitchat.android.services.AppStateStore.setNickname(nickname) } fun setIsConnected(connected: Boolean) { @@ -228,6 +229,7 @@ class ChatState( fun setSelectedPrivateChatPeer(peerID: String?) { _selectedPrivateChatPeer.value = peerID + com.bitchat.android.services.AppStateStore.setSelectedPrivateChatPeer(peerID) } fun setUnreadPrivateMessages(unread: Set) { diff --git a/app/src/main/java/com/bitchat/android/ui/ChatViewModel.kt b/app/src/main/java/com/bitchat/android/ui/ChatViewModel.kt index 9cf0f614..ba5775a0 100644 --- a/app/src/main/java/com/bitchat/android/ui/ChatViewModel.kt +++ b/app/src/main/java/com/bitchat/android/ui/ChatViewModel.kt @@ -178,8 +178,6 @@ class ChatViewModel( application = application, state = state, messageManager = messageManager, - privateChatManager = privateChatManager, - meshDelegateHandler = meshDelegateHandler, dataManager = dataManager, notificationManager = notificationManager ) @@ -445,6 +443,8 @@ class ChatViewModel( } override fun onCleared() { + geohashViewModel.shutdownUiSubscriptions() + com.bitchat.android.services.AppStateStore.setSelectedPrivateChatPeer(null) super.onCleared() // Note: Mesh service lifecycle is now managed by MainActivity } diff --git a/app/src/main/java/com/bitchat/android/ui/DataManager.kt b/app/src/main/java/com/bitchat/android/ui/DataManager.kt index 4f37d47a..8517ec15 100644 --- a/app/src/main/java/com/bitchat/android/ui/DataManager.kt +++ b/app/src/main/java/com/bitchat/android/ui/DataManager.kt @@ -198,8 +198,10 @@ class DataManager(private val context: Context) { // MARK: - Blocked Users Management + @Synchronized fun loadBlockedUsers() { val savedBlockedUsers = prefs.getStringSet("blocked_users", emptySet()) ?: emptySet() + _blockedUsers.clear() _blockedUsers.addAll(savedBlockedUsers) } @@ -217,6 +219,7 @@ class DataManager(private val context: Context) { saveBlockedUsers() } + @Synchronized fun isUserBlocked(fingerprint: String): Boolean { return _blockedUsers.contains(fingerprint) } @@ -226,8 +229,10 @@ class DataManager(private val context: Context) { private val _geohashBlockedUsers = mutableSetOf() // Set of nostr pubkey hex val geohashBlockedUsers: Set get() = _geohashBlockedUsers.toSet() + @Synchronized fun loadGeohashBlockedUsers() { val savedGeohashBlockedUsers = prefs.getStringSet("geohash_blocked_users", emptySet()) ?: emptySet() + _geohashBlockedUsers.clear() _geohashBlockedUsers.addAll(savedGeohashBlockedUsers) } @@ -245,6 +250,7 @@ class DataManager(private val context: Context) { saveGeohashBlockedUsers() } + @Synchronized fun isGeohashUserBlocked(pubkeyHex: String): Boolean { return _geohashBlockedUsers.contains(pubkeyHex) } diff --git a/app/src/main/java/com/bitchat/android/ui/GeohashViewModel.kt b/app/src/main/java/com/bitchat/android/ui/GeohashViewModel.kt index 8a2e31d8..d2b706ec 100644 --- a/app/src/main/java/com/bitchat/android/ui/GeohashViewModel.kt +++ b/app/src/main/java/com/bitchat/android/ui/GeohashViewModel.kt @@ -12,7 +12,7 @@ import com.bitchat.android.geohash.GeohashNostrPrivacyPolicy import com.bitchat.android.geohash.LiveLocationPrivacyGate import com.bitchat.android.nostr.GeohashMessageHandler import com.bitchat.android.nostr.GeohashRepository -import com.bitchat.android.nostr.NostrDirectMessageHandler +import com.bitchat.android.nostr.NostrBackgroundRuntime import com.bitchat.android.nostr.NostrIdentityBridge import com.bitchat.android.nostr.NostrProtocol import com.bitchat.android.nostr.NostrRelayManager @@ -20,69 +20,48 @@ import com.bitchat.android.nostr.NostrSubscriptionManager import com.bitchat.android.nostr.PoWPreferenceManager import com.bitchat.android.nostr.GeohashAliasRegistry import com.bitchat.android.nostr.GeohashConversationRegistry -import kotlinx.coroutines.CancellationException import kotlinx.coroutines.Job import kotlinx.coroutines.delay import kotlinx.coroutines.flow.StateFlow -import kotlinx.coroutines.flow.combine import kotlinx.coroutines.launch import java.util.Date -import kotlinx.coroutines.flow.collectLatest import kotlinx.coroutines.isActive -import kotlinx.coroutines.Dispatchers -import java.security.SecureRandom import java.util.UUID -import kotlin.random.asKotlinRandom class GeohashViewModel( application: Application, private val state: ChatState, private val messageManager: MessageManager, - private val privateChatManager: PrivateChatManager, - private val meshDelegateHandler: MeshDelegateHandler, private val dataManager: DataManager, private val notificationManager: NotificationManager ) : AndroidViewModel(application), DefaultLifecycleObserver { - companion object { - private const val TAG = "GeohashViewModel" - private val secureRandom = SecureRandom().asKotlinRandom() - } + companion object { private const val TAG = "GeohashViewModel" } private val repo = GeohashRepository(application, state, dataManager) - private val subscriptionManager = NostrSubscriptionManager(application, viewModelScope) + private val uiSubscriptionOwner = "geohash-ui-${UUID.randomUUID()}" + private val subscriptionManager = NostrSubscriptionManager( + application, + owner = uiSubscriptionOwner + ) private val geohashMessageHandler = GeohashMessageHandler( application = application, - state = state, - messageManager = messageManager, repo = repo, scope = viewModelScope, - dataManager = dataManager - ) - private val dmHandler = NostrDirectMessageHandler( - application = application, - state = state, - privateChatManager = privateChatManager, - meshDelegateHandler = meshDelegateHandler, - scope = viewModelScope, - repo = repo, - dataManager = dataManager + dataManager = dataManager, + addChannelMessage = messageManager::addChannelMessage ) - // Live channel message stream (kind 20000). Low-volume; kept alive in the background. - private var currentGeohashMsgSubId: String? = null // Presence heartbeat firehose (kind 20001). High-volume; paused while backgrounded. private var currentGeohashPresenceSubId: String? = null - private var currentDmSubId: String? = null - private var currentDmGeohash: String? = null private var geoTimer: Job? = null - private var globalPresenceJob: Job? = null private var locationChannelManager: com.bitchat.android.geohash.LocationChannelManager? = null private val activeSamplingGeohashes = mutableSetOf() private val samplingSubscriptionIds = mutableMapOf() private val liveSamplingSubscriptionGeohashes = mutableSetOf() private var requestedLiveSamplingGeohashes: Set = emptySet() private var requestedUserSamplingGeohashes: Set = emptySet() + private var uiSubscriptionsShutdown = false private val liveLocationRevocationListener: () -> Unit = { val revokedLiveGeohashes = liveSamplingSubscriptionGeohashes.toSet() revokedLiveGeohashes.forEach { geohash -> @@ -105,21 +84,10 @@ class GeohashViewModel( } fun initialize() { - subscriptionManager.connect() // Observe process lifecycle to manage background sampling kotlin.runCatching { ProcessLifecycleOwner.get().lifecycle.addObserver(this) } - val identity = NostrIdentityBridge.getCurrentNostrIdentity(getApplication()) - if (identity != null) { - // Use global chat-messages only for full account DMs (mesh context). For geohash DMs, subscribe per-geohash below. - subscriptionManager.subscribeGiftWraps( - pubkey = identity.publicKeyHex, - sinceMs = System.currentTimeMillis() - 172800000L, - id = "chat-messages", - handler = { event -> dmHandler.onGiftWrap(event, "", identity) } // geohash="" means global account DM (not geohash identity) - ) - } try { locationChannelManager = com.bitchat.android.geohash.LocationChannelManager.getInstance(getApplication()) viewModelScope.launch { @@ -133,10 +101,6 @@ class GeohashViewModel( state.setIsTeleported(teleported) } } - - // Start global presence heartbeat loop - startGlobalPresenceHeartbeat() - } catch (e: Exception) { Log.e(TAG, "Failed to initialize location channel state: ${e.message}") state.setSelectedLocationChannel(com.bitchat.android.geohash.ChannelID.Mesh) @@ -144,102 +108,17 @@ class GeohashViewModel( } } - private fun startGlobalPresenceHeartbeat() { - globalPresenceJob?.cancel() - globalPresenceJob = viewModelScope.launch(kotlinx.coroutines.Dispatchers.IO) { - val manager = locationChannelManager ?: return@launch - combine( - manager.availableChannels, - LiveLocationPrivacyGate.enabled - ) { channels, enabled -> - GeohashNostrPrivacyPolicy.livePresenceTargets(channels, enabled) - }.collectLatest { targetGeohashes -> - - if (targetGeohashes.isNotEmpty()) { - // Enter heartbeat loop for this set of channels - // If channels change (e.g. user moves), collectLatest cancels this loop and starts a new one immediately - while (true) { - // Randomize loop interval (40-80s, average 60s) - val loopInterval = secureRandom.nextLong(40000L, 80000L) - var timeSpent = 0L - - try { - Log.v(TAG, "💓 Broadcasting global presence to ${targetGeohashes.size} channels") - targetGeohashes.forEach { geohash -> - // Decorrelate individual broadcasts with random delay (1s-5s) - val stepDelay = secureRandom.nextLong(1000L, 10000L) - delay(stepDelay) - timeSpent += stepDelay - - broadcastLiveLocationPresence(geohash) - } - } catch (e: CancellationException) { - throw e - } catch (e: Exception) { - Log.w(TAG, "Global presence heartbeat error: ${e.message}") - } - - // Wait remaining time to satisfy target average cadence - val remaining = loopInterval - timeSpent - if (remaining > 0) { - delay(remaining) - } else { - delay(10000L) // Minimum guard delay - } - } - } - } - } - } - fun panicReset() { repo.clearAll() GeohashAliasRegistry.clear() GeohashConversationRegistry.clear() - subscriptionManager.disconnect() - currentGeohashMsgSubId = null + subscriptionManager.unsubscribeAllOwned() currentGeohashPresenceSubId = null - currentDmSubId = null - currentDmGeohash = null activeChannelGeohash = null geoTimer?.cancel() geoTimer = null - globalPresenceJob?.cancel() - globalPresenceJob = null try { NostrIdentityBridge.clearAllAssociations(getApplication()) } catch (_: Exception) {} - initialize() - } - - private suspend fun broadcastLiveLocationPresence(geohash: String) { - val manager = locationChannelManager ?: return - val token = LiveLocationPrivacyGate.captureToken() ?: return - val isCurrentLiveTarget = GeohashNostrPrivacyPolicy.livePresenceTargets( - manager.availableChannels.value, - liveLocationEnabled = true - ).contains(geohash) - if (!isCurrentLiveTarget || !LiveLocationPrivacyGate.accepts(token)) return - - try { - var identity: com.bitchat.android.nostr.NostrIdentity? = null - LiveLocationPrivacyGate.runIfAllowed(token) { - identity = NostrIdentityBridge.deriveIdentity(geohash, getApplication()) - } - val preparedIdentity = identity ?: return - if (!LiveLocationPrivacyGate.accepts(token)) return - val event = NostrProtocol.createGeohashPresenceEvent(geohash, preparedIdentity) - LiveLocationPrivacyGate.runIfAllowed(token) { - val relayManager = NostrRelayManager.getInstance(getApplication()) - relayManager.sendEventToGeohash( - event, - geohash, - includeDefaults = false, - nRelays = 5, - liveLocationToken = token - ) - } - } catch (e: Exception) { - Log.w(TAG, "Failed to send live-location presence") - } + NostrBackgroundRuntime.resetSubscriptions() } fun sendGeohashMessage(content: String, channel: com.bitchat.android.geohash.GeohashChannel, myPeerID: String, nickname: String?) { @@ -424,27 +303,36 @@ class GeohashViewModel( } fun ensureGeohashDMSubscriptionForConversation(conversationKey: String) { - val geohash = repo.getConversationGeohash(conversationKey) ?: return - if (currentDmGeohash == geohash && currentDmSubId != null) return - - currentDmSubId?.let(subscriptionManager::unsubscribe) - currentDmSubId = null - currentDmGeohash = null - subscribeChannelDM(geohash) + val geohash = repo.getConversationGeohash(conversationKey) + ?: NostrBackgroundRuntime.conversationGeohash(conversationKey) + ?: return + NostrBackgroundRuntime.ensureConversationDm(geohash) } - fun displayNameForNostrPubkeyUI(pubkeyHex: String): String = repo.displayNameForNostrPubkeyUI(pubkeyHex) - fun displayNameForGeohashConversation(pubkeyHex: String, sourceGeohash: String): String = repo.displayNameForGeohashConversation(pubkeyHex, sourceGeohash) + fun displayNameForNostrPubkeyUI(pubkeyHex: String): String { + val foregroundName = repo.displayNameForNostrPubkeyUI(pubkeyHex) + return foregroundName.takeUnless { it == "anon" } + ?: NostrBackgroundRuntime.displayNameForNostrPubkey(pubkeyHex) + ?: foregroundName + } + + fun displayNameForGeohashConversation(pubkeyHex: String, sourceGeohash: String): String { + val foregroundName = repo.displayNameForGeohashConversation(pubkeyHex, sourceGeohash) + return foregroundName.takeUnless { it == "anon" } + ?: NostrBackgroundRuntime.displayNameForGeohashConversation(pubkeyHex, sourceGeohash) + ?: foregroundName + } + + fun conversationGeohash(conversationKey: String): String? = + repo.getConversationGeohash(conversationKey) + ?: NostrBackgroundRuntime.conversationGeohash(conversationKey) fun peerIdentityForNostrPubkey(pubkeyHex: String): PeerIdentity = PeerIdentity.nostr(pubkeyHex) private fun switchLocationChannel(channel: com.bitchat.android.geohash.ChannelID?) { geoTimer?.cancel(); geoTimer = null - currentGeohashMsgSubId?.let { subscriptionManager.unsubscribe(it); currentGeohashMsgSubId = null } currentGeohashPresenceSubId?.let { subscriptionManager.unsubscribe(it); currentGeohashPresenceSubId = null } - currentDmSubId?.let { subscriptionManager.unsubscribe(it); currentDmSubId = null } - currentDmGeohash = null when (channel) { is com.bitchat.android.geohash.ChannelID.Mesh -> { @@ -481,17 +369,11 @@ class GeohashViewModel( val liveLocationToken = locationChannelManager ?.liveLocationTokenForSelectedChannel(channel.channel) - // Chat message stream (kind 20000) is low-volume; keep it alive even when - // backgrounded so geohash messages still arrive. - subscribeChannelMessages(channel.channel.geohash, liveLocationToken) // Presence heartbeat firehose (kind 20001) is the high-volume data hog; only // run it in the foreground. It is restored in onStart() and torn down in onStop(). if (isAppInForeground()) { subscribeChannelPresence(channel.channel.geohash, liveLocationToken) } - // Gift-wrap DM subscription is lightweight (filtered to our pubkey) and is - // kept alive in the background so geohash DMs still arrive. - subscribeChannelDM(channel.channel.geohash) } null -> { Log.d(TAG, "📡 No channel selected") @@ -501,25 +383,6 @@ class GeohashViewModel( } } - /** - * Subscribe to the chat message stream (kind 20000) for a geohash channel. - * Low-volume; kept alive in the background so messages keep arriving. - */ - private fun subscribeChannelMessages( - geohash: String, - liveLocationToken: Long? - ) { - val subId = "geohash-${UUID.randomUUID()}"; currentGeohashMsgSubId = subId - subscriptionManager.subscribeGeohashMessages( - geohash = geohash, - sinceMs = System.currentTimeMillis() - 3600000L, - limit = 200, - id = subId, - handler = { event -> geohashMessageHandler.onEvent(event, geohash) }, - liveLocationToken = liveLocationToken - ) - } - /** * Subscribe to the presence heartbeat firehose (kind 20001) for a geohash channel. * High-volume; only used to refresh the participant list, so it is torn down in @@ -540,25 +403,6 @@ class GeohashViewModel( ) } - /** - * Subscribe to gift-wrap DMs for a geohash channel's derived identity. - * Lightweight (filtered to our pubkey); kept alive in the background. - */ - private fun subscribeChannelDM(geohash: String) { - val dmIdentity = NostrIdentityBridge.deriveIdentity(geohash, getApplication()) - val dmSubId = "geo-dm-${UUID.randomUUID()}" - currentDmSubId = dmSubId - currentDmGeohash = geohash - subscriptionManager.subscribeGiftWraps( - pubkey = dmIdentity.publicKeyHex, - sinceMs = System.currentTimeMillis() - 172800000L, - id = dmSubId, - handler = { event -> dmHandler.onGiftWrap(event, geohash, dmIdentity) } - ) - // Also register alias in global registry for routing convenience - GeohashAliasRegistry.put("nostr_${dmIdentity.publicKeyHex.take(16)}", dmIdentity.publicKeyHex) - } - private fun startGeoParticipantsTimer() { geoTimer = viewModelScope.launch { while (repo.getCurrentGeohash() != null) { @@ -569,7 +413,16 @@ class GeohashViewModel( } override fun onCleared() { + shutdownUiSubscriptions() super.onCleared() + } + + fun shutdownUiSubscriptions() { + if (uiSubscriptionsShutdown) return + uiSubscriptionsShutdown = true + subscriptionManager.unsubscribeAllOwned() + geoTimer?.cancel() + geoTimer = null kotlin.runCatching { ProcessLifecycleOwner.get().lifecycle.removeObserver(this) } @@ -602,10 +455,6 @@ class GeohashViewModel( if (repo.getCurrentGeohash() != null && geoTimer?.isActive != true) { startGeoParticipantsTimer() } - // Resume the global presence heartbeat - if (globalPresenceJob?.isActive != true) { - startGlobalPresenceHeartbeat() - } } override fun onStop(owner: LifecycleOwner) { @@ -615,12 +464,9 @@ class GeohashViewModel( currentGeohashPresenceSubId?.let { subscriptionManager.unsubscribe(it); currentGeohashPresenceSubId = null } // Drop geohash sampling subscriptions activeSamplingGeohashes.forEach(::unsubscribeSampling) - // Stop broadcasting presence heartbeats - globalPresenceJob?.cancel(); globalPresenceJob = null // Stop participant-refresh polling geoTimer?.cancel(); geoTimer = null - // NOTE: gift-wrap DM subscriptions (per-geohash + global "chat-messages") are intentionally - // left active so direct messages still arrive while backgrounded. + // Process-owned message and DM subscriptions remain active. } private fun performSubscribeSampling(geohash: String) { diff --git a/app/src/main/java/com/bitchat/android/ui/PeerColorSeed.kt b/app/src/main/java/com/bitchat/android/ui/PeerColorSeed.kt new file mode 100644 index 00000000..b6eac7cd --- /dev/null +++ b/app/src/main/java/com/bitchat/android/ui/PeerColorSeed.kt @@ -0,0 +1,33 @@ +package com.bitchat.android.ui + +import com.bitchat.android.model.BitchatMessage +import java.util.Locale + +/** + * Stable, presentation-neutral identity used to derive a peer hue. + * + * ViewModels may expose this value, but only the UI theme resolves it to a rendered color. + */ +@JvmInline +value class PeerColorSeed(val value: String) + +fun meshPeerColorSeed(peerID: String): PeerColorSeed = + PeerColorSeed("noise:${peerID.lowercase(Locale.ROOT)}") + +fun nostrPeerColorSeed(pubkeyHex: String): PeerColorSeed = + PeerColorSeed("nostr:${pubkeyHex.lowercase(Locale.ROOT)}") + +fun peerColorSeedForMessage(message: BitchatMessage): PeerColorSeed { + val value = when { + message.senderPeerID?.startsWith("nostr:") == true || + message.senderPeerID?.startsWith("nostr_") == true -> { + "nostr:${message.senderPeerID.lowercase(Locale.ROOT)}" + } + message.senderPeerID?.length == 16 || message.senderPeerID?.length == 64 -> { + "noise:${message.senderPeerID.lowercase(Locale.ROOT)}" + } + else -> message.sender.lowercase(Locale.ROOT) + } + + return PeerColorSeed(value) +} diff --git a/app/src/main/java/com/bitchat/android/ui/PrivateChatManager.kt b/app/src/main/java/com/bitchat/android/ui/PrivateChatManager.kt index bd1b51bf..ee28a534 100644 --- a/app/src/main/java/com/bitchat/android/ui/PrivateChatManager.kt +++ b/app/src/main/java/com/bitchat/android/ui/PrivateChatManager.kt @@ -34,6 +34,7 @@ class PrivateChatManager( private val messageManager: MessageManager, private val dataManager: DataManager, private val noiseSessionDelegate: NoiseSessionDelegate, + private val trackUnreadMessages: Boolean = true, private val hasReadReceiptBeenSent: (messageID: String) -> Boolean = { false }, private val markMessageReadLocally: (messageID: String) -> Unit = {} ) { @@ -343,7 +344,7 @@ class PrivateChatManager( // Nostr messages originate here and must be added explicitly, even after their // sender alias has canonicalized to a contact_* conversation ID. if (origin == PrivateMessageOrigin.NOSTR) { - if (suppressUnread) { + if (suppressUnread || !trackUnreadMessages) { messageManager.addPrivateMessageNoUnread(conversationID, message) } else { messageManager.addPrivateMessage(conversationID, message) @@ -351,7 +352,10 @@ class PrivateChatManager( } // Track as unread for read receipt purposes if not focused - if (!suppressUnread && state.getSelectedPrivateChatPeerValue() != conversationID) { + if (trackUnreadMessages && + !suppressUnread && + state.getSelectedPrivateChatPeerValue() != conversationID + ) { val unreadList = unreadReceivedMessages.getOrPut(conversationID) { mutableListOf() } unreadList.add(message) Log.d(TAG, "Queued unread from $conversationID (count=${unreadList.size})") diff --git a/app/src/main/java/com/bitchat/android/util/AppConstants.kt b/app/src/main/java/com/bitchat/android/util/AppConstants.kt index ad2325ed..0210f58c 100644 --- a/app/src/main/java/com/bitchat/android/util/AppConstants.kt +++ b/app/src/main/java/com/bitchat/android/util/AppConstants.kt @@ -92,7 +92,7 @@ object AppConstants { const val SCAN_OFF_DURATION_ULTRA_LOW_MS: Long = 29_000L const val MAX_CONNECTIONS_NORMAL: Int = 8 const val MAX_CONNECTIONS_POWER_SAVE: Int = 8 - const val MAX_CONNECTIONS_ULTRA_LOW: Int = 4 + const val MAX_CONNECTIONS_ULTRA_LOW: Int = 8 } object Nostr { diff --git a/app/src/main/java/com/bitchat/android/wifi-aware/SyncedSocket.kt b/app/src/main/java/com/bitchat/android/wifi-aware/SyncedSocket.kt index c3039327..9688f04f 100644 --- a/app/src/main/java/com/bitchat/android/wifi-aware/SyncedSocket.kt +++ b/app/src/main/java/com/bitchat/android/wifi-aware/SyncedSocket.kt @@ -24,10 +24,9 @@ class SyncedSocket( private val outputStream: DataOutputStream companion object { - // Both peers exchange keep-alive frames every ~2s while connected, so a read that - // stalls well beyond that means the link is dead (half-open). Time out so the read - // loop can detect it and trigger disconnection instead of blocking forever. - const val DEFAULT_READ_TIMEOUT_MS = 15_000 + // Critical-background peers send at most every 30s. Tolerate three missed frames so a + // faster local profile cannot drop a healthy, slower remote peer merely for being idle. + const val DEFAULT_READ_TIMEOUT_MS = 90_000 } init { diff --git a/app/src/main/java/com/bitchat/android/wifi-aware/WifiAwareController.kt b/app/src/main/java/com/bitchat/android/wifi-aware/WifiAwareController.kt index f2ec4cfa..1c2b7e5f 100644 --- a/app/src/main/java/com/bitchat/android/wifi-aware/WifiAwareController.kt +++ b/app/src/main/java/com/bitchat/android/wifi-aware/WifiAwareController.kt @@ -6,12 +6,10 @@ import android.util.Log import kotlinx.coroutines.CoroutineScope import kotlinx.coroutines.Dispatchers import kotlinx.coroutines.SupervisorJob -import kotlinx.coroutines.cancel import kotlinx.coroutines.delay import kotlinx.coroutines.flow.MutableStateFlow import kotlinx.coroutines.flow.StateFlow import kotlinx.coroutines.flow.asStateFlow -import kotlinx.coroutines.isActive import kotlinx.coroutines.launch import java.util.concurrent.atomic.AtomicBoolean @@ -68,24 +66,20 @@ object WifiAwareController { Log.i(TAG, "Wi-Fi Aware unsupported: ${status.reason}") } setEnabled(enabledByDefault) - // Start background poller for debug surfacing - scope.launch { - while (isActive) { - try { - val s = service - if (s != null) { - _connectedPeers.value = s.getDeviceAddressToPeerMapping() // peerID -> ip - _knownPeers.value = s.getPeerNicknames() - _discoveredPeers.value = s.getDiscoveredPeerIds() - } else { - _connectedPeers.value = emptyMap() - _knownPeers.value = emptyMap() - _discoveredPeers.value = emptySet() - } - } catch (_: Exception) { } - delay(1000) - } - } + } + + internal fun publishDebugSnapshot( + connected: Map, + known: Map, + discovered: Set + ) { + _connectedPeers.value = connected + _knownPeers.value = known + _discoveredPeers.value = discovered + } + + private fun clearDebugSnapshot() { + publishDebugSnapshot(emptyMap(), emptyMap(), emptySet()) } fun setEnabled(value: Boolean) { @@ -202,9 +196,7 @@ object WifiAwareController { } try { stopped?.stopServices() } catch (_: Exception) { } try { com.bitchat.android.services.AppStateStore.clearTransportPeers("WIFI") } catch (_: Exception) { } - _connectedPeers.value = emptyMap() - _knownPeers.value = emptyMap() - _discoveredPeers.value = emptySet() + clearDebugSnapshot() try { com.bitchat.android.ui.debug.DebugSettingsManager.getInstance().addDebugMessage(com.bitchat.android.ui.debug.DebugMessage.SystemMessage("Wi‑Fi Aware stopped")) } catch (_: Exception) {} } @@ -214,9 +206,7 @@ object WifiAwareController { service = null _running.value = false try { com.bitchat.android.services.AppStateStore.clearTransportPeers("WIFI") } catch (_: Exception) { } - _connectedPeers.value = emptyMap() - _knownPeers.value = emptyMap() - _discoveredPeers.value = emptySet() + clearDebugSnapshot() } } diff --git a/app/src/main/java/com/bitchat/android/wifi-aware/WifiAwareMeshService.kt b/app/src/main/java/com/bitchat/android/wifi-aware/WifiAwareMeshService.kt index d2be00e2..85669aba 100644 --- a/app/src/main/java/com/bitchat/android/wifi-aware/WifiAwareMeshService.kt +++ b/app/src/main/java/com/bitchat/android/wifi-aware/WifiAwareMeshService.kt @@ -76,10 +76,6 @@ class WifiAwareMeshService(private val context: Context) : MeshService, Transpor private const val CLIENT_SOCKET_RETRY_DELAY_MS = 750L private const val CLIENT_SOCKET_ATTEMPTS = 3 private const val CLIENT_ROLE_REVERSAL_FAILURES = 3 - // Discovery freshness window for reconnection maintenance - private const val DISCOVERY_STALE_MS = 5L * 60 * 1000 - private const val DISCOVERY_IDLE_REFRESH_MS = 2L * 60 * 1000 - private const val DISCOVERY_SESSION_REFRESH_MIN_INTERVAL_MS = 90L * 1000 private const val ROLE_REVERSAL_PREFIX = "ROLE_SERVER:" } @@ -89,6 +85,7 @@ class WifiAwareMeshService(private val context: Context) : MeshService, Transpor // Peer ID must match BluetoothMeshService: first 16 hex chars of identity fingerprint (8 bytes) override val myPeerID: String = encryptionService.getIdentityFingerprint().take(16) private val serviceScope = CoroutineScope(Dispatchers.IO + SupervisorJob()) + private val powerManager = com.bitchat.android.mesh.PowerManager.getInstance(context.applicationContext) private val wifiTransport = WifiAwareTransport() private lateinit var meshCore: MeshCore private lateinit var fragmentingSender: FragmentingPacketSender @@ -171,6 +168,7 @@ class WifiAwareMeshService(private val context: Context) : MeshService, Transpor hooks = MeshCore.Hooks( onMessageReceived = { message -> handleMessageReceived(message) }, onAnnounceProcessed = { routed, _ -> + publishControllerDebugSnapshot() routed.peerID?.let { pid -> DirectLinkAnnouncementPolicy.observationFor(routed, MAX_TTL) ?.let(::observeDirectIngressLink) @@ -626,6 +624,7 @@ class WifiAwareMeshService(private val context: Context) : MeshService, Transpor val now = System.currentTimeMillis() discoveredTimestamps[peerId] = now lastDiscoveryActivityAt.set(now) + publishControllerDebugSnapshot() } @RequiresApi(Build.VERSION_CODES.Q) @@ -642,7 +641,8 @@ class WifiAwareMeshService(private val context: Context) : MeshService, Transpor if (!com.bitchat.android.wifiaware.WifiAwareController.enabled.value) return false val lastRefresh = lastDiscoveryRefreshAt.get() - if ((now - lastRefresh) < DISCOVERY_SESSION_REFRESH_MIN_INTERVAL_MS) return false + val minRefreshMs = powerManager.profile.value.wifiAware.discoverySessionRefreshMinMs + if ((now - lastRefresh) < minRefreshMs) return false if (!lastDiscoveryRefreshAt.compareAndSet(lastRefresh, now)) return false Log.i(TAG, "Refreshing Wi-Fi Aware discovery sessions ($reason)") @@ -657,7 +657,8 @@ class WifiAwareMeshService(private val context: Context) : MeshService, Transpor serviceScope.launch { while (isActive) { try { - delay(15_000) // Check every 15 seconds + val schedule = powerManager.profile.value.wifiAware + delay(schedule.connectionMaintenanceMs) if (!isActive) break val now = System.currentTimeMillis() @@ -665,18 +666,19 @@ class WifiAwareMeshService(private val context: Context) : MeshService, Transpor // 0. Prune stale discovery entries. PeerHandles become invalid when the // discovery sessions restart, so we must not keep pinging old handles forever. val staleIds = discoveredTimestamps.filter { (id, ts) -> - (now - ts) >= DISCOVERY_STALE_MS && !connectionTracker.isConnected(id) + (now - ts) >= schedule.discoveryStaleMs && !connectionTracker.isConnected(id) }.keys.toSet() if (staleIds.isNotEmpty()) { staleIds.forEach { discoveredTimestamps.remove(it) } handleToPeerId.entries.removeIf { it.value in staleIds } staleIds.forEach { subscribeHandles.remove(it) } staleIds.forEach { publishHandles.remove(it) } + publishControllerDebugSnapshot() } // 1. Identify peers that are discovered (recently seen) but not currently connected val recentDiscovered = discoveredTimestamps.filter { (id, ts) -> - (now - ts) < DISCOVERY_STALE_MS // Seen in last 5 minutes + (now - ts) < schedule.discoveryStaleMs }.keys // 2. Filter out those who are already connected @@ -726,7 +728,7 @@ class WifiAwareMeshService(private val context: Context) : MeshService, Transpor disconnectedPeers.isNotEmpty() && missingUsableHandle && !attemptedReconnect -> { refreshDiscoverySessions("missing peer handle", now) } - recentDiscovered.isEmpty() && idleFor >= DISCOVERY_IDLE_REFRESH_MS -> { + recentDiscovered.isEmpty() && idleFor >= schedule.discoveryIdleRefreshMs -> { refreshDiscoverySessions("idle discovery", now) } } @@ -869,6 +871,7 @@ class WifiAwareMeshService(private val context: Context) : MeshService, Transpor val synced = SyncedSocket(client) activeSocket = synced connectionTracker.onClientConnected(peerId, synced) + publishControllerDebugSnapshot() // We only ever accept a single data socket per server request. Close the // listening ServerSocket now so it can't block a future re-serve (its // presence makes hasOpenServerSocket() true for the life of the process) @@ -936,30 +939,8 @@ class WifiAwareMeshService(private val context: Context) : MeshService, Transpor pubSession: PublishDiscoverySession, peerHandle: PeerHandle ) { - // TCP keep-alive pings - serviceScope.launch { - try { - while (connectionTracker.isConnected(peerId)) { - // write empty byte array effectively sends [4 bytes length=0] which is our ping - try { - client.write(ByteArray(0)) - } catch (_: IOException) { - // The write side is dead. Don't just stop pinging: actively tear down so the - // half-open socket stops counting as "connected" and maintenance can retry. - handlePeerDisconnection(peerId, client) - break - } - delay(2_000) - } - } catch (_: Exception) {} - } - // Discovery keep-alive - serviceScope.launch { - var msgId = 0 - while (connectionTracker.isConnected(peerId)) { - try { pubSession.sendMessage(peerHandle, msgId++, ByteArray(0)) } catch (_: Exception) { break } - delay(20_000) - } + startProfiledKeepAlive(client, peerId) { msgId -> + pubSession.sendMessage(peerHandle, msgId, ByteArray(0)) } } @@ -1127,6 +1108,7 @@ class WifiAwareMeshService(private val context: Context) : MeshService, Transpor val synced = SyncedSocket(sock) activeSocket = synced connectionTracker.onClientConnected(peerId, synced) + publishControllerDebugSnapshot() clientSocketFailures.remove(peerId) try { meshCore.addOrUpdatePeer(peerId, peerId) } catch (_: Exception) {} listenerExec.execute { listenToPeer(synced, peerId) } @@ -1170,28 +1152,50 @@ class WifiAwareMeshService(private val context: Context) : MeshService, Transpor peerId: String, peerHandle: PeerHandle ) { - // TCP keep-alive + startProfiledKeepAlive(sock, peerId) { msgId -> + subscribeSession?.sendMessage(peerHandle, msgId, ByteArray(0)) + } + } + + /** + * One coroutine per peer schedules both TCP and discovery deadlines. This avoids two + * independently waking jobs while still adapting after every send to the latest profile. + */ + private fun startProfiledKeepAlive( + socket: SyncedSocket, + peerId: String, + sendDiscovery: (Int) -> Unit + ) { serviceScope.launch { - try { - while (connectionTracker.isConnected(peerId)) { + var discoveryMessageId = 0 + var nextTcpAt = 0L + var nextDiscoveryAt = 0L + while (connectionTracker.isConnected(peerId)) { + val now = android.os.SystemClock.elapsedRealtime() + val schedule = powerManager.profile.value.wifiAware + + if (now >= nextTcpAt) { try { - sock.write(ByteArray(0)) + socket.write(ByteArray(0)) } catch (_: IOException) { - // The write side is dead. Tear down so the half-open socket stops counting - // as "connected" and maintenance can retry instead of silently stalling. - handlePeerDisconnection(peerId, sock) + handlePeerDisconnection(peerId, socket) break } - delay(2_000) + nextTcpAt = now + schedule.tcpKeepAliveMs } - } catch (_: Exception) {} - } - // Discovery keep-alive - serviceScope.launch { - var msgId = 0 - while (connectionTracker.isConnected(peerId)) { - try { subscribeSession?.sendMessage(peerHandle, msgId++, ByteArray(0)) } catch (_: Exception) { break } - delay(20_000) + + if (now >= nextDiscoveryAt) { + try { + sendDiscovery(discoveryMessageId++) + } catch (_: Exception) { + // The TCP data path remains usable; discovery maintenance will recover its + // session independently without dropping this idle connection. + } + nextDiscoveryAt = now + schedule.discoveryKeepAliveMs + } + + val nextWakeAt = minOf(nextTcpAt, nextDiscoveryAt) + delay((nextWakeAt - android.os.SystemClock.elapsedRealtime()).coerceAtLeast(250L)) } } } @@ -1269,6 +1273,10 @@ class WifiAwareMeshService(private val context: Context) : MeshService, Transpor ) } catch (_: Exception) { } try { meshCore.setDirectConnection(observation.peerID, true) } catch (_: Exception) { } + try { + meshCore.gossipSyncManager.scheduleInitialSyncToPeer(observation.peerID, 1_000) + } catch (_: Exception) { } + publishControllerDebugSnapshot() Log.i( TAG, @@ -1355,6 +1363,7 @@ class WifiAwareMeshService(private val context: Context) : MeshService, Transpor } } // Else: socket replaced or inactive; do not remove peer/session, as a new socket has likely taken over + publishControllerDebugSnapshot() } } @@ -1565,6 +1574,16 @@ class WifiAwareMeshService(private val context: Context) : MeshService, Transpor fun getDiscoveredPeerIds(): Set = (handleToPeerId.values + discoveredTimestamps.keys).filter { it.isNotBlank() }.toSet() + private fun publishControllerDebugSnapshot() { + try { + com.bitchat.android.wifiaware.WifiAwareController.publishDebugSnapshot( + connected = getDeviceAddressToPeerMapping(), + known = getPeerNicknames(), + discovered = getDiscoveredPeerIds() + ) + } catch (_: Exception) { } + } + /** * Utility for logs/UI: pretty-prints one peer-to-address mapping per line. */ diff --git a/app/src/test/kotlin/com/bitchat/android/mesh/PowerProfileResolverTest.kt b/app/src/test/kotlin/com/bitchat/android/mesh/PowerProfileResolverTest.kt new file mode 100644 index 00000000..8d4df343 --- /dev/null +++ b/app/src/test/kotlin/com/bitchat/android/mesh/PowerProfileResolverTest.kt @@ -0,0 +1,102 @@ +package com.bitchat.android.mesh + +import org.junit.Assert.assertEquals +import org.junit.Assert.assertFalse +import org.junit.Assert.assertTrue +import org.junit.Test + +class PowerProfileResolverTest { + + @Test + fun `background normal without peers scans one second per minute`() { + val profile = resolve(battery = 80, background = true, peers = false) + + assertEquals(PowerManager.BatteryBand.NORMAL, profile.batteryBand) + assertEquals(1_000L, profile.ble.scanOnMs) + assertEquals(59_000L, profile.ble.scanOffMs) + assertFalse(profile.ble.continuousScan) + assertEquals(8, profile.ble.maxConnections) + assertEquals(60_000L, profile.meshAnnouncementIntervalMs) + } + + @Test + fun `any direct transport peer accelerates background BLE discovery to thirty seconds`() { + val profile = resolve(battery = 80, background = true, peers = true) + + assertEquals(1_000L, profile.ble.scanOnMs) + assertEquals(29_000L, profile.ble.scanOffMs) + } + + @Test + fun `background low profile uses selected wifi and nostr cadence`() { + val profile = resolve(battery = 20, background = true) + + assertEquals(PowerManager.BatteryBand.LOW, profile.batteryBand) + assertEquals(20_000L, profile.wifiAware.tcpKeepAliveMs) + assertEquals(120_000L, profile.wifiAware.discoveryKeepAliveMs) + assertEquals(120_000L, profile.wifiAware.connectionMaintenanceMs) + assertEquals(600_000L, profile.wifiAware.discoveryIdleRefreshMs) + assertEquals(900_000L, profile.wifiAware.discoveryStaleMs) + assertEquals(600_000L, profile.nostr.presenceHeartbeatMinMs) + assertEquals(600_000L, profile.nostr.subscriptionValidationMs) + assertEquals(120_000L, profile.meshAnnouncementIntervalMs) + } + + @Test + fun `background critical profile uses slowest wifi and nostr cadence without reducing connections`() { + val profile = resolve(battery = 10, background = true, peers = true) + + assertEquals(PowerManager.BatteryBand.CRITICAL, profile.batteryBand) + assertEquals(30_000L, profile.wifiAware.tcpKeepAliveMs) + assertEquals(180_000L, profile.wifiAware.discoveryKeepAliveMs) + assertEquals(180_000L, profile.wifiAware.connectionMaintenanceMs) + assertEquals(900_000L, profile.wifiAware.discoveryIdleRefreshMs) + assertEquals(1_200_000L, profile.wifiAware.discoveryStaleMs) + assertEquals(900_000L, profile.nostr.presenceHeartbeatMinMs) + assertEquals(900_000L, profile.nostr.subscriptionValidationMs) + assertEquals(300_000L, profile.meshAnnouncementIntervalMs) + assertEquals(8, profile.ble.maxConnections) + } + + @Test + fun `foreground normal preserves balanced BLE cadence`() { + val profile = resolve(battery = 80, background = false) + + assertEquals(PowerManager.PowerMode.BALANCED, profile.mode) + assertEquals(8_000L, profile.ble.scanOnMs) + assertEquals(2_000L, profile.ble.scanOffMs) + assertEquals(5_000L, profile.wifiAware.tcpKeepAliveMs) + assertEquals(30_000L, profile.meshAnnouncementIntervalMs) + } + + @Test + fun `foreground charging is performance but background charging stays power limited`() { + val foreground = PowerProfileResolver.resolve(80, true, false, false) + val background = PowerProfileResolver.resolve(80, true, true, false) + + assertEquals(PowerManager.PowerMode.PERFORMANCE, foreground.mode) + assertTrue(foreground.ble.continuousScan) + assertEquals(PowerManager.PowerMode.POWER_SAVER, background.mode) + assertFalse(background.ble.continuousScan) + assertEquals(59_000L, background.ble.scanOffMs) + } + + @Test + fun `battery thresholds are inclusive`() { + assertEquals(PowerManager.BatteryBand.NORMAL, resolve(21, true).batteryBand) + assertEquals(PowerManager.BatteryBand.LOW, resolve(20, true).batteryBand) + assertEquals(PowerManager.BatteryBand.LOW, resolve(11, true).batteryBand) + assertEquals(PowerManager.BatteryBand.CRITICAL, resolve(10, true).batteryBand) + } + + private fun resolve( + battery: Int, + background: Boolean, + peers: Boolean = false + ) = PowerProfileResolver.resolve( + batteryLevel = battery, + isCharging = false, + isBackground = background, + hasDirectPeers = peers + ) +} diff --git a/app/src/test/kotlin/com/bitchat/android/nostr/NostrBackgroundEventProcessorTest.kt b/app/src/test/kotlin/com/bitchat/android/nostr/NostrBackgroundEventProcessorTest.kt new file mode 100644 index 00000000..d7fffc2c --- /dev/null +++ b/app/src/test/kotlin/com/bitchat/android/nostr/NostrBackgroundEventProcessorTest.kt @@ -0,0 +1,65 @@ +package com.bitchat.android.nostr + +import android.app.Application +import androidx.test.core.app.ApplicationProvider +import com.bitchat.android.services.AppStateStore +import kotlinx.coroutines.CoroutineScope +import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.SupervisorJob +import kotlinx.coroutines.cancel +import kotlinx.coroutines.runBlocking +import kotlinx.coroutines.withTimeout +import org.junit.After +import org.junit.Assert.assertEquals +import org.junit.Before +import org.junit.Test +import org.junit.runner.RunWith +import org.robolectric.RobolectricTestRunner + +@RunWith(RobolectricTestRunner::class) +class NostrBackgroundEventProcessorTest { + private lateinit var scope: CoroutineScope + + @Before + fun setUp() { + AppStateStore.clear() + scope = CoroutineScope(SupervisorJob() + Dispatchers.IO) + } + + @After + fun tearDown() { + scope.cancel() + AppStateStore.clear() + } + + @Test + fun `cold start processes more events than the removed handoff queue capacity`() = runBlocking { + val application = ApplicationProvider.getApplicationContext() + val processor = NostrBackgroundEventProcessor(application, scope) + + repeat(300) { index -> + processor.onGeohashMessage( + event = NostrEvent( + id = "cold-start-$index", + pubkey = index.toString(16).padStart(64, '0'), + createdAt = 1, + kind = NostrKind.EPHEMERAL_EVENT, + tags = listOf(listOf("g", "u4pruy")), + content = "message-$index" + ), + geohash = "u4pruy" + ) + } + + withTimeout(5_000) { + while (AppStateStore.channelMessages.value["geo:u4pruy"].orEmpty().size < 300) { + kotlinx.coroutines.yield() + } + } + + assertEquals( + (0 until 300).map { "cold-start-$it" }.toSet(), + AppStateStore.channelMessages.value["geo:u4pruy"].orEmpty().map { it.id }.toSet() + ) + } +} diff --git a/app/src/test/kotlin/com/bitchat/android/nostr/NostrDirectMessageHandlerTest.kt b/app/src/test/kotlin/com/bitchat/android/nostr/NostrDirectMessageHandlerTest.kt index a9ded74d..f756447a 100644 --- a/app/src/test/kotlin/com/bitchat/android/nostr/NostrDirectMessageHandlerTest.kt +++ b/app/src/test/kotlin/com/bitchat/android/nostr/NostrDirectMessageHandlerTest.kt @@ -5,7 +5,6 @@ import com.bitchat.android.services.AppStateStore import com.bitchat.android.services.SeenMessageStore import com.bitchat.android.ui.ChatState import com.bitchat.android.ui.DataManager -import com.bitchat.android.ui.MeshDelegateHandler import com.bitchat.android.ui.MessageManager import com.bitchat.android.ui.NoiseSessionDelegate import com.bitchat.android.ui.PrivateChatManager @@ -72,7 +71,7 @@ class NostrDirectMessageHandlerTest { application = application, state = state, privateChatManager = privateChatManager, - meshDelegateHandler = mock(), + updateDeliveryStatus = { _, _ -> }, scope = scope, repo = GeohashRepository(application, state, dataManager), dataManager = dataManager, diff --git a/app/src/test/kotlin/com/bitchat/android/nostr/NostrPendingEventQueueTest.kt b/app/src/test/kotlin/com/bitchat/android/nostr/NostrPendingEventQueueTest.kt new file mode 100644 index 00000000..8954cbe3 --- /dev/null +++ b/app/src/test/kotlin/com/bitchat/android/nostr/NostrPendingEventQueueTest.kt @@ -0,0 +1,81 @@ +package com.bitchat.android.nostr + +import org.junit.Assert.assertEquals +import org.junit.Assert.assertNotEquals +import org.junit.Assert.assertNull +import org.junit.Test + +class NostrPendingEventQueueTest { + @Test + fun `empty relay set is not queued`() { + val queue = NostrPendingEventQueue(capacity = 2) + + assertNull(queue.enqueue(event("empty"), emptyList(), liveLocationToken = null)) + assertEquals(0, queue.size()) + } + + @Test + fun `capacity evicts the oldest publish`() { + val queue = NostrPendingEventQueue(capacity = 2) + queue.enqueue(event("one"), listOf("relay"), liveLocationToken = null) + queue.enqueue(event("two"), listOf("relay"), liveLocationToken = null) + queue.enqueue(event("three"), listOf("relay"), liveLocationToken = null) + + assertEquals( + listOf("two", "three"), + queue.pendingForRelay("relay").map { it.event.content } + ) + } + + @Test + fun `duplicate event publishes retain independent delivery state`() { + val queue = NostrPendingEventQueue(capacity = 4) + val signedEvent = event("same") + val firstId = requireNotNull( + queue.enqueue(signedEvent, listOf("relay-a", "relay-b"), liveLocationToken = null) + ) + val secondId = requireNotNull( + queue.enqueue(signedEvent, listOf("relay-a"), liveLocationToken = null) + ) + assertNotEquals(firstId, secondId) + + queue.markDelivered(firstId, "relay-a") + + assertEquals( + listOf(secondId), + queue.pendingForRelay("relay-a").map { it.queueId } + ) + assertEquals( + listOf(firstId), + queue.pendingForRelay("relay-b").map { it.queueId } + ) + + queue.markDelivered(firstId, "relay-b") + assertEquals(1, queue.size()) + } + + @Test + fun `privacy purge retains non-live publishes`() { + val queue = NostrPendingEventQueue(capacity = 4) + queue.enqueue(event("manual"), listOf("relay"), liveLocationToken = null) + queue.enqueue(event("live"), listOf("relay"), liveLocationToken = 42L) + + queue.removeLiveLocationEvents() + + assertEquals( + listOf("manual"), + queue.pendingForRelay("relay").map { it.event.content } + ) + } + + private fun event(content: String): NostrEvent { + val privateKey = "0".repeat(63) + "1" + return NostrEvent( + pubkey = NostrCrypto.derivePublicKey(privateKey), + createdAt = 1, + kind = NostrKind.TEXT_NOTE, + tags = emptyList(), + content = content + ).sign(privateKey) + } +} diff --git a/app/src/test/kotlin/com/bitchat/android/nostr/NostrRelayManagerLifecycleSmokeTest.kt b/app/src/test/kotlin/com/bitchat/android/nostr/NostrRelayManagerLifecycleSmokeTest.kt index 7bb4aa30..de53a579 100644 --- a/app/src/test/kotlin/com/bitchat/android/nostr/NostrRelayManagerLifecycleSmokeTest.kt +++ b/app/src/test/kotlin/com/bitchat/android/nostr/NostrRelayManagerLifecycleSmokeTest.kt @@ -9,6 +9,34 @@ import org.robolectric.RobolectricTestRunner @RunWith(RobolectricTestRunner::class) class NostrRelayManagerLifecycleSmokeTest { + @Test + fun `owner teardown is synchronous and preserves other subscription owners`() { + val manager = NostrRelayManager.shared + manager.disconnect() + manager.clearAllSubscriptions() + val filter = NostrFilter(kinds = listOf(NostrKind.TEXT_NOTE)) + + manager.subscribe( + filter = filter, + id = "background-contract", + handler = {}, + targetRelayUrls = emptyList(), + owner = NostrRelayManager.OWNER_BACKGROUND + ) + manager.subscribe( + filter = filter, + id = "ui-contract", + handler = {}, + targetRelayUrls = emptyList(), + owner = "test-ui" + ) + + manager.unsubscribeOwner("test-ui") + + assertEquals(setOf("background-contract"), manager.getActiveSubscriptions().keys) + manager.clearAllSubscriptions() + } + @Test fun `disconnected manager maintains subscription and empty publish invariants locally`() { val manager = NostrRelayManager.shared diff --git a/app/src/test/kotlin/com/bitchat/android/services/AppStateStoreTest.kt b/app/src/test/kotlin/com/bitchat/android/services/AppStateStoreTest.kt index afce826f..90ab0034 100644 --- a/app/src/test/kotlin/com/bitchat/android/services/AppStateStoreTest.kt +++ b/app/src/test/kotlin/com/bitchat/android/services/AppStateStoreTest.kt @@ -79,6 +79,10 @@ class AppStateStoreTest { setOf("ble-1", "wifi-1", "shared"), AppStateStore.getDirectPeers() ) + assertEquals( + setOf("ble-1", "wifi-1", "shared"), + AppStateStore.directPeers.value + ) } @Test diff --git a/app/src/test/kotlin/com/bitchat/android/ui/PrivateChatManagerTest.kt b/app/src/test/kotlin/com/bitchat/android/ui/PrivateChatManagerTest.kt index 354fdc62..d8371690 100644 --- a/app/src/test/kotlin/com/bitchat/android/ui/PrivateChatManagerTest.kt +++ b/app/src/test/kotlin/com/bitchat/android/ui/PrivateChatManagerTest.kt @@ -69,6 +69,34 @@ class PrivateChatManagerTest { assertEquals(listOf(message), state.getPrivateChatsValue()[conversationID]) } + @Test + fun `headless Nostr processing stores messages without retaining UI unread work`() { + val headlessManager = PrivateChatManager( + state = state, + messageManager = MessageManager(state), + dataManager = DataManager(RuntimeEnvironment.getApplication()), + noiseSessionDelegate = mock(), + trackUnreadMessages = false + ) + val message = BitchatMessage( + id = "background-nostr-message", + sender = "alice", + content = "background", + timestamp = Date(1), + isPrivate = true, + senderPeerID = "nostr_background" + ) + + headlessManager.handleIncomingPrivateMessage( + message = message, + suppressUnread = false, + origin = PrivateMessageOrigin.NOSTR + ) + + assertEquals(listOf(message), AppStateStore.privateMessages.value["nostr_background"]) + assertTrue(state.getUnreadPrivateMessagesValue().isEmpty()) + } + @Test fun `canonical conversation sends read receipt through live mesh peer id`() { val noiseKey = ByteArray(32) { 9 }