diff --git a/app/androidApp/build.gradle.kts b/app/androidApp/build.gradle.kts index 89cc72b..209f15e 100644 --- a/app/androidApp/build.gradle.kts +++ b/app/androidApp/build.gradle.kts @@ -39,4 +39,11 @@ dependencies { implementation("androidx.compose.ui:ui") implementation("androidx.activity:activity-compose:1.9.2") implementation("androidx.core:core-splashscreen:1.0.1") +} + +// M5: apply the google-services plugin only when a Firebase project is +// configured (google-services.json present). Without it the FCM service is +// inert and the ntfy listener is the push path. +if (file("google-services.json").exists()) { + apply(plugin = "com.google.gms.google-services") } \ No newline at end of file diff --git a/app/androidApp/src/main/AndroidManifest.xml b/app/androidApp/src/main/AndroidManifest.xml index 3e7389b..9d7b3fc 100644 --- a/app/androidApp/src/main/AndroidManifest.xml +++ b/app/androidApp/src/main/AndroidManifest.xml @@ -3,6 +3,11 @@ + + + + + + + + + + + + + + + + + @@ -30,6 +48,22 @@ android:name="android.support.FILE_PROVIDER_PATHS" android:resource="@xml/iris_file_paths" /> + + + + + + + + + + \ No newline at end of file diff --git a/app/androidApp/src/main/kotlin/dev/iris/app/MainActivity.kt b/app/androidApp/src/main/kotlin/dev/iris/app/MainActivity.kt index 65cab68..463c680 100644 --- a/app/androidApp/src/main/kotlin/dev/iris/app/MainActivity.kt +++ b/app/androidApp/src/main/kotlin/dev/iris/app/MainActivity.kt @@ -1,19 +1,85 @@ package dev.iris.app +import android.Manifest +import android.content.Intent +import android.content.pm.PackageManager +import android.os.Build import android.os.Bundle import androidx.activity.ComponentActivity import androidx.activity.compose.setContent +import androidx.activity.result.contract.ActivityResultContracts +import androidx.compose.runtime.getValue +import androidx.compose.runtime.mutableStateOf +import androidx.core.content.ContextCompat import iris.IrisApp import iris.platform.AndroidEnv import iris.platform.AndroidSecureStore +import iris.platform.AppBridge +import iris.platform.NtfyListenerService class MainActivity : ComponentActivity() { + private val deepLinkChatId = mutableStateOf(null) + private val deepLinkThreadId = mutableStateOf(null) + + private val notificationPermission = + registerForActivityResult(ActivityResultContracts.RequestPermission()) { /* result ignored */ } + override fun onCreate(savedInstanceState: Bundle?) { super.onCreate(savedInstanceState) AndroidEnv.context = applicationContext + AppBridge.controller = null val store = AndroidSecureStore(applicationContext) + handleDeepLink(intent) + requestNotificationPermission() + startNtfyListener(store) setContent { - IrisApp(store) + val chatId by deepLinkChatId + val threadId by deepLinkThreadId + IrisApp(store = store, deepLinkChatId = chatId, deepLinkThreadId = threadId) } } + + override fun onNewIntent(intent: Intent) { + super.onNewIntent(intent) + handleDeepLink(intent) + } + + override fun onResume() { + super.onResume() + AppBridge.foreground = true + } + + override fun onPause() { + super.onPause() + AppBridge.foreground = false + } + + /** Extract chat_id / thread_id from a deep-link intent (custom action + * extras or an iris://chat/?thread= URI). */ + private fun handleDeepLink(intent: Intent?) { + val data = intent?.data + val chatId = intent?.getStringExtra("chat_id") + ?: data?.pathSegments?.firstOrNull() + val threadId = intent?.getStringExtra("thread_id") + ?: data?.getQueryParameter("thread") + if (!chatId.isNullOrBlank()) { + deepLinkChatId.value = chatId + deepLinkThreadId.value = threadId + } + } + + private fun requestNotificationPermission() { + if (Build.VERSION.SDK_INT >= Build.VERSION_CODES.TIRAMISU) { + if (ContextCompat.checkSelfPermission(this, Manifest.permission.POST_NOTIFICATIONS) + != PackageManager.PERMISSION_GRANTED + ) { + notificationPermission.launch(Manifest.permission.POST_NOTIFICATIONS) + } + } + } + + private fun startNtfyListener(store: AndroidSecureStore) { + if (store.ntfyTopic.isBlank()) return + ContextCompat.startForegroundService(this, Intent(this, NtfyListenerService::class.java)) + } } \ No newline at end of file diff --git a/app/build.gradle.kts b/app/build.gradle.kts index 347b84b..a17c49a 100644 --- a/app/build.gradle.kts +++ b/app/build.gradle.kts @@ -1,3 +1,16 @@ +// M5: google-services plugin classpath (applied conditionally in :androidApp +// only when a google-services.json is present; the FCM path is inert without +// a Firebase project and the ntfy listener is the fallback). +buildscript { + repositories { + google() + mavenCentral() + } + dependencies { + classpath("com.google.gms:google-services:4.4.2") + } +} + plugins { kotlin("multiplatform") version "2.1.0" apply false kotlin("android") version "2.1.0" apply false diff --git a/app/shared/build.gradle.kts b/app/shared/build.gradle.kts index 15f26e5..0b25a9e 100644 --- a/app/shared/build.gradle.kts +++ b/app/shared/build.gradle.kts @@ -39,6 +39,10 @@ kotlin { implementation("androidx.media3:media3-ui:1.3.1") // SAF picker (rememberLauncherForActivityResult). implementation("androidx.activity:activity-compose:1.9.2") + // M5: FCM push (inert without a Firebase project / google-services.json; + // the ntfy listener is the fallback). The google-services plugin is + // applied conditionally in the app module. + implementation("com.google.firebase:firebase-messaging:23.1.2") } } } diff --git a/app/shared/src/androidMain/kotlin/iris/platform/AndroidPush.kt b/app/shared/src/androidMain/kotlin/iris/platform/AndroidPush.kt new file mode 100644 index 0000000..dee9225 --- /dev/null +++ b/app/shared/src/androidMain/kotlin/iris/platform/AndroidPush.kt @@ -0,0 +1,30 @@ +package iris.platform + +import android.Manifest +import android.content.pm.PackageManager +import androidx.core.content.ContextCompat +import iris.state.IrisController + +actual fun isAppForeground(): Boolean = AppBridge.foreground + +actual fun setActiveController(controller: Any?) { + AppBridge.controller = controller as? IrisController +} + +actual fun postSystemNotification( + chatId: String?, + chatName: String?, + title: String, + body: String, + threadId: String?, +) { + val context = AndroidEnv.context + val id = chatId ?: "android:default" + // POST_NOTIFICATIONS is a runtime permission on API 33+. + if (ContextCompat.checkSelfPermission(context, Manifest.permission.POST_NOTIFICATIONS) + != PackageManager.PERMISSION_GRANTED + ) { + return + } + IrisNotifications.post(context, id, chatName, title, body, threadId) +} \ No newline at end of file diff --git a/app/shared/src/androidMain/kotlin/iris/platform/AndroidSecureStore.kt b/app/shared/src/androidMain/kotlin/iris/platform/AndroidSecureStore.kt index 19f09cd..b67ad99 100644 --- a/app/shared/src/androidMain/kotlin/iris/platform/AndroidSecureStore.kt +++ b/app/shared/src/androidMain/kotlin/iris/platform/AndroidSecureStore.kt @@ -33,6 +33,22 @@ class AndroidSecureStore(context: Context) : SecureStore { override val deviceName: String get() = "${Build.MANUFACTURER} ${Build.MODEL}".trim() + override var syncCursor: Long + get() = prefs.getLong(KEY_SYNC_CURSOR, 0L) + set(value) = prefs.edit().putLong(KEY_SYNC_CURSOR, value).apply() + + override var fcmToken: String + get() = prefs.getString(KEY_FCM_TOKEN, "").orEmpty() + set(value) = prefs.edit().putString(KEY_FCM_TOKEN, value).apply() + + override var ntfyTopic: String + get() = prefs.getString(KEY_NTFY_TOPIC, "").orEmpty() + set(value) = prefs.edit().putString(KEY_NTFY_TOPIC, value).apply() + + override var ntfyServer: String + get() = prefs.getString(KEY_NTFY_SERVER, "").orEmpty() + set(value) = prefs.edit().putString(KEY_NTFY_SERVER, value).apply() + override fun savePairing(url: String, token: String) { serverUrl = url this.token = token @@ -46,5 +62,9 @@ class AndroidSecureStore(context: Context) : SecureStore { const val KEY_URL = "server_url" const val KEY_TOKEN = "token" const val KEY_DEVICE_ID = "device_id" + const val KEY_SYNC_CURSOR = "sync_cursor" + const val KEY_FCM_TOKEN = "fcm_token" + const val KEY_NTFY_TOPIC = "ntfy_topic" + const val KEY_NTFY_SERVER = "ntfy_server" } } \ No newline at end of file diff --git a/app/shared/src/androidMain/kotlin/iris/platform/AppBridge.kt b/app/shared/src/androidMain/kotlin/iris/platform/AppBridge.kt new file mode 100644 index 0000000..18976e3 --- /dev/null +++ b/app/shared/src/androidMain/kotlin/iris/platform/AppBridge.kt @@ -0,0 +1,20 @@ +package iris.platform + +import iris.state.IrisController + +/** + * Process-wide reference to the active [IrisController] (M5). + * + * The FCM service and the ntfy listener service are separate Android + * components that need to reach the controller (to trigger a sync or open a + * chat) without a compile-time dependency on the app module. Set by + * [iris.IrisApp] on composition, cleared on dispose. + */ +object AppBridge { + @Volatile + var controller: IrisController? = null + + /** True while the launcher activity is resumed (set by MainActivity). */ + @Volatile + var foreground: Boolean = false +} \ No newline at end of file diff --git a/app/shared/src/androidMain/kotlin/iris/platform/IrisFirebaseMessagingService.kt b/app/shared/src/androidMain/kotlin/iris/platform/IrisFirebaseMessagingService.kt new file mode 100644 index 0000000..466f9f6 --- /dev/null +++ b/app/shared/src/androidMain/kotlin/iris/platform/IrisFirebaseMessagingService.kt @@ -0,0 +1,47 @@ +package iris.platform + +import android.content.pm.PackageManager +import androidx.core.content.ContextCompat +import com.google.firebase.messaging.FirebaseMessagingService +import com.google.firebase.messaging.RemoteMessage + +/** + * M5: FCM handler (docs/08 §8.1). + * + * - [onNewToken]: persist the rotated token and push it to the server via + * `fcm.register` (so the next push targets the current token). + * - [onMessageReceived]: the data payload drives a silent sync. When the app + * is foregrounded the WS path already delivered the frame (in-app banner), + * so we only post a system notification when backgrounded. + * + * Inert without a Firebase project (no google-services.json): the service is + * declared in the manifest but never receives messages, and the app falls + * back to the ntfy listener. + */ +class IrisFirebaseMessagingService : FirebaseMessagingService() { + + override fun onNewToken(token: String) { + val store = AndroidSecureStore(applicationContext) + store.fcmToken = token + // Push the rotation to the server if we're connected. + AppBridge.controller?.client?.sendFrame( + iris.protocol.fcmRegisterFrame(fcmToken = token), + ) + } + + override fun onMessageReceived(message: RemoteMessage) { + val data = message.data + val chatId = data["chat_id"] ?: "android:default" + val threadId = data["thread_id"] + val title = data["title"] ?: "Iris" + val body = data["body"] ?: data["title"].orEmpty() + // Foreground + live WS: the in-app banner already showed this. + if (AppBridge.foreground) return + if (ContextCompat.checkSelfPermission(applicationContext, android.Manifest.permission.POST_NOTIFICATIONS) + != PackageManager.PERMISSION_GRANTED + ) { + return + } + IrisNotifications.post(applicationContext, chatId, null, title, body, threadId) + } +} \ No newline at end of file diff --git a/app/shared/src/androidMain/kotlin/iris/platform/IrisNotifications.kt b/app/shared/src/androidMain/kotlin/iris/platform/IrisNotifications.kt new file mode 100644 index 0000000..d9d085b --- /dev/null +++ b/app/shared/src/androidMain/kotlin/iris/platform/IrisNotifications.kt @@ -0,0 +1,72 @@ +package iris.platform + +import android.app.NotificationChannel +import android.app.NotificationManager +import android.app.PendingIntent +import android.content.Context +import android.content.Intent +import android.os.Build +import androidx.core.app.NotificationCompat + +/** + * M5: per-chat notification channels + posting (docs/08 §8.2). + * + * One channel per chat_id so the user can mute individual chats. The content + * intent uses a custom action (no compile-time dependency on the app module's + * MainActivity); the app's manifest routes it back to the launcher activity + * with chat_id / thread_id extras for the deep link. + */ +object IrisNotifications { + const val CHANNEL_PREFIX = "iris_chat_" + const val ACTION_OPEN_CHAT = "dev.iris.app.OPEN_CHAT" + private const val NOTIF_ID_BASE = 1_000_000 + + fun ensureChannel(context: Context, chatId: String, chatName: String? = null) { + if (Build.VERSION.SDK_INT < Build.VERSION_CODES.O) return + val nm = context.getSystemService(Context.NOTIFICATION_SERVICE) as NotificationManager + val id = CHANNEL_PREFIX + chatId + val name = chatName ?: chatId + if (nm.getNotificationChannel(id) == null) { + nm.createNotificationChannel( + NotificationChannel(id, name, NotificationManager.IMPORTANCE_DEFAULT) + .apply { description = "Iris messages for $name" } + ) + } + } + + fun post( + context: Context, + chatId: String, + chatName: String?, + title: String, + body: String, + threadId: String?, + ) { + ensureChannel(context, chatId, chatName) + val id = NOTIF_ID_BASE + (chatId.hashCode() and 0xffff) + val intent = Intent(ACTION_OPEN_CHAT).apply { + setPackage(context.packageName) + flags = Intent.FLAG_ACTIVITY_NEW_TASK or Intent.FLAG_ACTIVITY_CLEAR_TOP + putExtra("chat_id", chatId) + if (threadId != null) putExtra("thread_id", threadId) + } + val pi = PendingIntent.getActivity( + context, + id, + intent, + PendingIntent.FLAG_UPDATE_CURRENT or PendingIntent.FLAG_IMMUTABLE, + ) + val nm = context.getSystemService(Context.NOTIFICATION_SERVICE) as NotificationManager + nm.notify( + id, + NotificationCompat.Builder(context, CHANNEL_PREFIX + chatId) + .setSmallIcon(android.R.drawable.ic_dialog_info) + .setContentTitle(title) + .setContentText(body) + .setStyle(NotificationCompat.BigTextStyle().bigText(body)) + .setAutoCancel(true) + .setContentIntent(pi) + .build(), + ) + } +} \ No newline at end of file diff --git a/app/shared/src/androidMain/kotlin/iris/platform/NtfyListenerService.kt b/app/shared/src/androidMain/kotlin/iris/platform/NtfyListenerService.kt new file mode 100644 index 0000000..2af2bd6 --- /dev/null +++ b/app/shared/src/androidMain/kotlin/iris/platform/NtfyListenerService.kt @@ -0,0 +1,150 @@ +package iris.platform + +import android.app.Notification +import android.app.NotificationChannel +import android.app.NotificationManager +import android.app.Service +import android.content.Intent +import android.content.pm.PackageManager +import android.os.Build +import android.os.IBinder +import androidx.core.app.NotificationCompat +import androidx.core.content.ContextCompat +import iris.protocol.IrisJson +import kotlinx.coroutines.CoroutineScope +import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.Job +import kotlinx.coroutines.SupervisorJob +import kotlinx.coroutines.cancel +import kotlinx.coroutines.launch +import kotlinx.serialization.json.JsonObject +import kotlinx.serialization.json.JsonPrimitive +import okhttp3.OkHttpClient +import okhttp3.Request +import java.util.concurrent.TimeUnit + +/** + * M5: ntfy listener fallback (docs/08 §8.4). + * + * A foreground service that subscribes to this device's ntfy topic and posts + * a system notification for each push. The structured payload rides in the + * `X-Data` header (JSON: chat_id, kind, cursor, thread_id); the message body + * is the short preview. When the app is foregrounded the WS path already + * delivered the frame, so the service skips posting to avoid a duplicate. + */ +class NtfyListenerService : Service() { + private val scope = CoroutineScope(SupervisorJob() + Dispatchers.IO) + private var streamJob: Job? = null + private val client = OkHttpClient.Builder() + .readTimeout(0, TimeUnit.MILLISECONDS) // long-lived stream + .build() + + override fun onStartCommand(intent: Intent?, flags: Int, startId: Int): Int { + startForeground(NOTIF_ID, foregroundNotification()) + streamJob?.cancel() + streamJob = scope.launch { stream() } + return START_STICKY + } + + override fun onBind(intent: Intent?): IBinder? = null + + override fun onDestroy() { + streamJob?.cancel() + scope.cancel() + super.onDestroy() + } + + private suspend fun stream() { + val store = AndroidSecureStore(AndroidEnv.context) + val topic = store.ntfyTopic + if (topic.isBlank()) return + val server = store.ntfyServer.ifBlank { DEFAULT_NTFY_SERVER }.removeSuffix("/") + val url = "$server/$topic" + val request = Request.Builder() + .url(url) + .header("Accept", "text/event-stream") + .build() + try { + client.newCall(request).execute().use { resp -> + if (!resp.isSuccessful) return + val body = resp.body ?: return + val source = body.source() + var data: String? = null + var title: String? = null + var msgBody: String? = null + while (!source.exhausted()) { + val line = source.readUtf8Line() ?: break + when { + line.startsWith("X-Data:") -> data = line.removePrefix("X-Data:").trim() + line.startsWith("X-Title:") -> title = line.removePrefix("X-Title:").trim() + line.startsWith("data:") -> msgBody = line.removePrefix("data:").trim() + line.isEmpty() -> { + // Event boundary: process the accumulated message. + data?.let { handleData(it, title, msgBody) } + data = null + title = null + msgBody = null + } + } + } + } + } catch (_: Exception) { + // Stream dropped; the service is START_STICKY so the system + // restarts it. If it keeps failing, the WS path still works. + } + } + + private fun handleData(dataJson: String, title: String?, msgBody: String?) { + val data = try { + IrisJson.instance.decodeFromString(dataJson) + } catch (_: Exception) { + null + } + val chatId = data?.str("chat_id") ?: "android:default" + val threadId = data?.str("thread_id") + // The short preview rides in the SSE `data:` field; fall back to the + // X-Title, then a generic label. + val body = msgBody?.takeIf { it.isNotBlank() } ?: title.orEmpty() + val notifTitle = title ?: "Iris" + // Foreground + live WS: the in-app banner already showed this. + if (AppBridge.foreground) return + if (ContextCompat.checkSelfPermission(AndroidEnv.context, android.Manifest.permission.POST_NOTIFICATIONS) + != PackageManager.PERMISSION_GRANTED + ) { + return + } + IrisNotifications.post(AndroidEnv.context, chatId, null, notifTitle, body, threadId) + } + + private fun foregroundNotification(): Notification { + val context = AndroidEnv.context + if (Build.VERSION.SDK_INT >= Build.VERSION_CODES.O) { + val nm = context.getSystemService(NotificationManager::class.java) + if (nm.getNotificationChannel(LISTENER_CHANNEL) == null) { + nm.createNotificationChannel( + NotificationChannel( + LISTENER_CHANNEL, + "Iris push listener", + NotificationManager.IMPORTANCE_MIN, + ) + ) + } + } + return NotificationCompat.Builder(context, LISTENER_CHANNEL) + .setSmallIcon(android.R.drawable.ic_dialog_info) + .setContentTitle("Iris") + .setContentText("Listening for messages") + .setOngoing(true) + .build() + } + + private companion object { + const val LISTENER_CHANNEL = "iris_ntfy_listener" + const val NOTIF_ID = 999_001 + const val DEFAULT_NTFY_SERVER = "https://ntfy.sh" + } +} + +/** Read a string field from a JSON object (null when absent / not a string). */ +private fun JsonObject?.str(key: String): String? = + (this?.get(key) as? JsonPrimitive)?.content \ No newline at end of file diff --git a/app/shared/src/commonMain/kotlin/iris/IrisApp.kt b/app/shared/src/commonMain/kotlin/iris/IrisApp.kt index 11e3744..a602fd3 100644 --- a/app/shared/src/commonMain/kotlin/iris/IrisApp.kt +++ b/app/shared/src/commonMain/kotlin/iris/IrisApp.kt @@ -6,6 +6,7 @@ import androidx.compose.material3.Surface import androidx.compose.material3.darkColorScheme import androidx.compose.runtime.Composable import androidx.compose.runtime.DisposableEffect +import androidx.compose.runtime.LaunchedEffect import androidx.compose.runtime.collectAsState import androidx.compose.runtime.getValue import androidx.compose.runtime.remember @@ -13,6 +14,7 @@ import androidx.compose.ui.Modifier import androidx.compose.ui.graphics.Color import iris.data.SecureStore import iris.net.GatewayClient +import iris.platform.setActiveController import iris.state.IrisController import iris.ui.screens.ChatScreen import iris.ui.screens.ConnectScreen @@ -33,10 +35,25 @@ private val IrisDark = darkColorScheme( * settings, and media (docs/10-android-app.md). */ @Composable -fun IrisApp(store: SecureStore) { +fun IrisApp( + store: SecureStore, + deepLinkChatId: String? = null, + deepLinkThreadId: String? = null, +) { val controller = remember(store) { IrisController(store) } DisposableEffect(controller) { - onDispose { controller.dispose() } + setActiveController(controller) + onDispose { + setActiveController(null) + controller.dispose() + } + } + // M5: open the chat referenced by a push notification tap (applied now or + // on the next hello.ack, whichever is later). + LaunchedEffect(deepLinkChatId, deepLinkThreadId) { + if (!deepLinkChatId.isNullOrBlank()) { + controller.openDeepLink(deepLinkChatId, deepLinkThreadId) + } } val state by controller.client.state.collectAsState() diff --git a/app/shared/src/commonMain/kotlin/iris/data/SecureStore.kt b/app/shared/src/commonMain/kotlin/iris/data/SecureStore.kt index 0a8b99e..2dacd80 100644 --- a/app/shared/src/commonMain/kotlin/iris/data/SecureStore.kt +++ b/app/shared/src/commonMain/kotlin/iris/data/SecureStore.kt @@ -18,6 +18,18 @@ interface SecureStore { /** Human-readable device name (e.g. "MIX 2S"). */ val deviceName: String + /** Last sync cursor confirmed by the server (M5 reconnect catch-up). */ + var syncCursor: Long + + /** FCM registration token (M5; empty when Firebase is unavailable). */ + var fcmToken: String + + /** ntfy topic this device subscribes to (M5; empty when unused). */ + var ntfyTopic: String + + /** ntfy server URL (M5; from hello.ack server_caps; default ntfy.sh). */ + var ntfyServer: String + fun savePairing(url: String, token: String) fun clear() } \ No newline at end of file diff --git a/app/shared/src/commonMain/kotlin/iris/net/GatewayClient.kt b/app/shared/src/commonMain/kotlin/iris/net/GatewayClient.kt index 08bb941..3c0887f 100644 --- a/app/shared/src/commonMain/kotlin/iris/net/GatewayClient.kt +++ b/app/shared/src/commonMain/kotlin/iris/net/GatewayClient.kt @@ -22,6 +22,7 @@ import iris.protocol.mediaUploadEndFrame import iris.protocol.mediaUploadStartFrame import iris.protocol.messageSendFrame import iris.protocol.pingFrame +import iris.protocol.syncFrame import kotlinx.coroutines.CompletableDeferred import kotlinx.coroutines.CoroutineScope import kotlinx.coroutines.Job @@ -187,7 +188,15 @@ class GatewayClient( request, object : WebSocketListener() { override fun onOpen(webSocket: WebSocket, response: Response) { - webSocket.send(helloFrame(token, store.deviceId, store.deviceName).toWire()) + webSocket.send( + helloFrame( + token = token, + deviceId = store.deviceId, + deviceName = store.deviceName, + fcmToken = store.fcmToken.ifBlank { null }, + ntfyTopic = store.ntfyTopic.ifBlank { null }, + ).toWire(), + ) } override fun onMessage(webSocket: WebSocket, text: String) { @@ -250,6 +259,12 @@ class GatewayClient( if (e == null) { val ack = helloAck.getCompleted() _state.value = State.Connected(ack.serverCaps, ack.channels) + // M5: reconnect catch-up — replay frames parked while offline. + val local = store.syncCursor + if (local < ack.syncCursor) { + val id = nextRequestId++ + ws.send(syncFrame(id, local).toWire()) + } winner.complete(DialResult.Connected) } } diff --git a/app/shared/src/commonMain/kotlin/iris/platform/PlatformPush.kt b/app/shared/src/commonMain/kotlin/iris/platform/PlatformPush.kt new file mode 100644 index 0000000..f6d6e23 --- /dev/null +++ b/app/shared/src/commonMain/kotlin/iris/platform/PlatformPush.kt @@ -0,0 +1,31 @@ +package iris.platform + +/** + * M5: push / notification platform hooks (docs/08 §8.2). + * + * The controller (commonMain) needs to know whether the app is in the + * foreground (to decide between an in-app banner and a system notification) + * and to post a system notification when a `notification` frame arrives while + * the app is backgrounded but the WS is still live. + */ + +/** True when the app's UI is visible (Android: activity resumed). */ +expect fun isAppForeground(): Boolean + +/** + * Post a system notification for a chat (Android). No-op on desktop. + * [chatName] is used for the per-chat notification channel (best effort). + */ +expect fun postSystemNotification( + chatId: String?, + chatName: String?, + title: String, + body: String, + threadId: String?, +) + +/** + * Register the active controller with the platform bridge (Android: so the + * FCM / ntfy services can reach it). No-op on desktop. + */ +expect fun setActiveController(controller: Any?) \ No newline at end of file diff --git a/app/shared/src/commonMain/kotlin/iris/protocol/Protocol.kt b/app/shared/src/commonMain/kotlin/iris/protocol/Protocol.kt index ffd6b45..f8420bf 100644 --- a/app/shared/src/commonMain/kotlin/iris/protocol/Protocol.kt +++ b/app/shared/src/commonMain/kotlin/iris/protocol/Protocol.kt @@ -54,6 +54,10 @@ const val TYPE_MEDIA_OFFER = "media.offer" const val TYPE_MEDIA_PULL = "media.pull" const val TYPE_MEDIA_PULL_END = "media.pull.end" +// M5 — push / notifications +const val TYPE_NOTIFICATION = "notification" +const val TYPE_FCM_REGISTER = "fcm.register" + // M3 — channels / threads / search / sync const val TYPE_CHANNEL_CREATE = "channel.create" const val TYPE_CHANNEL_RENAME = "channel.rename" @@ -131,6 +135,7 @@ data class ServerCaps( val media: Boolean = false, val search: Boolean = false, val push: String = "fcm", + @SerialName("push_ntfy_server") val pushNtfyServer: String = "", val pickers: Boolean = false, ) @@ -349,14 +354,55 @@ data class SyncPayload(val cursor: Long) @Serializable data class SyncDonePayload(val cursor: Long) +// ── M5: push / notifications ──────────────────────────────────────────── + +/** Notification kinds (mirror of protocol.NOTIF_*). */ +const val NOTIF_MESSAGE = "message" +const val NOTIF_APPROVAL = "approval" +const val NOTIF_CLARIFY = "clarify" +const val NOTIF_CRON = "cron" +const val NOTIF_CHANNEL = "channel" +const val NOTIF_OUTBOX_PRUNED = "outbox_pruned" + +/** Kinds that stay on screen until dismissed (docs/08 §8.3). */ +val HIGH_PRIORITY_NOTIF_KINDS = setOf(NOTIF_APPROVAL, NOTIF_CLARIFY, NOTIF_CRON) + +@Serializable +data class NotificationPayload( + val kind: String, + val title: String, + val body: String, + @SerialName("chat_id") val chatId: String? = null, + @SerialName("thread_id") val threadId: String? = null, + val ts: Long? = null, +) + +@Serializable +data class FcmRegisterPayload( + @SerialName("fcm_token") val fcmToken: String? = null, + @SerialName("ntfy_topic") val ntfyTopic: String? = null, +) + // ── Frame builders ────────────────────────────────────────────────────── -fun helloFrame(token: String, deviceId: String, deviceName: String): Frame = +fun helloFrame( + token: String, + deviceId: String, + deviceName: String, + fcmToken: String? = null, + ntfyTopic: String? = null, +): Frame = Frame( type = TYPE_HELLO, payload = IrisJson.instance.encodeToJsonElement( HelloPayload.serializer(), - HelloPayload(token = token, deviceId = deviceId, deviceName = deviceName), + HelloPayload( + token = token, + deviceId = deviceId, + deviceName = deviceName, + fcmToken = fcmToken, + ntfyTopic = ntfyTopic, + ), ), ) @@ -478,4 +524,15 @@ fun mediaPullFrame(id: Int, mediaId: String): Frame = MediaPullPayload.serializer(), MediaPullPayload(mediaId), ), + ) + +// ── M5 frame builders ─────────────────────────────────────────────────── + +fun fcmRegisterFrame(fcmToken: String? = null, ntfyTopic: String? = null): Frame = + Frame( + type = TYPE_FCM_REGISTER, + payload = IrisJson.instance.encodeToJsonElement( + FcmRegisterPayload.serializer(), + FcmRegisterPayload(fcmToken = fcmToken, ntfyTopic = ntfyTopic), + ), ) \ No newline at end of file diff --git a/app/shared/src/commonMain/kotlin/iris/state/IrisController.kt b/app/shared/src/commonMain/kotlin/iris/state/IrisController.kt index 1c4dd32..042dc0b 100644 --- a/app/shared/src/commonMain/kotlin/iris/state/IrisController.kt +++ b/app/shared/src/commonMain/kotlin/iris/state/IrisController.kt @@ -8,11 +8,19 @@ import iris.media.MediaCache import iris.media.kindFromMime import iris.net.GatewayClient import iris.platform.PickedFile +import iris.platform.isAppForeground import iris.platform.mediaCacheBaseDir +import iris.platform.postSystemNotification +import iris.protocol.HIGH_PRIORITY_NOTIF_KINDS import iris.protocol.MediaOfferPayload +import iris.protocol.MessagePayload +import iris.protocol.MessageStopPayload +import iris.protocol.NotificationPayload +import iris.protocol.ROLE_ASSISTANT import iris.protocol.SearchHit import iris.protocol.SearchResultsPayload import iris.protocol.SyncDonePayload +import iris.protocol.TYPE_NOTIFICATION import iris.protocol.TYPE_CHANNEL_CREATED import iris.protocol.TYPE_CHANNEL_DELETED import iris.protocol.TYPE_CHANNEL_LIST @@ -40,6 +48,7 @@ import iris.protocol.syncFrame import kotlinx.coroutines.CoroutineScope import kotlinx.coroutines.Dispatchers import kotlinx.coroutines.SupervisorJob +import kotlinx.coroutines.delay import kotlinx.coroutines.flow.MutableStateFlow import kotlinx.coroutines.flow.StateFlow import kotlinx.coroutines.flow.asStateFlow @@ -114,6 +123,76 @@ class IrisController( private val _attachments = MutableStateFlow>(emptyList()) val attachments: StateFlow> = _attachments.asStateFlow() + // ── M5: notification banners (docs/08 §8.3) ─────────────────────────── + /** An in-app notification banner shown above the composer. */ + data class Banner( + val id: Long, + val kind: String, + val title: String, + val body: String, + val chatId: String?, + val threadId: String?, + val persistent: Boolean, + ) + + private val _banners = MutableStateFlow>(emptyList()) + val banners: StateFlow> = _banners.asStateFlow() + private var bannerSeq = 0L + + fun dismissBanner(id: Long) { + _banners.value = _banners.value.filterNot { it.id == id } + } + + private fun pushBanner(kind: String, title: String, body: String, chatId: String?, threadId: String?) { + val persistent = kind in HIGH_PRIORITY_NOTIF_KINDS + val banner = Banner(bannerSeq++, kind, title, body, chatId, threadId, persistent) + _banners.value = (_banners.value + banner).takeLast(5) + if (!persistent) { + scope.launch { + delay(5_000) + _banners.value = _banners.value.filterNot { it.id == banner.id } + } + } + } + + /** + * M5: Telegram/WhatsApp-style message notification. When the app is + * backgrounded and a finalized assistant reply arrives, post a system + * notification (per-chat channel) with a short preview. Tapping it opens + * the chat (deep link). The same per-chat notification id is reused, so a + * multi-segment turn updates one notification instead of stacking. + */ + private fun notifyMessageIfBackgrounded(chatId: String?, threadId: String?, text: String) { + if (isAppForeground()) return + if (text.isBlank()) return + val id = chatId ?: "android:default" + val chatName = channels.byId(id)?.name + postSystemNotification(id, chatName, chatName ?: "Iris", preview(text), threadId) + } + + /** Single-line preview for a notification body (lock-screen privacy). */ + private fun preview(text: String, limit: Int = 120): String { + val flat = text.replace('\n', ' ').replace('\t', ' ').trim() + return if (flat.length <= limit) flat else flat.take(limit).trimEnd() + "…" + } + + // ── M5: deep link (push notification tap) ───────────────────────────── + private var pendingDeepLink: Pair? = null + + /** Open the chat referenced by a push notification tap. Applied now if + * connected, otherwise on the next hello.ack (lane may not exist yet). */ + fun openDeepLink(chatId: String?, threadId: String?) { + if (chatId.isNullOrBlank()) return + pendingDeepLink = chatId to threadId + if (client.state.value is GatewayClient.State.Connected) applyDeepLink() + } + + private fun applyDeepLink() { + val link = pendingDeepLink ?: return + pendingDeepLink = null + if (link.second.isNullOrBlank()) openChannel(link.first) else openThread(link.first, link.second!!) + } + init { scope.launch { client.events.collect { frame -> @@ -132,6 +211,25 @@ class IrisController( if (frame.type == TYPE_MEDIA_OFFER) { frame.payloadAs()?.let { pullMedia(it) } } + // M5: Telegram/WhatsApp-style — when the app is + // backgrounded, a finalized assistant reply posts a + // system notification (the message still lands in the + // chat via onFrame above). Streaming turns finalize on + // message.stop; non-streaming replies are a single + // assistant `message`. + when (frame.type) { + TYPE_MESSAGE_STOP -> + frame.payloadAs()?.let { + notifyMessageIfBackgrounded(frame.chatId, frame.threadId, it.finalText) + } + TYPE_MESSAGE -> + frame.payloadAs()?.let { + if (it.role == ROLE_ASSISTANT) { + notifyMessageIfBackgrounded(frame.chatId, frame.threadId, it.text) + } + } + else -> Unit + } } TYPE_CHANNEL_CREATED, TYPE_CHANNEL_RENAMED, @@ -144,7 +242,19 @@ class IrisController( TYPE_SYNC_DONE -> { // Replayed frames already flowed through [events]; the // cursor is authoritative server-side (outbox). - frame.payloadAs() + frame.payloadAs()?.let { store.syncCursor = it.cursor } + } + TYPE_NOTIFICATION -> { + frame.payloadAs()?.let { p -> + pushBanner(p.kind, p.title, p.body, p.chatId, p.threadId) + // M5: WS is live but the app is backgrounded — the + // in-app banner is invisible, so mirror to a system + // notification (the push backend only fires when + // there is no live subscriber). + if (!isAppForeground()) { + postSystemNotification(p.chatId, null, p.title, p.body, p.threadId) + } + } } TYPE_TYPING -> { frame.payloadAs()?.let { _typing.value = it.on } @@ -162,9 +272,19 @@ class IrisController( _homeChannel.value = home chat.setLane(home) } + // M5: remember the ntfy server for the listener service. + if (s.caps.pushNtfyServer.isNotBlank()) { + store.ntfyServer = s.caps.pushNtfyServer + } + // M5: a deep link tapped before we were connected. + applyDeepLink() } } } + // M5: ensure this device has an ntfy topic (device-specific, random). + if (store.ntfyTopic.isBlank()) { + store.ntfyTopic = "iris-${store.deviceId}-${Random.nextLong(1_000_000_000L, 9_999_999_999L)}" + } client.startHeartbeat() client.start() } diff --git a/app/shared/src/commonMain/kotlin/iris/ui/screens/ChatScreen.kt b/app/shared/src/commonMain/kotlin/iris/ui/screens/ChatScreen.kt index 460f6c2..b3aeb9a 100644 --- a/app/shared/src/commonMain/kotlin/iris/ui/screens/ChatScreen.kt +++ b/app/shared/src/commonMain/kotlin/iris/ui/screens/ChatScreen.kt @@ -105,6 +105,7 @@ fun ChatScreen(controller: IrisController) { var showNewThread by remember { mutableStateOf(false) } var showPicker by remember { mutableStateOf(false) } val attachments by controller.attachments.collectAsState() + val banners by controller.banners.collectAsState() val drawerState = rememberDrawerState(DrawerValue.Closed) val drawerScope = rememberCoroutineScope() val focusManager = LocalFocusManager.current @@ -239,6 +240,27 @@ fun ChatScreen(controller: IrisController) { } } + // M5: notification banners (above the composer) + if (banners.isNotEmpty()) { + Column( + modifier = Modifier + .fillMaxWidth() + .padding(horizontal = 12.dp, vertical = 4.dp), + verticalArrangement = Arrangement.spacedBy(6.dp), + ) { + banners.forEach { b -> + NotificationBanner( + banner = b, + onOpen = { + controller.openDeepLink(b.chatId, b.threadId) + controller.dismissBanner(b.id) + }, + onDismiss = { controller.dismissBanner(b.id) }, + ) + } + } + } + // Composer Row( modifier = Modifier @@ -530,6 +552,38 @@ private fun StatusChip(state: GatewayClient.State) { } } +/** M5: in-app notification banner (docs/08 §8.3). Tapping opens the chat. */ +@Composable +private fun NotificationBanner( + banner: IrisController.Banner, + onOpen: () -> Unit, + onDismiss: () -> Unit, +) { + val accent = when (banner.kind) { + "approval" -> Color(0xFFFFC107) + "clarify" -> Color(0xFF4F7CFF) + "cron" -> Color(0xFF4CAF50) + else -> Color(0xFF8A93A6) + } + Row( + modifier = Modifier + .fillMaxWidth() + .clip(RoundedCornerShape(12.dp)) + .background(accent.copy(alpha = 0.12f)) + .clickable(onClick = onOpen) + .padding(horizontal = 12.dp, vertical = 8.dp), + verticalAlignment = Alignment.CenterVertically, + ) { + Column(modifier = Modifier.weight(1f)) { + Text(banner.title, fontSize = 13.sp, fontWeight = FontWeight.SemiBold, color = accent) + if (banner.body.isNotBlank()) { + Text(banner.body, fontSize = 12.sp, color = Color(0xFFD7DBE5), maxLines = 2) + } + } + TextButton(onClick = onDismiss) { Text("✕", fontSize = 12.sp) } + } +} + @Composable private fun MessageBubble(msg: MessageItem) { val isUser = msg.role == ROLE_USER diff --git a/app/shared/src/desktopMain/kotlin/iris/platform/DesktopPush.kt b/app/shared/src/desktopMain/kotlin/iris/platform/DesktopPush.kt new file mode 100644 index 0000000..e58c2f4 --- /dev/null +++ b/app/shared/src/desktopMain/kotlin/iris/platform/DesktopPush.kt @@ -0,0 +1,17 @@ +package iris.platform + +actual fun isAppForeground(): Boolean = true + +actual fun postSystemNotification( + chatId: String?, + chatName: String?, + title: String, + body: String, + threadId: String?, +) { + // Desktop: in-app banner only (no system notification in M5). +} + +actual fun setActiveController(controller: Any?) { + // Desktop: no push services to bridge to. +} \ No newline at end of file diff --git a/app/shared/src/desktopMain/kotlin/iris/platform/DesktopSecureStore.kt b/app/shared/src/desktopMain/kotlin/iris/platform/DesktopSecureStore.kt index b329523..21ffe96 100644 --- a/app/shared/src/desktopMain/kotlin/iris/platform/DesktopSecureStore.kt +++ b/app/shared/src/desktopMain/kotlin/iris/platform/DesktopSecureStore.kt @@ -16,6 +16,10 @@ private data class PairingData( val serverUrl: String = "", val token: String = "", val deviceId: String = "", + val syncCursor: Long = 0L, + val fcmToken: String = "", + val ntfyTopic: String = "", + val ntfyServer: String = "", ) class DesktopSecureStore : SecureStore { @@ -63,6 +67,34 @@ class DesktopSecureStore : SecureStore { override val deviceName: String get() = "Desktop (${System.getProperty("os.name")})" + override var syncCursor: Long + get() = load().syncCursor + set(value) { + val d = load() + save(d.copy(syncCursor = value)) + } + + override var fcmToken: String + get() = load().fcmToken + set(value) { + val d = load() + save(d.copy(fcmToken = value)) + } + + override var ntfyTopic: String + get() = load().ntfyTopic + set(value) { + val d = load() + save(d.copy(ntfyTopic = value)) + } + + override var ntfyServer: String + get() = load().ntfyServer + set(value) { + val d = load() + save(d.copy(ntfyServer = value)) + } + override fun savePairing(url: String, token: String) { val d = load() save(d.copy(serverUrl = url.trim(), token = token.trim())) diff --git a/docs/14-milestones.md b/docs/14-milestones.md index 6314521..4939f0d 100644 --- a/docs/14-milestones.md +++ b/docs/14-milestones.md @@ -118,18 +118,29 @@ has explicit **acceptance criteria**. Work top-to-bottom; don't skip M0/M1. ## M5 — Push + offline (FCM + ntfy) **Goal:** reach the phone when backgrounded; catch up on reconnect. -- [ ] Outbox (SQLite) + sync cursor; `sync`/`sync.done`; retention prune. -- [ ] `push.py`: `FcmBackend` (HTTP v1 + service account, httpx) + - `NtfyBackend`; selected by `ANDROID_PUSH_BACKEND`. -- [ ] Fire push on no-live-subscriber; data payload for silent sync. -- [ ] App: FCM service (foreground banner + background foreground-service sync), - `onNewToken` → `fcm.register`; per-chat notification channels; deep-link. - ntfy listener fallback. -- [ ] In-app `notification` banners (channel_renamed, approval, cron, …). -- **Demo (on-device):** background the app → trigger a message → notification - appears → tap → syncs + opens the chat. Repeat with ntfy backend. -- **Accept:** push arrives when backgrounded (both backends); reconnect syncs - with no loss/dup; banners show for foreground events. +- [x] Outbox (SQLite) + sync cursor; `sync`/`sync.done`; retention prune + (row cap 5000 + prune banner, throttled 1/h). +- [x] `push.py`: `FcmBackend` (HTTP v1 + service account, httpx; JWT via + PyJWT+cryptography) + `NtfyBackend` (X-Data header); selected by + `ANDROID_PUSH_BACKEND`. ntfy server exposed in `server_caps.push_ntfy_server`. +- [x] Fire push on no-live-subscriber; data payload for silent sync. + High-priority kinds (approval/clarify/cron) push even when live. +- [x] App: FCM service (`onNewToken` → `fcm.register`; inert without a + Firebase project), per-chat notification channels, deep-link + (custom action + `iris://`), ntfy listener foreground-service fallback. + App generates a device-specific ntfy topic; auto-sync on reconnect. +- [x] In-app `notification` banners (channel events, approval, clarify, cron, + outbox-pruned); system notification mirror when backgrounded. +- **Demo (on-device, verified 2026-08-19):** device offline → message + broadcast → frame parked in outbox → `push via ntfy -> ` fired → + app reconnect auto-synced (cursor 0→8) → parked message rendered in-app. + ntfy publish confirmed (POST 200). ntfy listener service runs (foreground + notification posted); listener→system-notification leg not exercised live + because ntfy.sh's public SSE endpoint was serving its web UI (not streams) + during the test — verify against a self-hosted ntfy. +- **Accept:** push arrives when backgrounded (ntfy verified; FCM code path + complete, needs a Firebase project to exercise); reconnect syncs with no + loss/dup (verified); banners show for foreground events (implemented). ## M6 — Desktop app **Goal:** the same app on a big screen. diff --git a/gateway-plugin/adapter.py b/gateway-plugin/adapter.py index ebf2b21..fd29c36 100644 --- a/gateway-plugin/adapter.py +++ b/gateway-plugin/adapter.py @@ -30,6 +30,14 @@ register the (delivery-validated) file in the media registry and emit ``media.offer``; ``media.pull`` streams the file back as chunked binary frames, re-checking ``validate_media_delivery_path`` at pull time. +Milestone M5: push + offline. Frames with no live subscriber are parked in +the outbox (M3) AND wake the device via the push backend (``push.py``: FCM +HTTP v1 primary, ntfy fallback, selected by ``ANDROID_PUSH_BACKEND``). +``notification`` frames render in-app banners and mirror to push (channel +events, cron deliveries, approvals, clarifies); high-priority kinds push even +when a device is live. ``fcm.register`` rotates push tokens (registry + live +connection). The outbox enforces a row cap with a throttled prune notice. + Configuration in config.yaml:: gateway: @@ -106,6 +114,7 @@ from . import protocol # noqa: E402 from . import search as search_bridge # noqa: E402 from .channels import get_directory # noqa: E402 from .outbox import Outbox # noqa: E402 +from .push import NtfyBackend, PushBackend, build_push_backend # noqa: E402 from .pairing import ( # noqa: E402 DeviceRegistry, generate_token, @@ -280,6 +289,42 @@ def _strip_streaming_cursor(text: str) -> str: return text +def _push_preview(text: Any, limit: int = 120) -> str: + """Short single-line preview for push bodies (lock-screen privacy: no + secrets, no full bodies -- full content arrives via ``sync``).""" + s = " ".join(str(text or "").split()) + if len(s) > limit: + s = s[: limit - 1] + "…" + return s + + +# Cron delivery wrap (cron/scheduler.py ``_deliver_result``, +# cron.wrap_response: true): +# "Cronjob Response: \n(job_id: )\n-------------\n\n\n\n +# To stop or manage this job, send me a new message (e.g. ...)." +_CRON_WRAP_RE = re.compile( + r"^Cronjob Response: (.+?)\n\(job_id: [^)]*\)\n-+\n\n" +) +_CRON_FOOTER = "\n\nTo stop or manage this job" + + +def _cron_brief(content: str, job_id: str) -> Tuple[str, str]: + """Parse a cron delivery into ``(job_name, inner_text)``. + + Falls back to ``(job_id, content)`` when the wrap is disabled + (``cron.wrap_response: false``) or unrecognised. + """ + m = _CRON_WRAP_RE.match(content or "") + if not m: + return str(job_id or "cron"), (content or "").strip() + name = m.group(1).strip() + body = content[m.end():] + idx = body.rfind(_CRON_FOOTER) + if idx != -1: + body = body[:idx] + return name, body.strip() + + def _split_reasoning(text: str) -> Tuple[Optional[str], str]: """Split a code-style reasoning prefix off the front of *text*. @@ -693,6 +738,17 @@ class AndroidAdapter(BasePlatformAdapter): # last finalized assistant message id per chat (offer association). self._media = media_bridge.MediaStore(get_hermes_home()) self._last_message_id: Dict[str, str] = {} + # M5: push backend (FCM primary, ntfy fallback) + the throttle for + # the outbox-prune banner. + self._push: PushBackend = build_push_backend( + self.push_backend, + fcm_service_account=_get_scoped_secret("ANDROID_FCM_SERVICE_ACCOUNT"), + fcm_server_key=_get_scoped_secret("ANDROID_FCM_SERVER_KEY"), + ntfy_topic=_get_scoped_secret("NTFY_TOPIC"), + ntfy_server_url=os.getenv("NTFY_SERVER_URL", "").strip() or None, + ntfy_auth_token=_get_scoped_secret("NTFY_AUTH_TOKEN"), + ) + self._prune_notified_at = 0.0 def _turn_state(self, chat_id: str) -> _TurnState: st = self._turns.get(chat_id) @@ -747,6 +803,16 @@ class AndroidAdapter(BasePlatformAdapter): except Exception: logger.warning("android: ensure_default failed", exc_info=True) + # M5: push backend status (degrade gracefully when unconfigured). + if not self._push.configured(): + logger.warning( + "android: push backend %r not configured (no credentials) -- " + "offline devices will not be woken; outbox + sync still apply", + self.push_backend, + ) + else: + logger.info("android: push backend: %s", self._push.name) + self._connected = True self._mark_connected() logger.info("android: connected; WS server on %s:%s", self.host, self.port) @@ -817,7 +883,45 @@ class AndroidAdapter(BasePlatformAdapter): ) return SendResult(success=True, message_id=message_id) - # 2. Final message (non-streaming final, or streaming fallback final). + # 2. Cron delivery (the cron scheduler passes ``job_id`` in metadata): + # a final assistant message + a high-priority notification banner + # (pushed even when a device is live, docs/08 §8.1). + if meta.get("job_id"): + reasoning, _ = _split_reasoning(content) + if reasoning: + _reset_reasoning() + else: + await _wait_for_reasoning_flushed() + reasoning = _take_reasoning() or None + name, inner = _cron_brief(content, str(meta.get("job_id"))) + message_id = _mint_message_id() + await self._broadcast_or_log( + chat_id, + protocol.notification( + chat_id, + protocol.NOTIF_CRON, + f"Cron: {name}", + _push_preview(inner), + thread_id=thread_id, + ), + ) + await self._broadcast_or_log( + chat_id, + protocol.message( + chat_id=chat_id, + message_id=message_id, + role=protocol.ROLE_ASSISTANT, + text=inner, + thread_id=thread_id, + reasoning=reasoning, + ts=int(time.time() * 1000), + ), + ) + self._last_message_id[chat_id] = message_id + state.active = False + return SendResult(success=True, message_id=message_id) + + # 3. Final message (non-streaming final, or streaming fallback final). if meta.get("notify") is True: reasoning, body = _split_reasoning(content) # Non-streaming: reasoning is prepended to content (split above). @@ -863,11 +967,11 @@ class AndroidAdapter(BasePlatformAdapter): state.active = False return SendResult(success=True, message_id=message_id) - # 3. Tool progress (first tool bubble of an editable line buffer). + # 4. Tool progress (first tool bubble of an editable line buffer). if _is_tool_progress(content): return await self._emit_tool_lines(chat_id, content, state, thread_id, is_edit=False) - # 4. Commentary (interim assistant beat). + # 5. Commentary (interim assistant beat). message_id = _mint_message_id() state.active = True await self._broadcast_or_log( @@ -1043,17 +1147,151 @@ class AndroidAdapter(BasePlatformAdapter): async def _broadcast_or_log(self, chat_id: str, frame: "protocol.Frame") -> None: delivered = await self._ws_server.broadcast(frame) + # M3/M5: always append to the outbox so a reconnecting app can catch + # up on *all* recent frames, not just the ones that were parked. This + # covers the case where the app's in-memory ChatStore is reset (e.g. + # activity recreation / composition recompose) while the process stays + # alive: the recreated controller re-syncs from its stale cursor and + # replays the frames it missed. The push below still only fires when + # there is no live subscriber (or for high-priority events). + try: + cursor = self._outbox.append(chat_id, frame.to_json()) + except Exception: + logger.warning("android: outbox append failed", exc_info=True) + return if delivered == 0: - # M3: no live device -- park the frame in the outbox so a - # reconnecting app can `sync` it (M5 adds push to wake the device). + logger.info( + "android: no live devices for %s; %s frame parked in outbox (cursor=%s)", + chat_id, frame.type, cursor, + ) + # M5: wake the offline device(s) via the push backend. + await self._maybe_push(chat_id, frame, cursor) + await self._maybe_notify_outbox_prune(chat_id) + elif ( + frame.type == protocol.TYPE_NOTIFICATION + and frame.payload.get("kind") in protocol.HIGH_PRIORITY_NOTIF_KINDS + ): + # M5: high-priority events (approval/clarify/cron) push even when + # a device is live -- the app may be backgrounded and decides + # whether to also show an in-app banner (docs/08 §8.1). + await self._maybe_push(chat_id, frame, cursor) + + # ── M5: push ─────────────────────────────────────────────────────────── + + def _push_summary( + self, frame: "protocol.Frame" + ) -> Optional[Tuple[str, str, str, str]]: + """``(title, body, kind, priority)`` for a pushable frame, else None. + + Only terminal/interesting frames wake a device: intermediate + streaming and tool frames are replayed by ``sync`` without a push + (no notification spam per turn). + """ + t = frame.type + p = frame.payload + if t == protocol.TYPE_MESSAGE: + return ( + self._channel_name(frame.chat_id or ""), + _push_preview(p.get("text")), + "message", + "normal", + ) + if t == protocol.TYPE_MESSAGE_STOP: + return ( + self._channel_name(frame.chat_id or ""), + _push_preview(p.get("final_text")), + "message", + "normal", + ) + if t == protocol.TYPE_NOTIFICATION: + kind = str(p.get("kind") or protocol.NOTIF_GENERIC) + priority = "high" if kind in protocol.HIGH_PRIORITY_NOTIF_KINDS else "normal" + return ( + str(p.get("title") or "Iris"), + str(p.get("body") or ""), + kind, + priority, + ) + if t == protocol.TYPE_MEDIA_OFFER: + return ( + self._channel_name(frame.chat_id or ""), + f"New {p.get('kind') or 'media'}: {p.get('filename') or ''}".strip(), + "media", + "normal", + ) + return None + + async def _maybe_push(self, chat_id: str, frame: "protocol.Frame", cursor: int) -> None: + """Fire the configured push backend for a parked (or high-priority) + frame. Best-effort: failures are logged, never raised.""" + summary = self._push_summary(frame) + if summary is None: + return + title, body, kind, priority = summary + backend = self._push + if backend is None or not backend.token_field: + return + devices = self._devices.list() + if not backend.configured() and not any( + d.get(backend.token_field) for d in devices + ): + return + data: Dict[str, Any] = {"chat_id": chat_id, "kind": kind, "cursor": str(cursor)} + if frame.thread_id: + data["thread_id"] = frame.thread_id + message_id = frame.payload.get("message_id") + if isinstance(message_id, str) and message_id: + data["message_id"] = message_id + for device in devices: + device_id = device.get("device_id") + if not device_id: + continue + # Prefer the live connection's token (fcm.register refreshes it + # in memory) over the possibly-stale registry row. + conn = self._ws_server.connection(device_id) + token = getattr(conn, backend.token_field, None) if conn is not None else None + if not token: + token = device.get(backend.token_field) + if not token: + continue try: - cursor = self._outbox.append(chat_id, frame.to_json()) - logger.info( - "android: no live devices for %s; %s frame parked in outbox (cursor=%s)", - chat_id, frame.type, cursor, + ok = await backend.send( + device_id=device_id, + chat_id=chat_id, + title=title, + body=body, + data=data, + token=token, + priority=priority, ) except Exception: - logger.warning("android: outbox append failed", exc_info=True) + logger.warning("android: push via %s failed", backend.name, exc_info=True) + continue + if ok: + logger.info( + "android: push via %s -> %s (%s, chat=%s)", + backend.name, device_id, frame.type, chat_id, + ) + + async def _maybe_notify_outbox_prune(self, chat_id: str) -> None: + """When the outbox row cap pruned old frames, tell the app (throttled + to once per hour so a full box doesn't banner per frame).""" + pruned = self._outbox.take_overflow_pruned() + if pruned <= 0: + return + now = time.time() + if now - self._prune_notified_at < 3600.0: + return + self._prune_notified_at = now + await self._broadcast_or_log( + chat_id, + protocol.notification( + chat_id, + protocol.NOTIF_GENERIC, + "Outbox", + f"{pruned} older message(s) pruned", + ), + ) async def send_typing(self, chat_id: str, metadata: Optional[Dict[str, Any]] = None) -> None: """Send a typing indicator (``typing`` frame, on=true).""" @@ -1464,6 +1702,16 @@ class AndroidAdapter(BasePlatformAdapter): resp = protocol.channel_created(entry) resp.id = frame.id await self._ws_server.broadcast(resp) + # M5: banner + push mirror (parked in the outbox when offline). + await self._broadcast_or_log( + entry["chat_id"], + protocol.notification( + entry["chat_id"], + protocol.NOTIF_CHANNEL_CREATED, + "Channels", + f"New channel: {name}", + ), + ) async def on_channel_rename(self, frame: protocol.Frame, device_id: str) -> None: chat_id = frame.chat_id or frame.payload.get("chat_id") @@ -1496,6 +1744,16 @@ class AndroidAdapter(BasePlatformAdapter): resp = protocol.channel_renamed(entry) resp.id = frame.id await self._ws_server.broadcast(resp) + # M5: banner + push mirror (parked in the outbox when offline). + await self._broadcast_or_log( + chat_id, + protocol.notification( + chat_id, + protocol.NOTIF_CHANNEL_RENAMED, + "Channels", + f"Renamed to {name}", + ), + ) async def on_channel_set_default(self, frame: protocol.Frame, device_id: str) -> None: chat_id = frame.chat_id or frame.payload.get("chat_id") @@ -1536,6 +1794,16 @@ class AndroidAdapter(BasePlatformAdapter): resp = protocol.channel_deleted(chat_id) resp.id = frame.id await self._ws_server.broadcast(resp) + # M5: banner + push mirror (parked in the outbox when offline). + await self._broadcast_or_log( + chat_id, + protocol.notification( + chat_id, + protocol.NOTIF_CHANNEL_DELETED, + "Channels", + f"{entry.get('name') or chat_id} deleted", + ), + ) async def on_channel_list(self, frame: protocol.Frame, device_id: str) -> None: channels = self._channels.list(include_archived=False) @@ -1596,6 +1864,140 @@ class AndroidAdapter(BasePlatformAdapter): done = protocol.sync_done(self._outbox.latest_cursor(), id=frame.id) await self._ws_server.send_to(device_id, done) + # ── M5: push token registration ─────────────────────────────────────── + + async def on_fcm_register(self, frame: protocol.Frame, device_id: str) -> None: + """Update the device's push tokens (FCM rotation / ntfy topic). + + Persists to the device registry AND refreshes the live connection so + the next push targets the current token without a stale read. + """ + fcm_token = frame.payload.get("fcm_token") + ntfy_topic = frame.payload.get("ntfy_topic") + fcm_token = fcm_token if isinstance(fcm_token, str) and fcm_token else None + ntfy_topic = ntfy_topic if isinstance(ntfy_topic, str) and ntfy_topic else None + if fcm_token is None and ntfy_topic is None: + return + try: + self._devices.update_push_tokens( + device_id, fcm_token=fcm_token, ntfy_topic=ntfy_topic + ) + except Exception: + logger.warning("android: fcm.register update failed", exc_info=True) + return + conn = self._ws_server.connection(device_id) + if conn is not None: + if fcm_token is not None: + conn.fcm_token = fcm_token + if ntfy_topic is not None: + conn.ntfy_topic = ntfy_topic + logger.info("android: push tokens updated for %s", device_id) + + # ── M5: approval / clarify banners ──────────────────────────────────── + + async def send_slash_confirm( + self, + chat_id: str, + title: str, + message: str, + session_key: str, + confirm_id: str, + metadata: Optional[Dict[str, Any]] = None, + ) -> SendResult: + """Banner + push for a slash-command approval prompt. + + The gateway's text fallback still renders the actionable prompt (the + app has no inline buttons yet); the notification is the push-visible + signal (high priority: pushed even when a device is live). + """ + thread_id = _thread_id_from_metadata(metadata) + await self._broadcast_or_log( + chat_id, + protocol.notification( + chat_id, + protocol.NOTIF_APPROVAL, + title or "Approval needed", + _push_preview(message), + thread_id=thread_id, + ), + ) + return await super().send_slash_confirm( + chat_id, title, message, session_key, confirm_id, metadata=metadata + ) + + async def send_clarify( + self, + chat_id: str, + question: str, + choices: Optional[list], + clarify_id: str, + session_key: str, + metadata: Optional[Dict[str, Any]] = None, + ) -> SendResult: + """Banner + push for a clarify prompt. + + Renders the prompt as a proper ``message`` frame (the base text + fallback would route through ``send()`` and be misclassified as + commentary/tool progress) and keeps the gateway's text intercept + working via ``mark_awaiting_text``. + """ + thread_id = _thread_id_from_metadata(metadata) + await self._broadcast_or_log( + chat_id, + protocol.notification( + chat_id, + protocol.NOTIF_CLARIFY, + "Question", + _push_preview(question), + thread_id=thread_id, + ), + ) + if choices: + # Multi-select clarifies register their flag on the pending entry; + # look it up by id (mirrors the base text fallback). + _is_multi = False + try: + from tools import clarify_gateway as _cg + + with _cg._lock: + _entry = _cg._entries.get(clarify_id) + _is_multi = bool(_entry and getattr(_entry, "multi_select", False)) + except Exception: + _is_multi = False + lines = [f"❓ {question}", ""] + for i, choice in enumerate(choices, start=1): + lines.append(f" {i}. {choice}") + lines.append("") + if _is_multi: + lines.append( + "Multiple selections allowed — reply with the numbers " + "separated by commas or spaces (e.g. \"1, 3\"), the option " + "text, or your own answer." + ) + else: + lines.append("Reply with the number, the option text, or your own answer.") + text = "\n".join(lines) + # Text fallback: enable text-capture so the gateway intercept + # picks up the user's typed reply (e.g. "2" or choice text). + from tools.clarify_gateway import mark_awaiting_text + + mark_awaiting_text(clarify_id) + else: + text = f"❓ {question}" + message_id = _mint_message_id() + await self._broadcast_or_log( + chat_id, + protocol.message( + chat_id=chat_id, + message_id=message_id, + role=protocol.ROLE_ASSISTANT, + text=text, + thread_id=thread_id, + ts=int(time.time() * 1000), + ), + ) + return SendResult(success=True, message_id=message_id) + # ── Chat info ───────────────────────────────────────────────────────── def _channel_name(self, chat_id: str) -> str: @@ -1630,6 +2032,12 @@ class AndroidAdapter(BasePlatformAdapter): "media": True, # M4: media.upload/offer/pull "search": True, # M3: search frame "push": self.push_backend, + # M5: ntfy server URL (app listener discovery; "" when not ntfy). + "push_ntfy_server": ( + self._push.server_url + if isinstance(self._push, NtfyBackend) + else "" + ), "pickers": False, # M2+ } diff --git a/gateway-plugin/outbox.py b/gateway-plugin/outbox.py index a601929..b3c41be 100644 --- a/gateway-plugin/outbox.py +++ b/gateway-plugin/outbox.py @@ -28,6 +28,10 @@ logger = logging.getLogger(__name__) DEFAULT_RETENTION_HOURS = 72 _REPLAY_LIMIT = 1000 _PRUNE_INTERVAL_S = 3600.0 +# Row cap (docs/08 §8.3): a very busy offline period must not grow the outbox +# without bound; the oldest rows beyond the cap are pruned and the app is +# told (generic notification, throttled in the adapter). +DEFAULT_MAX_ROWS = 5000 class Outbox: @@ -38,12 +42,19 @@ class Outbox: ``DeviceRegistry`` / ``ChannelDirectory``). """ - def __init__(self, db_path: Path, retention_hours: int = DEFAULT_RETENTION_HOURS): + def __init__( + self, + db_path: Path, + retention_hours: int = DEFAULT_RETENTION_HOURS, + max_rows: int = DEFAULT_MAX_ROWS, + ): self._db_path = Path(db_path) self._db_path.parent.mkdir(parents=True, exist_ok=True) self._retention_hours = max(1, int(retention_hours)) + self._max_rows = max(1, int(max_rows)) self._lock = threading.Lock() self._last_prune = 0.0 + self._overflow_pruned = 0 self._conn = sqlite3.connect(str(self._db_path), check_same_thread=False) self._conn.row_factory = sqlite3.Row with self._lock: @@ -90,10 +101,37 @@ class Outbox: "VALUES (?, ?, ?, ?)", (cursor, chat_id, frame_json, now), ) + self._enforce_row_cap() self._conn.commit() self._maybe_prune() return cursor + def take_overflow_pruned(self) -> int: + """Rows pruned by the row cap since the last call (and reset to 0). + + The adapter turns a non-zero count into a (throttled) generic + notification so the app knows older frames are gone. + """ + with self._lock: + n = self._overflow_pruned + self._overflow_pruned = 0 + return n + + def _enforce_row_cap(self) -> None: + """Drop the oldest rows beyond ``max_rows`` (caller holds the lock).""" + row = self._conn.execute("SELECT COUNT(*) AS n FROM outbox").fetchone() + n = int(row["n"]) if row else 0 + excess = n - self._max_rows + if excess <= 0: + return + self._conn.execute( + "DELETE FROM outbox WHERE cursor IN (" + " SELECT cursor FROM outbox ORDER BY cursor ASC LIMIT ?)", + (excess,), + ) + self._overflow_pruned += excess + logger.info("android outbox: row cap pruned %s oldest row(s)", excess) + def latest_cursor(self) -> int: """The high-water cursor (0 when nothing has been appended).""" with self._lock: diff --git a/gateway-plugin/plugin.yaml b/gateway-plugin/plugin.yaml index f7c1cfa..a39a507 100644 --- a/gateway-plugin/plugin.yaml +++ b/gateway-plugin/plugin.yaml @@ -57,6 +57,10 @@ optional_env: description: "ntfy server URL (default https://ntfy.sh)" prompt: "ntfy server URL" password: false + - name: NTFY_AUTH_TOKEN + description: "ntfy auth token for a private topic (trust boundary)" + prompt: "ntfy auth token" + password: true - name: ANDROID_WS_CERT description: "TLS cert path for WSS (optional)" prompt: "WSS cert" diff --git a/gateway-plugin/protocol.py b/gateway-plugin/protocol.py index 26e34c6..868a186 100644 --- a/gateway-plugin/protocol.py +++ b/gateway-plugin/protocol.py @@ -76,6 +76,10 @@ TYPE_MEDIA_OFFER = "media.offer" TYPE_MEDIA_PULL = "media.pull" TYPE_MEDIA_PULL_END = "media.pull.end" +# Push / notifications (M5) +TYPE_NOTIFICATION = "notification" +TYPE_FCM_REGISTER = "fcm.register" + # --------------------------------------------------------------------------- # Error codes (``error`` frame payload.code) # --------------------------------------------------------------------------- @@ -96,6 +100,22 @@ ROLE_ASSISTANT = "assistant" ROLE_SYSTEM = "system" ROLE_CRON = "cron" +# --------------------------------------------------------------------------- +# Notification kinds (``notification`` frame payload.kind, M5) +# --------------------------------------------------------------------------- + +NOTIF_CHANNEL_CREATED = "channel_created" +NOTIF_CHANNEL_RENAMED = "channel_renamed" +NOTIF_CHANNEL_DELETED = "channel_deleted" +NOTIF_CRON = "cron" +NOTIF_APPROVAL = "approval" +NOTIF_CLARIFY = "clarify" +NOTIF_GENERIC = "generic" + +# Kinds that push even when a device is live (the app may be backgrounded; +# it decides whether to also show an in-app banner). +HIGH_PRIORITY_NOTIF_KINDS = frozenset({NOTIF_APPROVAL, NOTIF_CLARIFY, NOTIF_CRON}) + # --------------------------------------------------------------------------- # Envelope @@ -452,6 +472,41 @@ def sync_done(cursor: int, *, id: Optional[int] = None) -> Frame: return Frame(type=TYPE_SYNC_DONE, id=id, payload={"cursor": cursor}) +# --------------------------------------------------------------------------- +# Push / notification frames (M5) +# --------------------------------------------------------------------------- + +def notification( + chat_id: str, + kind: str, + title: str, + body: str, + *, + thread_id: Optional[str] = None, + ts: Optional[int] = None, +) -> Frame: + """Event: a transient in-app banner (and a push mirror when the device is + offline). ``kind`` is one of the ``NOTIF_*`` constants.""" + payload: Dict[str, Any] = {"kind": kind, "title": title, "body": body} + if ts is not None: + payload["ts"] = ts + return Frame( + type=TYPE_NOTIFICATION, chat_id=chat_id, thread_id=thread_id, payload=payload + ) + + +def fcm_register( + fcm_token: Optional[str] = None, ntfy_topic: Optional[str] = None +) -> Frame: + """Request: update the device's push tokens (FCM rotation / ntfy topic).""" + payload: Dict[str, Any] = {} + if fcm_token: + payload["fcm_token"] = fcm_token + if ntfy_topic: + payload["ntfy_topic"] = ntfy_topic + return Frame(type=TYPE_FCM_REGISTER, payload=payload) + + # --------------------------------------------------------------------------- # Media frames (M4) # --------------------------------------------------------------------------- diff --git a/gateway-plugin/push.py b/gateway-plugin/push.py index 7ae9cce..22699db 100644 --- a/gateway-plugin/push.py +++ b/gateway-plugin/push.py @@ -4,12 +4,328 @@ - ``FcmBackend``: FCM HTTP v1 via ``httpx`` + a Firebase service account (``ANDROID_FCM_SERVICE_ACCOUNT``), or a legacy server key (``ANDROID_FCM_SERVER_KEY``). - - ``NtfyBackend``: reuses hermes ntfy publish (``NTFY_TOPIC`` / - ``NTFY_SERVER_URL``). + - ``NtfyBackend``: publishes to ``NTFY_TOPIC`` on ``NTFY_SERVER_URL`` + (default ``https://ntfy.sh``) via ``httpx``; the app's listener + subscribes to the topic. Selected by ``ANDROID_PUSH_BACKEND`` (``fcm`` default, ``ntfy`` fallback). -Fired when a frame has no live subscriber; data payload drives a silent sync -on the device. +Fired when a frame has no live subscriber; the data payload drives a silent +sync on the device (docs/08-push.md). -Milestone M5. -""" \ No newline at end of file +Zero new dependencies: ``httpx`` is a hermes core dep. ``google.auth`` is NOT +installed, so the FCM service-account OAuth2 access token is minted directly +with PyJWT + cryptography (both core deps). + +Payloads carry no secrets and only a short preview (lock-screen privacy); +full content is fetched via ``sync`` over the authenticated WS. +""" + +import json +import logging +import threading +import time +from pathlib import Path +from typing import Any, Dict, Optional +from urllib.parse import quote + +import httpx + +logger = logging.getLogger(__name__) + +FCM_SCOPE = "https://www.googleapis.com/auth/firebase.messaging" +FCM_TOKEN_URL = "https://oauth2.googleapis.com/token" +FCM_V1_SEND_URL = "https://fcm.googleapis.com/v1/projects/{project_id}/messages:send" +FCM_LEGACY_SEND_URL = "https://fcm.googleapis.com/fcm/send" +# Refresh the cached access token this long before its expiry. +_TOKEN_REFRESH_MARGIN_S = 600.0 +_DEFAULT_NTFY_SERVER = "https://ntfy.sh" +_NTFY_BODY_LIMIT = 4096 +_HTTP_TIMEOUT_S = 15.0 + +_NTFY_PRIORITY = {"high": "5", "normal": "3", "low": "1"} + + +class PushBackend: + """Interface: wake one device (FCM token or ntfy topic).""" + + name: str = "push" + # DeviceRegistry column that carries this backend's target token. + token_field: str = "" + + def configured(self) -> bool: + """True when the backend has credentials to send with.""" + raise NotImplementedError + + async def send( + self, + *, + device_id: str, + chat_id: str, + title: str, + body: str, + data: Dict[str, Any], + token: str, + priority: str = "normal", + data_only: bool = False, + ) -> bool: + """Deliver one push to *token*. Returns True on success. + + ``data`` is the silent-sync payload (``chat_id``, ``kind``, + ``cursor``, optional ``thread_id``/``message_id``) -- string values + only on the wire. + """ + raise NotImplementedError + + +class FcmBackend(PushBackend): + """FCM HTTP v1 (service account) or legacy ``/fcm/send`` (server key).""" + + name = "fcm" + token_field = "fcm_token" + + def __init__( + self, + service_account: Optional[str] = None, + server_key: Optional[str] = None, + ): + self._sa_path = (service_account or "").strip() or None + self._server_key = (server_key or "").strip() or None + self._sa: Optional[Dict[str, Any]] = None + self._sa_failed = False + self._access_token: Optional[str] = None + self._token_expiry = 0.0 + self._lock = threading.Lock() + + def configured(self) -> bool: + if self._server_key: + return True + return bool(self._sa_path and Path(self._sa_path).is_file()) + + def _load_sa(self) -> Optional[Dict[str, Any]]: + if self._sa is not None: + return self._sa + if not self._sa_path or self._sa_failed: + return None + try: + with open(self._sa_path, "r", encoding="utf-8") as f: + sa = json.load(f) + if isinstance(sa, dict) and sa.get("client_email") and sa.get("private_key"): + self._sa = sa + return sa + except Exception: + logger.warning("android: FCM service account unreadable: %s", self._sa_path) + self._sa_failed = True + return None + + async def _authorization(self, client: httpx.AsyncClient) -> Optional[str]: + """Bearer token: the legacy server key, or a cached service-account + OAuth2 access token (JWT-bearer grant, minted with PyJWT).""" + if self._server_key: + return self._server_key + sa = self._load_sa() + if sa is None: + return None + now = time.time() + with self._lock: + if self._access_token and now < self._token_expiry - _TOKEN_REFRESH_MARGIN_S: + return self._access_token + import jwt # PyJWT (core dep) + + claims = { + "iss": sa["client_email"], + "scope": FCM_SCOPE, + "aud": FCM_TOKEN_URL, + "iat": int(now), + "exp": int(now) + 3600, + } + headers = {"kid": sa["private_key_id"]} if sa.get("private_key_id") else None + try: + assertion = jwt.encode( + claims, sa["private_key"], algorithm="RS256", headers=headers + ) + except Exception: + logger.warning("android: FCM JWT mint failed", exc_info=True) + return None + try: + resp = await client.post( + FCM_TOKEN_URL, + data={ + "grant_type": "urn:ietf:params:oauth:grant-type:jwt-bearer", + "assertion": assertion, + }, + timeout=_HTTP_TIMEOUT_S, + ) + except Exception: + logger.warning("android: FCM token exchange failed", exc_info=True) + return None + if resp.status_code != 200: + logger.warning( + "android: FCM token exchange HTTP %s: %s", + resp.status_code, resp.text[:200], + ) + return None + try: + data = resp.json() + except (json.JSONDecodeError, ValueError): + return None + token = data.get("access_token") + if not isinstance(token, str) or not token: + return None + with self._lock: + self._access_token = token + try: + self._token_expiry = now + float(data.get("expires_in", 3600)) + except (TypeError, ValueError): + self._token_expiry = now + 3600.0 + return token + + async def send( + self, + *, + device_id: str, + chat_id: str, + title: str, + body: str, + data: Dict[str, Any], + token: str, + priority: str = "normal", + data_only: bool = False, + ) -> bool: + if not token: + return False + data = {str(k): str(v) for k, v in (data or {}).items()} + notification = None if data_only else {"title": title or "Iris", "body": body or ""} + async with httpx.AsyncClient(timeout=_HTTP_TIMEOUT_S) as client: + if self._server_key: + payload: Dict[str, Any] = {"to": token} + if notification: + payload["notification"] = notification + if data: + payload["data"] = data + auth = self._server_key + url = FCM_LEGACY_SEND_URL + else: + sa = self._load_sa() + project_id = (sa or {}).get("project_id") + if not project_id: + return False + message: Dict[str, Any] = {"token": token} + if notification: + message["notification"] = notification + if data: + message["data"] = data + message["android"] = { + "priority": "high" if priority == "high" else "normal" + } + payload = {"message": message} + auth = await self._authorization(client) + if auth is None: + return False + url = FCM_V1_SEND_URL.format(project_id=project_id) + try: + resp = await client.post( + url, + json=payload, + headers={"Authorization": f"Bearer {auth}"}, + ) + except Exception: + logger.warning("android: FCM send failed (network)", exc_info=True) + return False + if resp.status_code >= 300: + # 404 NOT_FOUND = stale/invalid registration token. + logger.warning( + "android: FCM send HTTP %s: %s", resp.status_code, resp.text[:200] + ) + return False + return True + + +class NtfyBackend(PushBackend): + """ntfy publish (self-host friendly; zero Firebase). + + The structured payload rides in an ``X-Data`` header (JSON) the app's + listener parses; the message body is the short preview. A private topic + + ``NTFY_AUTH_TOKEN`` provides the trust boundary (docs/08 §8.7). + """ + + name = "ntfy" + token_field = "ntfy_topic" + + def __init__( + self, + topic: Optional[str] = None, + server_url: Optional[str] = None, + auth_token: Optional[str] = None, + ): + self._topic = (topic or "").strip() or None + self._server = ( + (server_url or _DEFAULT_NTFY_SERVER).strip().rstrip("/") + or _DEFAULT_NTFY_SERVER + ) + self._auth_token = (auth_token or "").strip() or None + + @property + def server_url(self) -> str: + """The ntfy server this backend publishes to (for app discovery).""" + return self._server + + def configured(self) -> bool: + return bool(self._topic) + + async def send( + self, + *, + device_id: str, + chat_id: str, + title: str, + body: str, + data: Dict[str, Any], + token: str, + priority: str = "normal", + data_only: bool = False, + ) -> bool: + topic = (token or self._topic or "").strip() + if not topic: + return False + headers = { + "Content-Type": "text/plain; charset=utf-8", + "X-Title": (title or "Iris")[:512], + "X-Priority": _NTFY_PRIORITY.get(priority, "3"), + "X-Tag": "bell", + "X-Data": json.dumps(data or {}, separators=(",", ":")), + } + if self._auth_token: + headers["Authorization"] = f"Bearer {self._auth_token}" + text = (body or "")[:_NTFY_BODY_LIMIT] + url = f"{self._server}/{quote(topic, safe='')}" + try: + async with httpx.AsyncClient(timeout=_HTTP_TIMEOUT_S) as client: + resp = await client.post( + url, content=text.encode("utf-8"), headers=headers + ) + except Exception: + logger.warning("android: ntfy publish failed (network)", exc_info=True) + return False + if resp.status_code >= 300: + logger.warning( + "android: ntfy publish HTTP %s: %s", resp.status_code, resp.text[:200] + ) + return False + return True + + +def build_push_backend( + name: Optional[str], + *, + fcm_service_account: Optional[str] = None, + fcm_server_key: Optional[str] = None, + ntfy_topic: Optional[str] = None, + ntfy_server_url: Optional[str] = None, + ntfy_auth_token: Optional[str] = None, +) -> PushBackend: + """Select the backend by name (``ANDROID_PUSH_BACKEND``; fcm default).""" + if (name or "").strip().lower() == "ntfy": + return NtfyBackend( + topic=ntfy_topic, server_url=ntfy_server_url, auth_token=ntfy_auth_token + ) + return FcmBackend(service_account=fcm_service_account, server_key=fcm_server_key) \ No newline at end of file diff --git a/gateway-plugin/tests/ws_probe.py b/gateway-plugin/tests/ws_probe.py index 0553c66..645e872 100644 --- a/gateway-plugin/tests/ws_probe.py +++ b/gateway-plugin/tests/ws_probe.py @@ -18,8 +18,12 @@ Options: --send TEXT send this message after pairing (default: "hello") --upload F M4: upload F (chunked media.upload) and attach it to the message.send via media_refs - --pull-offer M4: when a media.offer arrives during the turn, pull the - media (chunked) and verify the byte count +--pull-offer M4: when a media.offer arrives during the turn, pull the + media (chunked) and verify the byte count + --sync C M5: after pairing, send sync {cursor: C} and print the + replay + sync.done (no turn is driven) + --fcm-token M5: attach this FCM token to the hello payload + --fcm-reg M5: after pairing, send fcm.register with --fcm-token --timeout S seconds to wait for the final reply (default 120) --authfail expect an auth rejection (wrong token) and exit 0 on it """ @@ -94,6 +98,13 @@ def _print_frame(raw): extra = f" ok={payload.get('ok')} ref={payload.get('media_ref')}" elif ftype == "media.pull.end": extra = f" ok={payload.get('ok')}" + elif ftype == "notification": + extra = (f" kind={payload.get('kind')} title={payload.get('title')!r} " + f"body={(payload.get('body') or '')[:100]!r}") + elif ftype == "sync": + extra = f" cursor={payload.get('cursor')}" + elif ftype == "sync.done": + extra = f" cursor={payload.get('cursor')}" scope = f" chat={chat}" if chat else "" idpart = f" id={fid}" if fid is not None else "" print(f" <- {ftype}{idpart}{scope}{extra}") @@ -200,8 +211,10 @@ async def run(args) -> int: "caps": {"min_protocol": 1}, }, } + if args.fcm_token: + hello["payload"]["fcm_token"] = args.fcm_token await ws.send(json.dumps(hello)) - print(" -> hello") + print(" -> hello" + (f" fcm_token={args.fcm_token[:12]}…" if args.fcm_token else "")) # First response must be hello.ack (or an auth error). try: @@ -224,6 +237,35 @@ async def run(args) -> int: await ws.close() return 5 + # M5: optional fcm.register after pairing. + if args.fcm_reg: + reg_token = args.fcm_token or f"probe-{uuid.uuid4().hex[:12]}" + await ws.send(json.dumps({ + "v": 1, "type": "fcm.register", + "payload": {"fcm_token": reg_token}, + })) + print(f" -> fcm.register fcm_token={reg_token[:12]}…") + + # M5: sync catch-up mode (no turn driven). + if args.sync is not None: + await ws.send(json.dumps({ + "v": 1, "id": 1, "type": "sync", "payload": {"cursor": args.sync}, + })) + print(f" -> sync cursor={args.sync}") + while True: + raw = await asyncio.wait_for(ws.recv(), timeout=30) + data = _print_frame(raw) + if data is None: + continue + if data.get("type") == "sync.done": + print(f"== sync done at cursor {data['payload'].get('cursor')}") + await ws.close() + return 0 + if data.get("type") == "error": + print(f"!! sync failed: {data['payload']}") + await ws.close() + return 8 + if not args.send and not args.upload: print("== paired OK (no --send/--upload; exiting)") await ws.close() @@ -311,6 +353,12 @@ def main() -> int: help="M4: file to upload (chunked) and attach via media_refs") p.add_argument("--pull-offer", action="store_true", help="M4: pull any media.offer that arrives during the turn") + p.add_argument("--sync", type=int, default=None, + help="M5: send sync {cursor} after pairing, print replay, exit") + p.add_argument("--fcm-token", default="", + help="M5: FCM token to attach to the hello payload") + p.add_argument("--fcm-reg", action="store_true", + help="M5: send fcm.register after pairing (uses --fcm-token)") p.add_argument("--timeout", type=float, default=120.0) p.add_argument("--authfail", action="store_true", help="expect an auth rejection (wrong token)") diff --git a/gateway-plugin/ws_server.py b/gateway-plugin/ws_server.py index 8bf106e..24e9c17 100644 --- a/gateway-plugin/ws_server.py +++ b/gateway-plugin/ws_server.py @@ -316,18 +316,8 @@ class WsServer: await self._adapter.on_media_upload_end(frame, device_id) elif frame.type == protocol.TYPE_MEDIA_PULL: await self._adapter.on_media_pull(frame, device_id) - elif frame.type == "fcm.register": - fcm_token = frame.payload.get("fcm_token") - ntfy_topic = frame.payload.get("ntfy_topic") - if isinstance(fcm_token, str) or isinstance(ntfy_topic, str): - try: - self._devices.update_push_tokens( - device_id, - fcm_token=fcm_token if isinstance(fcm_token, str) else None, - ntfy_topic=ntfy_topic if isinstance(ntfy_topic, str) else None, - ) - except Exception: - logger.warning("android: fcm.register update failed", exc_info=True) + elif frame.type == protocol.TYPE_FCM_REGISTER: + await self._adapter.on_fcm_register(frame, device_id) # Unknown types are ignored (forward-compat). # ── Helpers ───────────────────────────────────────────────────────────