diff --git a/app/shared/src/commonMain/kotlin/iris/data/ChatStore.kt b/app/shared/src/commonMain/kotlin/iris/data/ChatStore.kt index eb0c051..a8a19c7 100644 --- a/app/shared/src/commonMain/kotlin/iris/data/ChatStore.kt +++ b/app/shared/src/commonMain/kotlin/iris/data/ChatStore.kt @@ -444,6 +444,26 @@ class ChatStore { if (changed) _lanes.value = map } + /** + * Remove messages by id from every lane (a `message.deleted` frame). The + * server is authoritative: the frame carries no lane, and a message id is + * unique, so scanning all lanes is both safe and idempotent (a no-op when the + * id is absent, e.g. a local-only id another device never had). + */ + fun removeMessages(messageIds: Set) { + if (messageIds.isEmpty()) return + val map = _lanes.value.toMutableMap() + var changed = false + for ((lane, list) in map) { + val updated = list.filterNot { it.id in messageIds } + if (updated != list) { + map[lane] = updated + changed = true + } + } + if (changed) _lanes.value = map + } + /** M7: mark all pending user messages as failed (gateway error frame). */ fun failPending() { val map = _lanes.value.toMutableMap() diff --git a/app/shared/src/commonMain/kotlin/iris/protocol/Protocol.kt b/app/shared/src/commonMain/kotlin/iris/protocol/Protocol.kt index 05f6f0e..c7920ed 100644 --- a/app/shared/src/commonMain/kotlin/iris/protocol/Protocol.kt +++ b/app/shared/src/commonMain/kotlin/iris/protocol/Protocol.kt @@ -3,6 +3,7 @@ package iris.protocol import kotlinx.serialization.SerialName import kotlinx.serialization.Serializable import kotlinx.serialization.json.Json +import kotlinx.serialization.json.JsonArray import kotlinx.serialization.json.JsonElement import kotlinx.serialization.json.JsonObject import kotlinx.serialization.json.JsonPrimitive @@ -41,6 +42,10 @@ const val TYPE_TYPING = "typing" const val TYPE_MESSAGE_START = "message.start" const val TYPE_MESSAGE_UPDATE = "message.update" const val TYPE_MESSAGE_STOP = "message.stop" + +// Message deletion (app requests; broadcast to all devices) +const val TYPE_MESSAGE_DELETE = "message.delete" +const val TYPE_MESSAGE_DELETED = "message.deleted" const val TYPE_TOOL_START = "tool.start" const val TYPE_TOOL_PROGRESS = "tool.progress" const val TYPE_TOOL_END = "tool.end" @@ -385,6 +390,14 @@ data class HistoryPayload( @SerialName("oldest_message_id") val oldestMessageId: String? = null, ) +// ── message deletion ────────────────────────────────────────────────────── + +/** The message(s) removed from a chat/thread (response + broadcast). */ +@Serializable +data class MessageDeletedPayload( + @SerialName("message_ids") val messageIds: List = emptyList(), +) + // ── M5: push / notifications ──────────────────────────────────────────── /** Notification kinds (mirror of protocol.NOTIF_*). */ @@ -550,6 +563,24 @@ fun historyFrame( }, ) +/** Request deletion of the given message(s) in a chat/thread. The server + * removes them from the outbox and broadcasts `message.deleted`. */ +fun messageDeleteFrame( + id: Int, + chatId: String, + messageIds: List, + threadId: String? = null, +): Frame = + Frame( + id = id, + type = TYPE_MESSAGE_DELETE, + chatId = chatId, + threadId = threadId, + payload = buildJsonObject { + put("message_ids", JsonArray(messageIds.map { JsonPrimitive(it) })) + }, + ) + // ── M4 frame builders ──────────────────────────────────────────────────── fun mediaUploadStartFrame( diff --git a/app/shared/src/commonMain/kotlin/iris/state/IrisController.kt b/app/shared/src/commonMain/kotlin/iris/state/IrisController.kt index 32a0610..463ab5a 100644 --- a/app/shared/src/commonMain/kotlin/iris/state/IrisController.kt +++ b/app/shared/src/commonMain/kotlin/iris/state/IrisController.kt @@ -19,6 +19,7 @@ import iris.protocol.HIGH_PRIORITY_NOTIF_KINDS import iris.protocol.HistoryMessage import iris.protocol.HistoryPayload import iris.protocol.MediaOfferPayload +import iris.protocol.MessageDeletedPayload import iris.protocol.MessagePayload import iris.protocol.MessageStopPayload import iris.protocol.NotificationPayload @@ -38,6 +39,8 @@ import iris.protocol.TYPE_CHANNEL_RENAMED import iris.protocol.TYPE_COMMENTARY import iris.protocol.TYPE_MEDIA_OFFER import iris.protocol.TYPE_MESSAGE +import iris.protocol.TYPE_MESSAGE_DELETE +import iris.protocol.TYPE_MESSAGE_DELETED import iris.protocol.TYPE_MESSAGE_START import iris.protocol.TYPE_MESSAGE_STOP import iris.protocol.TYPE_MESSAGE_UPDATE @@ -56,6 +59,7 @@ import iris.protocol.channelListFrame import iris.protocol.channelRenameFrame import iris.protocol.channelSetDefaultFrame import iris.protocol.historyFrame +import iris.protocol.messageDeleteFrame import iris.protocol.searchFrame import iris.protocol.syncFrame import kotlinx.coroutines.CoroutineScope @@ -367,6 +371,13 @@ class IrisController( TYPE_READ_RECEIPT -> { frame.payloadAs()?.let { chat.markRead(it.messageId) } } + TYPE_MESSAGE_DELETED -> { + // A message was deleted (by this or another device). + // Idempotent: dropping an unknown id is a no-op. + frame.payloadAs()?.let { + chat.removeMessages(it.messageIds.toSet()) + } + } TYPE_STATUS -> { frame.payloadAs()?.let { _gatewayStatus.value = it.state } } @@ -539,6 +550,18 @@ class IrisController( ) } + /** Delete the given message(s) from the current lane (long-press select → + * delete). The server removes them from the outbox and broadcasts + * `message.deleted`; the local cache drops them on that frame (or + * immediately, below, for snappy UX). */ + fun deleteMessages(messageIds: List) { + if (messageIds.isEmpty()) return + val lane = chat.currentLane.value + val (chatId, threadId) = chat.parseLane(lane) + chat.removeMessages(messageIds.toSet()) + client.sendFrame(messageDeleteFrame(0, chatId, messageIds, threadId)) + } + // ── M4: attachments (pick -> upload -> send) ────────────────────────── /** Stage a picked file: upload it, then keep it as a pending attachment. */ 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 543731d..9bb0f5e 100644 --- a/app/shared/src/commonMain/kotlin/iris/ui/screens/ChatScreen.kt +++ b/app/shared/src/commonMain/kotlin/iris/ui/screens/ChatScreen.kt @@ -2,6 +2,7 @@ package iris.ui.screens import androidx.compose.foundation.ExperimentalFoundationApi import androidx.compose.foundation.background +import androidx.compose.foundation.border import androidx.compose.foundation.clickable import androidx.compose.foundation.combinedClickable import androidx.compose.foundation.focusable @@ -221,6 +222,34 @@ fun ChatScreen(controller: IrisController) { // Topic context menu (long-press / right-click): rename / delete a thread. var renameThread by remember { mutableStateOf(null) } var deleteThread by remember { mutableStateOf(null) } + // Message selection (long-press / right-click a bubble → multi-select → + // delete). The composer is replaced by a selection toolbar while active. + var selectionMode by remember { mutableStateOf(false) } + var selectedIds by remember { mutableStateOf>(emptySet()) } + var showDeleteConfirm by remember { mutableStateOf(false) } + + fun enterSelection(id: String) { + selectionMode = true + selectedIds = setOf(id) + } + + fun toggleSelect(id: String) { + selectedIds = if (id in selectedIds) selectedIds - id else selectedIds + id + if (selectedIds.isEmpty()) selectionMode = false + } + + fun exitSelection() { + selectionMode = false + selectedIds = emptySet() + } + + // Selection is scoped to the lane the messages were viewed in; leaving + // that lane (channel / thread switch) clears any in-progress selection + // so a later delete can't target the wrong lane. + LaunchedEffect(currentLane) { + if (selectionMode) exitSelection() + } + Column(modifier = Modifier.fillMaxSize()) { // Header (M7: avatar + two-line title pill, reference look). Wide // panes keep everything on one row; narrow phones move the utility @@ -355,12 +384,27 @@ fun ChatScreen(controller: IrisController) { when (row) { is ChatRow.Day -> DaySeparator(row.label) is ChatRow.Item -> when (val item = row.item) { - is MessageItem -> MessageBubble( - msg = item, - maxWidth = bubbleMaxWidth, - reasoningAutoCollapse = reasoningAutoCollapse, - onRetry = { controller.retrySend(item.id) }, - ) + is MessageItem -> { + // Only finalized messages are selectable: a + // streaming bubble has no final id yet and a + // pending echo isn't on the server to delete. + val selectable = !item.streaming && !item.pending + MessageBubble( + msg = item, + maxWidth = bubbleMaxWidth, + reasoningAutoCollapse = reasoningAutoCollapse, + onRetry = { controller.retrySend(item.id) }, + selectionMode = selectionMode && selectable, + selected = selectable && item.id in selectedIds, + onToggleSelect = { if (selectable) toggleSelect(item.id) }, + onLongPress = { + if (selectable) { + if (selectionMode) toggleSelect(item.id) + else enterSelection(item.id) + } + }, + ) + } is ToolItem -> if (toolDetail != ToolDetail.NOTHING) ToolCard(item, toolDetail) } } @@ -444,6 +488,15 @@ fun ChatScreen(controller: IrisController) { } } + // Selection toolbar (replaces the composer while messages are selected). + if (selectionMode) { + SelectionToolbar( + count = selectedIds.size, + onCancel = { exitSelection() }, + onDelete = { showDeleteConfirm = true }, + ) + } + // Composer (M7: rounded pill + accent circular send button). Sending is gated // on a live socket: sendMessage is a no-op while disconnected, so an // ungated send would show a pending bubble that never resolves. @@ -452,7 +505,7 @@ fun ChatScreen(controller: IrisController) { isConnected && (input.isNotBlank() || attachments.any { it.mediaRef != null && it.error == null }) val layoutDensity = LocalDensity.current.density var textHeightPx by remember { mutableFloatStateOf(0f) } - Row( + if (!selectionMode) Row( modifier = Modifier .fillMaxWidth() .padding(horizontal = 12.dp, vertical = 8.dp) @@ -575,6 +628,30 @@ fun ChatScreen(controller: IrisController) { }, ) } + // Message selection: confirm deleting the selected message(s). + if (showDeleteConfirm) { + val n = selectedIds.size + AlertDialog( + onDismissRequest = { showDeleteConfirm = false }, + title = { Text(if (n == 1) "Delete message?" else "Delete $n messages?") }, + text = { + Text( + if (n == 1) "This message will be deleted for all devices." + else "These messages will be deleted for all devices.", + ) + }, + confirmButton = { + TextButton(onClick = { + showDeleteConfirm = false + controller.deleteMessages(selectedIds.toList()) + exitSelection() + }) { Text("Delete") } + }, + dismissButton = { + TextButton(onClick = { showDeleteConfirm = false }) { Text("Cancel") } + }, + ) + } } } @@ -1418,8 +1495,18 @@ private fun NotificationBanner( } } +@OptIn(ExperimentalFoundationApi::class) @Composable -private fun MessageBubble(msg: MessageItem, maxWidth: Dp, reasoningAutoCollapse: Boolean, onRetry: () -> Unit) { +private fun MessageBubble( + msg: MessageItem, + maxWidth: Dp, + reasoningAutoCollapse: Boolean, + onRetry: () -> Unit, + selectionMode: Boolean = false, + selected: Boolean = false, + onToggleSelect: () -> Unit = {}, + onLongPress: () -> Unit = {}, +) { val isUser = msg.role == ROLE_USER val isCommentary = msg.isCommentary val bubbleColor = when { @@ -1437,15 +1524,31 @@ private fun MessageBubble(msg: MessageItem, maxWidth: Dp, reasoningAutoCollapse: modifier = Modifier.fillMaxWidth(), horizontalArrangement = if (isUser) Arrangement.End else Arrangement.Start, ) { + // Selection check sits on the side away from the bubble (user bubbles + // are right-aligned → check on the left; assistant on the left → right). + if (isUser && selectionMode) { + SelectionCheck(selected) + Spacer(modifier = Modifier.width(6.dp)) + } Column( modifier = Modifier .widthIn(max = maxWidth) .clip(RoundedCornerShape(14.dp)) - .background(bubbleColor) + .background(if (selected) IrisColors.chipSelected else bubbleColor) .then( - if (isUser && msg.status == MsgStatus.Failed) { - Modifier.clickable { onRetry() } - } else Modifier + // Long-press / right-click enters (or toggles within) message + // selection; a plain tap toggles while selecting, or retries a + // failed user send otherwise. Attached always so the long-press + // affordance exists even outside selection mode. + Modifier + .combinedClickable( + onClick = { + if (selectionMode) onToggleSelect() + else if (isUser && msg.status == MsgStatus.Failed) onRetry() + }, + onLongClick = { onLongPress() }, + ) + .rightClick { onLongPress() } ) .padding(horizontal = 12.dp, vertical = 8.dp), ) { @@ -1532,6 +1635,63 @@ private fun MessageBubble(msg: MessageItem, maxWidth: Dp, reasoningAutoCollapse: } } } + if (!isUser && selectionMode) { + Spacer(modifier = Modifier.width(6.dp)) + SelectionCheck(selected) + } + } +} + +/** Circular selection indicator shown beside a bubble in selection mode. */ +@Composable +private fun SelectionCheck(selected: Boolean) { + Box( + modifier = Modifier + .size(22.dp) + .clip(CircleShape) + .background(if (selected) IrisColors.primary else Color.Transparent) + .border( + width = if (selected) 0.dp else 1.5.dp, + color = IrisColors.textDim, + shape = CircleShape, + ), + contentAlignment = Alignment.Center, + ) { + if (selected) Text("✓", color = Color.White, fontSize = 12.sp) + } +} + +/** Toolbar shown in place of the composer while messages are selected. */ +@Composable +private fun SelectionToolbar(count: Int, onCancel: () -> Unit, onDelete: () -> Unit) { + Row( + modifier = Modifier + .fillMaxWidth() + .padding(horizontal = 12.dp, vertical = 8.dp) + .clip(RoundedCornerShape(28.dp)) + .background(IrisColors.surface) + .padding(horizontal = 8.dp, vertical = 4.dp), + verticalAlignment = Alignment.CenterVertically, + ) { + HeaderIconButton(onClick = onCancel) { + Text("✕", fontSize = 16.sp) + } + Spacer(modifier = Modifier.width(8.dp)) + Text( + if (count == 1) "1 selected" else "$count selected", + style = MaterialTheme.typography.bodyLarge, + modifier = Modifier.weight(1f), + ) + Box( + modifier = Modifier + .clip(RoundedCornerShape(20.dp)) + .background(if (count > 0) IrisColors.primary else IrisColors.primary.copy(alpha = 0.35f)) + .clickable(enabled = count > 0) { onDelete() } + .padding(horizontal = 16.dp, vertical = 8.dp), + contentAlignment = Alignment.Center, + ) { + Text("🗑 Delete", color = Color.White, fontSize = 14.sp) + } } } diff --git a/docs/04-wire-protocol.md b/docs/04-wire-protocol.md index 2767aed..a09ec6a 100644 --- a/docs/04-wire-protocol.md +++ b/docs/04-wire-protocol.md @@ -63,7 +63,17 @@ Streaming a bubble. `update` carries the **full** current text (app replaces). {"type":"message.start","chat_id":"…","payload":{"message_id":"m_9002","role":"assistant"}} {"type":"message.update","chat_id":"…","payload":{"message_id":"m_9002","text":"partial…"}} {"type":"message.stop","chat_id":"…","payload":{"message_id":"m_9002","final_text":"full…", - "reasoning":"…","model":"…","tokens":11}} + "reasoning":"…","model":"…","tokens":11}} +``` + +### `message.deleted` +The given message(s) were deleted from a chat/thread. Response to a +`message.delete` request (`id` set) **and** broadcast to every device so all of +them drop the message(s) from their cache. Also outboxed, so a device that was +offline learns of the deletion on its next `sync`. +```json +{"type":"message.deleted","id":30,"chat_id":"android:default","thread_id":null, + "payload":{"message_ids":["m_9001","m_9002"]}} ``` ### `commentary` @@ -300,6 +310,16 @@ Load a page of messages for a chat/thread (initial open, scroll-up pagination). `before_message_id` — return messages older than this (omit for newest page). `limit` — max messages (default 50, max 200). +### `message.delete` +Delete the given message(s) from a chat/thread. The server removes them from the +outbox (so `history`/`sync` no longer return them) and broadcasts +`message.deleted` to every device. Idempotent: a message already gone (pruned by +retention) still yields a `message.deleted` broadcast so live caches drop it. +```json +{"type":"message.delete","id":30,"chat_id":"android:default","thread_id":null, + "payload":{"message_ids":["m_9001","m_9002"]}} +``` + ### `commands.catalog` Fetch the full slash-command catalog (for the `Menü` bottom sheet). ```json @@ -346,10 +366,11 @@ Keepalive. `{"type":"ping","payload":{"ts":1724000000000}}` → `pong`. - Frames are ordered per connection (TCP/WS). Streaming `message.update` frames for a `message_id` are monotonic; the app may coalesce to the latest. -- Terminal frames (`message`, `message.stop`, `tool.end`, `notification`, - `picker.*`, `channel.*`, `agent.busy`, `agent.idle`, `history`, - `commands.catalog`, `commands.complete`) are **never dropped** under - backpressure; only intermediate `message.update`/`tool.progress` are coalesced. +- Terminal frames (`message`, `message.stop`, `message.deleted`, `tool.end`, + `notification`, `picker.*`, `channel.*`, `agent.busy`, `agent.idle`, + `history`, `commands.catalog`, `commands.complete`) are **never dropped** + under backpressure; only intermediate `message.update`/`tool.progress` are + coalesced. - Anything not delivered live goes to the **outbox** and is replayed by `sync`. - **Broadcast:** channel directory events (`channel.*`) and read-receipts are pushed to **all** connected devices for that gateway (no subscribe step). diff --git a/docs/10-android-app.md b/docs/10-android-app.md index a3ba7d2..14d845f 100644 --- a/docs/10-android-app.md +++ b/docs/10-android-app.md @@ -126,6 +126,20 @@ app/shared/src/ ### Intermediate messages - `commentary` frames → dimmed/smaller bubble, distinct from final answers. +### Message selection + delete +- **Long-press** a message bubble (touch) or **right-click** it (desktop mouse) + enters selection mode: the tapped message is selected (a circular check appears + beside each bubble) and the composer is replaced by a selection toolbar + (count + **Delete** + cancel ✕). +- While selecting, **tap** a bubble to toggle it; the ✕ (or deselecting the last + message) exits selection mode. Only finalized messages are selectable — a + streaming bubble has no final id yet and a pending echo isn't on the server. +- **Delete** → confirm dialog → `message.delete {message_ids:[…]}` for the + current lane. The server removes the message(s) from the outbox (so + `history`/`sync` no longer return them) and broadcasts `message.deleted` to + every device; each device drops them from its cache (the requesting device + also drops them locally for snappy UX). Deleting is idempotent. + ### Agent busy / stop / steer - `agent.busy` → show "thinking…" indicator in chat header (animated dots). - `agent.idle` → clear indicator. diff --git a/docs/protocol/frames.schema.json b/docs/protocol/frames.schema.json index 5214d72..7162d4f 100644 --- a/docs/protocol/frames.schema.json +++ b/docs/protocol/frames.schema.json @@ -42,6 +42,7 @@ "message.start": { "payload": { "message_id": { "type": "string" }, "role": { "type": "string" } } }, "message.update": { "payload": { "message_id": { "type": "string" }, "text": { "type": "string", "description": "Full current text (app replaces)." } } }, "message.stop": { "payload": { "message_id": { "type": "string" }, "final_text": { "type": "string" }, "reasoning": { "type": "string" }, "model": { "type": "string" }, "tokens": { "type": "integer" }, "ts": { "type": "integer" } } }, + "message.deleted": { "description": "The given message(s) were deleted from a chat/thread. Response to a message.delete request (id set) and broadcast to every device so all drop them from their cache; also outboxed so an offline device learns of the deletion on its next sync.", "payload": { "message_ids": { "type": "array", "items": { "type": "string" } } } }, "commentary": { "description": "Intermediate assistant beat.", "payload": { "message_id": { "type": "string" }, "text": { "type": "string" } } }, "tool.start": { "payload": { "index": { "type": "integer" }, "name": { "type": "string" }, "preview": { "type": "string" }, "args": { "type": "object" } } }, "tool.progress": { "payload": { "index": { "type": "integer" }, "name": { "type": "string" }, "note": { "type": "string" } } }, @@ -77,6 +78,7 @@ "search": { "payload": { "query": { "type": "string" }, "scope": { "type": "string", "enum": ["all", "chat"] }, "chat_id": { "type": "string" }, "thread_id": { "type": "string" }, "limit": { "type": "integer", "description": "Optional; server default 20." } } }, "sync": { "description": "Reconnect catch-up; replays undelivered outbox frames only (not full history).", "payload": { "cursor": { "type": "integer" } } }, "history": { "description": "Load a page of full message history for a chat/thread (initial open, scroll-up pagination).", "payload": { "before_message_id": { "type": "string", "description": "Return messages older than this (omit for newest page)." }, "limit": { "type": "integer", "description": "Max messages (default 50, max 200)." } } }, + "message.delete": { "description": "Delete the given message(s) from a chat/thread. The server removes them from the outbox (so history/sync no longer return them) and broadcasts message.deleted to every device. Idempotent: a message already gone (pruned) still yields a message.deleted broadcast.", "payload": { "message_ids": { "type": "array", "items": { "type": "string" }, "description": "One or more message_id values to delete." } } }, "fcm.register": { "payload": { "fcm_token": { "type": "string" }, "ntfy_topic": { "type": "string" } } }, "ping": { "payload": { "ts": { "type": "integer" } } } } @@ -105,7 +107,7 @@ ], "reliability": { "ordering": "Per-connection (TCP/WS). message.update for a message_id is monotonic; app may coalesce to latest.", - "never_dropped": ["message", "message.stop", "tool.end", "notification", "channel.*", "search.results", "error"], + "never_dropped": ["message", "message.stop", "message.deleted", "tool.end", "notification", "channel.*", "search.results", "error"], "coalescable_under_backpressure": ["message.update", "tool.progress"], "offline": "Undelivered frames go to the outbox; replayed by sync. Terminal frames always outboxed." } diff --git a/gateway-plugin/adapter.py b/gateway-plugin/adapter.py index 7cc0232..4a77a0c 100644 --- a/gateway-plugin/adapter.py +++ b/gateway-plugin/adapter.py @@ -2216,6 +2216,52 @@ class AndroidAdapter(BasePlatformAdapter): ) await self._ws_server.send_to(device_id, resp) + # ── Message deletion (app -> agent) ─────────────────────────────────── + + async def on_message_delete(self, frame: protocol.Frame, device_id: str) -> None: + """Handle an inbound ``message.delete`` request. + + Removes the requested message(s) from the outbox (so ``history`` and + ``sync`` no longer return them) and broadcasts ``message.deleted`` to + every device (outboxed too, so an offline device learns of the + deletion on its next ``sync``). Deleting is idempotent: a message that + is already gone (pruned by retention) simply yields 0 removed rows, + and the ``message.deleted`` broadcast is still emitted so live caches + drop it. + """ + payload = frame.payload + chat_id = frame.chat_id or payload.get("chat_id") + if not isinstance(chat_id, str) or not chat_id.strip(): + await self._ws_server.send_to( + device_id, + protocol.error(protocol.ERR_NOT_FOUND, "message.delete requires chat_id", id=frame.id), + ) + return + chat_id = chat_id.strip() + thread_id = frame.thread_id or payload.get("thread_id") + if not isinstance(thread_id, str) or not thread_id.strip(): + thread_id = None + message_ids = payload.get("message_ids") + if not isinstance(message_ids, list): + message_ids = [payload.get("message_id")] if payload.get("message_id") else [] + message_ids = [m for m in message_ids if isinstance(m, str) and m.strip()] + if not message_ids: + await self._ws_server.send_to( + device_id, + protocol.error(protocol.ERR_UNSUPPORTED, "message.delete requires message_ids", id=frame.id), + ) + return + removed = 0 + for mid in message_ids: + removed += self._outbox.delete_message(chat_id, mid, thread_id=thread_id) + logger.info( + "android: message.delete from %s chat_id=%r thread_id=%r ids=%s removed=%s", + device_id, chat_id, thread_id, message_ids, removed, + ) + resp = protocol.message_deleted(chat_id, message_ids, thread_id=thread_id) + resp.id = frame.id + await self._broadcast_or_log(chat_id, resp) + # ── M5: push token registration ─────────────────────────────────────── async def on_fcm_register(self, frame: protocol.Frame, device_id: str) -> None: diff --git a/gateway-plugin/outbox.py b/gateway-plugin/outbox.py index 0770b29..e9a7eec 100644 --- a/gateway-plugin/outbox.py +++ b/gateway-plugin/outbox.py @@ -277,6 +277,55 @@ class Outbox: "oldest_message_id": oldest_message_id, } + # ── message deletion ────────────────────────────────────────────────── + + def delete_message( + self, + chat_id: str, + message_id: str, + thread_id: Optional[str] = None, + ) -> int: + """Remove every outbox frame belonging to *message_id* in *chat_id*. + + A message can span several frames (``message`` / ``message.start`` / + ``message.update`` / ``message.stop`` / ``media.offer`` / + ``commentary``); all of them are removed so neither ``history`` nor a + ``sync`` replay can resurrect the message. The delete is scoped to the + exact lane: a flat-lane delete (``thread_id=None``) matches only frames + with no ``thread_id``, and a thread delete matches only that thread's + frames (a ``message_id`` is unique to one lane, so this is a safety + net, not a filter that drops real frames). Returns the number of rows + removed (0 when the message is not in the outbox — e.g. already pruned + by retention). + """ + if not message_id: + return 0 + with self._lock: + rows = self._conn.execute( + "SELECT cursor, frame FROM outbox WHERE chat_id = ?", (chat_id,) + ).fetchall() + cursors: List[int] = [] + for r in rows: + try: + frame = json.loads(r["frame"]) + except (json.JSONDecodeError, TypeError): + continue + if not isinstance(frame, dict): + continue + if frame.get("thread_id") != thread_id: + continue + payload = frame.get("payload") + if isinstance(payload, dict) and payload.get("message_id") == message_id: + cursors.append(int(r["cursor"])) + if not cursors: + return 0 + placeholders = ",".join("?" * len(cursors)) + self._conn.execute( + f"DELETE FROM outbox WHERE cursor IN ({placeholders})", cursors + ) + self._conn.commit() + return len(cursors) + # ── retention ───────────────────────────────────────────────────────── def _maybe_prune(self) -> None: diff --git a/gateway-plugin/protocol.py b/gateway-plugin/protocol.py index 5093fac..2e79c9b 100644 --- a/gateway-plugin/protocol.py +++ b/gateway-plugin/protocol.py @@ -42,6 +42,10 @@ TYPE_MESSAGE_START = "message.start" TYPE_MESSAGE_UPDATE = "message.update" TYPE_MESSAGE_STOP = "message.stop" +# Message deletion (app requests; broadcast to all devices) +TYPE_MESSAGE_DELETE = "message.delete" +TYPE_MESSAGE_DELETED = "message.deleted" + # Tool activity (M2) TYPE_TOOL_START = "tool.start" TYPE_TOOL_PROGRESS = "tool.progress" @@ -528,6 +532,33 @@ def history( ) +# --------------------------------------------------------------------------- +# Message deletion frames +# --------------------------------------------------------------------------- + +def message_deleted( + chat_id: str, + message_ids: List[str], + *, + thread_id: Optional[str] = None, + id: Optional[int] = None, +) -> Frame: + """Broadcast: the given message(s) were deleted from a chat/thread. + + Carried by the response to a ``message.delete`` request (``id`` set) and + broadcast to every device so all of them drop the message(s) from their + cache. Also outboxed, so a device that was offline learns of the deletion + on its next ``sync``. + """ + return Frame( + type=TYPE_MESSAGE_DELETED, + id=id, + chat_id=chat_id, + thread_id=thread_id, + payload={"message_ids": list(message_ids)}, + ) + + # --------------------------------------------------------------------------- # Push / notification frames (M5) # --------------------------------------------------------------------------- diff --git a/gateway-plugin/ws_server.py b/gateway-plugin/ws_server.py index 88485a2..4e5347a 100644 --- a/gateway-plugin/ws_server.py +++ b/gateway-plugin/ws_server.py @@ -380,6 +380,8 @@ class WsServer: await self._adapter.on_sync(frame, device_id) elif frame.type == protocol.TYPE_HISTORY: await self._adapter.on_history(frame, device_id) + elif frame.type == protocol.TYPE_MESSAGE_DELETE: + await self._adapter.on_message_delete(frame, device_id) elif frame.type == protocol.TYPE_MEDIA_UPLOAD_START: await self._adapter.on_media_upload_start(frame, device_id) elif frame.type == protocol.TYPE_MEDIA_UPLOAD_END: