Compare commits
2
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
9f3f9842c8 | ||
|
|
acd5fb4ad0 |
No files matched your search
@@ -15,7 +15,7 @@ import iris.IrisApp
|
|||||||
import iris.platform.AndroidEnv
|
import iris.platform.AndroidEnv
|
||||||
import iris.platform.AndroidSecureStore
|
import iris.platform.AndroidSecureStore
|
||||||
import iris.platform.AppBridge
|
import iris.platform.AppBridge
|
||||||
import iris.platform.NtfyListenerService
|
import iris.platform.syncNtfyListener
|
||||||
|
|
||||||
class MainActivity : ComponentActivity() {
|
class MainActivity : ComponentActivity() {
|
||||||
private val deepLinkChatId = mutableStateOf<String?>(null)
|
private val deepLinkChatId = mutableStateOf<String?>(null)
|
||||||
@@ -31,7 +31,9 @@ class MainActivity : ComponentActivity() {
|
|||||||
val store = AndroidSecureStore(applicationContext)
|
val store = AndroidSecureStore(applicationContext)
|
||||||
handleDeepLink(intent)
|
handleDeepLink(intent)
|
||||||
requestNotificationPermission()
|
requestNotificationPermission()
|
||||||
startNtfyListener(store)
|
// M5: the ntfy listener (and its permanent notification) only runs
|
||||||
|
// when the paired gateway pushes via ntfy; unknown ("") = not yet.
|
||||||
|
syncNtfyListener(store.pushBackend)
|
||||||
setContent {
|
setContent {
|
||||||
val chatId by deepLinkChatId
|
val chatId by deepLinkChatId
|
||||||
val threadId by deepLinkThreadId
|
val threadId by deepLinkThreadId
|
||||||
@@ -58,9 +60,11 @@ class MainActivity : ComponentActivity() {
|
|||||||
* extras or an iris://chat/<id>?thread=<tid> URI). */
|
* extras or an iris://chat/<id>?thread=<tid> URI). */
|
||||||
private fun handleDeepLink(intent: Intent?) {
|
private fun handleDeepLink(intent: Intent?) {
|
||||||
val data = intent?.data
|
val data = intent?.data
|
||||||
val chatId = intent?.getStringExtra("chat_id")
|
val chatId =
|
||||||
|
intent?.getStringExtra("chat_id")
|
||||||
?: data?.pathSegments?.firstOrNull()
|
?: data?.pathSegments?.firstOrNull()
|
||||||
val threadId = intent?.getStringExtra("thread_id")
|
val threadId =
|
||||||
|
intent?.getStringExtra("thread_id")
|
||||||
?: data?.getQueryParameter("thread")
|
?: data?.getQueryParameter("thread")
|
||||||
if (!chatId.isNullOrBlank()) {
|
if (!chatId.isNullOrBlank()) {
|
||||||
deepLinkChatId.value = chatId
|
deepLinkChatId.value = chatId
|
||||||
@@ -77,9 +81,4 @@ class MainActivity : ComponentActivity() {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
private fun startNtfyListener(store: AndroidSecureStore) {
|
|
||||||
if (store.ntfyTopic.isBlank()) return
|
|
||||||
ContextCompat.startForegroundService(this, Intent(this, NtfyListenerService::class.java))
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
@@ -1,6 +1,7 @@
|
|||||||
package iris.platform
|
package iris.platform
|
||||||
|
|
||||||
import android.Manifest
|
import android.Manifest
|
||||||
|
import android.content.Intent
|
||||||
import android.content.pm.PackageManager
|
import android.content.pm.PackageManager
|
||||||
import androidx.core.content.ContextCompat
|
import androidx.core.content.ContextCompat
|
||||||
import iris.state.IrisController
|
import iris.state.IrisController
|
||||||
@@ -11,6 +12,19 @@ actual fun setActiveController(controller: Any?) {
|
|||||||
AppBridge.controller = controller as? IrisController
|
AppBridge.controller = controller as? IrisController
|
||||||
}
|
}
|
||||||
|
|
||||||
|
actual fun syncNtfyListener(backend: String) {
|
||||||
|
val context = AndroidEnv.context
|
||||||
|
val intent = Intent(context, NtfyListenerService::class.java)
|
||||||
|
val store = AndroidSecureStore(context)
|
||||||
|
if (backend == "ntfy" && store.ntfyTopic.isNotBlank()) {
|
||||||
|
ContextCompat.startForegroundService(context, intent)
|
||||||
|
} else {
|
||||||
|
// FCM (or unknown) backend: no persistent listener, no permanent
|
||||||
|
// notification. stopService is a no-op when it isn't running.
|
||||||
|
context.stopService(intent)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
actual fun postSystemNotification(
|
actual fun postSystemNotification(
|
||||||
chatId: String?,
|
chatId: String?,
|
||||||
chatName: String?,
|
chatName: String?,
|
||||||
|
|||||||
@@ -110,6 +110,10 @@ class AndroidSecureStore(
|
|||||||
get() = prefs.getString(KEY_NTFY_SERVER, "").orEmpty()
|
get() = prefs.getString(KEY_NTFY_SERVER, "").orEmpty()
|
||||||
set(value) = prefs.edit().putString(KEY_NTFY_SERVER, value).apply()
|
set(value) = prefs.edit().putString(KEY_NTFY_SERVER, value).apply()
|
||||||
|
|
||||||
|
override var pushBackend: String
|
||||||
|
get() = prefs.getString(KEY_PUSH_BACKEND, "").orEmpty()
|
||||||
|
set(value) = prefs.edit().putString(KEY_PUSH_BACKEND, value).apply()
|
||||||
|
|
||||||
override var threadsEnabled: Boolean
|
override var threadsEnabled: Boolean
|
||||||
get() = prefs.getBoolean(KEY_THREADS_ENABLED, false)
|
get() = prefs.getBoolean(KEY_THREADS_ENABLED, false)
|
||||||
set(value) = prefs.edit().putBoolean(KEY_THREADS_ENABLED, value).apply()
|
set(value) = prefs.edit().putBoolean(KEY_THREADS_ENABLED, value).apply()
|
||||||
@@ -184,6 +188,7 @@ class AndroidSecureStore(
|
|||||||
const val KEY_FCM_TOKEN = "fcm_token"
|
const val KEY_FCM_TOKEN = "fcm_token"
|
||||||
const val KEY_NTFY_TOPIC = "ntfy_topic"
|
const val KEY_NTFY_TOPIC = "ntfy_topic"
|
||||||
const val KEY_NTFY_SERVER = "ntfy_server"
|
const val KEY_NTFY_SERVER = "ntfy_server"
|
||||||
|
const val KEY_PUSH_BACKEND = "push_backend"
|
||||||
const val KEY_THREADS_ENABLED = "threads_enabled"
|
const val KEY_THREADS_ENABLED = "threads_enabled"
|
||||||
const val KEY_TOOL_DETAIL = "tool_detail"
|
const val KEY_TOOL_DETAIL = "tool_detail"
|
||||||
const val KEY_STREAMING_ENABLED = "streaming_enabled"
|
const val KEY_STREAMING_ENABLED = "streaming_enabled"
|
||||||
|
|||||||
@@ -4,6 +4,7 @@ import android.content.pm.PackageManager
|
|||||||
import androidx.core.content.ContextCompat
|
import androidx.core.content.ContextCompat
|
||||||
import com.google.firebase.messaging.FirebaseMessagingService
|
import com.google.firebase.messaging.FirebaseMessagingService
|
||||||
import com.google.firebase.messaging.RemoteMessage
|
import com.google.firebase.messaging.RemoteMessage
|
||||||
|
import iris.net.GatewayClient
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* M5: FCM handler (docs/08 §8.1).
|
* M5: FCM handler (docs/08 §8.1).
|
||||||
@@ -30,13 +31,23 @@ class IrisFirebaseMessagingService : FirebaseMessagingService() {
|
|||||||
}
|
}
|
||||||
|
|
||||||
override fun onMessageReceived(message: RemoteMessage) {
|
override fun onMessageReceived(message: RemoteMessage) {
|
||||||
|
// Foreground + live WS: the in-app banner already showed this.
|
||||||
|
if (AppBridge.foreground) return
|
||||||
|
// Live WS: the frame arrives over the socket and the controller
|
||||||
|
// mirrors it to a system notification itself — posting here would
|
||||||
|
// duplicate it (docs/08 §8.7).
|
||||||
|
if (AppBridge.controller?.client?.state?.value is GatewayClient.State.Connected) return
|
||||||
|
// Backgrounded/killed: FCM already displayed the `notification`
|
||||||
|
// payload on our behalf (the data payload only carries sync
|
||||||
|
// metadata). Posting again would show a second notification with a
|
||||||
|
// different id. Data-only messages (no notification payload) are the
|
||||||
|
// exception: the app must display them itself.
|
||||||
|
if (message.notification != null) return
|
||||||
val data = message.data
|
val data = message.data
|
||||||
val chatId = data["chat_id"] ?: "android:default"
|
val chatId = data["chat_id"] ?: "android:default"
|
||||||
val threadId = data["thread_id"]
|
val threadId = data["thread_id"]
|
||||||
val title = data["title"] ?: "Iris"
|
val title = data["title"] ?: "Iris"
|
||||||
val body = data["body"] ?: data["title"].orEmpty()
|
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)
|
if (ContextCompat.checkSelfPermission(applicationContext, android.Manifest.permission.POST_NOTIFICATIONS)
|
||||||
!= PackageManager.PERMISSION_GRANTED
|
!= PackageManager.PERMISSION_GRANTED
|
||||||
) {
|
) {
|
||||||
|
|||||||
@@ -30,6 +30,11 @@ interface SecureStore {
|
|||||||
/** ntfy server URL (M5; from hello.ack server_caps; default ntfy.sh). */
|
/** ntfy server URL (M5; from hello.ack server_caps; default ntfy.sh). */
|
||||||
var ntfyServer: String
|
var ntfyServer: String
|
||||||
|
|
||||||
|
/** Push backend of the paired gateway ("fcm"/"ntfy"; from hello.ack
|
||||||
|
* server_caps; empty when unknown — decides whether the ntfy listener
|
||||||
|
* foreground service runs at all). */
|
||||||
|
var pushBackend: String
|
||||||
|
|
||||||
/** UI setting: show threads (topics) in the chat view. */
|
/** UI setting: show threads (topics) in the chat view. */
|
||||||
var threadsEnabled: Boolean
|
var threadsEnabled: Boolean
|
||||||
|
|
||||||
@@ -61,6 +66,10 @@ interface SecureStore {
|
|||||||
* system font scale). */
|
* system font scale). */
|
||||||
var fontSizeScale: Float
|
var fontSizeScale: Float
|
||||||
|
|
||||||
fun savePairing(url: String, token: String)
|
fun savePairing(
|
||||||
|
url: String,
|
||||||
|
token: String,
|
||||||
|
)
|
||||||
|
|
||||||
fun clear()
|
fun clear()
|
||||||
}
|
}
|
||||||
@@ -70,7 +70,14 @@ class GatewayClient(
|
|||||||
sealed interface State {
|
sealed interface State {
|
||||||
data object Disconnected : State
|
data object Disconnected : State
|
||||||
data object Connecting : State
|
data object Connecting : State
|
||||||
data class Connected(val caps: ServerCaps, val channels: List<ChannelInfo>) : State
|
data class Connected(
|
||||||
|
val caps: ServerCaps,
|
||||||
|
val channels: List<ChannelInfo>,
|
||||||
|
/** M5: highest outbox cursor already pushed to this device
|
||||||
|
* (from hello.ack; 0 = never). Sync-replayed frames at/below
|
||||||
|
* it must not re-post system notifications (docs/08 §8.7). */
|
||||||
|
val lastPushedCursor: Long = 0,
|
||||||
|
) : State
|
||||||
data object Reconnecting : State
|
data object Reconnecting : State
|
||||||
data class AuthFailed(val message: String) : State
|
data class AuthFailed(val message: String) : State
|
||||||
}
|
}
|
||||||
@@ -275,7 +282,7 @@ class GatewayClient(
|
|||||||
helloAck.invokeOnCompletion { e ->
|
helloAck.invokeOnCompletion { e ->
|
||||||
if (e == null) {
|
if (e == null) {
|
||||||
val ack = helloAck.getCompleted()
|
val ack = helloAck.getCompleted()
|
||||||
_state.value = State.Connected(ack.serverCaps, ack.channels)
|
_state.value = State.Connected(ack.serverCaps, ack.channels, ack.lastPushedCursor)
|
||||||
// M5: reconnect catch-up — replay frames parked while offline.
|
// M5: reconnect catch-up — replay frames parked while offline.
|
||||||
val local = store.syncCursor
|
val local = store.syncCursor
|
||||||
if (local < ack.syncCursor) {
|
if (local < ack.syncCursor) {
|
||||||
|
|||||||
@@ -29,3 +29,12 @@ expect fun postSystemNotification(
|
|||||||
* FCM / ntfy services can reach it). No-op on desktop.
|
* FCM / ntfy services can reach it). No-op on desktop.
|
||||||
*/
|
*/
|
||||||
expect fun setActiveController(controller: Any?)
|
expect fun setActiveController(controller: Any?)
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Align the ntfy listener with the paired gateway's push backend (Android:
|
||||||
|
* start/stop the listener foreground service — and with it the permanent
|
||||||
|
* "Listening for messages" notification; no-op on desktop). The listener
|
||||||
|
* only runs when the gateway actually pushes via ntfy; FCM gateways need
|
||||||
|
* no persistent listener.
|
||||||
|
*/
|
||||||
|
expect fun syncNtfyListener(backend: String)
|
||||||
@@ -115,6 +115,11 @@ data class Frame(
|
|||||||
val type: String,
|
val type: String,
|
||||||
@SerialName("chat_id") val chatId: String? = null,
|
@SerialName("chat_id") val chatId: String? = null,
|
||||||
@SerialName("thread_id") val threadId: String? = null,
|
@SerialName("thread_id") val threadId: String? = null,
|
||||||
|
/** Outbox cursor this frame was parked under. Set only on frames
|
||||||
|
* replayed by `sync` (live frames carry none) — the app skips
|
||||||
|
* re-notifying replayed frames with `cursor <= lastPushedCursor`
|
||||||
|
* (they already woke the device via push, docs/08 §8.7). */
|
||||||
|
val cursor: Long? = null,
|
||||||
val payload: JsonElement = JsonObject(emptyMap()),
|
val payload: JsonElement = JsonObject(emptyMap()),
|
||||||
) {
|
) {
|
||||||
/** Payload as a JSON object (the wire format); parse per-type with
|
/** Payload as a JSON object (the wire format); parse per-type with
|
||||||
@@ -185,6 +190,10 @@ data class HelloAckPayload(
|
|||||||
@SerialName("server_caps") val serverCaps: ServerCaps = ServerCaps(),
|
@SerialName("server_caps") val serverCaps: ServerCaps = ServerCaps(),
|
||||||
@SerialName("sync_cursor") val syncCursor: Long = 0,
|
@SerialName("sync_cursor") val syncCursor: Long = 0,
|
||||||
val channels: List<ChannelInfo> = emptyList(),
|
val channels: List<ChannelInfo> = emptyList(),
|
||||||
|
/** M5: highest outbox cursor already delivered to THIS device via the
|
||||||
|
* push backend (0 = never). Sync-replayed frames at/below it must not
|
||||||
|
* re-post system notifications (dedupe, docs/08 §8.7). */
|
||||||
|
@SerialName("last_pushed_cursor") val lastPushedCursor: Long = 0,
|
||||||
)
|
)
|
||||||
|
|
||||||
// ── message (server -> app) ─────────────────────────────────────────────
|
// ── message (server -> app) ─────────────────────────────────────────────
|
||||||
|
|||||||
@@ -13,6 +13,7 @@ import iris.platform.PickedFile
|
|||||||
import iris.platform.isAppForeground
|
import iris.platform.isAppForeground
|
||||||
import iris.platform.mediaCacheBaseDir
|
import iris.platform.mediaCacheBaseDir
|
||||||
import iris.platform.postSystemNotification
|
import iris.platform.postSystemNotification
|
||||||
|
import iris.platform.syncNtfyListener
|
||||||
import iris.protocol.ChannelDeletedPayload
|
import iris.protocol.ChannelDeletedPayload
|
||||||
import iris.protocol.ChannelInfo
|
import iris.protocol.ChannelInfo
|
||||||
import iris.protocol.CommandsCatalogPayload
|
import iris.protocol.CommandsCatalogPayload
|
||||||
@@ -135,6 +136,19 @@ class IrisController(
|
|||||||
private val _gatewayStatus = MutableStateFlow<String?>(null)
|
private val _gatewayStatus = MutableStateFlow<String?>(null)
|
||||||
val gatewayStatus: StateFlow<String?> = _gatewayStatus.asStateFlow()
|
val gatewayStatus: StateFlow<String?> = _gatewayStatus.asStateFlow()
|
||||||
|
|
||||||
|
/** M5: highest outbox cursor already delivered to this device via the
|
||||||
|
* push backend (from hello.ack; 0 = never). Sync-replayed frames with
|
||||||
|
* `cursor <= lastPushedCursor` already woke the device via push, so the
|
||||||
|
* app must not post a second system notification for them (docs/08
|
||||||
|
* §8.7). In-memory: hello.ack refreshes it on every (re)connect. */
|
||||||
|
@Volatile
|
||||||
|
private var lastPushedCursor: Long = 0
|
||||||
|
|
||||||
|
/** True while [frame] is a sync replay that already reached the device
|
||||||
|
* via push (live frames carry no cursor and are never suppressed). */
|
||||||
|
private fun isPushedReplay(frame: iris.protocol.Frame): Boolean =
|
||||||
|
frame.cursor?.let { it <= lastPushedCursor } ?: false
|
||||||
|
|
||||||
// ── M3: threads toggle (per-app for now; per-channel lands later) ─────
|
// ── M3: threads toggle (per-app for now; per-channel lands later) ─────
|
||||||
// Persisted (Settings → "Threads").
|
// Persisted (Settings → "Threads").
|
||||||
private val _threadsEnabled = MutableStateFlow(store.threadsEnabled)
|
private val _threadsEnabled = MutableStateFlow(store.threadsEnabled)
|
||||||
@@ -392,13 +406,15 @@ class IrisController(
|
|||||||
when (frame.type) {
|
when (frame.type) {
|
||||||
TYPE_MESSAGE_STOP -> {
|
TYPE_MESSAGE_STOP -> {
|
||||||
frame.payloadAs<MessageStopPayload>()?.let {
|
frame.payloadAs<MessageStopPayload>()?.let {
|
||||||
|
if (!isPushedReplay(frame)) {
|
||||||
notifyMessageIfBackgrounded(frame.chatId, frame.threadId, it.finalText)
|
notifyMessageIfBackgrounded(frame.chatId, frame.threadId, it.finalText)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
}
|
||||||
|
|
||||||
TYPE_MESSAGE -> {
|
TYPE_MESSAGE -> {
|
||||||
frame.payloadAs<MessagePayload>()?.let {
|
frame.payloadAs<MessagePayload>()?.let {
|
||||||
if (it.role == ROLE_ASSISTANT) {
|
if (it.role == ROLE_ASSISTANT && !isPushedReplay(frame)) {
|
||||||
notifyMessageIfBackgrounded(frame.chatId, frame.threadId, it.text)
|
notifyMessageIfBackgrounded(frame.chatId, frame.threadId, it.text)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -492,8 +508,10 @@ class IrisController(
|
|||||||
// M5: WS is live but the app is backgrounded — the
|
// M5: WS is live but the app is backgrounded — the
|
||||||
// in-app banner is invisible, so mirror to a system
|
// in-app banner is invisible, so mirror to a system
|
||||||
// notification (the push backend only fires when
|
// notification (the push backend only fires when
|
||||||
// there is no live subscriber).
|
// there is no live subscriber). Suppressed for
|
||||||
if (!isAppForeground()) {
|
// sync replays that already woke the device via
|
||||||
|
// push (docs/08 §8.7).
|
||||||
|
if (!isAppForeground() && !isPushedReplay(frame)) {
|
||||||
postSystemNotification(p.chatId, null, p.title, p.body, p.threadId)
|
postSystemNotification(p.chatId, null, p.title, p.body, p.threadId)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -539,6 +557,8 @@ class IrisController(
|
|||||||
prevState = s
|
prevState = s
|
||||||
if (s is GatewayClient.State.Connected) {
|
if (s is GatewayClient.State.Connected) {
|
||||||
channels.setAll(s.channels)
|
channels.setAll(s.channels)
|
||||||
|
// M5: refresh the push-dedupe watermark (docs/08 §8.7).
|
||||||
|
lastPushedCursor = s.lastPushedCursor
|
||||||
val home = s.channels.firstOrNull { it.isDefault }?.chatId
|
val home = s.channels.firstOrNull { it.isDefault }?.chatId
|
||||||
if (home != null) {
|
if (home != null) {
|
||||||
_homeChannel.value = home
|
_homeChannel.value = home
|
||||||
@@ -554,6 +574,11 @@ class IrisController(
|
|||||||
if (s.caps.pushNtfyServer.isNotBlank()) {
|
if (s.caps.pushNtfyServer.isNotBlank()) {
|
||||||
store.ntfyServer = s.caps.pushNtfyServer
|
store.ntfyServer = s.caps.pushNtfyServer
|
||||||
}
|
}
|
||||||
|
// M5: align the ntfy listener (and its permanent
|
||||||
|
// "Listening for messages" notification) with the
|
||||||
|
// gateway's push backend — only ntfy gateways need it.
|
||||||
|
store.pushBackend = s.caps.push
|
||||||
|
syncNtfyListener(s.caps.push)
|
||||||
// M5: a deep link tapped before we were connected.
|
// M5: a deep link tapped before we were connected.
|
||||||
applyDeepLink()
|
applyDeepLink()
|
||||||
// Slash-command catalog for the composer's "/" drawer
|
// Slash-command catalog for the composer's "/" drawer
|
||||||
|
|||||||
@@ -19,6 +19,10 @@ actual fun setActiveController(controller: Any?) {
|
|||||||
DesktopBridge.controller = controller as? IrisController
|
DesktopBridge.controller = controller as? IrisController
|
||||||
}
|
}
|
||||||
|
|
||||||
|
actual fun syncNtfyListener(backend: String) {
|
||||||
|
// Desktop has no ntfy listener service.
|
||||||
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* M6: OS notifications (docs/11 §11.2). Desktop has no FCM; the tray + the
|
* M6: OS notifications (docs/11 §11.2). Desktop has no FCM; the tray + the
|
||||||
* platform notifier cover the "backgrounded" leg. Shells out to
|
* platform notifier cover the "backgrounded" leg. Shells out to
|
||||||
@@ -26,20 +30,31 @@ actual fun setActiveController(controller: Any?) {
|
|||||||
* Best effort — failures are ignored (the in-app banner is the primary path).
|
* Best effort — failures are ignored (the in-app banner is the primary path).
|
||||||
*/
|
*/
|
||||||
object DesktopNotifier {
|
object DesktopNotifier {
|
||||||
fun post(title: String, body: String) {
|
fun post(
|
||||||
|
title: String,
|
||||||
|
body: String,
|
||||||
|
) {
|
||||||
try {
|
try {
|
||||||
val os = System.getProperty("os.name").lowercase()
|
val os = System.getProperty("os.name").lowercase()
|
||||||
val cmd = when {
|
val cmd =
|
||||||
os.contains("linux") ->
|
when {
|
||||||
|
os.contains("linux") -> {
|
||||||
listOf("notify-send", "-a", "Iris", "-c", "iris", title, body)
|
listOf("notify-send", "-a", "Iris", "-c", "iris", title, body)
|
||||||
os.contains("mac") ->
|
}
|
||||||
|
|
||||||
|
os.contains("mac") -> {
|
||||||
listOf(
|
listOf(
|
||||||
"osascript", "-e",
|
"osascript",
|
||||||
|
"-e",
|
||||||
"display notification \"${esc(body)}\" with title \"${esc(title)}\"",
|
"display notification \"${esc(body)}\" with title \"${esc(title)}\"",
|
||||||
)
|
)
|
||||||
else ->
|
}
|
||||||
|
|
||||||
|
else -> {
|
||||||
listOf(
|
listOf(
|
||||||
"powershell", "-NoProfile", "-Command",
|
"powershell",
|
||||||
|
"-NoProfile",
|
||||||
|
"-Command",
|
||||||
"Add-Type -AssemblyName System.Windows.Forms; " +
|
"Add-Type -AssemblyName System.Windows.Forms; " +
|
||||||
"\$n = New-Object System.Windows.Forms.NotifyIcon; " +
|
"\$n = New-Object System.Windows.Forms.NotifyIcon; " +
|
||||||
"\$n.Icon = [System.Drawing.SystemIcons]::Information; " +
|
"\$n.Icon = [System.Drawing.SystemIcons]::Information; " +
|
||||||
@@ -49,11 +64,11 @@ object DesktopNotifier {
|
|||||||
"Start-Sleep -Milliseconds 4500; \$n.Dispose()",
|
"Start-Sleep -Milliseconds 4500; \$n.Dispose()",
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
}
|
||||||
ProcessBuilder(cmd).redirectErrorStream(true).start()
|
ProcessBuilder(cmd).redirectErrorStream(true).start()
|
||||||
} catch (_: Exception) {
|
} catch (_: Exception) {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
private fun esc(s: String): String =
|
private fun esc(s: String): String = s.replace("\\", "\\\\").replace("\"", "\\\"").replace("\n", " ")
|
||||||
s.replace("\\", "\\\\").replace("\"", "\\\"").replace("\n", " ")
|
|
||||||
}
|
}
|
||||||
@@ -37,6 +37,7 @@ class DesktopSecureStore : SecureStore {
|
|||||||
val fcmToken: String = "",
|
val fcmToken: String = "",
|
||||||
val ntfyTopic: String = "",
|
val ntfyTopic: String = "",
|
||||||
val ntfyServer: String = "",
|
val ntfyServer: String = "",
|
||||||
|
val pushBackend: String = "",
|
||||||
val threadsEnabled: Boolean = false,
|
val threadsEnabled: Boolean = false,
|
||||||
val toolDetail: String = "truncated",
|
val toolDetail: String = "truncated",
|
||||||
val streamingEnabled: Boolean = true,
|
val streamingEnabled: Boolean = true,
|
||||||
@@ -153,6 +154,13 @@ class DesktopSecureStore : SecureStore {
|
|||||||
save(d.copy(ntfyServer = value))
|
save(d.copy(ntfyServer = value))
|
||||||
}
|
}
|
||||||
|
|
||||||
|
override var pushBackend: String
|
||||||
|
get() = load().pushBackend
|
||||||
|
set(value) {
|
||||||
|
val d = load()
|
||||||
|
save(d.copy(pushBackend = value))
|
||||||
|
}
|
||||||
|
|
||||||
override var threadsEnabled: Boolean
|
override var threadsEnabled: Boolean
|
||||||
get() = load().threadsEnabled
|
get() = load().threadsEnabled
|
||||||
set(value) {
|
set(value) {
|
||||||
|
|||||||
@@ -21,6 +21,10 @@ Every frame:
|
|||||||
- `v` — protocol version (currently `1`). Server rejects unknown major versions.
|
- `v` — protocol version (currently `1`). Server rejects unknown major versions.
|
||||||
- `id` — request id (client-chosen). Responses/acks echo it. Events have no `id`.
|
- `id` — request id (client-chosen). Responses/acks echo it. Events have no `id`.
|
||||||
- `chat_id` / `thread_id` — top-level for convenience; may also be in `payload`.
|
- `chat_id` / `thread_id` — top-level for convenience; may also be in `payload`.
|
||||||
|
- `cursor` — outbox cursor the frame was parked under. Present **only** on
|
||||||
|
frames replayed by `sync` (live frames carry none). The app compares it
|
||||||
|
against `last_pushed_cursor` from `hello.ack` to skip re-notifying frames
|
||||||
|
that already woke the device via push (`08-push.md` §8.7).
|
||||||
- Unknown `type`s are ignored (forward-compat); unknown `payload` fields ignored.
|
- Unknown `type`s are ignored (forward-compat); unknown `payload` fields ignored.
|
||||||
|
|
||||||
**Binary media frames** are not JSON. A media transfer is: one JSON header frame
|
**Binary media frames** are not JSON. A media transfer is: one JSON header frame
|
||||||
@@ -36,9 +40,14 @@ Pairing succeeded.
|
|||||||
"server_caps":{"streaming":true,"reasoning":true,"tools":true,"media":true,
|
"server_caps":{"streaming":true,"reasoning":true,"tools":true,"media":true,
|
||||||
"search":true,"push":"fcm","pickers":true},
|
"search":true,"push":"fcm","pickers":true},
|
||||||
"sync_cursor":1042,
|
"sync_cursor":1042,
|
||||||
|
"last_pushed_cursor":1040,
|
||||||
"channels":[{"chat_id":"android:default","name":"Default","kind":"default","is_default":true}]
|
"channels":[{"chat_id":"android:default","name":"Default","kind":"default","is_default":true}]
|
||||||
}}
|
}}
|
||||||
```
|
```
|
||||||
|
`last_pushed_cursor` is the highest outbox cursor already delivered to THIS
|
||||||
|
device via the push backend (0 = never). The app skips system notifications
|
||||||
|
for sync-replayed frames with `cursor <= last_pushed_cursor` — they already
|
||||||
|
woke the device via push (dedupe, `08-push.md` §8.7).
|
||||||
|
|
||||||
### `message`
|
### `message`
|
||||||
A final / standalone message.
|
A final / standalone message.
|
||||||
|
|||||||
@@ -103,3 +103,33 @@ persist until acted on.
|
|||||||
over the authenticated WS.
|
over the authenticated WS.
|
||||||
- ntfy: use a **private topic + auth token** for any real trust boundary (hermes
|
- ntfy: use a **private topic + auth token** for any real trust boundary (hermes
|
||||||
ntfy adapter guidance).
|
ntfy adapter guidance).
|
||||||
|
|
||||||
|
## 8.8 Notification dedupe (push vs. sync)
|
||||||
|
|
||||||
|
A message sent while the device is offline is notified **twice** by naive
|
||||||
|
design: once by the push (FCM displays the `notification` payload), and again
|
||||||
|
when the app reconnects, syncs the outbox, and mirrors the replayed frames to
|
||||||
|
system notifications (the background-mirror path, §8.5). The fix is a per-
|
||||||
|
device push watermark:
|
||||||
|
|
||||||
|
- **Gateway** records the highest outbox cursor delivered to each device via
|
||||||
|
the push backend (`devices.last_pushed_cursor`, advanced only on a
|
||||||
|
*successful* send) and returns it in `hello.ack` as `last_pushed_cursor`.
|
||||||
|
- **Gateway** coalesces back-to-back pushes per chat (5 s window): a cron
|
||||||
|
delivery parks a notification frame AND a message frame, and only the first
|
||||||
|
pushes — the second reaches the app via sync (tap the first notification).
|
||||||
|
- **Gateway** tags every `sync`-replayed frame with its outbox cursor in the
|
||||||
|
frame envelope (`cursor`; live frames carry none).
|
||||||
|
- **App** skips system notifications for replayed frames with
|
||||||
|
`cursor <= lastPushedCursor` (they already woke the device). Live frames are
|
||||||
|
never suppressed — that is exactly the case where no push fired and the app
|
||||||
|
must notify itself.
|
||||||
|
- **App** FCM handler (`onMessageReceived`) posts nothing when the WS is
|
||||||
|
connected (the background-mirror path handles it) and nothing when the
|
||||||
|
message carried a `notification` payload (FCM already displayed it);
|
||||||
|
data-only messages are the exception (the app must display them itself).
|
||||||
|
|
||||||
|
Residual edge: FCM is at-least-once, so a lost device ack can still produce a
|
||||||
|
duplicate *system-displayed* notification (two `FCM-Notification:*` ids). The
|
||||||
|
designed evolution is the data-only push option (§8.2.1), which moves display
|
||||||
|
into the app and lets it use a stable per-message notification id.
|
||||||
@@ -12,6 +12,7 @@
|
|||||||
"type": { "type": "string", "description": "Frame type (see frame_types)." },
|
"type": { "type": "string", "description": "Frame type (see frame_types)." },
|
||||||
"chat_id": { "type": "string", "description": "Optional chat scope (e.g. android:default, android:chan_7)." },
|
"chat_id": { "type": "string", "description": "Optional chat scope (e.g. android:default, android:chan_7)." },
|
||||||
"thread_id": { "type": "string", "description": "Optional thread scope within a chat_id." },
|
"thread_id": { "type": "string", "description": "Optional thread scope within a chat_id." },
|
||||||
|
"cursor": { "type": "integer", "description": "Outbox cursor the frame was parked under. Present ONLY on frames replayed by sync (live frames carry none). The app skips re-notifying replayed frames with cursor <= last_pushed_cursor (docs/08 §8.7)." },
|
||||||
"payload": { "type": "object", "description": "Type-specific payload." }
|
"payload": { "type": "object", "description": "Type-specific payload." }
|
||||||
}
|
}
|
||||||
},
|
},
|
||||||
@@ -22,6 +23,7 @@
|
|||||||
"payload": {
|
"payload": {
|
||||||
"server_caps": { "type": "object", "properties": { "streaming": {"type":"boolean"}, "reasoning": {"type":"boolean"}, "tools": {"type":"boolean"}, "media": {"type":"boolean"}, "search": {"type":"boolean"}, "push": {"type":"string","enum":["fcm","ntfy","none"]}, "push_ntfy_server": {"type":"string","description":"ntfy server URL for the app's listener; empty string when the backend is not ntfy."}, "pickers": {"type":"boolean"} } },
|
"server_caps": { "type": "object", "properties": { "streaming": {"type":"boolean"}, "reasoning": {"type":"boolean"}, "tools": {"type":"boolean"}, "media": {"type":"boolean"}, "search": {"type":"boolean"}, "push": {"type":"string","enum":["fcm","ntfy","none"]}, "push_ntfy_server": {"type":"string","description":"ntfy server URL for the app's listener; empty string when the backend is not ntfy."}, "pickers": {"type":"boolean"} } },
|
||||||
"sync_cursor": { "type": "integer" },
|
"sync_cursor": { "type": "integer" },
|
||||||
|
"last_pushed_cursor": { "type": "integer", "description": "Highest outbox cursor already delivered to THIS device via the push backend (0 = never). The app skips system notifications for sync-replayed frames at/below it (dedupe, docs/08 §8.7)." },
|
||||||
"channels": { "type": "array", "items": { "$ref": "#/definitions/channel" } }
|
"channels": { "type": "array", "items": { "$ref": "#/definitions/channel" } }
|
||||||
}
|
}
|
||||||
},
|
},
|
||||||
|
|||||||
@@ -439,6 +439,11 @@ def _strip_streaming_cursor(text: str) -> str:
|
|||||||
return text
|
return text
|
||||||
|
|
||||||
|
|
||||||
|
# M5: coalesce back-to-back pushes for the same chat (a cron delivery parks
|
||||||
|
# a notification frame AND a message frame; only the first should push).
|
||||||
|
_PUSH_COALESCE_S = 5.0
|
||||||
|
|
||||||
|
|
||||||
def _push_preview(text: Any, limit: int = 120) -> str:
|
def _push_preview(text: Any, limit: int = 120) -> str:
|
||||||
"""Short single-line preview for push bodies (lock-screen privacy: no
|
"""Short single-line preview for push bodies (lock-screen privacy: no
|
||||||
secrets, no full bodies -- full content arrives via ``sync``)."""
|
secrets, no full bodies -- full content arrives via ``sync``)."""
|
||||||
@@ -1019,6 +1024,9 @@ class AndroidAdapter(BasePlatformAdapter):
|
|||||||
ntfy_auth_token=_get_scoped_secret("NTFY_AUTH_TOKEN"),
|
ntfy_auth_token=_get_scoped_secret("NTFY_AUTH_TOKEN"),
|
||||||
)
|
)
|
||||||
self._prune_notified_at = 0.0
|
self._prune_notified_at = 0.0
|
||||||
|
# M5: per-chat push throttle (epoch seconds of the last successful
|
||||||
|
# push). The coalesced frame still reaches the app via sync.
|
||||||
|
self._last_push_at: Dict[str, float] = {}
|
||||||
|
|
||||||
def _turn_state(self, chat_id: str) -> _TurnState:
|
def _turn_state(self, chat_id: str) -> _TurnState:
|
||||||
st = self._turns.get(chat_id)
|
st = self._turns.get(chat_id)
|
||||||
@@ -1537,6 +1545,16 @@ class AndroidAdapter(BasePlatformAdapter):
|
|||||||
summary = self._push_summary(frame)
|
summary = self._push_summary(frame)
|
||||||
if summary is None:
|
if summary is None:
|
||||||
return
|
return
|
||||||
|
# M5: coalesce back-to-back pushes for the same chat (cron delivery
|
||||||
|
# = notification frame + message frame). The suppressed frame is
|
||||||
|
# still synced when the app reconnects.
|
||||||
|
now = time.time()
|
||||||
|
if now - self._last_push_at.get(chat_id, 0.0) < _PUSH_COALESCE_S:
|
||||||
|
logger.info(
|
||||||
|
"android: push coalesced for %s (%s frame within %.0fs of last push)",
|
||||||
|
chat_id, frame.type, _PUSH_COALESCE_S,
|
||||||
|
)
|
||||||
|
return
|
||||||
title, body, kind, priority = summary
|
title, body, kind, priority = summary
|
||||||
backend = self._push
|
backend = self._push
|
||||||
if backend is None or not backend.token_field:
|
if backend is None or not backend.token_field:
|
||||||
@@ -1578,6 +1596,15 @@ class AndroidAdapter(BasePlatformAdapter):
|
|||||||
logger.warning("android: push via %s failed", backend.name, exc_info=True)
|
logger.warning("android: push via %s failed", backend.name, exc_info=True)
|
||||||
continue
|
continue
|
||||||
if ok:
|
if ok:
|
||||||
|
# M5: remember that this cursor reached the device via push,
|
||||||
|
# so the app can dedupe it on the next sync replay.
|
||||||
|
try:
|
||||||
|
self._devices.update_push_cursor(device_id, cursor)
|
||||||
|
except Exception:
|
||||||
|
logger.warning(
|
||||||
|
"android: push cursor update failed for %s", device_id, exc_info=True
|
||||||
|
)
|
||||||
|
self._last_push_at[chat_id] = time.time()
|
||||||
logger.info(
|
logger.info(
|
||||||
"android: push via %s -> %s (%s, chat=%s)",
|
"android: push via %s -> %s (%s, chat=%s)",
|
||||||
backend.name, device_id, frame.type, chat_id,
|
backend.name, device_id, frame.type, chat_id,
|
||||||
@@ -2362,6 +2389,10 @@ class AndroidAdapter(BasePlatformAdapter):
|
|||||||
id=raw.get("id") if isinstance(raw.get("id"), int) else None,
|
id=raw.get("id") if isinstance(raw.get("id"), int) else None,
|
||||||
chat_id=raw.get("chat_id") if isinstance(raw.get("chat_id"), str) else e.get("chat_id"),
|
chat_id=raw.get("chat_id") if isinstance(raw.get("chat_id"), str) else e.get("chat_id"),
|
||||||
thread_id=raw.get("thread_id") if isinstance(raw.get("thread_id"), str) else None,
|
thread_id=raw.get("thread_id") if isinstance(raw.get("thread_id"), str) else None,
|
||||||
|
# M5: tag replayed frames with their outbox cursor so the app
|
||||||
|
# can skip re-notifying frames that already woke the device
|
||||||
|
# via push (cursor <= last_pushed_cursor, docs/08 §8.7).
|
||||||
|
cursor=e.get("cursor"),
|
||||||
v=raw.get("v") if isinstance(raw.get("v"), int) else protocol.PROTOCOL_VERSION,
|
v=raw.get("v") if isinstance(raw.get("v"), int) else protocol.PROTOCOL_VERSION,
|
||||||
)
|
)
|
||||||
await self._ws_server.send_to(device_id, replayed)
|
await self._ws_server.send_to(device_id, replayed)
|
||||||
|
|||||||
+53
-10
@@ -17,7 +17,7 @@ import sqlite3
|
|||||||
import threading
|
import threading
|
||||||
import time
|
import time
|
||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
from typing import Any, Dict, List, Optional
|
from typing import Any
|
||||||
from urllib.parse import quote
|
from urllib.parse import quote
|
||||||
|
|
||||||
logger = logging.getLogger(__name__)
|
logger = logging.getLogger(__name__)
|
||||||
@@ -31,7 +31,7 @@ def generate_token() -> str:
|
|||||||
return secrets.token_hex(TOKEN_BYTES)
|
return secrets.token_hex(TOKEN_BYTES)
|
||||||
|
|
||||||
|
|
||||||
def verify_token(provided: Optional[str], expected: Optional[str]) -> bool:
|
def verify_token(provided: str | None, expected: str | None) -> bool:
|
||||||
"""Constant-time token comparison (never time-leaks the token)."""
|
"""Constant-time token comparison (never time-leaks the token)."""
|
||||||
if not provided or not expected:
|
if not provided or not expected:
|
||||||
return False
|
return False
|
||||||
@@ -65,6 +65,7 @@ def pairing_url(host: str, port: int, secure: bool = False) -> str:
|
|||||||
# Device registry (SQLite)
|
# Device registry (SQLite)
|
||||||
# ---------------------------------------------------------------------------
|
# ---------------------------------------------------------------------------
|
||||||
|
|
||||||
|
|
||||||
class DeviceRegistry:
|
class DeviceRegistry:
|
||||||
"""Persistent device registry under ``get_hermes_home()/"android"``.
|
"""Persistent device registry under ``get_hermes_home()/"android"``.
|
||||||
|
|
||||||
@@ -88,20 +89,32 @@ class DeviceRegistry:
|
|||||||
caps TEXT NOT NULL DEFAULT '{}',
|
caps TEXT NOT NULL DEFAULT '{}',
|
||||||
fcm_token TEXT,
|
fcm_token TEXT,
|
||||||
ntfy_topic TEXT,
|
ntfy_topic TEXT,
|
||||||
|
last_pushed_cursor INTEGER NOT NULL DEFAULT 0,
|
||||||
last_seen REAL NOT NULL DEFAULT 0,
|
last_seen REAL NOT NULL DEFAULT 0,
|
||||||
created REAL NOT NULL DEFAULT 0
|
created REAL NOT NULL DEFAULT 0
|
||||||
)
|
)
|
||||||
"""
|
"""
|
||||||
)
|
)
|
||||||
|
# M5: migrate pre-push-cursor databases (the column carries the
|
||||||
|
# highest outbox cursor already delivered to the device via the
|
||||||
|
# push backend; hello.ack returns it for notification dedupe).
|
||||||
|
cols = {
|
||||||
|
r["name"]
|
||||||
|
for r in self._conn.execute("PRAGMA table_info(devices)").fetchall()
|
||||||
|
}
|
||||||
|
if "last_pushed_cursor" not in cols:
|
||||||
|
self._conn.execute(
|
||||||
|
"ALTER TABLE devices ADD COLUMN last_pushed_cursor INTEGER NOT NULL DEFAULT 0"
|
||||||
|
)
|
||||||
self._conn.commit()
|
self._conn.commit()
|
||||||
|
|
||||||
def upsert(
|
def upsert(
|
||||||
self,
|
self,
|
||||||
device_id: str,
|
device_id: str,
|
||||||
name: str,
|
name: str,
|
||||||
caps: Optional[Dict[str, Any]] = None,
|
caps: dict[str, Any] | None = None,
|
||||||
fcm_token: Optional[str] = None,
|
fcm_token: str | None = None,
|
||||||
ntfy_topic: Optional[str] = None,
|
ntfy_topic: str | None = None,
|
||||||
) -> None:
|
) -> None:
|
||||||
now = time.time()
|
now = time.time()
|
||||||
caps_json = json.dumps(caps or {}, separators=(",", ":"))
|
caps_json = json.dumps(caps or {}, separators=(",", ":"))
|
||||||
@@ -125,8 +138,8 @@ class DeviceRegistry:
|
|||||||
def update_push_tokens(
|
def update_push_tokens(
|
||||||
self,
|
self,
|
||||||
device_id: str,
|
device_id: str,
|
||||||
fcm_token: Optional[str] = None,
|
fcm_token: str | None = None,
|
||||||
ntfy_topic: Optional[str] = None,
|
ntfy_topic: str | None = None,
|
||||||
) -> None:
|
) -> None:
|
||||||
with self._lock:
|
with self._lock:
|
||||||
self._conn.execute(
|
self._conn.execute(
|
||||||
@@ -149,14 +162,43 @@ class DeviceRegistry:
|
|||||||
)
|
)
|
||||||
self._conn.commit()
|
self._conn.commit()
|
||||||
|
|
||||||
def get(self, device_id: str) -> Optional[Dict[str, Any]]:
|
def update_push_cursor(self, device_id: str, cursor: int) -> None:
|
||||||
|
"""Advance the device's last-pushed cursor (monotonic; never
|
||||||
|
regresses). Called after a successful push send."""
|
||||||
|
try:
|
||||||
|
cursor = max(0, int(cursor or 0))
|
||||||
|
except (TypeError, ValueError):
|
||||||
|
return
|
||||||
|
with self._lock:
|
||||||
|
self._conn.execute(
|
||||||
|
"""
|
||||||
|
UPDATE devices SET last_pushed_cursor = MAX(last_pushed_cursor, ?)
|
||||||
|
WHERE device_id = ?
|
||||||
|
""",
|
||||||
|
(cursor, device_id),
|
||||||
|
)
|
||||||
|
self._conn.commit()
|
||||||
|
|
||||||
|
def last_pushed_cursor(self, device_id: str) -> int:
|
||||||
|
"""Highest outbox cursor pushed to this device (0 = never/unknown)."""
|
||||||
|
with self._lock:
|
||||||
|
row = self._conn.execute(
|
||||||
|
"SELECT last_pushed_cursor FROM devices WHERE device_id = ?",
|
||||||
|
(device_id,),
|
||||||
|
).fetchone()
|
||||||
|
try:
|
||||||
|
return int(row["last_pushed_cursor"]) if row else 0
|
||||||
|
except (TypeError, ValueError, KeyError, IndexError):
|
||||||
|
return 0
|
||||||
|
|
||||||
|
def get(self, device_id: str) -> dict[str, Any] | None:
|
||||||
with self._lock:
|
with self._lock:
|
||||||
row = self._conn.execute(
|
row = self._conn.execute(
|
||||||
"SELECT * FROM devices WHERE device_id = ?", (device_id,)
|
"SELECT * FROM devices WHERE device_id = ?", (device_id,)
|
||||||
).fetchone()
|
).fetchone()
|
||||||
return _row_to_device(row) if row else None
|
return _row_to_device(row) if row else None
|
||||||
|
|
||||||
def list(self) -> List[Dict[str, Any]]:
|
def list(self) -> list[dict[str, Any]]:
|
||||||
with self._lock:
|
with self._lock:
|
||||||
rows = self._conn.execute(
|
rows = self._conn.execute(
|
||||||
"SELECT * FROM devices ORDER BY last_seen DESC"
|
"SELECT * FROM devices ORDER BY last_seen DESC"
|
||||||
@@ -171,7 +213,7 @@ class DeviceRegistry:
|
|||||||
pass
|
pass
|
||||||
|
|
||||||
|
|
||||||
def _row_to_device(row: sqlite3.Row) -> Dict[str, Any]:
|
def _row_to_device(row: sqlite3.Row) -> dict[str, Any]:
|
||||||
try:
|
try:
|
||||||
caps = json.loads(row["caps"] or "{}")
|
caps = json.loads(row["caps"] or "{}")
|
||||||
if not isinstance(caps, dict):
|
if not isinstance(caps, dict):
|
||||||
@@ -184,6 +226,7 @@ def _row_to_device(row: sqlite3.Row) -> Dict[str, Any]:
|
|||||||
"caps": caps,
|
"caps": caps,
|
||||||
"fcm_token": row["fcm_token"],
|
"fcm_token": row["fcm_token"],
|
||||||
"ntfy_topic": row["ntfy_topic"],
|
"ntfy_topic": row["ntfy_topic"],
|
||||||
|
"last_pushed_cursor": row["last_pushed_cursor"] or 0,
|
||||||
"last_seen": row["last_seen"],
|
"last_seen": row["last_seen"],
|
||||||
"created": row["created"],
|
"created": row["created"],
|
||||||
}
|
}
|
||||||
+116
-77
@@ -17,7 +17,7 @@ Milestone M5: notification, fcm.register, read.receipt, status.
|
|||||||
|
|
||||||
import json
|
import json
|
||||||
from dataclasses import dataclass, field
|
from dataclasses import dataclass, field
|
||||||
from typing import Any, Dict, List, Optional
|
from typing import Any, Optional
|
||||||
|
|
||||||
PROTOCOL_VERSION = 1
|
PROTOCOL_VERSION = 1
|
||||||
|
|
||||||
@@ -146,30 +146,38 @@ STATUS_DEGRADED = "degraded"
|
|||||||
# Envelope
|
# Envelope
|
||||||
# ---------------------------------------------------------------------------
|
# ---------------------------------------------------------------------------
|
||||||
|
|
||||||
|
|
||||||
@dataclass
|
@dataclass
|
||||||
class Frame:
|
class Frame:
|
||||||
"""One wire frame.
|
"""One wire frame.
|
||||||
|
|
||||||
``v`` is always serialised; ``id``/``chat_id``/``thread_id`` are
|
``v`` is always serialised; ``id``/``chat_id``/``thread_id``/``cursor``
|
||||||
omitted when ``None`` (events carry no ``id``; chat-scoped frames carry
|
are omitted when ``None`` (events carry no ``id``; chat-scoped frames
|
||||||
``chat_id``/``thread_id`` at the top level for convenience).
|
carry ``chat_id``/``thread_id`` at the top level for convenience).
|
||||||
|
|
||||||
|
``cursor`` is set only on frames replayed by ``sync``: the outbox cursor
|
||||||
|
the frame was parked under. The app uses it to skip re-notifying frames
|
||||||
|
that already woke the device via push (docs/08 §8.7).
|
||||||
"""
|
"""
|
||||||
|
|
||||||
type: str
|
type: str
|
||||||
payload: Dict[str, Any] = field(default_factory=dict)
|
payload: dict[str, Any] = field(default_factory=dict)
|
||||||
id: Optional[int] = None
|
id: int | None = None
|
||||||
chat_id: Optional[str] = None
|
chat_id: str | None = None
|
||||||
thread_id: Optional[str] = None
|
thread_id: str | None = None
|
||||||
|
cursor: int | None = None
|
||||||
v: int = PROTOCOL_VERSION
|
v: int = PROTOCOL_VERSION
|
||||||
|
|
||||||
def to_dict(self) -> Dict[str, Any]:
|
def to_dict(self) -> dict[str, Any]:
|
||||||
d: Dict[str, Any] = {"v": self.v, "type": self.type}
|
d: dict[str, Any] = {"v": self.v, "type": self.type}
|
||||||
if self.id is not None:
|
if self.id is not None:
|
||||||
d["id"] = self.id
|
d["id"] = self.id
|
||||||
if self.chat_id is not None:
|
if self.chat_id is not None:
|
||||||
d["chat_id"] = self.chat_id
|
d["chat_id"] = self.chat_id
|
||||||
if self.thread_id is not None:
|
if self.thread_id is not None:
|
||||||
d["thread_id"] = self.thread_id
|
d["thread_id"] = self.thread_id
|
||||||
|
if self.cursor is not None:
|
||||||
|
d["cursor"] = self.cursor
|
||||||
d["payload"] = self.payload
|
d["payload"] = self.payload
|
||||||
return d
|
return d
|
||||||
|
|
||||||
@@ -202,17 +210,29 @@ class Frame:
|
|||||||
thread_id = data.get("thread_id")
|
thread_id = data.get("thread_id")
|
||||||
if not isinstance(thread_id, str):
|
if not isinstance(thread_id, str):
|
||||||
thread_id = None
|
thread_id = None
|
||||||
return cls(type=ftype, payload=payload, id=fid, chat_id=chat_id, thread_id=thread_id)
|
cursor = data.get("cursor")
|
||||||
|
if not isinstance(cursor, int) or isinstance(cursor, bool):
|
||||||
|
cursor = None
|
||||||
|
return cls(
|
||||||
|
type=ftype,
|
||||||
|
payload=payload,
|
||||||
|
id=fid,
|
||||||
|
chat_id=chat_id,
|
||||||
|
thread_id=thread_id,
|
||||||
|
cursor=cursor,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
# ---------------------------------------------------------------------------
|
# ---------------------------------------------------------------------------
|
||||||
# Frame constructors (server -> app)
|
# Frame constructors (server -> app)
|
||||||
# ---------------------------------------------------------------------------
|
# ---------------------------------------------------------------------------
|
||||||
|
|
||||||
|
|
||||||
def hello_ack(
|
def hello_ack(
|
||||||
server_caps: Dict[str, Any],
|
server_caps: dict[str, Any],
|
||||||
sync_cursor: int = 0,
|
sync_cursor: int = 0,
|
||||||
channels: Optional[list] = None,
|
channels: list | None = None,
|
||||||
|
last_pushed_cursor: int = 0,
|
||||||
) -> Frame:
|
) -> Frame:
|
||||||
return Frame(
|
return Frame(
|
||||||
type=TYPE_HELLO_ACK,
|
type=TYPE_HELLO_ACK,
|
||||||
@@ -220,6 +240,11 @@ def hello_ack(
|
|||||||
"server_caps": server_caps,
|
"server_caps": server_caps,
|
||||||
"sync_cursor": sync_cursor,
|
"sync_cursor": sync_cursor,
|
||||||
"channels": channels or [],
|
"channels": channels or [],
|
||||||
|
# M5: highest outbox cursor already delivered to THIS device via
|
||||||
|
# the push backend (0 = never). The app skips system
|
||||||
|
# notifications for sync-replayed frames at/below it (dedupe,
|
||||||
|
# docs/08 §8.7).
|
||||||
|
"last_pushed_cursor": last_pushed_cursor,
|
||||||
},
|
},
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -230,15 +255,15 @@ def message(
|
|||||||
role: str,
|
role: str,
|
||||||
text: str,
|
text: str,
|
||||||
*,
|
*,
|
||||||
thread_id: Optional[str] = None,
|
thread_id: str | None = None,
|
||||||
reasoning: Optional[str] = None,
|
reasoning: str | None = None,
|
||||||
media: Optional[list] = None,
|
media: list | None = None,
|
||||||
reply_to: Optional[str] = None,
|
reply_to: str | None = None,
|
||||||
model: Optional[str] = None,
|
model: str | None = None,
|
||||||
tokens: Optional[int] = None,
|
tokens: int | None = None,
|
||||||
ts: Optional[int] = None,
|
ts: int | None = None,
|
||||||
) -> Frame:
|
) -> Frame:
|
||||||
payload: Dict[str, Any] = {
|
payload: dict[str, Any] = {
|
||||||
"message_id": message_id,
|
"message_id": message_id,
|
||||||
"role": role,
|
"role": role,
|
||||||
"text": text,
|
"text": text,
|
||||||
@@ -255,10 +280,12 @@ def message(
|
|||||||
payload["tokens"] = tokens
|
payload["tokens"] = tokens
|
||||||
if ts is not None:
|
if ts is not None:
|
||||||
payload["ts"] = ts
|
payload["ts"] = ts
|
||||||
return Frame(type=TYPE_MESSAGE, chat_id=chat_id, thread_id=thread_id, payload=payload)
|
return Frame(
|
||||||
|
type=TYPE_MESSAGE, chat_id=chat_id, thread_id=thread_id, payload=payload
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
def typing(chat_id: str, on: bool = True, *, thread_id: Optional[str] = None) -> Frame:
|
def typing(chat_id: str, on: bool = True, *, thread_id: str | None = None) -> Frame:
|
||||||
return Frame(
|
return Frame(
|
||||||
type=TYPE_TYPING,
|
type=TYPE_TYPING,
|
||||||
chat_id=chat_id,
|
chat_id=chat_id,
|
||||||
@@ -271,12 +298,13 @@ def typing(chat_id: str, on: bool = True, *, thread_id: Optional[str] = None) ->
|
|||||||
# Streaming frames (M2)
|
# Streaming frames (M2)
|
||||||
# ---------------------------------------------------------------------------
|
# ---------------------------------------------------------------------------
|
||||||
|
|
||||||
|
|
||||||
def message_start(
|
def message_start(
|
||||||
chat_id: str,
|
chat_id: str,
|
||||||
message_id: str,
|
message_id: str,
|
||||||
role: str = ROLE_ASSISTANT,
|
role: str = ROLE_ASSISTANT,
|
||||||
*,
|
*,
|
||||||
thread_id: Optional[str] = None,
|
thread_id: str | None = None,
|
||||||
) -> Frame:
|
) -> Frame:
|
||||||
"""Open a streaming bubble."""
|
"""Open a streaming bubble."""
|
||||||
return Frame(
|
return Frame(
|
||||||
@@ -292,7 +320,7 @@ def message_update(
|
|||||||
message_id: str,
|
message_id: str,
|
||||||
text: str,
|
text: str,
|
||||||
*,
|
*,
|
||||||
thread_id: Optional[str] = None,
|
thread_id: str | None = None,
|
||||||
) -> Frame:
|
) -> Frame:
|
||||||
"""Replace the live bubble text (full snapshot)."""
|
"""Replace the live bubble text (full snapshot)."""
|
||||||
return Frame(
|
return Frame(
|
||||||
@@ -308,14 +336,14 @@ def message_stop(
|
|||||||
message_id: str,
|
message_id: str,
|
||||||
final_text: str,
|
final_text: str,
|
||||||
*,
|
*,
|
||||||
thread_id: Optional[str] = None,
|
thread_id: str | None = None,
|
||||||
reasoning: Optional[str] = None,
|
reasoning: str | None = None,
|
||||||
model: Optional[str] = None,
|
model: str | None = None,
|
||||||
tokens: Optional[int] = None,
|
tokens: int | None = None,
|
||||||
ts: Optional[int] = None,
|
ts: int | None = None,
|
||||||
) -> Frame:
|
) -> Frame:
|
||||||
"""Finalize a streaming bubble."""
|
"""Finalize a streaming bubble."""
|
||||||
payload: Dict[str, Any] = {
|
payload: dict[str, Any] = {
|
||||||
"message_id": message_id,
|
"message_id": message_id,
|
||||||
"final_text": final_text,
|
"final_text": final_text,
|
||||||
}
|
}
|
||||||
@@ -339,16 +367,17 @@ def message_stop(
|
|||||||
# Tool activity frames (M2)
|
# Tool activity frames (M2)
|
||||||
# ---------------------------------------------------------------------------
|
# ---------------------------------------------------------------------------
|
||||||
|
|
||||||
|
|
||||||
def tool_start(
|
def tool_start(
|
||||||
chat_id: str,
|
chat_id: str,
|
||||||
index: int,
|
index: int,
|
||||||
name: str,
|
name: str,
|
||||||
*,
|
*,
|
||||||
thread_id: Optional[str] = None,
|
thread_id: str | None = None,
|
||||||
preview: Optional[str] = None,
|
preview: str | None = None,
|
||||||
args: Optional[Dict[str, Any]] = None,
|
args: dict[str, Any] | None = None,
|
||||||
) -> Frame:
|
) -> Frame:
|
||||||
payload: Dict[str, Any] = {"index": index, "name": name}
|
payload: dict[str, Any] = {"index": index, "name": name}
|
||||||
if preview:
|
if preview:
|
||||||
payload["preview"] = preview
|
payload["preview"] = preview
|
||||||
if args:
|
if args:
|
||||||
@@ -366,10 +395,10 @@ def tool_progress(
|
|||||||
index: int,
|
index: int,
|
||||||
name: str,
|
name: str,
|
||||||
*,
|
*,
|
||||||
thread_id: Optional[str] = None,
|
thread_id: str | None = None,
|
||||||
note: Optional[str] = None,
|
note: str | None = None,
|
||||||
) -> Frame:
|
) -> Frame:
|
||||||
payload: Dict[str, Any] = {"index": index, "name": name}
|
payload: dict[str, Any] = {"index": index, "name": name}
|
||||||
if note:
|
if note:
|
||||||
payload["note"] = note
|
payload["note"] = note
|
||||||
return Frame(
|
return Frame(
|
||||||
@@ -385,12 +414,12 @@ def tool_end(
|
|||||||
index: int,
|
index: int,
|
||||||
name: str,
|
name: str,
|
||||||
*,
|
*,
|
||||||
thread_id: Optional[str] = None,
|
thread_id: str | None = None,
|
||||||
ok: bool = True,
|
ok: bool = True,
|
||||||
duration: Optional[float] = None,
|
duration: float | None = None,
|
||||||
output_preview: Optional[str] = None,
|
output_preview: str | None = None,
|
||||||
) -> Frame:
|
) -> Frame:
|
||||||
payload: Dict[str, Any] = {"index": index, "name": name, "ok": ok}
|
payload: dict[str, Any] = {"index": index, "name": name, "ok": ok}
|
||||||
if duration is not None:
|
if duration is not None:
|
||||||
payload["duration"] = duration
|
payload["duration"] = duration
|
||||||
if output_preview:
|
if output_preview:
|
||||||
@@ -407,12 +436,13 @@ def tool_end(
|
|||||||
# Commentary frame (M2)
|
# Commentary frame (M2)
|
||||||
# ---------------------------------------------------------------------------
|
# ---------------------------------------------------------------------------
|
||||||
|
|
||||||
|
|
||||||
def commentary(
|
def commentary(
|
||||||
chat_id: str,
|
chat_id: str,
|
||||||
message_id: str,
|
message_id: str,
|
||||||
text: str,
|
text: str,
|
||||||
*,
|
*,
|
||||||
thread_id: Optional[str] = None,
|
thread_id: str | None = None,
|
||||||
) -> Frame:
|
) -> Frame:
|
||||||
"""An intermediate assistant beat (between tool iterations)."""
|
"""An intermediate assistant beat (between tool iterations)."""
|
||||||
return Frame(
|
return Frame(
|
||||||
@@ -427,9 +457,10 @@ def commentary(
|
|||||||
# Channel directory frames (M3)
|
# Channel directory frames (M3)
|
||||||
# ---------------------------------------------------------------------------
|
# ---------------------------------------------------------------------------
|
||||||
|
|
||||||
def _channel_payload(entry: Dict[str, Any]) -> Dict[str, Any]:
|
|
||||||
|
def _channel_payload(entry: dict[str, Any]) -> dict[str, Any]:
|
||||||
"""Project a directory entry onto the wire shape."""
|
"""Project a directory entry onto the wire shape."""
|
||||||
payload: Dict[str, Any] = {
|
payload: dict[str, Any] = {
|
||||||
"chat_id": entry.get("chat_id"),
|
"chat_id": entry.get("chat_id"),
|
||||||
"name": entry.get("name"),
|
"name": entry.get("name"),
|
||||||
"kind": entry.get("kind", "channel"),
|
"kind": entry.get("kind", "channel"),
|
||||||
@@ -451,7 +482,7 @@ def _channel_payload(entry: Dict[str, Any]) -> Dict[str, Any]:
|
|||||||
return payload
|
return payload
|
||||||
|
|
||||||
|
|
||||||
def channel_created(entry: Dict[str, Any], auto: bool = False) -> Frame:
|
def channel_created(entry: dict[str, Any], auto: bool = False) -> Frame:
|
||||||
"""Broadcast: a channel/thread was created.
|
"""Broadcast: a channel/thread was created.
|
||||||
|
|
||||||
``auto=True`` marks a thread the gateway minted itself for an incoming
|
``auto=True`` marks a thread the gateway minted itself for an incoming
|
||||||
@@ -465,7 +496,7 @@ def channel_created(entry: Dict[str, Any], auto: bool = False) -> Frame:
|
|||||||
return Frame(type=TYPE_CHANNEL_CREATED, payload=payload)
|
return Frame(type=TYPE_CHANNEL_CREATED, payload=payload)
|
||||||
|
|
||||||
|
|
||||||
def channel_renamed(entry: Dict[str, Any]) -> Frame:
|
def channel_renamed(entry: dict[str, Any]) -> Frame:
|
||||||
"""Broadcast: a channel/thread was renamed."""
|
"""Broadcast: a channel/thread was renamed."""
|
||||||
return Frame(type=TYPE_CHANNEL_RENAMED, payload=_channel_payload(entry))
|
return Frame(type=TYPE_CHANNEL_RENAMED, payload=_channel_payload(entry))
|
||||||
|
|
||||||
@@ -475,7 +506,7 @@ def channel_deleted(chat_id: str) -> Frame:
|
|||||||
return Frame(type=TYPE_CHANNEL_DELETED, payload={"chat_id": chat_id})
|
return Frame(type=TYPE_CHANNEL_DELETED, payload={"chat_id": chat_id})
|
||||||
|
|
||||||
|
|
||||||
def channel_list(channels: List[Dict[str, Any]]) -> Frame:
|
def channel_list(channels: list[dict[str, Any]]) -> Frame:
|
||||||
"""Full directory (response to a ``channel.list`` request)."""
|
"""Full directory (response to a ``channel.list`` request)."""
|
||||||
return Frame(
|
return Frame(
|
||||||
type=TYPE_CHANNEL_LIST,
|
type=TYPE_CHANNEL_LIST,
|
||||||
@@ -487,12 +518,13 @@ def channel_list(channels: List[Dict[str, Any]]) -> Frame:
|
|||||||
# Search frames (M3)
|
# Search frames (M3)
|
||||||
# ---------------------------------------------------------------------------
|
# ---------------------------------------------------------------------------
|
||||||
|
|
||||||
|
|
||||||
def search_results(
|
def search_results(
|
||||||
query: str,
|
query: str,
|
||||||
scope: str,
|
scope: str,
|
||||||
hits: List[Dict[str, Any]],
|
hits: list[dict[str, Any]],
|
||||||
*,
|
*,
|
||||||
id: Optional[int] = None,
|
id: int | None = None,
|
||||||
) -> Frame:
|
) -> Frame:
|
||||||
"""Response to a ``search`` request.
|
"""Response to a ``search`` request.
|
||||||
|
|
||||||
@@ -509,7 +541,8 @@ def search_results(
|
|||||||
# Slash-command catalog frames
|
# Slash-command catalog frames
|
||||||
# ---------------------------------------------------------------------------
|
# ---------------------------------------------------------------------------
|
||||||
|
|
||||||
def commands_catalog(commands: List[Dict[str, Any]], *, id: Optional[int] = None) -> Frame:
|
|
||||||
|
def commands_catalog(commands: list[dict[str, Any]], *, id: int | None = None) -> Frame:
|
||||||
"""Response to a ``commands.catalog`` request: the gateway's slash-command
|
"""Response to a ``commands.catalog`` request: the gateway's slash-command
|
||||||
catalog for the app's ``/`` drawer.
|
catalog for the app's ``/`` drawer.
|
||||||
|
|
||||||
@@ -527,7 +560,8 @@ def commands_catalog(commands: List[Dict[str, Any]], *, id: Optional[int] = None
|
|||||||
# Sync frames (M3 outbox)
|
# Sync frames (M3 outbox)
|
||||||
# ---------------------------------------------------------------------------
|
# ---------------------------------------------------------------------------
|
||||||
|
|
||||||
def sync_done(cursor: int, *, id: Optional[int] = None) -> Frame:
|
|
||||||
|
def sync_done(cursor: int, *, id: int | None = None) -> Frame:
|
||||||
"""Terminal frame of a ``sync`` replay: the new cursor to persist."""
|
"""Terminal frame of a ``sync`` replay: the new cursor to persist."""
|
||||||
return Frame(type=TYPE_SYNC_DONE, id=id, payload={"cursor": cursor})
|
return Frame(type=TYPE_SYNC_DONE, id=id, payload={"cursor": cursor})
|
||||||
|
|
||||||
@@ -536,20 +570,21 @@ def sync_done(cursor: int, *, id: Optional[int] = None) -> Frame:
|
|||||||
# History frame (full message history for a chat/thread)
|
# History frame (full message history for a chat/thread)
|
||||||
# ---------------------------------------------------------------------------
|
# ---------------------------------------------------------------------------
|
||||||
|
|
||||||
|
|
||||||
def history(
|
def history(
|
||||||
chat_id: str,
|
chat_id: str,
|
||||||
messages: List[Dict[str, Any]],
|
messages: list[dict[str, Any]],
|
||||||
has_more: bool,
|
has_more: bool,
|
||||||
*,
|
*,
|
||||||
thread_id: Optional[str] = None,
|
thread_id: str | None = None,
|
||||||
oldest_message_id: Optional[str] = None,
|
oldest_message_id: str | None = None,
|
||||||
id: Optional[int] = None,
|
id: int | None = None,
|
||||||
) -> Frame:
|
) -> Frame:
|
||||||
"""Response to a ``history`` request: a page of final messages for a
|
"""Response to a ``history`` request: a page of final messages for a
|
||||||
chat/thread, ordered oldest → newest. ``has_more`` signals older pages
|
chat/thread, ordered oldest → newest. ``has_more`` signals older pages
|
||||||
exist; ``oldest_message_id`` is the ``before_message_id`` for the next
|
exist; ``oldest_message_id`` is the ``before_message_id`` for the next
|
||||||
(older) page."""
|
(older) page."""
|
||||||
payload: Dict[str, Any] = {
|
payload: dict[str, Any] = {
|
||||||
"messages": messages,
|
"messages": messages,
|
||||||
"has_more": has_more,
|
"has_more": has_more,
|
||||||
}
|
}
|
||||||
@@ -568,12 +603,13 @@ def history(
|
|||||||
# Message deletion frames
|
# Message deletion frames
|
||||||
# ---------------------------------------------------------------------------
|
# ---------------------------------------------------------------------------
|
||||||
|
|
||||||
|
|
||||||
def message_deleted(
|
def message_deleted(
|
||||||
chat_id: str,
|
chat_id: str,
|
||||||
message_ids: List[str],
|
message_ids: list[str],
|
||||||
*,
|
*,
|
||||||
thread_id: Optional[str] = None,
|
thread_id: str | None = None,
|
||||||
id: Optional[int] = None,
|
id: int | None = None,
|
||||||
) -> Frame:
|
) -> Frame:
|
||||||
"""Broadcast: the given message(s) were deleted from a chat/thread.
|
"""Broadcast: the given message(s) were deleted from a chat/thread.
|
||||||
|
|
||||||
@@ -595,18 +631,19 @@ def message_deleted(
|
|||||||
# Push / notification frames (M5)
|
# Push / notification frames (M5)
|
||||||
# ---------------------------------------------------------------------------
|
# ---------------------------------------------------------------------------
|
||||||
|
|
||||||
|
|
||||||
def notification(
|
def notification(
|
||||||
chat_id: str,
|
chat_id: str,
|
||||||
kind: str,
|
kind: str,
|
||||||
title: str,
|
title: str,
|
||||||
body: str,
|
body: str,
|
||||||
*,
|
*,
|
||||||
thread_id: Optional[str] = None,
|
thread_id: str | None = None,
|
||||||
ts: Optional[int] = None,
|
ts: int | None = None,
|
||||||
) -> Frame:
|
) -> Frame:
|
||||||
"""Event: a transient in-app banner (and a push mirror when the device is
|
"""Event: a transient in-app banner (and a push mirror when the device is
|
||||||
offline). ``kind`` is one of the ``NOTIF_*`` constants."""
|
offline). ``kind`` is one of the ``NOTIF_*`` constants."""
|
||||||
payload: Dict[str, Any] = {"kind": kind, "title": title, "body": body}
|
payload: dict[str, Any] = {"kind": kind, "title": title, "body": body}
|
||||||
if ts is not None:
|
if ts is not None:
|
||||||
payload["ts"] = ts
|
payload["ts"] = ts
|
||||||
return Frame(
|
return Frame(
|
||||||
@@ -614,11 +651,9 @@ def notification(
|
|||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
def fcm_register(
|
def fcm_register(fcm_token: str | None = None, ntfy_topic: str | None = None) -> Frame:
|
||||||
fcm_token: Optional[str] = None, ntfy_topic: Optional[str] = None
|
|
||||||
) -> Frame:
|
|
||||||
"""Request: update the device's push tokens (FCM rotation / ntfy topic)."""
|
"""Request: update the device's push tokens (FCM rotation / ntfy topic)."""
|
||||||
payload: Dict[str, Any] = {}
|
payload: dict[str, Any] = {}
|
||||||
if fcm_token:
|
if fcm_token:
|
||||||
payload["fcm_token"] = fcm_token
|
payload["fcm_token"] = fcm_token
|
||||||
if ntfy_topic:
|
if ntfy_topic:
|
||||||
@@ -640,6 +675,7 @@ def read_receipt(chat_id: str, message_id: str) -> Frame:
|
|||||||
# Gateway health frame (M5)
|
# Gateway health frame (M5)
|
||||||
# ---------------------------------------------------------------------------
|
# ---------------------------------------------------------------------------
|
||||||
|
|
||||||
|
|
||||||
def status(state: str) -> Frame:
|
def status(state: str) -> Frame:
|
||||||
"""Gateway health state (``state`` is one of the ``STATUS_*`` constants)."""
|
"""Gateway health state (``state`` is one of the ``STATUS_*`` constants)."""
|
||||||
return Frame(type=TYPE_STATUS, payload={"state": state})
|
return Frame(type=TYPE_STATUS, payload={"state": state})
|
||||||
@@ -649,6 +685,7 @@ def status(state: str) -> Frame:
|
|||||||
# Media frames (M4)
|
# Media frames (M4)
|
||||||
# ---------------------------------------------------------------------------
|
# ---------------------------------------------------------------------------
|
||||||
|
|
||||||
|
|
||||||
def media_offer(
|
def media_offer(
|
||||||
media_id: str,
|
media_id: str,
|
||||||
kind: str,
|
kind: str,
|
||||||
@@ -656,16 +693,16 @@ def media_offer(
|
|||||||
size: int,
|
size: int,
|
||||||
filename: str,
|
filename: str,
|
||||||
*,
|
*,
|
||||||
chat_id: Optional[str] = None,
|
chat_id: str | None = None,
|
||||||
thread_id: Optional[str] = None,
|
thread_id: str | None = None,
|
||||||
message_id: Optional[str] = None,
|
message_id: str | None = None,
|
||||||
) -> Frame:
|
) -> Frame:
|
||||||
"""Event: the agent produced media the app can fetch via ``media.pull``.
|
"""Event: the agent produced media the app can fetch via ``media.pull``.
|
||||||
|
|
||||||
``message_id`` (optional) associates the offer with the assistant message
|
``message_id`` (optional) associates the offer with the assistant message
|
||||||
it belongs to (the app falls back to the lane's last assistant message).
|
it belongs to (the app falls back to the lane's last assistant message).
|
||||||
"""
|
"""
|
||||||
payload: Dict[str, Any] = {
|
payload: dict[str, Any] = {
|
||||||
"media_id": media_id,
|
"media_id": media_id,
|
||||||
"kind": kind,
|
"kind": kind,
|
||||||
"mime": mime,
|
"mime": mime,
|
||||||
@@ -674,15 +711,17 @@ def media_offer(
|
|||||||
}
|
}
|
||||||
if message_id:
|
if message_id:
|
||||||
payload["message_id"] = message_id
|
payload["message_id"] = message_id
|
||||||
return Frame(type=TYPE_MEDIA_OFFER, chat_id=chat_id, thread_id=thread_id, payload=payload)
|
return Frame(
|
||||||
|
type=TYPE_MEDIA_OFFER, chat_id=chat_id, thread_id=thread_id, payload=payload
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
def media_pull_end(ok: bool, *, id: Optional[int] = None) -> Frame:
|
def media_pull_end(ok: bool, *, id: int | None = None) -> Frame:
|
||||||
"""Terminal frame of a ``media.pull`` binary stream."""
|
"""Terminal frame of a ``media.pull`` binary stream."""
|
||||||
return Frame(type=TYPE_MEDIA_PULL_END, id=id, payload={"ok": ok})
|
return Frame(type=TYPE_MEDIA_PULL_END, id=id, payload={"ok": ok})
|
||||||
|
|
||||||
|
|
||||||
def media_upload_ack(ok: bool, media_ref: str, *, id: Optional[int] = None) -> Frame:
|
def media_upload_ack(ok: bool, media_ref: str, *, id: int | None = None) -> Frame:
|
||||||
"""Response to ``media.upload.end``: the ref is cached and may be used in
|
"""Response to ``media.upload.end``: the ref is cached and may be used in
|
||||||
a ``message.send`` ``media_refs``. Failures use ``error`` frames instead."""
|
a ``message.send`` ``media_refs``. Failures use ``error`` frames instead."""
|
||||||
return Frame(
|
return Frame(
|
||||||
@@ -692,12 +731,12 @@ def media_upload_ack(ok: bool, media_ref: str, *, id: Optional[int] = None) -> F
|
|||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
def error(code: str, message: str, *, id: Optional[int] = None) -> Frame:
|
def error(code: str, message: str, *, id: int | None = None) -> Frame:
|
||||||
return Frame(type=TYPE_ERROR, id=id, payload={"code": code, "message": message})
|
return Frame(type=TYPE_ERROR, id=id, payload={"code": code, "message": message})
|
||||||
|
|
||||||
|
|
||||||
def pong(ts: Optional[int] = None) -> Frame:
|
def pong(ts: int | None = None) -> Frame:
|
||||||
payload: Dict[str, Any] = {}
|
payload: dict[str, Any] = {}
|
||||||
if ts is not None:
|
if ts is not None:
|
||||||
payload["ts"] = ts
|
payload["ts"] = ts
|
||||||
return Frame(type=TYPE_PONG, payload=payload)
|
return Frame(type=TYPE_PONG, payload=payload)
|
||||||
@@ -290,6 +290,9 @@ class WsServer:
|
|||||||
server_caps=self._adapter.server_caps(),
|
server_caps=self._adapter.server_caps(),
|
||||||
sync_cursor=self._adapter._outbox.latest_cursor(),
|
sync_cursor=self._adapter._outbox.latest_cursor(),
|
||||||
channels=self._adapter.channel_list(),
|
channels=self._adapter.channel_list(),
|
||||||
|
# M5: lets the app dedupe sync-replayed notifications that
|
||||||
|
# already woke this device via push (docs/08 §8.7).
|
||||||
|
last_pushed_cursor=self._adapter._devices.last_pushed_cursor(device_id),
|
||||||
)
|
)
|
||||||
try:
|
try:
|
||||||
await ws.send(ack.to_json())
|
await ws.send(ack.to_json())
|
||||||
|
|||||||
Reference in new issue
Block a user