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: '' };
}
}