6 Commits
Author SHA1 Message Date
ARIA 70282dfb65 fix(release): use Forgejo-style /assets and /tags API routes
CI / Gateway plugin tests (push) Successful in 4m55s
CI / Kotlin tests (android host + desktop) (push) Successful in 6m48s
The server (gitea.zephyre.one) exposes a Forgejo-compatible API:
- release attachments live at POST /releases/{id}/assets, not /attachments
- tag deletion is DELETE /tags/{tag}, not DELETE /git/refs/tags/{tag}
(verified against the live API: /attachments 404s, /assets and /tags exist)
2026-08-23 19:53:17 +02:00
ARIA 44e8c7322e fix(release): delete leftover tag on re-run; support \\n in changelog input
CI / Gateway plugin tests (push) Successful in 5m0s
CI / Kotlin tests (android host + desktop) (push) Successful in 6m43s
- Gitea's DELETE /releases/:id does not remove the tag, so re-running the
  workflow for the same version failed with 409 (curl exit 22). The
  re-run safety block now also deletes the tag via git/refs/tags.
- Replace curl -sf with an api() wrapper that prints Gitea's error body
  on HTTP >= 400 instead of failing silently.
- The workflow_dispatch changelog input is single-line (Gitea has no
  multiline input type); convert literal \\n to real newlines and
  document it in the input description.
2026-08-23 19:03:59 +02:00
ARIA 29d0c1a73f Fix two gateway test failures: outbox lane scoping + SSE teardown race
CI / Gateway plugin tests (push) Successful in 5m46s
CI / Kotlin tests (android host + desktop) (push) Successful in 6m57s
- outbox: delete_message/message_info now match the exact lane first
  (a flat-lane delete/lookup with thread_id=None sees only frames with
  no thread_id) and fall back to the message_id across all lanes only
  when the exact lane matches nothing. Previously lane=None meant
  'any lane' in the first pass, so a flat-lane delete also removed
  same-id frames from threads (test expected 3 removed, got 4).

- http_server: the SSE live loop skipped queued frames when stop() set
  sub.closed before the handler thread reached the loop (descheduled
  under load between the initial hello/status writes and the loop).
  The loop now drains frames queued before the close, so the
  status{restarting} teardown broadcast always reaches the client
  before EOF (test_disconnect_broadcasts_status_restarting was flaky
  ~70% under CPU load).
