Auto-threading, history pagination, streaming toggle + tool/reasoning display settings

- Auto-threading (Telegram topic-mode workflow): message.send {auto_thread}
  mints a fresh AI-named thread (instant derived title, LLM upgrade via
  channel.renamed); channel.created {auto:true}; the app jumps into the new
  thread and relocates the optimistic pending bubble.
- history frame: paged full message history for initial channel open /
  scroll-up pagination (reconstructed from the outbox log).
- Streaming on/off: gateway side (display.platforms.android.streaming) plus a
  per-device app toggle (Settings → Streaming); reasoning/model/tokens carried
  on message frames.
- Context menu: long-press (touch) / right-click (desktop) thread affordances
  via a KMP rightClick expect/actual.
- Settings → Reasoning: auto-collapse long reasoning blocks (default on).
- Tool detail: the gateway now always supplies full tool data — it forces
  verbose tool progress (full args → tool.start.args) and captures each
  completed call via the post_tool_call hook (output/duration/ok → tool.end).
  The app reveals the full call + output on expand (Truncated) and
  auto-expands cards in Everything mode.
This commit is contained in:
ARIA committed 2026-08-20 16:26:52 +02:00
1 parent 1d67de5f12
commit 60ec2b44a7
24 files changed
+1115 -81

No files matched your search

