diff --git a/CHANGELOG.md b/CHANGELOG.md index 952db61d..c6c7bcb3 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,5 +1,32 @@ # Change log +## 1.2.0-rc.4 + +* Fixed: `host()` resolves with the live connection state instead of always `null` +* Fixed: `onConnection(false)` fires when a retired or background connection closes +* Fixed: notification taps are recorded after launching the app, so tap-opened screens land on top +* Fixed: `host()` reports `false` when there was nothing to subscribe instead of a spurious open + +## 1.2.0-rc.3 + +* Added: `getInitialNotification()`, `onNotificationOpened()`, and the `PushNotificationOpened` type (topic + data) +* Added: live push connection open/close events surfaced via `onOpen`/`onClose` +* Fixed: notification taps delivered reliably once, via a `PushTaps` queue +* Fixed: foreground notifications de-duplicated per title +* Updated: `resume` leaves a live Push host's connection untouched + +## 1.2.0-rc.2 + +* Added: `Push.getInitialNotification()` returns the notification whose tap launched the app +* Added: `Push.onNotificationOpened(callback)` fires on each notification tap +* Added: `Push.backgroundStatus()` reports Android exact-alarm, battery, and foreground-service state +* Added: `Push.requestExactAlarms()` and `Push.requestIgnoreBatteryOptimizations()` helpers +* Added: `SubscribeOptions.notifyInForeground` to post notifications while the app is visible +* Added: `PushBackgroundStatus` and `PushNotificationOpened` exported types +* Fixed: reliable Android background wake-ups via a periodic watchdog and alternating jobs +* Fixed: `resume` keeps subscriptions only for the same signed-in user +* Fixed: notification images are size-capped to prevent out-of-memory crashes + ## 1.2.0-rc.1 * Added: background push notifications render the server `notification` title, body, and image diff --git a/README.md b/README.md index b2579f08..0ba50f82 100644 --- a/README.md +++ b/README.md @@ -38,8 +38,8 @@ A subscription with `background: true` keeps delivering after the app is backgro the device restarts, until it is unsubscribed or `push.close()` is called (do this on sign-out). The SDK's native Android module (autolinked) saves the subscription, and a scheduled job and alarm wake the app every 15 to 60 seconds to reconnect; the broker replays what was sent in -between (`retry: true`). Messages no in-app callback receives are posted as notifications that -open the app. It reconnects with the credential saved at subscribe time, so use a session rather +between (`retry: true`). While the app is not on screen, each message is posted as a +notification that opens the app. It reconnects with the credential saved at subscribe time, so use a session rather than a short-lived JWT. On Android 13 and later, the first background subscription asks the user for the @@ -58,6 +58,86 @@ const sub = await push.subscribe('news', (message) => console.log(message.data), await push.setForeground(true); ``` +Notifications are posted while the app is backgrounded or closed. While it is on screen your +callback shows the message, so none is posted unless the subscription passes +`notifyInForeground: true`. + +#### Delivery while the app is closed + +Messages sent while the app is closed arrive at the next scheduled wake-up. While the device is +awake that is about every 15 seconds, or about every 60 seconds once exact alarms are allowed; +without exact alarms the wake-ups are inexact, so battery saver can defer them further. In Doze +(screen off and idle for a while) Android limits background alarms, exact ones included, to about +one every nine minutes, so a closed app can take several minutes to receive a message: allowing +exact alarms makes wake-ups punctual, it does not lift Doze. For immediate delivery, also in +Doze, use `push.setForeground(true)` (see above). + +The SDK uses exact alarms on its own whenever the app may schedule them. To allow it: + +1. Declare the permissions in your app's `AndroidManifest.xml`. Both are optional and subject to + Google Play policy: `SCHEDULE_EXACT_ALARM` needs a declaration in the Play Console, and + `USE_EXACT_ALARM` is reserved for alarm, clock and calendar apps (the SDK does not use it). + + ```xml + + + ``` + +With Expo, list them under `android.permissions` in `app.json` instead. + +2. On Android 13 and later the user has to allow exact alarms, under Settings > Apps > Special app + access > Alarms & reminders. Android 12 grants a declared `SCHEDULE_EXACT_ALARM` + automatically, and older versions need nothing. Check with `await push.backgroundStatus()`: when `bestEffort` is + true, explain why to the user, then from a user action open that screen with `push.requestExactAlarms()`, or + ask for the battery-optimisation exemption with `push.requestIgnoreBatteryOptimizations()`. Both return false when there is + nothing to ask, including when the permission is not declared. The SDK never opens these + screens on its own. + +If the user force-stops the app (Settings > Force stop, and on some devices swiping it away from +recents), Android cancels its alarms and jobs: nothing is delivered until the app is opened +again, and the broker then replays what was sent meanwhile. + +Set the notification icon with +`` in your +``; without it a generic icon is used. + +```js +const status = await push.backgroundStatus(); // null outside Android +if (status?.bestEffort) { + // Explain why, then from a button press: + await push.requestExactAlarms(); +} +``` + +Saved background delivery follows the app's current session, also while the app is closed: each +background run re-reads the session cookie, so a rotated session of the same user replaces the +saved one, and signing out (no session) or signing in as someone else stops background delivery. +If the broker still refuses the credential, delivery stops and `onError` reports it the next time +the app registers one. Still call `push.close()` on sign-out. + +#### Upgrading from an earlier release candidate + +The Expo config plugin is gone: the native module now declares everything background delivery +needs. Remove `"react-native-appwrite"` from the `plugins` list in `app.json`, or `expo config` +fails to load it. + +#### Opening a tapped notification + +Read the `data` sent with `createPush` when the user taps a background notification: + +```js +// The tap that launched the app (reported once, so call it at startup). +const opened = await push.getInitialNotification(); +if (opened) { + openSale(opened.data.saleId); +} + +// Taps while the app is running, including in the background. +const stop = push.onNotificationOpened(({ topic, data }) => openSale(data.saleId)); +``` + +On Android the SDK's native module reports the taps. Elsewhere they come from `expo-notifications`. + Foreground mode runs a `remoteMessaging` foreground service, which Google Play asks apps to declare in the Play Console. Apps that never enable it can remove the service from their merged manifest with `tools:node="remove"` on `io.appwrite.services.PushService` and diff --git a/android/src/main/AndroidManifest.xml b/android/src/main/AndroidManifest.xml index 44f894ed..d2b7d500 100644 --- a/android/src/main/AndroidManifest.xml +++ b/android/src/main/AndroidManifest.xml @@ -16,6 +16,13 @@ + Unit)? = null + override fun getName(): String = NAME + // Resolves with the tapped notification that launched the app as JSON (topic and payload), + // once, or null. + @ReactMethod + fun getInitialNotification(promise: Promise) = settle(promise) { + PushTaps.take()?.let { JSONObject(opened(it)).toString() } + } + + // Emits each tap while JS listens for them; until then a tap waits for getInitialNotification. + @ReactMethod + fun listenOpened(listening: Boolean, promise: Promise) = settle(promise) { + openedListeners = (openedListeners + if (listening) 1 else -1).coerceAtLeast(0) + if (openedListeners > 0 && stopOpened == null) { + stopOpened = PushTaps.listen { tap -> emit(OPENED_EVENT, opened(tap)) } + } else if (openedListeners == 0) { + stopOpened?.invoke() + stopOpened = null + } + null + } + + override fun invalidate() { + stopOpened?.invoke() + stopOpened = null + super.invalidate() + } + // Resolves once the connection is up and every filter is subscribed, or rejects with why not. @ReactMethod fun host(config: String, subscriptions: String, promise: Promise) { try { bridge.host(config, subscriptions) { error -> if (error == null) { - promise.resolve(null) + promise.resolve(bridge.isConnected()) } else { promise.reject(ERROR_CODE, error) } @@ -91,11 +126,20 @@ class AppwritePushModule internal constructor( fun hasSaved(promise: Promise) = settle(promise) { bridge.hasSaved() } @ReactMethod - fun resume(promise: Promise) = settle(promise) { - bridge.resume() + fun resume(authMethod: String?, credential: String?, signedOutWhenMissing: Boolean, promise: Promise) = settle(promise) { + bridge.resume(authMethod, credential, signedOutWhenMissing) null } + @ReactMethod + fun backgroundStatus(promise: Promise) = settle(promise) { bridge.backgroundStatus() } + + @ReactMethod + fun requestExactAlarms(promise: Promise) = settle(promise) { bridge.requestExactAlarms() } + + @ReactMethod + fun requestIgnoreBatteryOptimizations(promise: Promise) = settle(promise) { bridge.requestIgnoreBatteryOptimizations() } + @ReactMethod fun setErrorCallback(registered: Boolean, promise: Promise) = settle(promise) { bridge.setErrorCallback(registered) } @@ -110,6 +154,8 @@ class AppwritePushModule internal constructor( @ReactMethod fun removeListeners(count: Double) = Unit + private fun opened(tap: PushTap): Map = mapOf("topic" to tap.topic, "payload" to tap.payload) + private fun settle(promise: Promise, block: () -> Any?) { try { promise.resolve(block()) @@ -122,6 +168,8 @@ class AppwritePushModule internal constructor( const val NAME = "AppwritePush" const val MESSAGE_EVENT = "AppwritePushMessage" const val ERROR_EVENT = "AppwritePushError" + const val OPENED_EVENT = "AppwritePushOpened" + const val CONNECTION_EVENT = "AppwritePushConnection" private const val ERROR_CODE = "appwrite_push" } } diff --git a/android/src/main/java/io/appwrite/services/PushBackground.kt b/android/src/main/java/io/appwrite/services/PushBackground.kt index a550a555..e06dfabe 100644 --- a/android/src/main/java/io/appwrite/services/PushBackground.kt +++ b/android/src/main/java/io/appwrite/services/PushBackground.kt @@ -1,5 +1,7 @@ package io.appwrite.services +import android.Manifest +import android.app.ActivityManager import android.app.AlarmManager import android.app.NotificationChannel import android.app.NotificationManager @@ -14,11 +16,14 @@ import android.graphics.Bitmap import android.graphics.BitmapFactory import android.net.ConnectivityManager import android.net.Network +import android.net.Uri import android.os.Build import android.os.PowerManager import android.os.SystemClock +import android.provider.Settings import android.util.AtomicFile import android.util.Log +import android.webkit.CookieManager import androidx.core.app.NotificationCompat import androidx.core.app.NotificationManagerCompat import androidx.core.content.ContextCompat @@ -30,8 +35,10 @@ import com.hivemq.client.mqtt.mqtt5.message.publish.Mqtt5Publish import io.appwrite.exceptions.AppwriteException import org.json.JSONArray import org.json.JSONObject +import java.io.ByteArrayOutputStream import java.io.File import java.io.IOException +import java.io.InputStream import java.net.HttpURLConnection import java.net.URL import java.util.UUID @@ -53,6 +60,7 @@ internal data class PushEntry( val filter: String, val retry: Boolean, val title: String?, + val notifyInForeground: Boolean = false, ) /** @@ -96,7 +104,11 @@ internal object PushBackground { const val EXTRA_PAYLOAD = "io.appwrite.push.PAYLOAD" private const val SERVICE_CHANNEL_ID = "appwrite-push-service" - private const val JOB_ID = 0x50555348 + const val JOB_ID = 0x50555348 + const val NEXT_JOB_ID = 0x5055534A + const val EXPEDITED_JOB_ID = 0x5055534B + const val WATCHDOG_JOB_ID = 0x50555349 + private const val WATCHDOG_INTERVAL_MS = 15 * 60 * 1_000L private const val INTERVAL_MS = 15_000L private const val PRIVILEGED_INTERVAL_MS = 60_000L private const val RESTART_DELAY_MS = 3_000L @@ -116,6 +128,15 @@ internal object PushBackground { private const val HEARTBEAT_RESET_MS = 3 * 24 * 60 * 60 * 1_000L private const val REQUEST_TIMEOUT_SECONDS = 10L private const val IMAGE_TIMEOUT_MS = 5_000 + + // How long a scheduled run stays up after connecting so the broker can replay what was missed: + // until nothing arrived for DRAIN_QUIET_MS and nothing is being delivered, at most the run's cap. + const val JOB_DRAIN_MS = 15_000L + const val RECEIVER_DRAIN_MS = 5_000L + private const val DRAIN_QUIET_MS = 2_000L + private const val DRAIN_POLL_MS = 200L + private const val IMAGE_MAX_BYTES = 5 * 1024 * 1024 + private const val IMAGE_MAX_PX = 1_024 private const val CONNECT_TIMEOUT_SECONDS = 20L // How long a message waits for a listener that acknowledges it itself (the React Native and @@ -139,6 +160,16 @@ internal object PushBackground { @Volatile var errorHandler: ((Throwable) -> Boolean)? = null + // Tells the Push instance that hosts its subscriptions here when the connection opens (true) + // or closes (false). + @Volatile + var connectionHandler: ((Boolean) -> Unit)? = null + + // Whether the live connection (not a retired one) is open. + @Volatile + var isConnected = false + private set + @Volatile var serviceRunning = false @@ -164,6 +195,15 @@ internal object PushBackground { Thread(runnable, "AppwritePushBackground").apply { isDaemon = true } } private var mqtt: Mqtt5AsyncClient? = null + + // Set when [mqtt] is replaced or closed, so that client stops reconnecting and delivering. + private var mqttRetired: AtomicBoolean? = null + + // Messages being delivered, and when one last arrived or finished, for a run's drain window. + private val delivering = AtomicInteger(0) + + @Volatile + private var lastDelivery = 0L private var connectedConfig: PushConfig? = null private var connectedNetwork: Network? = null private var lastHeartbeat = 0L @@ -264,23 +304,120 @@ internal object PushBackground { } } - /** Resume saved background delivery when the app starts, without waiting for the next run. */ - fun resume(context: Context) { + /** + * Resume saved background delivery when the app starts, without waiting for the next run, with + * the app's current credential. A rotated session of the same user replaces the saved one; a + * different user drops the saved subscriptions, as does no credential when + * [signedOutWhenMissing] (the caller read the session store, so none means signed out). + * While a live Push hosts subscriptions here, it does nothing: that Push owns the connection. + */ + fun resume(context: Context, authMethod: String? = null, credential: String? = null, signedOutWhenMissing: Boolean = false) { val app = context.applicationContext - val active = synchronized(lock) { + val resume = synchronized(lock) { load(app) - isActive() + val current = config + when { + !isActive() || current == null -> null + // A live Push hosts here and owns the connection and its credential; resuming only + // reconciles saved delivery, so it leaves a live host alone. + listeners.isNotEmpty() -> null + // No credential: signed out when the caller can tell, else possibly not set yet. + authMethod == null || credential == null -> if (signedOutWhenMissing) false else null + current.authMethod == authMethod && current.credential == credential -> true + // A new credential for the same user (a rotated session) replaces the saved one, + // so the stale one is never sent and refused. + sameUser(current.authMethod, current.credential, authMethod, credential) -> { + config = current.copy(authMethod = authMethod, credential = credential) + persist(app) + true + } + // Another user signed in: the saved subscriptions were not theirs. + else -> false + } } - if (active) { - tick(app) {} + when (resume) { + true -> tick(app) {} + false -> stop(app) + null -> Unit } } + private enum class CredentialRefresh { UNCHANGED, CHANGED, GONE } + + // Bring the saved credential up to date with the session cookie the app + // signs in with, when it comes from one (see [PushConfig.sessionCookieUrl]): a rotated session of + // the same user replaces it, and no session or another user's is GONE. UNCHANGED when there is + // no cookie source or it cannot be read. + private fun refreshCredential(context: Context): CredentialRefresh { + val current = synchronized(lock) { config } ?: return CredentialRefresh.UNCHANGED + val url = current.sessionCookieUrl ?: return CredentialRefresh.UNCHANGED + val session = readSessionCookie(url, current.project) ?: return CredentialRefresh.UNCHANGED + return when { + session.isEmpty() -> CredentialRefresh.GONE + current.authMethod == "appwrite-session" && session == current.credential -> CredentialRefresh.UNCHANGED + sameUser(current.authMethod, current.credential, "appwrite-session", session) -> { + synchronized(lock) { + config = current.copy(authMethod = "appwrite-session", credential = session) + persist(context.applicationContext) + } + CredentialRefresh.CHANGED + } + else -> CredentialRefresh.GONE + } + } + + // The credential to send for a connection made with [connected]: the session re-read from the + // cookie store when one backs it (so an automatic reconnect uses a rotated session), else the + // one it connected with. Only a newer credential of the same user, for the same connection + // (broker, client id, project), is borrowed: a connection being replaced by [host] keeps its + // own. Null when the app signed out or another user signed in: authentication is aborted and + // delivery stops. + private fun currentCredential(context: Context, connected: PushConfig): String? { + if (refreshCredential(context) == CredentialRefresh.GONE) { + worker.execute { stop(context) } + return null + } + val current = synchronized(lock) { config } ?: return connected.credential + val sameConnection = current.host == connected.host && current.port == connected.port && + current.clientId == connected.clientId && current.project == connected.project && + current.authMethod == connected.authMethod + return if (sameConnection && sameUser(connected.authMethod, connected.credential, current.authMethod, current.credential)) { + current.credential + } else { + connected.credential + } + } + + // The `a_session_` cookie for [url] in the WebView cookie store, decoded; "" when there + // is none, and null when the store cannot be read. + private fun readSessionCookie(url: String, project: String): String? = runCatching { + val name = "a_session_$project" + CookieManager.getInstance().getCookie(url) + ?.split(";") + ?.map { it.trim() } + ?.firstOrNull { it.substringBefore("=") == name } + ?.substringAfter("=") + ?.let { Uri.decode(it) } + ?: "" + }.getOrNull() + + private fun sameUser(savedMethod: String, saved: String, method: String, credential: String): Boolean { + fun userId(method: String, credential: String) = + if (method == "appwrite-session") userIdFromSession(credential) else userIdFromJwt(credential) + val user = userId(method, credential) + return user.isNotEmpty() && user == userId(savedMethod, saved) + } + /** * One scheduled run: under a wakelock, reconnect if needed, send a heartbeat when one is due, * re-arm the next run, and call [done] once finished. Shuts down when nothing is saved. + * + * A run of a chained job ([jobId]) arms the next run under the other chained id: scheduling the + * id that is running would stop the run in progress. With [drainMs], the run stays up after + * connecting while the broker replays what was missed (at most that long), so a process woken + * from the background is not cached and frozen before the messages are delivered. */ - fun tick(context: Context, done: () -> Unit) { + fun tick(context: Context, jobId: Int? = null, drainMs: Long = 0, done: () -> Unit) { val app = context.applicationContext val wakeLock = acquireWakeLock(app) worker.execute { @@ -291,12 +428,17 @@ internal object PushBackground { } if (active) { // Arm the next run first, so a run cut short (killed, or a slow broker outlasting - // the wakelock) still leaves the chain armed. A refused credential in converge() - // shuts down, cancelling it again. - schedule(app, INTERVAL_MS) + // the wakelock) still leaves the chain armed; the watchdog job restarts a chain + // that was dropped anyway. A refused credential in converge() shuts down, + // cancelling them again. + schedule(app, INTERVAL_MS, runningJobId = jobId) + ensureWatchdog(app) ensureService(app) registerNetworkCallback(app) converge(app, heartbeat = true) + if (drainMs > 0 && mqtt != null) { + awaitDeliveries(drainMs) + } } else { shutdown(app) } @@ -309,9 +451,49 @@ internal object PushBackground { } } + // Worker thread only. Wait until no message arrived for DRAIN_QUIET_MS and none is being + // delivered, at most [maxMs]. Timed with the real monotonic clock ([nowMs]), which keeps moving + // while this thread sleeps. + private fun awaitDeliveries(maxMs: Long) { + val start = nowMs() + if (lastDelivery < start) { + lastDelivery = start + } + while (nowMs() - start < maxMs) { + if (delivering.get() == 0 && nowMs() - lastDelivery >= DRAIN_QUIET_MS) { + return + } + Thread.sleep(DRAIN_POLL_MS) + } + } + + private fun nowMs(): Long = System.nanoTime() / 1_000_000 + + /** + * Android 12+: hand a run to an expedited job, which starts at once with network access and + * keeps the process runnable until it finishes. A broadcast receiver woken in a cached process + * runs at background priority and can be frozen before it connects. Returns false when the + * job cannot be scheduled (older Android, or the expedited quota is used up), so the caller + * runs the tick itself. + */ + fun runExpedited(context: Context): Boolean { + if (Build.VERSION.SDK_INT < Build.VERSION_CODES.S) { + return false + } + return try { + val job = JobInfo.Builder(EXPEDITED_JOB_ID, ComponentName(context, PushJobService::class.java)) + .setExpedited(true) + .setRequiredNetworkType(JobInfo.NETWORK_TYPE_ANY) + .build() + context.getSystemService(JobScheduler::class.java)?.schedule(job) == JobScheduler.RESULT_SUCCESS + } catch (e: Exception) { + false + } + } + /** Run once more shortly after the task was removed, so an OEM kill that follows is undone. */ fun scheduleRestart(context: Context) { - schedule(context.applicationContext, RESTART_DELAY_MS, jobToo = false) + schedule(context.applicationContext, RESTART_DELAY_MS, alarmOnly = true) } /** @@ -326,6 +508,7 @@ internal object PushBackground { loaded = false } errorHandler = null + connectionHandler = null unregisterNetworkCallback(context.applicationContext) worker.submit { disconnect() }.get() } @@ -361,9 +544,14 @@ internal object PushBackground { // Must hold [lock]. Every background subscription, live or saved, one per filter and title. // Subscriptions sharing a filter and title are saved once, with QoS 1 when any of them wants it. private fun entries(): List = - (saved + listeners.filter { it.background }.flatMap { listener -> listener.topics.map { PushEntry(it, listener.retry, listener.title) } }) + (saved + liveEntries()) .groupBy { it.filter to it.title } - .map { (key, group) -> PushEntry(key.first, group.any { it.retry }, key.second) } + .map { (key, group) -> PushEntry(key.first, group.any { it.retry }, key.second, group.any { it.notifyInForeground }) } + + // Must hold [lock]. + private fun liveEntries(): List = listeners.filter { it.background }.flatMap { listener -> + listener.topics.map { PushEntry(it, listener.retry, listener.title, listener.notifyInForeground) } + } // Must hold [lock]. private fun persist(context: Context) { @@ -397,6 +585,7 @@ internal object PushBackground { private fun activate(context: Context, done: ((Throwable?) -> Unit)? = null) { registerNetworkCallback(context) schedule(context, INTERVAL_MS) + ensureWatchdog(context) ensureService(context) worker.execute { converge(context, heartbeat = false, done) } } @@ -437,6 +626,11 @@ internal object PushBackground { } private fun convergeOnce(context: Context, heartbeat: Boolean) { + // Signed out, or signed in as someone else, since the subscriptions were saved. + if (refreshCredential(context) == CredentialRefresh.GONE) { + stop(context) + return + } val (current, wanted) = synchronized(lock) { config to wanted() } if (current == null || wanted.isEmpty()) { disconnect() @@ -477,24 +671,50 @@ internal object PushBackground { // A failed first connect of a waited-for converge goes to the waiting caller only. val quiet = waiting val connected = AtomicBoolean(false) + val retired = AtomicBoolean(false) val client = buildPushClient( config.copy(keepAlive = KEEP_ALIVE_SECONDS), + onConnected = { + if (!retired.get()) { + isConnected = true + connectionHandler?.invoke(true) + } + }, + onDisconnected = { + if (!retired.get()) { + isConnected = false + connectionHandler?.invoke(false) + } + }, onError = { error -> if (!quiet || connected.get()) { report(error) } }, onAuthRefused = { error -> worker.execute { refused(context, error) } }, + credential = { currentCredential(context, config) }, + retired = retired, ) // Registered before connecting, so the backlog the broker replays right after each // SUBSCRIBE reaches the app however the subscription was made (including the ones the // client restores itself after an automatic reconnect). Acknowledged by hand, once the // message was delivered (see [deliver]). - client.publishes(MqttGlobalPublishFilter.ALL, { publish -> deliver(context, publish) }, true) + client.publishes( + MqttGlobalPublishFilter.ALL, + { publish -> + if (!retired.get()) { + deliver(context, publish) + } + }, + true, + ) try { - client.sendPushConnect(config.copy(keepAlive = KEEP_ALIVE_SECONDS), CONNECT_TIMEOUT_SECONDS) + client.sendPushConnect(config.copy(keepAlive = KEEP_ALIVE_SECONDS), CONNECT_TIMEOUT_SECONDS) { + currentCredential(context, config) + } } catch (e: Exception) { val cause = (e as? ExecutionException)?.cause ?: e + retired.set(true) client.disconnect() val error = connectError(cause) when { @@ -508,6 +728,7 @@ internal object PushBackground { } connected.set(true) mqtt = client + mqttRetired = retired connectedConfig = config connectedNetwork = activeNetwork(context) lastHeartbeat = SystemClock.elapsedRealtime() @@ -575,6 +796,13 @@ internal object PushBackground { // Worker thread only. private fun disconnect() { + mqttRetired?.set(true) + mqttRetired = null + // A retired client's own disconnect is not reported, so report the live one closing here. + if (isConnected) { + isConnected = false + connectionHandler?.invoke(false) + } mqtt?.disconnect() mqtt = null connectedConfig = null @@ -586,6 +814,21 @@ internal object PushBackground { // registered the message is also saved, and the next onError registration receives it. A // waiting caller receives it instead. private fun refused(context: Context, error: Throwable) { + // The saved credential may be a session the app has since rotated: retry once with the + // current one. Refused again, it is unchanged and delivery stops below. + if (!waiting) { + when (refreshCredential(context)) { + CredentialRefresh.CHANGED -> { + converge(context, heartbeat = false) + return + } + CredentialRefresh.GONE -> { + stop(context) + return + } + CredentialRefresh.UNCHANGED -> Unit + } + } if (waiting) { fail(error) } else if (!report(error)) { @@ -596,11 +839,24 @@ internal object PushBackground { // Deliver a message, then acknowledge it: right away, or once every listener that acknowledges // messages itself has done so (at most ACK_TIMEOUT_SECONDS later). + // A message counts as being delivered until it is acknowledged: for the React Native and Flutter + // bridges, once their callback has run (or the acknowledgement times out), not when the event is + // handed over, so a run's drain window waits for it. private fun deliver(context: Context, publish: Mqtt5Publish) { + delivering.incrementAndGet() + lastDelivery = nowMs() + deliverAcknowledged(context, publish) { + lastDelivery = nowMs() + delivering.decrementAndGet() + } + } + + private fun deliverAcknowledged(context: Context, publish: Mqtt5Publish, settled: () -> Unit) { val acknowledged = AtomicBoolean(false) val acknowledge = { if (acknowledged.compareAndSet(false, true)) { runCatching { publish.acknowledge() } + settled() } } val pending = AtomicInteger(1) @@ -624,7 +880,7 @@ internal object PushBackground { settle() } if (!acknowledged.get()) { - worker.schedule({ acknowledge() }, ACK_TIMEOUT_SECONDS, TimeUnit.SECONDS) + ackTimer.schedule({ acknowledge() }, ACK_TIMEOUT_SECONDS, TimeUnit.SECONDS) } } @@ -656,8 +912,13 @@ internal object PushBackground { // The app's PushReceiver gets only messages no live callback received. val handled = matching.isEmpty() && deliverToReceivers(context, message) val content = notificationContent(message) - val titles = matching.filter { it.background }.map { it.title ?: message.topic } + - if (handled) emptyList() else entries.map { it.title ?: message.topic } + // Background subscriptions notify while the app is not visible; on screen, only those + // that opted in with notifyInForeground do, since the app shows the message itself. A + // message no live callback received still notifies, as the app has not shown it. + val foreground = appInForeground() + val shown = foreground && matching.isNotEmpty() + val titles = matching.filter { it.background && (!foreground || it.notifyInForeground) }.map { it.title ?: message.topic } + + if (handled) emptyList() else entries.filter { !shown || it.notifyInForeground }.map { it.title ?: message.topic } // A title the server sent replaces every subscription's, so one notification is posted. titles.map { content.title ?: it }.distinct().forEach { notify(context, message, it, content) } } finally { @@ -712,14 +973,13 @@ internal object PushBackground { builder.setLargeIcon(image) .setStyle(NotificationCompat.BigPictureStyle().bigPicture(image).bigLargeIcon(null as Bitmap?).setSummaryText(body)) } - context.packageManager.getLaunchIntentForPackage(context.packageName)?.let { launch -> - launch.addFlags(Intent.FLAG_ACTIVITY_NEW_TASK or Intent.FLAG_ACTIVITY_SINGLE_TOP) - .putExtra(EXTRA_TOPIC, message.topic) - .putExtra(EXTRA_PAYLOAD, message.data) - builder.setContentIntent( - PendingIntent.getActivity(context, id, launch, PendingIntent.FLAG_UPDATE_CURRENT or PendingIntent.FLAG_IMMUTABLE), - ) - } + val open = Intent(context, PushOpenActivity::class.java) + .addFlags(Intent.FLAG_ACTIVITY_NEW_TASK or Intent.FLAG_ACTIVITY_NO_ANIMATION) + .putExtra(EXTRA_TOPIC, message.topic) + .putExtra(EXTRA_PAYLOAD, message.data) + builder.setContentIntent( + PendingIntent.getActivity(context, id, open, PendingIntent.FLAG_UPDATE_CURRENT or PendingIntent.FLAG_IMMUTABLE), + ) try { manager.notify(id, builder.build()) } catch (e: SecurityException) { @@ -727,6 +987,11 @@ internal object PushBackground { } } + // Whether one of the app's activities is visible: a foreground service alone does not count. + private fun appInForeground(): Boolean = ActivityManager.RunningAppProcessInfo() + .also { ActivityManager.getMyMemoryState(it) } + .importance == ActivityManager.RunningAppProcessInfo.IMPORTANCE_FOREGROUND + /** The server's `notification` block in [message]: nulls when the payload has none or is not JSON. */ fun notificationContent(message: PushMessage): PushNotificationContent { val notification = runCatching { JSONObject(message.data) }.getOrNull()?.optJSONObject("notification") @@ -736,6 +1001,11 @@ internal object PushBackground { } // Downloads notification images, so a slow one is abandoned without holding up delivery. + // Acknowledgement timeouts, apart from [worker] so they fire while a run's drain window holds it. + private val ackTimer = Executors.newSingleThreadScheduledExecutor { runnable -> + Thread(runnable, "AppwritePushAck").apply { isDaemon = true } + } + private val imageLoader = Executors.newCachedThreadPool { runnable -> Thread(runnable, "AppwritePushImage").apply { isDaemon = true } } @@ -746,7 +1016,7 @@ internal object PushBackground { val connection = runCatching { URL(url).openConnection() as HttpURLConnection }.getOrNull() ?: return null connection.connectTimeout = IMAGE_TIMEOUT_MS connection.readTimeout = IMAGE_TIMEOUT_MS - val download = imageLoader.submit { connection.inputStream.use { BitmapFactory.decodeStream(it) } } + val download = imageLoader.submit { connection.inputStream.use { readLimited(it, IMAGE_MAX_BYTES) }?.let { decodeImage(it) } } return try { download.get(IMAGE_TIMEOUT_MS.toLong(), TimeUnit.MILLISECONDS) } catch (e: Exception) { @@ -757,6 +1027,37 @@ internal object PushBackground { } } + // At most [limit] bytes of [input], or null when it holds more. + private fun readLimited(input: InputStream, limit: Int): ByteArray? { + val out = ByteArrayOutputStream() + val buffer = ByteArray(8_192) + while (true) { + val read = input.read(buffer) + if (read < 0) { + return out.toByteArray() + } + if (out.size() + read > limit) { + return null + } + out.write(buffer, 0, read) + } + } + + // Decode [bytes] downsampled so neither side exceeds IMAGE_MAX_PX, so a large photo cannot + // exhaust memory while a notification is posted. + private fun decodeImage(bytes: ByteArray): Bitmap? { + val bounds = BitmapFactory.Options().apply { inJustDecodeBounds = true } + BitmapFactory.decodeByteArray(bytes, 0, bytes.size, bounds) + if (bounds.outWidth <= 0 || bounds.outHeight <= 0) { + return null + } + var sample = 1 + while (bounds.outWidth / sample > IMAGE_MAX_PX || bounds.outHeight / sample > IMAGE_MAX_PX) { + sample *= 2 + } + return BitmapFactory.decodeByteArray(bytes, 0, bytes.size, BitmapFactory.Options().apply { inSampleSize = sample }) + } + /** The ongoing notification the foreground service shows, on its own quiet channel. */ fun ongoingNotification(context: Context) = NotificationCompat.Builder(context, SERVICE_CHANNEL_ID) .also { createChannels(context) } @@ -815,19 +1116,16 @@ internal object PushBackground { } // Arm the next run in [delayMs]: the job, and an alarm just after it in case the job is late. - private fun schedule(context: Context, delayMs: Long, jobToo: Boolean = true) { - val interval = if (delayMs == INTERVAL_MS && privileged(context)) PRIVILEGED_INTERVAL_MS else delayMs - if (jobToo) { - try { - val job = JobInfo.Builder(JOB_ID, ComponentName(context, PushJobService::class.java)) - .setMinimumLatency(interval) - .setRequiredNetworkType(JobInfo.NETWORK_TYPE_ANY) - .setPersisted(true) - .build() - context.getSystemService(JobScheduler::class.java)?.schedule(job) - } catch (e: Exception) { - report(e) + // The chained job alternates between two ids, so a run ([runningJobId]) never schedules its + // own id, which would stop it; the other id is cancelled only when it is not the one running. + private fun schedule(context: Context, delayMs: Long, alarmOnly: Boolean = false, runningJobId: Int? = null) { + val interval = interval(context, delayMs) + if (!alarmOnly) { + val next = if (runningJobId == JOB_ID) NEXT_JOB_ID else JOB_ID + if (runningJobId == null) { + context.getSystemService(JobScheduler::class.java)?.cancel(if (next == JOB_ID) NEXT_JOB_ID else JOB_ID) } + scheduleJob(context, delayMs, next) } val alarms = context.getSystemService(AlarmManager::class.java) ?: return val at = SystemClock.elapsedRealtime() + interval + 1_000L @@ -843,8 +1141,48 @@ internal object PushBackground { } } + private fun interval(context: Context, delayMs: Long): Long = + if (delayMs == INTERVAL_MS && privileged(context)) PRIVILEGED_INTERVAL_MS else delayMs + + private fun scheduleJob(context: Context, delayMs: Long, id: Int) { + try { + val job = JobInfo.Builder(id, ComponentName(context, PushJobService::class.java)) + .setMinimumLatency(interval(context, delayMs)) + .setRequiredNetworkType(JobInfo.NETWORK_TYPE_ANY) + .setPersisted(true) + .build() + context.getSystemService(JobScheduler::class.java)?.schedule(job) + } catch (e: Exception) { + report(e) + } + } + + // A periodic job, independent of the chain of runs, that restarts the chain when a run was + // dropped (the process killed between runs, a deferred alarm, a throttled job). Scheduled once, + // so its period is not reset by every run. + private fun ensureWatchdog(context: Context) { + val jobs = context.getSystemService(JobScheduler::class.java) ?: return + try { + if (jobs.allPendingJobs.any { it.id == WATCHDOG_JOB_ID }) { + return + } + jobs.schedule( + JobInfo.Builder(WATCHDOG_JOB_ID, ComponentName(context, PushJobService::class.java)) + .setPeriodic(WATCHDOG_INTERVAL_MS) + .setRequiredNetworkType(JobInfo.NETWORK_TYPE_ANY) + .setPersisted(true) + .build(), + ) + } catch (e: Exception) { + report(e) + } + } + private fun cancelSchedule(context: Context) { context.getSystemService(JobScheduler::class.java)?.cancel(JOB_ID) + context.getSystemService(JobScheduler::class.java)?.cancel(NEXT_JOB_ID) + context.getSystemService(JobScheduler::class.java)?.cancel(EXPEDITED_JOB_ID) + context.getSystemService(JobScheduler::class.java)?.cancel(WATCHDOG_JOB_ID) context.getSystemService(AlarmManager::class.java)?.cancel(tickIntent(context)) } @@ -855,6 +1193,63 @@ internal object PushBackground { PendingIntent.FLAG_UPDATE_CURRENT or PendingIntent.FLAG_IMMUTABLE, ) + /** What background delivery can rely on now; see [PushBackgroundStatus]. */ + fun backgroundStatus(context: Context): PushBackgroundStatus { + val app = context.applicationContext + return PushBackgroundStatus( + exactAlarms = Build.VERSION.SDK_INT < Build.VERSION_CODES.S || + app.getSystemService(AlarmManager::class.java)?.canScheduleExactAlarms() == true, + ignoringBatteryOptimizations = app.getSystemService(PowerManager::class.java)?.isIgnoringBatteryOptimizations(app.packageName) == true, + foregroundService = PushStore.foreground(app), + ) + } + + /** + * Open the system screen where the user allows exact alarms (Android 12+). Returns false when + * there is nothing to ask: already allowed, older Android, or the app does not declare + * `SCHEDULE_EXACT_ALARM`. + */ + fun requestExactAlarms(context: Context): Boolean { + if (Build.VERSION.SDK_INT < Build.VERSION_CODES.S || backgroundStatus(context).exactAlarms) { + return false + } + return openSettings(context, Settings.ACTION_REQUEST_SCHEDULE_EXACT_ALARM, Manifest.permission.SCHEDULE_EXACT_ALARM) + } + + /** + * Ask the user to exempt the app from battery optimisation. Returns false when there is + * nothing to ask: already exempt, or the app does not declare + * `REQUEST_IGNORE_BATTERY_OPTIMIZATIONS`. + */ + fun requestIgnoreBatteryOptimizations(context: Context): Boolean { + if (backgroundStatus(context).ignoringBatteryOptimizations) { + return false + } + return openSettings( + context, + Settings.ACTION_REQUEST_IGNORE_BATTERY_OPTIMIZATIONS, + Manifest.permission.REQUEST_IGNORE_BATTERY_OPTIMIZATIONS, + ) + } + + // Open a settings screen for this package when the app declares [permission]. + private fun openSettings(context: Context, action: String, permission: String): Boolean { + val app = context.applicationContext + val declared = runCatching { + @Suppress("DEPRECATION") + app.packageManager.getPackageInfo(app.packageName, PackageManager.GET_PERMISSIONS).requestedPermissions + }.getOrNull()?.contains(permission) == true + if (!declared) { + return false + } + return try { + app.startActivity(Intent(action, Uri.parse("package:${app.packageName}")).addFlags(Intent.FLAG_ACTIVITY_NEW_TASK)) + true + } catch (e: Exception) { + false + } + } + // Allowed exact alarms or exempt from battery optimisation: the alarm is reliable, so the // runs can be spaced further apart. private fun privileged(context: Context): Boolean = @@ -945,6 +1340,7 @@ internal object PushStore { authMethod = it.getString("authMethod"), credential = it.getString("credential"), project = it.getString("project"), + sessionCookieUrl = if (it.isNull("sessionCookieUrl")) null else it.optString("sessionCookieUrl").ifEmpty { null }, ) } val list = json.getJSONArray("entries") @@ -954,6 +1350,7 @@ internal object PushStore { filter = entry.getString("filter"), retry = entry.getBoolean("retry"), title = if (entry.isNull("title")) null else entry.getString("title"), + notifyInForeground = entry.optBoolean("notifyInForeground", false), ) } config to entries @@ -973,11 +1370,12 @@ internal object PushStore { .put("keepAlive", config.keepAlive) .put("authMethod", config.authMethod) .put("credential", config.credential) - .put("project", config.project), + .put("project", config.project) + .put("sessionCookieUrl", config.sessionCookieUrl ?: JSONObject.NULL), ) .put( "entries", - JSONArray(entries.map { JSONObject().put("filter", it.filter).put("retry", it.retry).put("title", it.title ?: JSONObject.NULL) }), + JSONArray(entries.map { JSONObject().put("filter", it.filter).put("retry", it.retry).put("title", it.title ?: JSONObject.NULL).put("notifyInForeground", it.notifyInForeground) }), ) write(context, STATE_FILE, json) } diff --git a/android/src/main/java/io/appwrite/services/PushBridge.kt b/android/src/main/java/io/appwrite/services/PushBridge.kt index d7b79b5d..e5cb5aae 100644 --- a/android/src/main/java/io/appwrite/services/PushBridge.kt +++ b/android/src/main/java/io/appwrite/services/PushBridge.kt @@ -24,6 +24,8 @@ class PushBridge(context: Context, private val events: Events) { fun onMessage(subscriptionId: String, message: PushMessage, ackToken: String) fun onError(message: String) + + fun onConnection(connected: Boolean) } private val appContext = context.applicationContext @@ -37,13 +39,17 @@ class PushBridge(context: Context, private val events: Events) { @Volatile private var errorCallback = false + private var connected: Boolean? = null + /** * Host the subscriptions in [subscriptionsJson] on the background connection described by * [configJson]. * * Config: `{"host", "port", "tls", "tlsInsecure", "clientId" (optional; per user and install - * when empty), "authMethod", "credential", "project"}`. Subscriptions: an array of - * `{"id", "topic", "background", "title" (optional), "retry"}`. + * when empty), "authMethod", "credential", "project", "sessionCookieUrl" (optional: the endpoint + * whose session cookie in the WebView cookie store holds the credential, so background runs + * follow a rotated session)}`. Subscriptions: an array of + * `{"id", "topic", "background", "title" (optional), "retry", "notifyInForeground" (optional)}`. * * [done] is called once the connection is up and every filter is subscribed, or with the * failure's message; the SDK rejects its subscribe with it and reports it itself. @@ -62,6 +68,7 @@ class PushBridge(context: Context, private val events: Events) { authMethod = authMethod, credential = credential, project = json.optString("project"), + sessionCookieUrl = if (json.isNull("sessionCookieUrl")) null else json.optString("sessionCookieUrl").ifEmpty { null }, ) val list = JSONArray(subscriptionsJson) val listeners = (0 until list.length()).map { index -> @@ -80,6 +87,7 @@ class PushBridge(context: Context, private val events: Events) { background = subscription.optBoolean("background", false), title = if (subscription.isNull("title")) null else subscription.optString("title").ifEmpty { null }, retry = subscription.optBoolean("retry", true), + notifyInForeground = subscription.optBoolean("notifyInForeground", false), ) } PushBackground.errorHandler = { error -> @@ -88,21 +96,54 @@ class PushBridge(context: Context, private val events: Events) { } errorCallback } + PushBackground.connectionHandler = { connection(it) } PushBackground.host(appContext, config, listeners) { error -> + if (error == null && PushBackground.isConnected) { + connection(true) + } done(error?.let { it.message ?: it.toString() }) } } + /** Whether the background connection is open now. */ + fun isConnected(): Boolean = PushBackground.isConnected + + private fun connection(open: Boolean) { + synchronized(this) { + if (connected == open) { + return + } + connected = open + } + events.onConnection(open) + } + /** Acknowledge a message once the SDK's callback for it has run. */ fun ack(token: String) { acks.remove(token)?.second?.invoke() } /** Stop hosting the live subscriptions, keeping the ones saved by an earlier run. */ - fun release() = PushBackground.release(appContext) + fun release() { + forgetConnection() + PushBackground.release(appContext) + } /** Stop background delivery and forget every saved subscription (sign-out). */ - fun stop() = PushBackground.stop(appContext) + fun stop() { + forgetConnection() + PushBackground.stop(appContext) + } + + private fun forgetConnection() { + PushBackground.connectionHandler = null + val wasOpen = synchronized(this) { + (connected == true).also { connected = null } + } + if (wasOpen) { + events.onConnection(false) + } + } /** Run background delivery in a foreground service; saved across restarts. */ fun setForeground(enabled: Boolean) = PushBackground.setForeground(appContext, enabled) @@ -110,8 +151,14 @@ class PushBridge(context: Context, private val events: Events) { /** Whether an earlier run saved background subscriptions that are still delivered. */ fun hasSaved(): Boolean = PushBackground.hasSaved(appContext) - /** Resume saved background delivery now instead of at the next scheduled run. */ - fun resume() = PushBackground.resume(appContext) + /** + * Resume saved background delivery now instead of at the next scheduled run, with the app's + * current credential (null when it has none). A rotated session of the same user replaces + * the saved one; another user, or none when [signedOutWhenMissing], drops the saved + * subscriptions without reporting an error. + */ + fun resume(authMethod: String?, credential: String?, signedOutWhenMissing: Boolean) = + PushBackground.resume(appContext, authMethod, credential, signedOutWhenMissing) /** * Record whether the app has an onError callback. Returns the refusal that stopped background @@ -122,6 +169,22 @@ class PushBridge(context: Context, private val events: Events) { return if (registered) PushStore.takeStoppedError(appContext) else null } + /** What background delivery can rely on, as `{"exactAlarms", "ignoringBatteryOptimizations", "foregroundService", "bestEffort"}`. */ + fun backgroundStatus(): String = PushBackground.backgroundStatus(appContext).let { + JSONObject() + .put("exactAlarms", it.exactAlarms) + .put("ignoringBatteryOptimizations", it.ignoringBatteryOptimizations) + .put("foregroundService", it.foregroundService) + .put("bestEffort", it.bestEffort) + .toString() + } + + /** Open the system screen that allows exact alarms; false when there is nothing to ask. */ + fun requestExactAlarms(): Boolean = PushBackground.requestExactAlarms(appContext) + + /** Ask to exempt the app from battery optimisation; false when there is nothing to ask. */ + fun requestIgnoreBatteryOptimizations(): Boolean = PushBackground.requestIgnoreBatteryOptimizations(appContext) + /** The default client id for a credential, so a foreground connection shares its session. */ fun defaultClientId(authMethod: String, credential: String): String = defaultPushClientId(appContext, authMethod, credential) diff --git a/android/src/main/java/io/appwrite/services/PushCore.kt b/android/src/main/java/io/appwrite/services/PushCore.kt index 20c56242..15f86b6a 100644 --- a/android/src/main/java/io/appwrite/services/PushCore.kt +++ b/android/src/main/java/io/appwrite/services/PushCore.kt @@ -12,6 +12,7 @@ import com.hivemq.client.mqtt.mqtt5.Mqtt5ClientConfig import com.hivemq.client.mqtt.mqtt5.auth.Mqtt5EnhancedAuthMechanism import com.hivemq.client.mqtt.mqtt5.exceptions.Mqtt5ConnAckException import com.hivemq.client.mqtt.mqtt5.exceptions.Mqtt5DisconnectException +import com.hivemq.client.mqtt.mqtt5.lifecycle.Mqtt5ClientReconnector import com.hivemq.client.mqtt.mqtt5.message.auth.Mqtt5Auth import com.hivemq.client.mqtt.mqtt5.message.auth.Mqtt5AuthBuilder import com.hivemq.client.mqtt.mqtt5.message.auth.Mqtt5EnhancedAuthBuilder @@ -30,6 +31,7 @@ import java.util.concurrent.CompletableFuture import java.util.concurrent.ExecutionException import java.util.concurrent.TimeUnit import java.util.concurrent.atomic.AtomicBoolean +import java.util.concurrent.atomic.AtomicReference import javax.net.ssl.ManagerFactoryParameters import javax.net.ssl.TrustManager import javax.net.ssl.TrustManagerFactory @@ -49,6 +51,24 @@ data class PushMessage( get() = String(payload, Charsets.UTF_8) } +/** + * What background delivery can rely on. Without exact alarms, the scheduled wake-ups are inexact + * and Doze can defer them, so delivery is [bestEffort] unless foreground mode keeps the + * connection open. + */ +data class PushBackgroundStatus( + /** The app may schedule exact alarms (`SCHEDULE_EXACT_ALARM`, granted). */ + val exactAlarms: Boolean, + /** The app is exempt from battery optimisation. */ + val ignoringBatteryOptimizations: Boolean, + /** Foreground mode (`setForeground(true)`) keeps the connection open in a service. */ + val foregroundService: Boolean, +) { + /** Wake-ups may be deferred by Doze, so messages can arrive late while the app is closed. */ + val bestEffort: Boolean + get() = !exactAlarms && !foregroundService +} + /** Connection parameters shared between the in-process client and the foreground service. */ internal data class PushConfig( val host: String, @@ -60,6 +80,10 @@ internal data class PushConfig( val authMethod: String, val credential: String, val project: String, + // The endpoint whose `a_session_` cookie in the WebView cookie store holds the app's + // session, when its credential comes from there (React Native). Background runs read it, so + // they keep up with a rotated session while no app code runs. + val sessionCookieUrl: String? = null, ) /** @@ -110,6 +134,8 @@ internal fun buildPushClient( onDisconnected: (() -> Unit)? = null, onError: ((Throwable) -> Unit)? = null, onAuthRefused: ((Throwable) -> Unit)? = null, + credential: () -> String? = { config.credential }, + retired: AtomicBoolean? = null, ): Mqtt5AsyncClient { // config.clientId is stable per user (see buildConfig), so a reconnect rebuild reuses it // and resumes the same session. @@ -130,11 +156,25 @@ internal fun buildPushClient( // Whether this client has connected yet: until it has, a refused CONNECT is final. val connected = AtomicBoolean(false) + val client = AtomicReference() builder = builder.addConnectedListener { + // A reconnect that was already under way when the client was retired disconnects at once, + // so it never competes with its replacement for the broker session. + if (retired?.get() == true) { + client.get()?.disconnect() + return@addConnectedListener + } connected.set(true) onConnected?.invoke() } builder = builder.addDisconnectedListener { context -> + // A client its owner replaced or closed ([retired]) never reconnects: one still retrying + // in the background would otherwise come back as a second connection for the same id. + if (retired?.get() == true) { + context.reconnector.reconnect(false) + onDisconnected?.invoke() + return@addDisconnectedListener + } // A refused first CONNECT (bad credential, rate limit) is final: stop the automatic // reconnect so pushConnect fails with the broker's reason instead of retrying forever. val refused = !connected.get() && context.cause is Mqtt5ConnAckException @@ -142,8 +182,16 @@ internal fun buildPushClient( // re-auth) is final too, instead of being retried with the same credential forever. val authRefused = onAuthRefused != null && connected.get() && context.source != MqttDisconnectSource.USER && isAuthRefusal(context.cause) - if (refused || authRefused) { + // Authentication was aborted because the app has no current credential (signed out): there + // is nothing to reconnect with, so the attempt fails now instead of retrying until timeout. + val noCredential = generateSequence(context.cause) { it.cause }.any { it is PushNoCredentialException } + if (refused || authRefused || noCredential) { context.reconnector.reconnect(false) + } else if (context.reconnector.isReconnect) { + // HiveMQ's automatic reconnect sends a CONNECT it builds itself, without user + // properties: the broker then misses the project and refuses the credential. Reconnect + // with the same CONNECT this client connects with instead. + (context.reconnector as? Mqtt5ClientReconnector)?.connect(pushConnectMessage(config, credential)) } onDisconnected?.invoke() if (!refused) { @@ -168,7 +216,7 @@ internal fun buildPushClient( } } - return builder.buildAsync() + return builder.buildAsync().also { client.set(it) } } /** @@ -207,10 +255,31 @@ internal fun isAuthRefusal(error: Throwable?): Boolean = when (error) { else -> false } -/** Send the CONNECT and wait for the CONNACK, at most [timeoutSeconds] when given. */ -internal fun Mqtt5AsyncClient.sendPushConnect(config: PushConfig, timeoutSeconds: Long? = null) { - val connAck = connectWith() - .enhancedAuth(PushAuthMechanism(config.authMethod, config.credential)) +/** + * Send the CONNECT and wait for the CONNACK, at most [timeoutSeconds] when given. [credential] + * supplies the credential each time one is sent, including on the client's automatic reconnects + * and re-auths, so they can use a session rotated since this CONNECT; null aborts the attempt. + */ +internal fun Mqtt5AsyncClient.sendPushConnect( + config: PushConfig, + timeoutSeconds: Long? = null, + credential: () -> String? = { config.credential }, +) { + val connAck = connect(pushConnectMessage(config, credential)) + if (timeoutSeconds == null) { + connAck.get() + } else { + connAck.get(timeoutSeconds, TimeUnit.SECONDS) + } +} + +/** + * The CONNECT for [config]: enhanced auth with the credential from [credential], and the project + * as a user property, which the broker resolves the credential against. + */ +internal fun pushConnectMessage(config: PushConfig, credential: () -> String?): Mqtt5Connect = + Mqtt5Connect.builder() + .enhancedAuth(PushAuthMechanism(config.authMethod, config.project, credential)) // Clean start is always off, so the broker keeps this client's session and can // redeliver missed messages to QoS-1 subscriptions on reconnect. .cleanStart(false) @@ -218,13 +287,7 @@ internal fun Mqtt5AsyncClient.sendPushConnect(config: PushConfig, timeoutSeconds .userProperties() .add("projectId", config.project) .applyUserProperties() - .send() - if (timeoutSeconds == null) { - connAck.get() - } else { - connAck.get(timeoutSeconds, TimeUnit.SECONDS) - } -} + .build() /** * The error a lost connection reports, or null when it is not an error (a disconnect the app @@ -266,11 +329,14 @@ internal fun Mqtt5Publish.toPushMessage(): PushMessage = /** * Single-step MQTT 5 enhanced auth: the method and credential are placed in the CONNECT - * packet and the broker accepts them in the CONNACK, with no AUTH round-trip. + * packet and the broker accepts them in the CONNACK, with no AUTH round-trip. The credential is + * read from [credential] each time, so an automatic reconnect or re-auth sends the current one; + * when it has none, the attempt fails instead of sending a credential known to be stale. */ internal class PushAuthMechanism( private val method: String, - private val credential: String, + private val project: String, + private val credential: () -> String?, ) : Mqtt5EnhancedAuthMechanism { override fun getMethod(): MqttUtf8String = MqttUtf8String.of(method) @@ -282,7 +348,8 @@ internal class PushAuthMechanism( connect: Mqtt5Connect, authBuilder: Mqtt5EnhancedAuthBuilder, ): CompletableFuture { - authBuilder.data(credential.toByteArray(Charsets.UTF_8)) + val current = credential() ?: return failedAuth() + authBuilder.data(current.toByteArray(Charsets.UTF_8)) return CompletableFuture.completedFuture(null) } @@ -290,7 +357,12 @@ internal class PushAuthMechanism( clientConfig: Mqtt5ClientConfig, authBuilder: Mqtt5AuthBuilder, ): CompletableFuture { - authBuilder.data(credential.toByteArray(Charsets.UTF_8)) + val current = credential() ?: return failedAuth() + // The broker resolves a re-auth against the project in its user properties too. + authBuilder.data(current.toByteArray(Charsets.UTF_8)) + .userProperties() + .add("projectId", project) + .applyUserProperties() return CompletableFuture.completedFuture(null) } @@ -317,8 +389,14 @@ internal class PushAuthMechanism( override fun onAuthError(clientConfig: Mqtt5ClientConfig, cause: Throwable) = Unit override fun onReAuthError(clientConfig: Mqtt5ClientConfig, cause: Throwable) = Unit + + private fun failedAuth(): CompletableFuture = + CompletableFuture().apply { completeExceptionally(PushNoCredentialException()) } } +/** Authentication aborted because there is no current credential to send (the app signed out). */ +internal class PushNoCredentialException : IllegalStateException("No current credential to authenticate with") + /** * A [TrustManagerFactory] that accepts any certificate — used when `tlsInsecure` is set * (for a broker whose certificate does not chain to a public root). @@ -353,6 +431,8 @@ internal class PushListener( // Called instead of [callback] when set, with the acknowledgement to run once the message was // handled: the React Native and Flutter bridges acknowledge after their callback has run. val acknowledgingCallback: ((PushMessage, () -> Unit) -> Unit)? = null, + // Also post a background notification while the app is visible. + val notifyInForeground: Boolean = false, ) /** MQTT topic-filter match with '+' (single level) and '#' (multi level). */ diff --git a/android/src/main/java/io/appwrite/services/PushOpenActivity.kt b/android/src/main/java/io/appwrite/services/PushOpenActivity.kt new file mode 100644 index 00000000..354ca6f0 --- /dev/null +++ b/android/src/main/java/io/appwrite/services/PushOpenActivity.kt @@ -0,0 +1,57 @@ +package io.appwrite.services + +import android.app.Activity +import android.content.Intent +import android.os.Bundle + +class PushOpenActivity : Activity() { + override fun onCreate(savedInstanceState: Bundle?) { + super.onCreate(savedInstanceState) + val topic = intent.getStringExtra(PushBackground.EXTRA_TOPIC) + val payload = intent.getStringExtra(PushBackground.EXTRA_PAYLOAD) + packageManager.getLaunchIntentForPackage(packageName)?.let { launch -> + launch.addFlags(Intent.FLAG_ACTIVITY_NEW_TASK or Intent.FLAG_ACTIVITY_SINGLE_TOP) + if (topic != null && payload != null) { + launch.putExtra(PushBackground.EXTRA_TOPIC, topic).putExtra(PushBackground.EXTRA_PAYLOAD, payload) + } + startActivity(launch) + } + // After opening the app, so a screen a tap callback opens lands on top of it. + if (topic != null && payload != null) { + PushTaps.record(PushTap(topic, payload)) + } + finish() + } +} + +internal data class PushTap(val topic: String, val payload: String) + +internal object PushTaps { + private var pending: PushTap? = null + private val listeners = mutableListOf<(PushTap) -> Unit>() + + fun record(tap: PushTap) { + val deliver = synchronized(this) { + listeners.toList().also { + if (it.isEmpty()) { + pending = tap + } + } + } + deliver.forEach { it(tap) } + } + + @Synchronized + fun take(): PushTap? = pending.also { pending = null } + + fun listen(onTap: (PushTap) -> Unit): () -> Unit { + synchronized(this) { + listeners.add(onTap) + } + return { + synchronized(this) { + listeners.remove(onTap) + } + } + } +} diff --git a/android/src/main/java/io/appwrite/services/PushWakeups.kt b/android/src/main/java/io/appwrite/services/PushWakeups.kt index 091fd528..c6d56165 100644 --- a/android/src/main/java/io/appwrite/services/PushWakeups.kt +++ b/android/src/main/java/io/appwrite/services/PushWakeups.kt @@ -12,11 +12,15 @@ import android.content.Intent */ class PushJobService : JobService() { override fun onStartJob(params: JobParameters): Boolean { - PushBackground.tick(this) { jobFinished(params, false) } + // A chained job arms the next under the other chained id; the watchdog and an expedited + // run (handed over by the alarm) only run once. Each stays up while missed messages arrive. + val chained = params.jobId.takeIf { it == PushBackground.JOB_ID || it == PushBackground.NEXT_JOB_ID } + PushBackground.tick(this, chained, PushBackground.JOB_DRAIN_MS) { jobFinished(params, false) } return true } - // The run re-arms the next one itself, so a stopped run needs no retry. + // The run re-arms the next one itself, and the watchdog repeats on its own, so a stopped run + // needs no retry. override fun onStopJob(params: JobParameters): Boolean = false } @@ -29,8 +33,13 @@ class PushAlarmReceiver : BroadcastReceiver() { if (intent.action != PushBackground.ACTION_TICK) { return } + // On Android 12+ the run goes to an expedited job, which keeps the process runnable while it + // connects and the broker replays; it runs here when the system refuses one. + if (PushBackground.runExpedited(context)) { + return + } val pending = goAsync() - PushBackground.tick(context) { pending.finish() } + PushBackground.tick(context, drainMs = PushBackground.RECEIVER_DRAIN_MS) { pending.finish() } } } diff --git a/package-lock.json b/package-lock.json index 0a128e84..33230e5f 100644 --- a/package-lock.json +++ b/package-lock.json @@ -1,12 +1,12 @@ { "name": "react-native-appwrite", - "version": "1.2.0-rc.1", + "version": "1.2.0-rc.4", "lockfileVersion": 3, "requires": true, "packages": { "": { "name": "react-native-appwrite", - "version": "1.2.0-rc.1", + "version": "1.2.0-rc.4", "license": "BSD-3-Clause", "dependencies": { "buffer": "6.0.3", diff --git a/package.json b/package.json index 16c14202..12b4f76a 100644 --- a/package.json +++ b/package.json @@ -2,7 +2,7 @@ "name": "react-native-appwrite", "homepage": "https://appwrite.io/support", "description": "Appwrite is an open-source self-hosted backend server that abstracts and simplifies complex and repetitive development tasks behind a very simple REST API", - "version": "1.2.0-rc.1", + "version": "1.2.0-rc.4", "license": "BSD-3-Clause", "main": "dist/cjs/sdk.js", "exports": { diff --git a/src/client.ts b/src/client.ts index 605f530c..4bc8575c 100644 --- a/src/client.ts +++ b/src/client.ts @@ -218,7 +218,7 @@ class Client { 'x-sdk-name': 'React Native', 'x-sdk-platform': 'client', 'x-sdk-language': 'reactnative', - 'x-sdk-version': '1.2.0-rc.1', + 'x-sdk-version': '1.2.0-rc.4', 'X-Appwrite-Response-Format': '2.3.0', }; diff --git a/src/index.ts b/src/index.ts index 18a4e75a..579ed73a 100644 --- a/src/index.ts +++ b/src/index.ts @@ -18,7 +18,9 @@ export { VectorsDB } from './services/vectors-db'; export { Realtime } from './services/realtime'; export { Push } from './services/push'; export type { + PushBackgroundStatus, PushMessage, + PushNotificationOpened, PushSubscription, SubscribeOptions, MessageCallback, diff --git a/src/react-native-shim.d.ts b/src/react-native-shim.d.ts index 914861ba..cd64d011 100644 --- a/src/react-native-shim.d.ts +++ b/src/react-native-shim.d.ts @@ -1,4 +1,8 @@ declare module 'react-native' { + export const AppState: { + readonly currentState: string | null; + }; + export const Platform: { readonly OS: string; readonly Version: number | string; @@ -107,6 +111,19 @@ declare module 'expo-notifications' { channelId: string, channel: { name: string; importance: number }, ): Promise; + export interface NotificationResponse { + notification: { + request: { content: { data?: Record | null } }; + }; + } + + export function getLastNotificationResponseAsync(): Promise; + // Missing from older expo-notifications versions. + export const clearLastNotificationResponseAsync: + (() => Promise) | undefined; + export function addNotificationResponseReceivedListener( + listener: (response: NotificationResponse) => void, + ): { remove(): void }; export function scheduleNotificationAsync(request: { content: { title: string; diff --git a/src/services/push.ts b/src/services/push.ts index 8ae70847..880e4f4c 100644 --- a/src/services/push.ts +++ b/src/services/push.ts @@ -12,7 +12,12 @@ import mqtt, { type IPublishPacket, type MqttClient, } from 'mqtt'; -import { NativeEventEmitter, NativeModules, Platform } from 'react-native'; +import { + AppState, + NativeEventEmitter, + NativeModules, + Platform, +} from 'react-native'; import { Client } from '../client'; import { createTcpStream } from '../lib/tcp-stream'; @@ -31,6 +36,26 @@ export interface PushMessage { export type MessageCallback = (message: PushMessage) => void | Promise; +/** What Android background delivery can rely on; see {@link Push.backgroundStatus}. */ +export interface PushBackgroundStatus { + /** The app may schedule exact alarms (`SCHEDULE_EXACT_ALARM`, granted). */ + exactAlarms: boolean; + /** The app is exempt from battery optimisation. */ + ignoringBatteryOptimizations: boolean; + /** Foreground mode ({@link Push.setForeground}) keeps the connection open in a service. */ + foregroundService: boolean; + /** Wake-ups may be deferred by Doze, so messages can arrive late while the app is closed. */ + bestEffort: boolean; +} + +/** A background notification the user tapped. */ +export interface PushNotificationOpened { + /** The topic the message was published to. */ + topic: string; + /** The `data` sent with the message (e.g. `createPush`), parsed; `{}` when it had none. */ + data: Record; +} + /** Per-subscription options passed to `subscribe` and `PushSubscription.update`. */ export interface SubscribeOptions { /** @@ -52,6 +77,12 @@ export interface SubscribeOptions { background?: boolean; /** Notification title for this subscription's messages. Defaults to the message topic. */ title?: string; + /** + * Also post the background notification while the app is on screen. Defaults to false: + * while the app is visible it shows the message itself through the callback, so + * notifications are posted only while it is backgrounded or closed. + */ + notifyInForeground?: boolean; /** * Retry delivery of messages missed while disconnected. `true` (the default) subscribes * at QoS 1 so the broker holds this topic's messages and redelivers them on reconnect; @@ -118,6 +149,7 @@ export class Push extends Service { private readonly native: NativePush | null = nativePush(); /** Whether the subscriptions currently live on the native module's connection. */ private nativeActive = false; + private nativeOpen = false; /** Bumped when the subscriptions move to the native host, superseding an open() in flight. */ private connectionEpoch = 0; /** Bumped by close(), so a subscribe that was waiting on a connection stops instead of reopening it. */ @@ -152,6 +184,7 @@ export class Push extends Service { callback: MessageCallback; background: boolean; title?: string; + notifyInForeground: boolean; qos: 0 | 1; } >(); @@ -201,8 +234,12 @@ export class Push extends Service { this.tlsInsecure = endpointInsecure; // Resume background delivery saved by an earlier run now, instead of at its next - // scheduled wake-up. - this.native?.resume().catch((err) => this.report(err)); + // scheduled wake-up, with the current session: a rotated session of the same user + // replaces the saved one, and signed out (no session cookie) drops the saved + // subscriptions. + if (this.native) { + this.resumeNative(this.native).catch((err) => this.report(err)); + } } // The broker subscription for a filter uses the highest QoS any local subscription on @@ -255,6 +292,119 @@ export class Push extends Service { await this.native?.setForeground(enabled); } + /** + * Android: what background delivery can rely on. When `bestEffort` is true the scheduled + * wake-ups are inexact and Doze can defer them, so messages may arrive late while the app is + * closed; explain why, then ask with {@link requestExactAlarms} or + * {@link requestIgnoreBatteryOptimizations}. Null elsewhere. + */ + async backgroundStatus(): Promise { + const json = await this.native?.backgroundStatus(); + return json ? (JSON.parse(json) as PushBackgroundStatus) : null; + } + + /** + * Android 12+: open the system screen where the user allows exact alarms, for punctual + * background wake-ups. Call it from a user action, never on its own. Resolves false when + * there is nothing to ask: already allowed, older Android, the app does not declare + * `SCHEDULE_EXACT_ALARM`, or not Android. + */ + async requestExactAlarms(): Promise { + return (await this.native?.requestExactAlarms()) ?? false; + } + + /** + * Android: ask the user to exempt the app from battery optimisation. Call it from a user + * action. Resolves false when there is nothing to ask: already exempt, the app does not + * declare `REQUEST_IGNORE_BATTERY_OPTIMIZATIONS`, or not Android. + */ + async requestIgnoreBatteryOptimizations(): Promise { + return ( + (await this.native?.requestIgnoreBatteryOptimizations()) ?? false + ); + } + + /** + * The background notification whose tap launched the app, or null. Call it once at startup: + * on Android a tap is reported only once. Elsewhere it reads `expo-notifications`. + * + * ```ts + * const opened = await push.getInitialNotification(); + * if (opened) openSale(opened.data.saleId); + * ``` + */ + async getInitialNotification(): Promise { + if (Platform.OS === 'android') { + const json = await this.native?.getInitialNotification(); + if (!json) { + return null; + } + const opened = JSON.parse(json) as NativeOpened; + return toOpened(opened.topic, opened.payload); + } + const notifications = expoNotifications(); + if (!notifications) { + return null; + } + const response = await notifications.getLastNotificationResponseAsync(); + const opened = response ? openedFromExpo(response) : null; + if (opened) { + await notifications.clearLastNotificationResponseAsync?.(); + } + return opened; + } + + /** + * Call [callback] each time the user taps a background notification while the app is + * running, including in the background. Returns a function that stops listening. The tap + * that launched the app comes from `getInitialNotification()` instead. + * + * ```ts + * const stop = push.onNotificationOpened(({ data }) => openSale(data.saleId)); + * ``` + */ + onNotificationOpened( + callback: (opened: PushNotificationOpened) => void, + ): () => void { + if (Platform.OS === 'android') { + if (!this.native) { + return () => undefined; + } + const native = this.native; + const emitter = new NativeEventEmitter(native); + const subscription = emitter.addListener( + 'AppwritePushOpened', + (event: NativeOpened) => + callback(toOpened(event.topic, event.payload)), + ); + native.listenOpened(true).catch((err) => this.report(err)); + let listening = true; + return () => { + if (!listening) { + return; + } + listening = false; + subscription.remove(); + native.listenOpened(false).catch((err) => this.report(err)); + }; + } + const notifications = expoNotifications(); + if (!notifications) { + return () => undefined; + } + const listener = ( + response: import('expo-notifications').NotificationResponse, + ): void => { + const opened = openedFromExpo(response); + if (opened) { + callback(opened); + } + }; + const subscription = + notifications.addNotificationResponseReceivedListener(listener); + return () => subscription.remove(); + } + /** * The signed-in user's own topic, `users/`, for a topic-less subscribe. The id comes * from the credential the connection authenticates with: the JWT when one is set (never @@ -436,6 +586,7 @@ export class Push extends Service { callback, background, title: options.title, + notifyInForeground: options.notifyInForeground ?? false, qos, }); return id; @@ -487,6 +638,7 @@ export class Push extends Service { callback, background, title: options.title, + notifyInForeground: options.notifyInForeground ?? false, qos, }); acks.push( @@ -616,6 +768,9 @@ export class Push extends Service { if (next.title !== undefined) { entry.title = next.title; } + if (next.notifyInForeground !== undefined) { + entry.notifyInForeground = next.notifyInForeground; + } if (next.retry !== undefined) { const q: 0 | 1 = next.retry ? 1 : 0; if (q !== entry.qos) { @@ -670,6 +825,9 @@ export class Push extends Service { if (next.title !== undefined) { entry.title = next.title; } + if (next.notifyInForeground !== undefined) { + entry.notifyInForeground = next.notifyInForeground; + } if (next.retry !== undefined) { entry.qos = next.retry ? 1 : 0; } @@ -722,18 +880,41 @@ export class Push extends Service { } } this.joinNative(native); - return enqueueNative(() => this.sendToNative(native)); + // Instances that joined an open connection hear onOpen now; the others when it opens. + return enqueueNative(() => this.sendToNative(native)).then((open) => { + if (open) { + for (const push of nativeHosts) { + push.nativeConnection(true); + } + } + }); } - /** Host every native host's subscriptions with this Push's credential. */ - private async sendToNative(native: NativePush): Promise { + /** This Push's view of the native host's connection, reported once per change. */ + private nativeConnection(open: boolean): void { + if (!this.nativeActive || this.nativeOpen === open) { + return; + } + this.nativeOpen = open; + if (open) { + this.onOpenCb?.(); + } else { + this.onCloseCb?.(); + } + } + + /** + * Host every native host's subscriptions with this Push's credential. Resolves with whether + * the connection is open (false when there was nothing to host). + */ + private async sendToNative(native: NativePush): Promise { // Host what is subscribed when this runs, so a close() or an unsubscribe queued // meanwhile is not undone by an older request. const subscriptions = [...nativeHosts].flatMap((push) => push.nativeEntries(), ); if (subscriptions.length === 0) { - return; + return false; } if ([...nativeHosts].some((push) => push.hasBackgroundSubs())) { requestNotificationPermission(); @@ -748,14 +929,21 @@ export class Push extends Service { authMethod, credential, project: this.client.config.project ?? '', + // A session from the cookie store: background runs re-read it there, so they follow + // a rotated session while no JS runs. + sessionCookieUrl: + credential === this.cookieSession && !this.client.config.jwt + ? (this.client.config.endpoint ?? '') + : '', }; await native.setErrorCallback( [...nativeHosts].some((push) => push.onErrorCb !== undefined), ); - await native.host( + const open = await native.host( JSON.stringify(config), JSON.stringify(subscriptions), ); + return open === true; } /** Move this Push's subscriptions onto the native host, closing its own connection. */ @@ -770,6 +958,9 @@ export class Push extends Service { this.everConnected = false; this.listenNative(native); nativeHosts.add(this); + if (!this.nativeActive) { + this.nativeOpen = false; + } this.nativeActive = true; } @@ -807,6 +998,7 @@ export class Push extends Service { background: sub.background, title: sub.title ?? null, retry: sub.qos === 1, + notifyInForeground: sub.notifyInForeground, })); } @@ -893,6 +1085,11 @@ export class Push extends Service { (event: { message: string }) => this.report(new Error(event.message)), ), + emitter.addListener( + 'AppwritePushConnection', + (event: { connected: boolean }) => + this.nativeConnection(event.connected), + ), ]; } @@ -1111,6 +1308,27 @@ export class Push extends Service { return this.connecting.then(() => this.mqtt!); } + /** + * Resume saved native delivery with the current credential. Only a cookie lookup that + * succeeded and found no session counts as signed out; a failed or impossible lookup keeps + * the saved subscriptions. + */ + private async resumeNative(native: NativePush): Promise { + const { jwt, session, endpoint, project } = this.client.config; + let signedOut = false; + if (!jwt && !session) { + const lookup = await lookupSessionCookie(endpoint, project); + this.cookieSession = lookup.session; + signedOut = lookup.found === false; + } + const current = this.currentCredential(); + await native.resume( + current?.authMethod ?? null, + current?.credential ?? null, + signedOut, + ); + } + private async loadCookieSession(): Promise { const { jwt, session, endpoint, project } = this.client.config; if (jwt || session) { @@ -1119,6 +1337,18 @@ export class Push extends Service { this.cookieSession = await sessionCookie(endpoint, project); } + /** The credential the connection would use, or null when there is none. */ + private currentCredential(): { + authMethod: AuthMethod; + credential: string; + } | null { + try { + return this.credential(); + } catch { + return null; + } + } + /** The credential set on the client (via Client.setJWT / setSession). */ private credential(): { authMethod: AuthMethod; credential: string } { if (this.client.config.jwt) { @@ -1345,38 +1575,62 @@ export class Push extends Service { qos: packet.qos, }; const serverTitle = notificationContent(message).title; - const titles = new Set(); + const postedTitles = new Set(); + // While the app is on screen it shows the message itself: only subscriptions that + // opted in with notifyInForeground also post a notification then. + const foreground = AppState.currentState === 'active'; for (const sub of this.subscriptions.values()) { if (matches(sub.topic, message.topic)) { await sub.callback(message); // Notification is per-subscription: only subs that opted in post one, each - // with its own title. A title the server sent replaces theirs, so one posts. - if (sub.background) { - titles.add(serverTitle ?? sub.title ?? message.topic); + // title once. A title the server sent replaces theirs, so it posts once. + const title = serverTitle ?? sub.title ?? message.topic; + if ( + sub.background && + (!foreground || sub.notifyInForeground) && + !postedTitles.has(title) + ) { + postedTitles.add(title); + await this.notify(message, title); } } } - for (const title of titles) { - await this.notify(message, title); - } } } /** The SDK's native Android module (see android/), which hosts background delivery. */ interface NativePush { - host(config: string, subscriptions: string): Promise; + /** Resolves with whether the connection is open once every filter is subscribed. */ + host(config: string, subscriptions: string): Promise; ack(token: string): void; release(): Promise; stop(): Promise; setForeground(enabled: boolean): Promise; hasSaved(): Promise; - resume(): Promise; + resume( + authMethod: string | null, + credential: string | null, + signedOutWhenMissing: boolean, + ): Promise; + /** {@link PushBackgroundStatus} as JSON. */ + backgroundStatus(): Promise; + requestExactAlarms(): Promise; + requestIgnoreBatteryOptimizations(): Promise; setErrorCallback(registered: boolean): Promise; defaultClientId(authMethod: string, credential: string): Promise; + /** The tapped notification that launched the app, as {@link NativeOpened} JSON. */ + getInitialNotification(): Promise; + listenOpened(listening: boolean): Promise; addListener(eventName: string): void; removeListeners(count: number): void; } +/** A tapped notification the native module reports: its topic and raw payload. */ +interface NativeOpened { + topic: string; + payload: string; +} + /** A message the native module delivers for a hosted subscription. */ interface NativeMessage { id: string; @@ -1424,6 +1678,44 @@ function requestNotificationPermission(): void { .catch(() => undefined); } +/** A tapped notification's topic and message `data`, parsed from its payload. */ +function toOpened(topic: string, payload: string): PushNotificationOpened { + let data: unknown; + try { + data = (JSON.parse(payload) as { data?: unknown })?.data; + } catch { + data = undefined; + } + return { + topic, + data: + typeof data === 'object' && data !== null && !Array.isArray(data) + ? (data as Record) + : {}, + }; +} + +/** A tap on a notification this SDK posted through expo-notifications, else null. */ +function openedFromExpo( + response: import('expo-notifications').NotificationResponse, +): PushNotificationOpened | null { + const data = response.notification.request.content.data; + return typeof data?.topic === 'string' && typeof data.payload === 'string' + ? toOpened(data.topic, data.payload) + : null; +} + +/** expo-notifications when the app has installed it, else null. */ +function expoNotifications(): typeof import('expo-notifications') | null { + try { + /* eslint-disable @typescript-eslint/no-require-imports */ + return require('expo-notifications'); + /* eslint-enable @typescript-eslint/no-require-imports */ + } catch { + return null; + } +} + /** What a message's `notification` block asks a background notification to show. */ interface NotificationContent { /** Whether the message has a `notification` block at all. */ @@ -1544,15 +1836,29 @@ async function sessionCookie( endpoint: string | undefined, project: string | undefined, ): Promise { + return (await lookupSessionCookie(endpoint, project)).session; +} + +/** + * Look up the session cookie. `found` is false only when the lookup succeeded and there is no + * session, and null when it could not be looked up (no cookie module, endpoint or project, or + * the lookup failed). + */ +async function lookupSessionCookie( + endpoint: string | undefined, + project: string | undefined, +): Promise<{ found: boolean | null; session: string }> { const cookies = NativeModules?.AppwriteCookies as NativeCookies | undefined; if (!cookies || !endpoint || !project) { - return ''; + return { found: null, session: '' }; } try { const value = await cookies.session(endpoint, project); - return value ? decodeURIComponent(value) : ''; + return value + ? { found: true, session: decodeURIComponent(value) } + : { found: false, session: '' }; } catch { - return ''; + return { found: null, session: '' }; } }