2026-08-23 14:45:09 +02:00
ARIA d801a18db5 fix FCM push notifications
CI / Gateway plugin tests (push) Failing after 6m27s
CI / Kotlin tests (android host + desktop) (push) Successful in 7m2s
2026-08-23 14:32:40 +02:00
Pakobbix 560d9c19b3 Merge pull request 'feat(app): copy messages via bubble context menu + selection toolbar' (#9) from feat/message-copy-context-menu into master
CI / Kotlin tests (android host + desktop) (push) Successful in 9m0s
CI / Gateway plugin tests (push) Failing after 9m18s
Reviewed-on: #9
2026-08-23 11:10:30 +00:00
ARIA 32db7fc4e8 feat(app): copy messages via bubble context menu + selection toolbar
CI / Kotlin tests (android host + desktop) (pull_request) Successful in 8m18s
CI / Gateway plugin tests (pull_request) Failing after 9m55s
Long-press (touch) / right-click (desktop) on a finalized message now
opens a context menu anchored to the bubble: Copy / Select messages /
Delete. The multi-select toolbar gains a Copy button that joins the
selected messages in display order. Copies raw markdown so formatting
survives pasting into other markdown apps; blank text is a no-op.
Cancelling a menu-opened delete confirm clears the staged selection so
the bubble doesn't stay highlighted.
2026-08-23 13:08:44 +02:00
7 changed files with 481 additions and 165 deletions

No files matched your search

+28 -8
View File
@@ -8,7 +8,7 @@ on:
required: true
type: string
changelog:
description: "Release notes (markdown, shown on the release page)"
description: "Release notes (markdown, shown on the release page). Single-line field — use literal \\n for line breaks."
required: false
type: string
@@ -178,19 +178,38 @@ jobs:
REPO="${GITEA_REPOSITORY:-$GITHUB_REPOSITORY}"
TOKEN="${RELEASE_TOKEN:-$GITHUB_TOKEN}"
VERSION=$(jq -r '.inputs.version' "$GITHUB_EVENT_PATH")
CHANGELOG=$(jq -r '.inputs.changelog // ""' "$GITHUB_EVENT_PATH")
# The dispatch input is a single-line field; turn literal \n into real newlines.
CHANGELOG=$(jq -r '.inputs.changelog // ""' "$GITHUB_EVENT_PATH" | sed 's/\\n/\n/g')
TAG="v$VERSION"
API="$SERVER/api/v1/repos/$REPO"
AUTH="Authorization: token $TOKEN"
# Re-run safety: drop a previous release (and its tag) for this version.
OLD_ID=$(curl -sf -H "$AUTH" "$API/releases/tags/$TAG" | jq -r '.id // empty')
# curl wrapper: on HTTP >= 400, print the response body (Gitea's error
# message) before failing — plain `curl -f` hides it (exit 22).
api() {
local code body
body=$(mktemp)
code=$(curl -s -o "$body" -w '%{http_code}' "$@") || { cat "$body"; rm -f "$body"; return 1; }
if [ "${code:0:1}" != "2" ]; then
echo "API error $code: $(cat "$body")" >&2
rm -f "$body"
return 1
fi
cat "$body"
rm -f "$body"
}
# Re-run safety: drop a previous release AND its tag for this version.
# (Gitea's DELETE /releases/:id does NOT remove the tag; a leftover tag
# makes the POST below fail with 409.)
OLD_ID=$(api -H "$AUTH" "$API/releases/tags/$TAG" | jq -r '.id // empty') || true
if [ -n "$OLD_ID" ]; then
curl -sf -X DELETE -H "$AUTH" "$API/releases/$OLD_ID" > /dev/null
api -X DELETE -H "$AUTH" "$API/releases/$OLD_ID" > /dev/null
fi
api -X DELETE -H "$AUTH" "$API/tags/$TAG" > /dev/null || true
# Gitea creates the tag at the default branch HEAD automatically.
RELEASE_ID=$(curl -sf -X POST -H "$AUTH" -H "Content-Type: application/json" \
RELEASE_ID=$(api -X POST -H "$AUTH" -H "Content-Type: application/json" \
"$API/releases" \
-d "$(jq -n --arg tag "$TAG" --arg title "Iris $VERSION" --arg body "$CHANGELOG" \
'{tag_name:$tag, title:$title, body:$body}')" \
@@ -200,7 +219,8 @@ jobs:
for f in "$GITHUB_WORKSPACE"/iris-android-v* "$GITHUB_WORKSPACE"/iris-desktop-*; do
[ -f "$f" ] || continue
echo "Uploading $(basename "$f")"
curl -sf -X POST -H "$AUTH" -F "attachment=@$f" \
"$API/releases/$RELEASE_ID/attachments" > /dev/null
# Forgejo-style API: release assets live under /assets, not /attachments.
api -X POST -H "$AUTH" -F "attachment=@$f" \
"$API/releases/$RELEASE_ID/assets" > /dev/null
done
echo "Done: $SERVER/$REPO/releases/tag/$TAG"
@@ -487,6 +487,22 @@ fun ChatScreen(controller: IrisController) {
var selectionMode by remember { mutableStateOf(false) }
var selectedIds by remember { mutableStateOf<Set<String>>(emptySet()) }
var showDeleteConfirm by remember { mutableStateOf(false) }
// True when the confirm dialog was opened from a bubble's context menu
// (which stages selectedIds WITHOUT entering selection mode) — cancel
// must then clear the staged selection so the bubble doesn't stay
// highlighted.
var deleteConfirmFromMenu by remember { mutableStateOf(false) }
val clipboard = LocalClipboardManager.current
// Copy raw markdown (formatting survives pasting into other markdown
// apps) with a toast, since the clipboard write itself is silent.
// Blank text (e.g. a media-only message) is a no-op: writing an empty
// string would wipe the clipboard for no reason.
fun copyText(text: String) {
if (text.isBlank()) return
clipboard.setText(AnnotatedString(text))
toastMessage = "Copied"
}
fun enterSelection(id: String) {
selectionMode = true
@@ -670,14 +686,30 @@ fun ChatScreen(controller: IrisController) {
selected = selectable && item.id in selectedIds,
onToggleSelect = { if (selectable) toggleSelect(item.id) },
onLongPress = {
if (selectable) {
if (selectionMode) {
toggleSelect(item.id)
} else {
enterSelection(item.id)
}
}
if (selectable && selectionMode) toggleSelect(item.id)
},
onCopy =
if (selectable) {
{ copyText(item.text) }
} else {
null
},
onDelete =
if (selectable) {
{
selectedIds = setOf(item.id)
deleteConfirmFromMenu = true
showDeleteConfirm = true
}
} else {
null
},
onSelect =
if (selectable) {
{ enterSelection(item.id) }
} else {
null
},
runtimeFooterEnabled = runtimeFooterEnabled,
runtimeFooterFields = runtimeFooterFields,
)
@@ -913,7 +945,20 @@ fun ChatScreen(controller: IrisController) {
SelectionToolbar(
count = selectedIds.size,
onCancel = { exitSelection() },
onDelete = { showDeleteConfirm = true },
// Multi-select copy: join the selected messages in
// display order, separated by a blank line.
onCopy = {
val text =
items
.filterIsInstance<MessageItem>()
.filter { it.id in selectedIds }
.joinToString("\n\n") { it.text }
copyText(text)
},
onDelete = {
deleteConfirmFromMenu = false
showDeleteConfirm = true
},
)
}
@@ -1062,8 +1107,16 @@ fun ChatScreen(controller: IrisController) {
// Message selection: confirm deleting the selected message(s).
if (showDeleteConfirm) {
val n = selectedIds.size
// Cancel: a menu-opened confirm staged selectedIds without
// entering selection mode — clear it so the bubble doesn't
// stay highlighted (toolbar path keeps its selection).
fun cancelDeleteConfirm() {
showDeleteConfirm = false
if (deleteConfirmFromMenu) selectedIds = emptySet()
}
AlertDialog(
onDismissRequest = { showDeleteConfirm = false },
onDismissRequest = { cancelDeleteConfirm() },
title = { Text(if (n == 1) "Delete message?" else "Delete $n messages?") },
text = {
Text(
@@ -1082,7 +1135,7 @@ fun ChatScreen(controller: IrisController) {
}) { Text("Delete") }
},
dismissButton = {
TextButton(onClick = { showDeleteConfirm = false }) { Text("Cancel") }
TextButton(onClick = { cancelDeleteConfirm() }) { Text("Cancel") }
},
)
}
@@ -2460,6 +2513,9 @@ private fun MessageBubble(
selected: Boolean = false,
onToggleSelect: () -> Unit = {},
onLongPress: () -> Unit = {},
onCopy: (() -> Unit)? = null,
onDelete: (() -> Unit)? = null,
onSelect: (() -> Unit)? = null,
runtimeFooterEnabled: Boolean = false,
runtimeFooterFields: List<String> = emptyList(),
) {
@@ -2489,141 +2545,178 @@ private fun MessageBubble(
SelectionCheck(selected)
Spacer(modifier = Modifier.width(6.dp))
}
Column(
modifier =
Modifier
.widthIn(max = maxWidth)
.clip(RoundedCornerShape(14.dp))
.background(if (selected) IrisColors.chipSelected else bubbleColor)
.then(
// 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()
}
// Context menu (long-press / right-click): copy / select / delete.
// Anchored to the bubble itself (Box), not the full-width row, so it
// pops up next to the bubble on both sides.
val hasMenu = onCopy != null || onDelete != null || onSelect != null
var menuOpen by remember { mutableStateOf(false) }
Box {
Column(
modifier =
Modifier
.widthIn(max = maxWidth)
.clip(RoundedCornerShape(14.dp))
.background(if (selected) IrisColors.chipSelected else bubbleColor)
.then(
// Long-press / right-click opens the context menu (or
// toggles selection while selecting); 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 = {
if (hasMenu && !selectionMode) menuOpen = true else onLongPress()
},
).rightClick {
if (hasMenu && !selectionMode) menuOpen = true else onLongPress()
},
onLongClick = { onLongPress() },
).rightClick { onLongPress() },
).padding(horizontal = 12.dp, vertical = 8.dp),
) {
// Reasoning block above the answer (assistant, non-commentary).
if (!isUser && !isCommentary && !msg.reasoning.isNullOrBlank()) {
ReasoningBlock(msg.reasoning!!, reasoningAutoCollapse, bubbleColor)
Spacer(modifier = Modifier.height(6.dp))
}
if (isUser) {
// M7: user text rendered as markdown (bold / italic / inline
// code), so a pasted snippet or emphasis shows up the same as
// in replies.
if (msg.text.isNotBlank() || msg.streaming) {
MarkdownText(
text =
msg.text.prepareForMarkdown().preserveNewlinesAsHardBreaks() +
if (msg.streaming) " ▉" else "",
color = textColor,
fontSize = 15.sp,
isStreaming = msg.streaming,
)
).padding(horizontal = 12.dp, vertical = 8.dp),
) {
// Reasoning block above the answer (assistant, non-commentary).
if (!isUser && !isCommentary && !msg.reasoning.isNullOrBlank()) {
ReasoningBlock(msg.reasoning!!, reasoningAutoCollapse, bubbleColor)
Spacer(modifier = Modifier.height(6.dp))
}
// M7: timestamp + delivery status on their own line at the
// bottom-right (Telegram look, same as agent bubbles).
Row(
modifier = Modifier.fillMaxWidth(),
horizontalArrangement = Arrangement.End,
verticalAlignment = Alignment.CenterVertically,
) {
if (time.isNotEmpty()) {
Text(time, fontSize = 10.sp, color = textColor.copy(alpha = 0.6f))
if (isUser) {
// M7: user text rendered as markdown (bold / italic / inline
// code), so a pasted snippet or emphasis shows up the same as
// in replies.
if (msg.text.isNotBlank() || msg.streaming) {
MarkdownText(
text =
msg.text.prepareForMarkdown().preserveNewlinesAsHardBreaks() +
if (msg.streaming) " ▉" else "",
color = textColor,
fontSize = 15.sp,
isStreaming = msg.streaming,
)
}
when (msg.status) {
MsgStatus.Pending -> {
Text(" sending…", fontSize = 10.sp, color = textColor.copy(alpha = 0.6f))
}
MsgStatus.Sent -> {
Text(" ✓", fontSize = 10.sp, color = textColor.copy(alpha = 0.8f))
}
MsgStatus.Read -> {
Text(" ✓✓", fontSize = 10.sp, color = textColor.copy(alpha = 0.8f))
}
MsgStatus.Failed -> {
Unit
}
}
}
if (msg.status == MsgStatus.Failed) {
// M7: timestamp + delivery status on their own line at the
// bottom-right (Telegram look, same as agent bubbles).
Row(
modifier = Modifier.fillMaxWidth(),
horizontalArrangement = Arrangement.End,
) {
Text("Failed to send — tap to retry", fontSize = 10.sp, color = IrisColors.errorText)
}
}
} else {
if (msg.text.isNotBlank() || msg.streaming) {
// M8: render the agent's reply as markdown (bold / italic /
// underscore, tables, highlighted code blocks). Leading
// newlines are stripped so the text hugs the top of the
// bubble; single newlines become hard breaks (models write
// status lines and wrapped text expecting a break per line,
// same as user input); the ▉ cursor is kept while streaming.
val displayText =
msg.text.prepareForMarkdown().preserveNewlinesAsHardBreaks() +
if (msg.streaming) " ▉" else ""
MarkdownText(
text = displayText,
color = textColor,
fontSize = if (isCommentary) 13.sp else 15.sp,
modifier = Modifier.fillMaxWidth(),
isStreaming = msg.streaming,
)
}
}
// M4: media attachments (image / player / document chip).
if (msg.media.isNotEmpty()) {
Spacer(modifier = Modifier.height(6.dp))
Column(verticalArrangement = Arrangement.spacedBy(6.dp)) {
msg.media.forEach { m -> MediaAttachment(m) }
}
}
if (!isUser) {
// Runtime-metadata footer + timestamp on ONE line (footer left,
// time right) — Telegram-style. The app controls whether/what
// the footer shows via Settings → Runtime footer.
val footerText =
if (!isCommentary && !msg.streaming && runtimeFooterEnabled) {
buildRuntimeFooterText(msg, runtimeFooterFields)
} else {
""
}
if (footerText.isNotEmpty() || time.isNotEmpty()) {
Spacer(modifier = Modifier.height(4.dp))
Row(
modifier = Modifier.fillMaxWidth(),
verticalAlignment = Alignment.CenterVertically,
) {
if (footerText.isNotEmpty()) {
Text(
footerText,
color = textColor.copy(alpha = 0.4f),
fontSize = 10.sp,
modifier = Modifier.weight(1f),
)
} else {
Spacer(modifier = Modifier.weight(1f))
}
if (time.isNotEmpty()) {
Text(time, fontSize = 10.sp, color = textColor.copy(alpha = 0.5f))
Text(time, fontSize = 10.sp, color = textColor.copy(alpha = 0.6f))
}
when (msg.status) {
MsgStatus.Pending -> {
Text(" sending…", fontSize = 10.sp, color = textColor.copy(alpha = 0.6f))
}
MsgStatus.Sent -> {
Text(" ✓", fontSize = 10.sp, color = textColor.copy(alpha = 0.8f))
}
MsgStatus.Read -> {
Text(" ✓✓", fontSize = 10.sp, color = textColor.copy(alpha = 0.8f))
}
MsgStatus.Failed -> {
Unit
}
}
}
if (msg.status == MsgStatus.Failed) {
Row(
modifier = Modifier.fillMaxWidth(),
horizontalArrangement = Arrangement.End,
) {
Text("Failed to send — tap to retry", fontSize = 10.sp, color = IrisColors.errorText)
}
}
} else {
if (msg.text.isNotBlank() || msg.streaming) {
// M8: render the agent's reply as markdown (bold / italic /
// underscore, tables, highlighted code blocks). Leading
// newlines are stripped so the text hugs the top of the
// bubble; single newlines become hard breaks (models write
// status lines and wrapped text expecting a break per line,
// same as user input); the ▉ cursor is kept while streaming.
val displayText =
msg.text.prepareForMarkdown().preserveNewlinesAsHardBreaks() +
if (msg.streaming) " ▉" else ""
MarkdownText(
text = displayText,
color = textColor,
fontSize = if (isCommentary) 13.sp else 15.sp,
modifier = Modifier.fillMaxWidth(),
isStreaming = msg.streaming,
)
}
}
// M4: media attachments (image / player / document chip).
if (msg.media.isNotEmpty()) {
Spacer(modifier = Modifier.height(6.dp))
Column(verticalArrangement = Arrangement.spacedBy(6.dp)) {
msg.media.forEach { m -> MediaAttachment(m) }
}
}
if (!isUser) {
// Runtime-metadata footer + timestamp on ONE line (footer left,
// time right) — Telegram-style. The app controls whether/what
// the footer shows via Settings → Runtime footer.
val footerText =
if (!isCommentary && !msg.streaming && runtimeFooterEnabled) {
buildRuntimeFooterText(msg, runtimeFooterFields)
} else {
""
}
if (footerText.isNotEmpty() || time.isNotEmpty()) {
Spacer(modifier = Modifier.height(4.dp))
Row(
modifier = Modifier.fillMaxWidth(),
verticalAlignment = Alignment.CenterVertically,
) {
if (footerText.isNotEmpty()) {
Text(
footerText,
color = textColor.copy(alpha = 0.4f),
fontSize = 10.sp,
modifier = Modifier.weight(1f),
)
} else {
Spacer(modifier = Modifier.weight(1f))
}
if (time.isNotEmpty()) {
Text(time, fontSize = 10.sp, color = textColor.copy(alpha = 0.5f))
}
}
}
}
}
if (hasMenu) {
DropdownMenu(expanded = menuOpen, onDismissRequest = { menuOpen = false }) {
onCopy?.let {
DropdownMenuItem(text = { Text("Copy") }, onClick = {
menuOpen = false
it()
})
}
onSelect?.let {
DropdownMenuItem(
text = { Text("Select messages") },
onClick = {
menuOpen = false
it()
},
)
}
onDelete?.let {
DropdownMenuItem(text = { Text("Delete") }, onClick = {
menuOpen = false
it()
})
}
}
}
@@ -2699,6 +2792,7 @@ private fun SelectionCheck(selected: Boolean) {
private fun SelectionToolbar(
count: Int,
onCancel: () -> Unit,
onCopy: () -> Unit,
onDelete: () -> Unit,
) {
Row(
@@ -2720,6 +2814,20 @@ private fun SelectionToolbar(
style = MaterialTheme.typography.bodyLarge,
modifier = Modifier.weight(1f),
)
// Copy the selection (secondary action → neutral chip, Delete stays
// the accent primary action).
Box(
modifier =
Modifier
.clip(RoundedCornerShape(20.dp))
.background(if (count > 0) IrisColors.chip else Color.Transparent)
.clickable(enabled = count > 0) { onCopy() }
.padding(horizontal = 16.dp, vertical = 8.dp),
contentAlignment = Alignment.Center,
) {
Text("⧉ Copy", color = if (count > 0) IrisColors.textBright else IrisColors.textDim, fontSize = 14.sp)
}
Spacer(modifier = Modifier.width(8.dp))
Box(
modifier =
Modifier
+37 -4
View File
@@ -5,13 +5,44 @@ relay. **Decision: FCM primary, ntfy fallback** (`IRIS_PUSH_BACKEND`).
## 8.1 When push fires
- A frame targets a `chat_id` whose device is **disconnected** (WS closed) →
drop to **outbox** + fire **push**.
- A frame targets a `chat_id` whose device is **disconnected** (no live
SSE/long-poll subscriber) → drop to **outbox** + fire **push** — with one
refinement, *turn-aware push* (below).
- Also fire push for high-priority foreground events the user should see even if
the app is backgrounded (approvals, clarifies, cron completions) — the app
decides whether to also show an in-app banner.
- If the device is **connected**, no push (the live frame is enough).
### 8.1.1 Turn-aware push (one push per turn, final answer as body)
An agent turn can span minutes and emit several completed status messages
("researching X…", "found Y…", "writing findings…", final answer). Pushing each
parked message would spam an offline user with the steps in between. So while
the agent's turn is **in flight** (hermes holds the typing indicator on from
turn start until the handler's `finally` at turn end), normal-priority
`message` / `message.stop` / `media.offer` frames that park with no live
device are **held back** per chat instead of pushing; the latest one is pushed
when the turn ends (`stop_typing`), so the offline user gets **one push with
the final answer**. Details:
- Turn state is tracked per `chat_id` from the typing indicator
(`send_typing` → in flight, `stop_typing` → ended; hermes fires
`stop_typing` in the handler's `finally`, after the final send, so the
flush always sees the final frame).
- The held-back frame is still parked in the outbox — sync catch-up is
unaffected; only the push is deferred.
- **High-priority notifications** (approval/clarify/cron) push immediately,
even mid-turn — they need user action.
- If the device **reconnects mid-turn** (SSE/long-poll open), the held-back
push is dropped: the app syncs the parked frames and must not get a
duplicate push at turn end.
- If the turn ends while the device is live, nothing is pushed (the frames
were delivered live / synced).
- Turns without a typing indicator (e.g. typing disabled in config) and
non-turn deliveries (cron, standalone sends) push immediately as before.
- Best-effort: a gateway crash mid-turn loses the held-back push (the frames
remain in the outbox and sync on reconnect).
## 8.2 `PushBackend` interface (`push.py`)
```python
@@ -26,9 +57,10 @@ class PushBackend(Protocol):
Selected at adapter init by `IRIS_PUSH_BACKEND` (`fcm` default, `ntfy`).
### 8.2.1 `FcmBackend` (primary)
- **FCM HTTP v1 API** via `httpx` (core dep). Auth = Firebase **service
account** (`IRIS_FCM_SERVICE_ACCOUNT` JSON path) → mint a short-lived
OAuth2 access token (cached, refreshed before expiry).
OAuth2 access token (cached, refreshed before expiry).
- Fallback: legacy **server key** (`IRIS_FCM_SERVER_KEY`) if no service
account (simpler, but legacy).
- Target = the device's **FCM token** (registered via `hello` /
@@ -43,6 +75,7 @@ Selected at adapter init by `IRIS_PUSH_BACKEND` (`fcm` default, `ntfy`).
devices).
### 8.2.2 `NtfyBackend` (fallback, self-host friendly)
- Reuses hermes's existing ntfy publish path (hermes ships an ntfy adapter).
- Publish to `NTFY_TOPIC` on `NTFY_SERVER_URL` (default `https://ntfy.sh`) via
`httpx` POST, with an `X-Title` / `X-Message` / `X-Tag` / `X-Priority` and a
@@ -132,4 +165,4 @@ device push watermark:
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.
into the app and lets it use a stable per-message notification id.
+60 -2
View File
@@ -1360,6 +1360,16 @@ class IrisAdapter(BasePlatformAdapter):
# 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] = {}
# M5: turn-aware push. chat_ids whose agent turn is in flight, per
# the typing indicator (hermes turns typing on at turn start and off
# in the handler's finally at turn end). While a turn is active and
# the device is offline, pushable message/media frames are held back
# in _pending_push instead of pushing per intermediate status
# message; the latest one is pushed when the turn ends, so an
# offline user gets ONE push with the final answer. High-priority
# notifications (approval/clarify/cron) always push immediately.
self._typing_turns: set[str] = set()
self._pending_push: dict[str, tuple[protocol.Frame, int]] = {}
# Live todo lists (lane key -> items): the agent's planning state,
# re-served as a snapshot when a device opens its event stream.
# In-memory only — a gateway restart drops it (the agent re-emits on
@@ -1948,8 +1958,27 @@ class IrisAdapter(BasePlatformAdapter):
frame.type,
cursor,
)
# M5: wake the offline device(s) via the push backend.
await self._maybe_push(chat_id, frame, cursor)
# M5: wake the offline device(s) via the push backend -- unless
# this is a normal-priority message/media frame while the
# agent's turn is still in flight: hold it back and push the
# latest one when the turn ends (stop_typing), so an offline
# user gets one push with the final answer instead of one per
# intermediate status message. High-priority notifications
# (approval/clarify/cron) still push immediately.
if chat_id in self._typing_turns and frame.type in (
protocol.TYPE_MESSAGE,
protocol.TYPE_MESSAGE_STOP,
protocol.TYPE_MEDIA_OFFER,
):
self._pending_push[chat_id] = (frame, cursor)
logger.info(
"iris: push deferred for %s (turn active; %s frame cursor=%s)",
chat_id,
frame.type,
cursor,
)
else:
await self._maybe_push(chat_id, frame, cursor)
await self._maybe_notify_outbox_prune(chat_id)
elif (
frame.type == protocol.TYPE_NOTIFICATION
@@ -2099,6 +2128,9 @@ class IrisAdapter(BasePlatformAdapter):
tid = metadata.get("thread_id")
if isinstance(tid, str) and tid:
thread_id = tid
# M5: typing on marks the chat's agent turn as in flight (hermes
# turns typing on at turn start, before the first output).
self._typing_turns.add(chat_id)
frame = protocol.typing(chat_id, True, thread_id=thread_id)
await self._http_server.fanout(frame, cursor=None)
@@ -2106,6 +2138,32 @@ class IrisAdapter(BasePlatformAdapter):
"""Clear the typing indicator (``typing`` frame, on=false)."""
frame = protocol.typing(chat_id, False)
await self._http_server.fanout(frame, cursor=None)
# M5: turn ended (hermes fires stop_typing in the handler's finally,
# after the final send). Flush the held-back push -- the final
# answer -- but only while the device is still offline; a live
# device already got the frames via its event stream / sync.
# Idempotent: hermes may call stop_typing more than once per turn.
self._typing_turns.discard(chat_id)
pending = self._pending_push.pop(chat_id, None)
if pending is None:
return
if self._http_server.has_devices():
logger.info("iris: deferred push dropped for %s (device back online)", chat_id)
return
held_frame, held_cursor = pending
await self._maybe_push(chat_id, held_frame, held_cursor)
def on_device_online(self) -> None:
"""A device opened its event stream (SSE/long-poll): it will sync
the outbox, so drop any held-back pushes -- flushing them later
would duplicate what the app already shows. Called from the HTTP
server's handler thread; dict.clear() is atomic under the GIL."""
if self._pending_push:
logger.info(
"iris: device online; dropping %d deferred push(es)",
len(self._pending_push),
)
self._pending_push.clear()
# ── M4: outbound media (agent -> app) ─────────────────────────────────
#
+26 -6
View File
@@ -275,6 +275,13 @@ class HttpServer:
def _add_sub(self, sub: _Subscriber) -> None:
with self._subs_lock:
self._subs.setdefault(sub.device_id, []).append(sub)
# M5: a live subscriber will sync the outbox -- tell the adapter to
# drop any held-back (deferred) pushes so the turn-end flush doesn't
# duplicate what the app already shows. getattr-guard: test doubles
# may use a bare adapter stub.
on_online = getattr(self._adapter, "on_device_online", None)
if on_online is not None:
on_online()
def _remove_sub(self, sub: _Subscriber) -> None:
with self._subs_lock:
@@ -516,12 +523,25 @@ class HttpServer:
for snap in self._adapter.todo_snapshot_frames():
self._write_sse(handler, "frame", None, snap.to_json())
# 3. Live frames (cursor=None frames have no id).
while not sub.closed.is_set():
try:
item = sub.q.get(timeout=SSE_HEARTBEAT_S)
except queue.Empty:
self._write_raw(handler, ": hb\n\n")
continue
while True:
if sub.closed.is_set():
# stop() can land between the initial writes above and
# this loop (the handler thread is descheduled under
# load): drain the frames queued before the close — e.g.
# the status{restarting} teardown broadcast — so the
# client sees them before EOF instead of losing them to
# the closed check.
try:
item = sub.q.get_nowait()
except queue.Empty:
reason = "stopped"
break
else:
try:
item = sub.q.get(timeout=SSE_HEARTBEAT_S)
except queue.Empty:
self._write_raw(handler, ": hb\n\n")
continue
if item is _STOP:
reason = "stopped"
break
+14 -10
View File
@@ -303,10 +303,11 @@ class Outbox:
frame of a streamed reply -- or ``None`` when the message is not in the
outbox (e.g. already pruned by retention).
The lane is matched exactly first; when that finds nothing the lookup
falls back to the ``message_id`` alone (it is a unique uuid4), so a
stale/missing ``thread_id`` on the request still resolves the row.
(``lane=None`` in the scan means "any lane".)
The lane is matched exactly first (a flat-lane lookup, ``thread_id
= None``, sees only frames with no ``thread_id``); when that finds
nothing the lookup falls back to the ``message_id`` alone across all
lanes (it is a unique uuid4), so a stale/missing ``thread_id`` on the
request still resolves the row.
"""
if not message_id:
return None
@@ -315,7 +316,7 @@ class Outbox:
"SELECT frame FROM outbox WHERE chat_id = ?", (chat_id,)
).fetchall()
def scan(lane: str | None) -> dict[str, Any] | None:
def scan(lane: str | None, exact: bool) -> dict[str, Any] | None:
msg_frame: dict[str, Any] | None = None
stop_frame: dict[str, Any] | None = None
for r in rows:
@@ -325,7 +326,7 @@ class Outbox:
continue
if not isinstance(frame, dict):
continue
if lane is not None and _frame_thread_id(frame) != lane:
if exact and _frame_thread_id(frame) != lane:
continue
payload = frame.get("payload")
if not isinstance(payload, dict) or payload.get("message_id") != message_id:
@@ -345,7 +346,7 @@ class Outbox:
}
return msg_frame or stop_frame
return scan(thread_id) or scan(None)
return scan(thread_id, exact=True) or scan(None, exact=False)
def delete_message(
self,
@@ -376,7 +377,7 @@ class Outbox:
"SELECT cursor, frame FROM outbox WHERE chat_id = ?", (chat_id,)
).fetchall()
def cursors_for(lane: str | None) -> list[int]:
def cursors_for(lane: str | None, exact: bool) -> list[int]:
out: list[int] = []
for r in rows:
try:
@@ -385,14 +386,17 @@ class Outbox:
continue
if not isinstance(frame, dict):
continue
if lane is not None and _frame_thread_id(frame) != lane:
if exact and _frame_thread_id(frame) != lane:
continue
payload = frame.get("payload")
if isinstance(payload, dict) and payload.get("message_id") == message_id:
out.append(int(r["cursor"]))
return out
cursors = cursors_for(thread_id) or cursors_for(None)
# Exact lane first (a flat-lane delete must not reach into
# threads); fall back to the message_id across all lanes only
# when the exact lane matches nothing (stale/missing thread_id).
cursors = cursors_for(thread_id, exact=True) or cursors_for(thread_id, exact=False)
if not cursors:
return 0
# One bound-parameter delete per cursor (a message spans only a few
+73
View File
@@ -1729,6 +1729,79 @@ async def test_push_skipped_when_backend_unconfigured(adapter):
assert fake.calls == []
@pytest.mark.asyncio
async def test_push_deferred_while_turn_active_flushed_on_stop_typing(adapter):
"""Turn-aware push: while the agent's turn is in flight (typing on) and
the device is offline, intermediate status messages park without a push;
the turn-end (stop_typing) flushes ONE push with the latest (final)
message. Regression: a long research turn pushed every status report."""
fake = _FakePush()
adapter._push = fake
adapter._devices.upsert(DEVICE_ID, "Test", {}, fcm_token="tok-1")
await adapter.send_typing("default") # turn starts
await adapter.send("default", "researching XXXX", metadata={"notify": True})
assert fake.calls == [] # intermediate: held back
await adapter.send(
"default", "found XXX, cross-referencing", metadata={"notify": True}
)
assert fake.calls == [] # still held back; latest replaces the pending
await adapter.send("default", "done, findings ready", metadata={"notify": True})
assert fake.calls == []
await adapter.stop_typing("default") # turn ends
assert len(fake.calls) == 1
call = fake.calls[0]
assert call["chat_id"] == "default"
assert call["data"]["kind"] == "message"
assert "done, findings ready" in call["body"]
# All frames are parked in the outbox for sync catch-up.
assert adapter._outbox.latest_cursor() == 3
@pytest.mark.asyncio
async def test_deferred_push_dropped_when_device_comes_online(adapter, ws_client):
"""If the device reconnects mid-turn it syncs the parked frames, so the
turn-end flush must not push a duplicate."""
ws, _ = ws_client
fake = _FakePush()
adapter._push = fake
adapter._devices.upsert(DEVICE_ID, "Test", {}, fcm_token="tok-1")
await adapter.send_typing("default")
await adapter.send("default", "intermediate", metadata={"notify": True})
assert fake.calls == []
# Device reconnects mid-turn (SSE open -> on_device_online clears the
# pending push); the parked frame arrives via the event stream.
frames = await recv_until(ws, lambda f: f.get("type") == "message")
assert frames[-1]["payload"]["text"] == "intermediate"
await adapter.stop_typing("default")
assert fake.calls == []
@pytest.mark.asyncio
async def test_high_priority_notification_pushes_immediately_during_turn(plugin, adapter):
"""Approval/clarify/cron notifications need user action: they push
immediately even while a turn is in flight and the device is offline."""
fake = _FakePush()
adapter._push = fake
adapter._devices.upsert(DEVICE_ID, "Test", {}, fcm_token="tok-1")
await adapter.send_typing("default")
await adapter._broadcast_or_log(
"default",
plugin.protocol.notification(
"default", plugin.protocol.NOTIF_APPROVAL, "Approve", "run rm -rf?"
),
)
assert len(fake.calls) == 1
assert fake.calls[0]["priority"] == "high"
assert fake.calls[0]["data"]["kind"] == "approval"
# ── M5: fcm.register ───────────────────────────────────────────────────────