Merge pull request #781 from permissionlesstech/codex/background-power-optimization

Centralize adaptive background power scheduling
This commit is contained in:
callebtc 2026-07-28 00:15:08 +02:00 committed by GitHub
commit 57d11299da
No known key found for this signature in database
GPG Key ID: B5690EEEBB952194
36 changed files with 1771 additions and 881 deletions

View File

@ -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) { }

View File

@ -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
}
/**

View File

@ -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)
}
}
}

View File

@ -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) { }

View File

@ -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<String>) {

View File

@ -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<RuntimePerformanceProfile> = _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
}
}

View File

@ -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<String>, channel: String?) {
when {
isBleEnabled() -> bluetooth.sendMessage(content, mentions, channel)

View File

@ -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)
}

View File

@ -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}")
}

View File

@ -28,14 +28,17 @@ class GeohashRepository(
// conversation key (e.g., "nostr_<pub16>") -> source geohash it belongs to
private val conversationGeohash: MutableMap<String, String> = 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<String, Int>()
@ -180,12 +193,15 @@ class GeohashRepository(
state.setGeohashParticipantCounts(counts)
}
@Synchronized
fun putNostrKeyMapping(tempKeyOrPeer: String, pubkeyHex: String) {
nostrKeyMapping[tempKeyOrPeer] = pubkeyHex
}
@Synchronized
fun getNostrKeyMapping(): Map<String, String> = 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)

View File

@ -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()
}
}

View File

@ -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
}
}

View File

@ -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 -> {

View File

@ -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<String>,
val liveLocationToken: Long?
)
private val lock = Any()
private val entries = ArrayDeque<Entry>()
private var nextQueueId = 1L
fun enqueue(
event: NostrEvent,
relayUrls: Collection<String>,
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<Delivery> = 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 }
}

View File

@ -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<Relay>()
private val connections = ConcurrentHashMap<String, WebSocket>()
private val reconnectJobs = ConcurrentHashMap<String, Job>()
private val desiredConnected = AtomicBoolean(false)
private val subscriptions = ConcurrentHashMap<String, Set<String>>() // relay URL -> subscription IDs
private val messageHandlers = ConcurrentHashMap<String, (NostrEvent) -> Unit>()
@ -99,28 +106,25 @@ class NostrRelayManager private constructor() {
val handler: (NostrEvent) -> Unit,
val targetRelayUrls: Set<String>? = 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<String>,
val liveLocationToken: Long? = null
)
private val messageQueue = mutableListOf<QueuedEvent>()
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<String>? = 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<String>? = 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)
}
}
}

View File

@ -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)
}

View File

@ -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()

View File

@ -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
}
}

View File

@ -17,6 +17,8 @@ object AppStateStore {
private val peerIdsByTransport = mutableMapOf<String, Set<String>>()
// Direct (single-hop) peer IDs per transport, used to gossip a unified neighbor set.
private val directPeerIdsByTransport = mutableMapOf<String, Set<String>>()
private val _directPeers = MutableStateFlow<Set<String>>(emptySet())
val directPeers: StateFlow<Set<String>> = _directPeers.asStateFlow()
// Connected peer IDs (mesh ephemeral IDs)
private val _peers = MutableStateFlow<List<String>>(emptyList())
val peers: StateFlow<List<String>> = _peers.asStateFlow()
@ -29,6 +31,12 @@ object AppStateStore {
private val _privateMessages = MutableStateFlow<Map<String, List<BitchatMessage>>>(emptyMap())
val privateMessages: StateFlow<Map<String, List<BitchatMessage>>> = _privateMessages.asStateFlow()
private val _nickname = MutableStateFlow("")
val nickname: StateFlow<String> = _nickname.asStateFlow()
private val _selectedPrivateChatPeer = MutableStateFlow<String?>(null)
val selectedPrivateChatPeer: StateFlow<String?> = _selectedPrivateChatPeer.asStateFlow()
// Channel messages by channel name
private val _channelMessages = MutableStateFlow<Map<String, List<BitchatMessage>>>(emptyMap())
val channelMessages: StateFlow<Map<String, List<BitchatMessage>>> = _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<String>) {
synchronized(this) {
peerIdsByTransport[transportId] = ids.toSet()
@ -69,20 +85,25 @@ object AppStateStore {
fun setTransportDirectPeers(transportId: String, ids: Collection<String>) {
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<String> {
synchronized(this) {
return directPeerIdsByTransport.values.flatten().toSet()
}
fun getDirectPeers(): Set<String> = _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
}
}

View File

@ -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<String>) {

View File

@ -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
}

View File

@ -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<String>() // Set of nostr pubkey hex
val geohashBlockedUsers: Set<String> 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)
}

View File

@ -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<String>()
private val samplingSubscriptionIds = mutableMapOf<String, String>()
private val liveSamplingSubscriptionGeohashes = mutableSetOf<String>()
private var requestedLiveSamplingGeohashes: Set<String> = emptySet()
private var requestedUserSamplingGeohashes: Set<String> = 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) {

View File

@ -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)
}

View File

@ -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})")

View File

@ -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 {

View File

@ -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 {

View File

@ -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<String, String>,
known: Map<String, String>,
discovered: Set<String>
) {
_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()
}
}

View File

@ -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<String> =
(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.
*/

View File

@ -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
)
}

View File

@ -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<Application>()
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()
)
}
}

View File

@ -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<MeshDelegateHandler>(),
updateDeliveryStatus = { _, _ -> },
scope = scope,
repo = GeohashRepository(application, state, dataManager),
dataManager = dataManager,

View File

@ -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)
}
}

View File

@ -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

View File

@ -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

View File

@ -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 }