+366 -18
View File
@@ -58,12 +58,14 @@ Or via environment variables (overrides config.yaml; secrets live in .env):
"""
import asyncio
import json
import logging
import os
import re
import threading
import time
import uuid
from collections import deque
from dataclasses import dataclass, field
from typing import Any, Dict, List, Optional, Tuple
@@ -185,6 +187,75 @@ def _reset_reasoning() -> None:
_reasoning_flushed.clear()
# ---------------------------------------------------------------------------
# M2 — tool-result capture (post_tool_call hook)
#
# The gateway renders tool *progress* lines to the platform but never streams
# the tool *output* (it is the agent's concern, persisted to history, not
# presentation). To let the app show the full call + result on demand
# (Settings → Tool detail), we capture each completed tool call via the
# ``post_tool_call`` hook and attach it to the ``tool.end`` frame.
#
# Global FIFO (like the reasoning buffer): a personal android gateway serves
# one active turn at a time, and records are matched to the open tool by name
# in completion order. Bounded so a runaway turn can't grow it without limit.
# ---------------------------------------------------------------------------
_tool_results: "deque[Dict[str, Any]]" = deque()
_tool_results_lock = threading.Lock()
_MAX_TOOL_RESULTS = 200
_MAX_OUTPUT_PREVIEW = 8000
def _on_post_tool_call(**kwargs: Any) -> None:
"""Plugin hook: capture a completed tool call's result + timing."""
result = kwargs.get("result")
record = {
"tool_name": kwargs.get("tool_name") or "",
"result": (str(result) if result is not None else "")[:_MAX_OUTPUT_PREVIEW],
"duration_ms": kwargs.get("duration_ms") or 0,
"status": kwargs.get("status") or "ok",
}
with _tool_results_lock:
_tool_results.append(record)
while len(_tool_results) > _MAX_TOOL_RESULTS:
_tool_results.popleft()
def _take_tool_result(tool_name: str) -> Optional[Dict[str, Any]]:
"""Pop the first completed record matching *tool_name* (FIFO), else None."""
if not tool_name:
return None
with _tool_results_lock:
for i, rec in enumerate(_tool_results):
if rec["tool_name"] == tool_name:
del _tool_results[i]
return rec
return None
def _reset_tool_results() -> None:
"""Clear captured records (turn boundary — drop anything unconsumed)."""
with _tool_results_lock:
_tool_results.clear()
def _tool_end_fields(tool_name: str) -> Dict[str, Any]:
"""Build the ``tool.end`` enrichment (ok/duration/output_preview) from the
captured hook record for *tool_name*; empty dict when none is available
(e.g. tool_progress off, or the call came from another session)."""
rec = _take_tool_result(tool_name)
if rec is None:
return {}
fields: Dict[str, Any] = {
"ok": rec["status"] == "ok",
"output_preview": rec["result"] or None,
}
if rec["duration_ms"]:
fields["duration"] = round(rec["duration_ms"] / 1000.0, 3)
return fields
# ---------------------------------------------------------------------------
# Defaults
# ---------------------------------------------------------------------------
@@ -283,6 +354,25 @@ def _thread_id_from_metadata(metadata: Optional[Dict[str, Any]]) -> Optional[str
return None
def _derive_thread_name(text: str) -> str:
"""Instant auto-thread name from the user's opening message (no model).
Reuses hermes' session-title derivation (``agent/title_generator.py``):
a deterministic slice of the user's own words, so the thread is named the
moment it is created. The LLM upgrade (``_schedule_thread_title_upgrade``)
replaces it moments later — the same two-stage titling hermes uses for
sessions (derived < llm < user).
"""
try:
from agent.title_generator import derive_title
title = derive_title(text)
except Exception:
logger.debug("Thread name derivation failed", exc_info=True)
title = None
return (title or "").strip() or "New thread"
def _strip_streaming_cursor(text: str) -> str:
if text and text.endswith(_STREAMING_CURSOR):
return text[: -len(_STREAMING_CURSOR)]
@@ -406,6 +496,50 @@ def _extract_code_block(content: str) -> Optional[str]:
return None
def _extract_verbose_args(line: str, content: str) -> Optional[Dict[str, Any]]:
"""Recover the full args dict from a verbose tool line, else ``None``.
In verbose mode the gateway renders ``<emoji> <name>(keys)`` on one line
and the full args JSON on the line that follows. When *line* is such a
header, return the parsed JSON object from the following line.
"""
parts = line.strip().split(None, 1)
if len(parts) < 2 or not _TOOL_NAME_ARGS_RE.match(parts[1]):
return None
lines = content.splitlines()
for i, ln in enumerate(lines):
if ln.strip() != line.strip():
continue
for follow in lines[i + 1:]:
follow = follow.strip()
if not follow:
continue
if follow.startswith("{"):
try:
obj = json.loads(follow)
return obj if isinstance(obj, dict) else None
except Exception:
return None
return None # next non-empty line is not the args JSON
return None
def _short_preview_from_args(args: Dict[str, Any], cap: int = 60) -> Optional[str]:
"""Derive a short one-line preview from a verbose args dict.
The verbose line carries no explicit preview, so the Truncated display
would otherwise lose its one-liner. Use the first non-empty string value
(whitespace-collapsed, capped) as a stand-in.
"""
if not isinstance(args, dict):
return None
for value in args.values():
if isinstance(value, str) and value.strip():
s = " ".join(value.split())
return s[: cap - 1] + "…" if len(s) > cap else s
return None
def _is_tool_progress(content: str) -> bool:
"""Heuristic: does *content* look like gateway tool-progress line(s)?
@@ -439,6 +573,9 @@ class _TurnState:
tool_index: int = 0
# Index of the most recently started tool (awaiting tool.end).
open_tool_index: Optional[int] = None
# Name of the most recently started tool (matches the post_tool_call
# record when the tool completes, so tool.end can carry its output).
open_tool_name: Optional[str] = None
# Tool lines already emitted as tool.start (dedup across edits).
seen_tool_lines: set = field(default_factory=set)
@@ -601,6 +738,46 @@ async def _standalone_send(
}
# ---------------------------------------------------------------------------
# Verbose tool progress (full args on the progress line)
# ---------------------------------------------------------------------------
def _ensure_verbose_tool_progress() -> None:
"""Ensure the android platform renders tool progress in ``verbose`` mode.
Verbose mode makes the gateway's tool-progress line carry the FULL
argument JSON (not just a ~40-char preview), which the adapter parses
into the ``tool.start`` frame's ``args`` field; the app then decides how
much to show (Settings → Tool detail). The tool *output* is captured
separately via the ``post_tool_call`` hook (verbose mode does not stream
it).
Best-effort and idempotent: writes
``display.platforms.android.tool_progress: verbose`` to config.yaml only
when it isn't already set. The gateway's config cache is mtime-keyed, so
the write takes effect on the next turn without a restart. Never raises.
"""
try:
from hermes_cli.config import load_config_readonly
cfg = load_config_readonly() or {}
display = cfg.get("display") or {}
platforms = display.get("platforms") or {}
android = platforms.get("android") or {}
if android.get("tool_progress") == "verbose":
return # already set
from utils import atomic_roundtrip_yaml_update
atomic_roundtrip_yaml_update(
get_hermes_home() / "config.yaml",
"display.platforms.android.tool_progress",
"verbose",
)
logger.info("android: set display.platforms.android.tool_progress=verbose")
except Exception:
logger.debug("android: could not ensure verbose tool_progress", exc_info=True)
# ---------------------------------------------------------------------------
# Interactive setup (hermes gateway setup flow)
# ---------------------------------------------------------------------------
@@ -654,6 +831,10 @@ def interactive_setup() -> None:
print_info(f"Pairing URL: {qr_payload(host or DEFAULT_HOST, _parse_port(port), token)}")
print_info(f"Server URL: {url}")
# Always render tool progress verbosely so the app receives the full tool
# call args (it decides how much to show via Settings → Tool detail).
_ensure_verbose_tool_progress()
print_success("Android configuration saved to ~/.hermes/.env")
print_info("Restart the gateway for changes to take effect: hermes gateway restart")
@@ -675,6 +856,10 @@ class AndroidAdapter(BasePlatformAdapter):
platform = Platform("android")
super().__init__(config=config, platform=platform)
# Ensure verbose tool progress (full args on the progress line) so the
# app can show the full tool call on demand. Best-effort; idempotent.
_ensure_verbose_tool_progress()
extra = getattr(config, "extra", {}) or {}
# Connection settings (env vars override config.yaml)
@@ -1094,38 +1279,59 @@ class AndroidAdapter(BasePlatformAdapter):
parsed = self._parse_tool_line_or_block(line, content)
if parsed is None:
continue
name, preview = parsed
# A new tool begins: close the previously-open one.
name, preview, args = parsed
# A new tool begins: close the previously-open one, attaching the
# output/duration/ok captured by the post_tool_call hook.
if state.open_tool_index is not None:
extra = _tool_end_fields(state.open_tool_name or "")
await self._broadcast_or_log(
chat_id,
protocol.tool_end(chat_id, state.open_tool_index, "", ok=True, thread_id=thread_id),
protocol.tool_end(
chat_id, state.open_tool_index, state.open_tool_name or "",
ok=extra.get("ok", True),
duration=extra.get("duration"),
output_preview=extra.get("output_preview"),
thread_id=thread_id,
),
)
state.tool_index += 1
state.open_tool_index = state.tool_index
state.open_tool_name = name
await self._broadcast_or_log(
chat_id,
protocol.tool_start(
chat_id, state.tool_index, name,
preview=preview, thread_id=thread_id,
preview=preview, args=args, thread_id=thread_id,
),
)
return SendResult(success=True, message_id=message_id)
@staticmethod
def _parse_tool_line_or_block(line: str, content: str) -> Optional[Tuple[str, Optional[str]]]:
"""Parse a tool line, expanding a terminal code block to its command."""
def _parse_tool_line_or_block(
line: str, content: str
) -> Optional[Tuple[str, Optional[str], Optional[Dict[str, Any]]]]:
"""Parse a tool line into ``(name, preview, args)``.
Expands a terminal code block to its command, and a verbose header
(``<emoji> <name>(keys)``) to its full args JSON (the JSON sits on the
following line). ``args`` is ``None`` unless the line is a verbose
header with a parseable JSON body.
"""
parsed = _parse_tool_line(line)
if parsed is not None:
name, preview = parsed
# Terminal code block: the command lives in the fenced lines that
# follow the "<emoji> terminal" head line.
if name == "terminal" and preview is None and "```" in content:
cmd = _extract_code_block(content)
if cmd:
return name, cmd
return parsed
return None
if parsed is None:
return None
name, preview = parsed
# Terminal code block: the command lives in the fenced lines that
# follow the "<emoji> terminal" head line.
if name == "terminal" and preview is None and "```" in content:
cmd = _extract_code_block(content)
if cmd:
return name, cmd, None
# Verbose mode: recover the full args from the following JSON line.
args = _extract_verbose_args(line, content)
if args is not None and preview is None:
preview = _short_preview_from_args(args)
return name, preview, args
async def _close_open_tool(
self, chat_id: str, state: _TurnState, thread_id: Optional[str]
@@ -1137,11 +1343,19 @@ class AndroidAdapter(BasePlatformAdapter):
tool it was waiting on has returned).
"""
if state.open_tool_index is not None:
extra = _tool_end_fields(state.open_tool_name or "")
await self._broadcast_or_log(
chat_id,
protocol.tool_end(chat_id, state.open_tool_index, "", ok=True, thread_id=thread_id),
protocol.tool_end(
chat_id, state.open_tool_index, state.open_tool_name or "",
ok=extra.get("ok", True),
duration=extra.get("duration"),
output_preview=extra.get("output_preview"),
thread_id=thread_id,
),
)
state.open_tool_index = None
state.open_tool_name = None
def _reset_tool_state(self, state: _TurnState) -> None:
"""Clear per-turn tool bookkeeping (called at turn finalization)."""
@@ -1149,6 +1363,10 @@ class AndroidAdapter(BasePlatformAdapter):
state.seen_tool_lines = set()
state.tool_index = 0
state.open_tool_index = None
state.open_tool_name = None
# Drop any captured tool results not consumed by a tool.end this turn
# (e.g. tool_progress off) so they can't leak into the next turn.
_reset_tool_results()
async def _broadcast_or_log(self, chat_id: str, frame: "protocol.Frame") -> None:
delivered = await self._ws_server.broadcast(frame)
@@ -1465,6 +1683,33 @@ class AndroidAdapter(BasePlatformAdapter):
if not isinstance(reply_to, str) or not reply_to.strip():
reply_to = None
# Auto-threading (the app's Threads setting, docs/06 §6.3): a message
# in a channel's flat lane gets its own fresh thread, the way Telegram
# topic mode mints a topic per new conversation. The thread is named
# instantly from the user's opening message (derived title) and the
# LLM upgrades the name in the background. The user echo, the agent
# turn, and all streaming frames then carry the new thread_id.
# Skipped for slash commands (session-scoped, not conversation
# starters) and replies (they continue where the user is).
auto_thread = payload.get("auto_thread") is True
if (
auto_thread
and thread_id is None
and text.strip()
and not text.lstrip().startswith("/")
and reply_to is None
):
entry = self._channels.create(
name=_derive_thread_name(text),
kind="thread",
parent_chat_id=chat_id,
)
thread_id = entry["chat_id"]
# Bare broadcast (like channel.create): the directory is
# re-served on hello.ack, so no outbox entry is needed.
await self._ws_server.broadcast(protocol.channel_created(entry, auto=True))
self._schedule_thread_title_upgrade(entry["chat_id"], text)
# M4: resolve media refs (single-use; unknown ref -> error).
media_urls: List[str] = []
media_types: List[str] = []
@@ -1494,6 +1739,12 @@ class AndroidAdapter(BasePlatformAdapter):
# Echo to all devices: the sender confirms (server-assigned id),
# other devices see the message too (single-user, multi-device).
# Routed through _broadcast_or_log (not a bare broadcast) so the echo
# is appended to the outbox: the app's ChatStore is in-memory only, so
# after a process death / activity recreation the only way the user's
# own message is restored is via the sync replay. Without this, user
# messages vanish on reconnect while bot messages (already parked)
# survive.
message_id = f"m_{uuid.uuid4().hex[:16]}"
echo = protocol.message(
chat_id=chat_id,
@@ -1505,7 +1756,7 @@ class AndroidAdapter(BasePlatformAdapter):
reply_to=reply_to,
ts=int(time.time() * 1000),
)
await self._ws_server.broadcast(echo)
await self._broadcast_or_log(chat_id, echo)
# Refs are consumed by this message (no replay).
for ref in media_refs:
self._media.pop_inbound(ref)
@@ -1554,6 +1805,48 @@ class AndroidAdapter(BasePlatformAdapter):
protocol.read_receipt(chat_id, message_id),
)
def _schedule_thread_title_upgrade(self, thread_id: str, text: str) -> None:
"""Upgrade an auto-created thread's name with the model's title.
Stage 2 of hermes' two-stage session titling (``agent/title_generator
.py``): the thread was created with an instant derived name; this
background call on the ``title_generation`` auxiliary task replaces it
with the model's title and broadcasts ``channel.renamed``. Best-effort
— any failure (config, model, network) leaves the derived name in
place, and a thread the user already renamed or archived is untouched.
"""
loop = asyncio.get_running_loop()
def _work() -> None:
try:
from agent.title_generator import generate_title
title = generate_title(text)
except Exception:
logger.debug("Thread title upgrade failed", exc_info=True)
return
if not title:
return
entry = self._channels.get(thread_id)
if entry is None or entry.get("archived"):
return
if (entry.get("name") or "") == title:
return
renamed = self._channels.rename(thread_id, title)
if renamed is None:
return
try:
asyncio.run_coroutine_threadsafe(
self._ws_server.broadcast(protocol.channel_renamed(renamed)),
loop,
)
except Exception:
logger.debug("Thread title rename broadcast failed", exc_info=True)
threading.Thread(
target=_work, daemon=True, name="android-thread-title"
).start()
# ── M4: inbound media (app -> agent) ──────────────────────────────────
#
# ``media.upload.start`` -> raw binary frames (one at a time per
@@ -1875,6 +2168,54 @@ class AndroidAdapter(BasePlatformAdapter):
done = protocol.sync_done(self._outbox.latest_cursor(), id=frame.id)
await self._ws_server.send_to(device_id, done)
# ── Full message history (initial channel open / scroll-up) ───────────
async def on_history(self, frame: protocol.Frame, device_id: str) -> None:
"""Handle an inbound ``history`` request.
``sync`` only replays the outbox delta since the device's cursor, so
after a process death the app's in-memory ChatStore is empty and the
delta does not cover older messages. ``history`` loads the full
message list for a chat/thread (reconstructed from the outbox log) so
the app can populate the view on first open / restart.
"""
payload = frame.payload
chat_id = frame.chat_id or payload.get("chat_id")
logger.info("android: history request from %s chat_id=%r", device_id, 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_UNSUPPORTED, "history requires a 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
before = payload.get("before_message_id")
if not isinstance(before, str) or not before.strip():
before = None
limit_raw = payload.get("limit")
try:
limit = int(limit_raw) if limit_raw is not None else 50
except (TypeError, ValueError):
limit = 50
page = self._outbox.history(
chat_id,
thread_id=thread_id,
before_message_id=before,
limit=limit,
)
resp = protocol.history(
chat_id,
page["messages"],
page["has_more"],
thread_id=thread_id,
oldest_message_id=page["oldest_message_id"],
id=frame.id,
)
await self._ws_server.send_to(device_id, resp)
# ── M5: push token registration ───────────────────────────────────────
async def on_fcm_register(self, frame: protocol.Frame, device_id: str) -> None:
@@ -2123,6 +2464,13 @@ def register(ctx):
ctx.register_hook("on_stream_delta", _on_stream_delta)
except Exception:
logger.debug("android: on_stream_delta hook registration failed", exc_info=True)
# M2: capture each completed tool call's result + timing so the tool.end
# frame can carry the output (the gateway never streams tool output to
# platforms). The app shows it on demand (Settings → Tool detail).
try:
ctx.register_hook("post_tool_call", _on_post_tool_call)
except Exception:
logger.debug("android: post_tool_call hook registration failed", exc_info=True)
ctx.register_platform(
name="android",
label="Android",
+108
View File
@@ -169,6 +169,114 @@ class Outbox:
)
return out
# ── history (full message history for a chat/thread) ──────────────────
def history(
self,
chat_id: str,
thread_id: Optional[str] = None,
before_message_id: Optional[str] = None,
limit: int = 50,
) -> Dict[str, Any]:
"""Final messages for a chat/thread, for the ``history`` frame.
Reconstructs the message list from the outbox log: a final message is
either a standalone ``message`` frame (user echo / non-streaming
assistant) or a ``message.stop`` frame (streaming assistant final).
Intermediate frames (``message.start``/``message.update``, tool,
commentary, notification, …) are skipped. Messages are deduplicated by
``message_id`` and returned oldest → newest.
Pagination: ``before_message_id`` returns the page of messages older
than that id (the newest page when omitted). Returns
``{messages, has_more, oldest_message_id}``.
"""
limit = max(1, min(int(limit or 50), 200))
with self._lock:
rows = self._conn.execute(
"SELECT cursor, frame FROM outbox WHERE chat_id = ? "
"ORDER BY cursor ASC",
(chat_id,),
).fetchall()
final: List[Dict[str, Any]] = []
for r in rows:
try:
frame = json.loads(r["frame"])
except (json.JSONDecodeError, TypeError):
continue
if not isinstance(frame, dict):
continue
if thread_id is not None and frame.get("thread_id") != thread_id:
continue
ftype = frame.get("type")
payload = frame.get("payload")
if not isinstance(payload, dict):
continue
if ftype == "message":
role = payload.get("role")
if role not in ("user", "assistant"):
continue
final.append(
{
"cursor": int(r["cursor"]),
"message_id": payload.get("message_id"),
"role": role,
"text": payload.get("text", ""),
"reasoning": payload.get("reasoning"),
"model": payload.get("model"),
"tokens": payload.get("tokens"),
"ts": payload.get("ts"),
"media": payload.get("media"),
}
)
elif ftype == "message.stop":
final.append(
{
"cursor": int(r["cursor"]),
"message_id": payload.get("message_id"),
"role": "assistant",
"text": payload.get("final_text", ""),
"reasoning": payload.get("reasoning"),
"model": payload.get("model"),
"tokens": payload.get("tokens"),
"ts": payload.get("ts"),
"media": None,
}
)
# Deduplicate by message_id (keep the latest occurrence), keep order.
by_id: Dict[str, Dict[str, Any]] = {}
for m in final:
mid = m.get("message_id")
if mid:
by_id[mid] = m
ordered = sorted(by_id.values(), key=lambda m: m["cursor"])
# Paginate: messages older than ``before_message_id`` (newest page when
# omitted / not found — a pruned anchor falls back to the newest page).
if before_message_id:
idx = next(
(i for i, m in enumerate(ordered) if m["message_id"] == before_message_id),
None,
)
pool = ordered[:idx] if idx is not None else ordered
else:
pool = ordered
has_more = len(pool) > limit
page = pool[-limit:] if has_more else list(pool)
oldest_message_id = page[0]["message_id"] if page else None
for m in page:
m.pop("cursor", None)
# Omit absent optional fields (the app's serializer treats a
# missing key as its default, but a JSON ``null`` for a
# non-nullable field like ``media`` would fail to parse).
for key in ("reasoning", "model", "tokens", "ts", "media"):
if m.get(key) is None:
m.pop(key, None)
return {
"messages": page,
"has_more": has_more,
"oldest_message_id": oldest_message_id,
}
# ── retention ─────────────────────────────────────────────────────────
def _maybe_prune(self) -> None:
+47 -3
View File
@@ -68,6 +68,9 @@ TYPE_SEARCH_RESULTS = "search.results"
TYPE_SYNC = "sync"
TYPE_SYNC_DONE = "sync.done"
# Full message history (initial channel open / scroll-up pagination)
TYPE_HISTORY = "history"
# Media (M4)
TYPE_MEDIA_UPLOAD_START = "media.upload.start"
TYPE_MEDIA_UPLOAD_END = "media.upload.end"
@@ -430,9 +433,18 @@ def _channel_payload(entry: Dict[str, Any]) -> Dict[str, Any]:
return payload
def channel_created(entry: Dict[str, Any]) -> Frame:
"""Broadcast: a channel/thread was created."""
return Frame(type=TYPE_CHANNEL_CREATED, payload=_channel_payload(entry))
def channel_created(entry: Dict[str, Any], auto: bool = False) -> Frame:
"""Broadcast: a channel/thread was created.
``auto=True`` marks a thread the gateway minted itself for an incoming
message (auto-threading, docs/06 §6.3): the app jumps into it and the
name is an instant derived title, upgraded by the LLM via a follow-up
``channel.renamed``.
"""
payload = _channel_payload(entry)
if auto:
payload["auto"] = True
return Frame(type=TYPE_CHANNEL_CREATED, payload=payload)
def channel_renamed(entry: Dict[str, Any]) -> Frame:
@@ -484,6 +496,38 @@ def sync_done(cursor: int, *, id: Optional[int] = None) -> Frame:
return Frame(type=TYPE_SYNC_DONE, id=id, payload={"cursor": cursor})
# ---------------------------------------------------------------------------
# History frame (full message history for a chat/thread)
# ---------------------------------------------------------------------------
def history(
chat_id: str,
messages: List[Dict[str, Any]],
has_more: bool,
*,
thread_id: Optional[str] = None,
oldest_message_id: Optional[str] = None,
id: Optional[int] = None,
) -> Frame:
"""Response to a ``history`` request: a page of final messages for a
chat/thread, ordered oldest → newest. ``has_more`` signals older pages
exist; ``oldest_message_id`` is the ``before_message_id`` for the next
(older) page."""
payload: Dict[str, Any] = {
"messages": messages,
"has_more": has_more,
}
if oldest_message_id:
payload["oldest_message_id"] = oldest_message_id
return Frame(
type=TYPE_HISTORY,
id=id,
chat_id=chat_id,
thread_id=thread_id,
payload=payload,
)
# ---------------------------------------------------------------------------
# Push / notification frames (M5)
# ---------------------------------------------------------------------------
+2
View File
@@ -378,6 +378,8 @@ class WsServer:
await self._adapter.on_search(frame, device_id)
elif frame.type == protocol.TYPE_SYNC:
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_MEDIA_UPLOAD_START:
await self._adapter.on_media_upload_start(frame, device_id)
elif frame.type == protocol.TYPE_MEDIA_UPLOAD_END: