test(e2e/platforms): address review — non-vacuous stream finalize, stranger refusal, redelivery transcript, known_failure gates

- KNOWN cells use tests.e2e.core._pending_fixes.known_failure (Rig.gate) around the final
  assertions only; no strict xfail (merge-order safe)
- streaming scenarios: private chat per scenario (Rig.fresh_user), same-chat barrier re-sent
  when dropped in the #121393 window; wait_idle docstring corrected; TURN_TIMEOUT 90->30
- stream_finalize_rejected now faults the real finalize per transport (Telegram draft final
  sendMessage / group editMessageText, Slack closing appendStream/stopStream, Discord edit) and
  asserts the fault fired; + stream_finalize_rejected_group, stream_reply_once
- Telegram stand-in: sendMessageDraft/sendRichMessageDraft, unknown methods 404, empty/too-long
  text and MarkdownV2 parse errors; Slack unknown_method + users.conversations; Discord empty
  content 50006, handler 4xx recorded as faulted
- approval_click: stranger clicks Deny and must be refused (answered) before the owner approves
- redelivery_one_reply asserts the inbound reached the transcript once
- dm_one_reply asserts no formatting rejection (plain-text fallback) for a punctuated reply
- rig: model.supports_vision so the heic control PNG reaches the model as an image
- Slack stream_finalize_rejected(_group) KNOWN #95430 (partial stream left beside re-post)
This commit is contained in:
teknium1
2026-09-24 05:24:59 -07:00
committed by Teknium
parent c5b8469ea2
commit e9330b2406
14 changed files with 462 additions and 135 deletions

View File

@@ -6,25 +6,36 @@ request Hermes sent the model, or persisted state (state.db). Tokens ``[in:<id>]
to the model turns it started, so duplicate or dropped turns are counted, never guessed.
Negative claims ("no second reply", "no turn for the unmentioned message") are settled by a BARRIER:
a later inbound in the same chat whose reply proves the gateway already processed everything
queued before it — never by sleeping.
a later inbound whose reply proves the gateway already processed everything queued before it — never
by sleeping. Streaming scenarios run each in a private chat of their own (``Rig.fresh_user``); their
same-chat barrier is re-sent if the first one is dropped while the adapter closes the streamed turn
(#121393: no platform-visible signal marks that window).
Open bugs are ``Rig.gate()`` blocks around a scenario's FINAL assertions only (see
``tests/e2e/core/_pending_fixes.known_failure``): the cell xfails only on that bug's own message,
fails on anything else, and simply passes once the fix lands.
"""
from __future__ import annotations
import contextlib
import html
import io
import itertools
import re
import time
from dataclasses import dataclass
from dataclasses import dataclass, field
from pathlib import Path
from typing import Any, Callable, Dict, List
from typing import Any, Callable, ContextManager, Dict, Iterator, List, Optional, Tuple
from tests.e2e.core._pending_fixes import known_failure
from tests.e2e.core.platforms._helpers import Director, GatewayUnderTest
from tests.fakes.fake_llm_provider import FakeLLMServer, Text, ToolCall
from tests.fakes.platforms._standin import Visible, wait_until
TURN_TIMEOUT = 90.0
# A turn takes ~3 s even under 8-way parallel load; the budget only bounds how long one red scenario
# can hold the file (the CI runner kills a file at 300 s, a module-scoped rig serves ~12 scenarios).
TURN_TIMEOUT = 30.0
@dataclass
@@ -33,10 +44,22 @@ class Rig:
drv: Any
director: Director
llm: FakeLLMServer
known: Optional[Tuple[str, str]] = None # (pattern, reason) for the running scenario's open bug
_users: Optional[Iterator[str]] = field(default=None, repr=False)
def ctx(self) -> str:
return f"{self.drv.describe()}\n{self.gw.tail(4000)}"
def gate(self) -> ContextManager[None]:
"""Wrap a scenario's final assertions: xfail only on the running scenario's KNOWN bug signature."""
return known_failure(*self.known) if self.known else contextlib.nullcontext()
def fresh_user(self) -> str:
"""An allowlisted user whose private chat no earlier scenario used (``drv.stream_users``)."""
if self._users is None:
self._users = itertools.cycle(self.drv.stream_users)
return next(self._users)
def head(aid: str) -> str:
return f"<<{aid}>>"
@@ -69,15 +92,30 @@ def wait_reply(rig: Rig, chat_id: str, aid: str, what: str = "", timeout: float
what or f"reply {aid} visible in {chat_id}", timeout=timeout, on_timeout=rig.ctx)
def barrier(rig: Rig, token: str, *, group: bool = False) -> None:
"""One more inbound in the same chat, answered: everything queued before it has been handled."""
aid = f"B-{token}"
# A reply can be visible while its turn is still closing; wait for the turn boundary first so
# the barrier is a new turn (an inbound in that closing window is tracked as #121393).
def barrier(rig: Rig, token: str, *, group: bool = False, user_id: Optional[str] = None, attempts: int = 1) -> None:
"""One more inbound in the same chat, answered: everything queued before it has been handled.
The gateway serializes a chat: an inbound arriving while a turn runs is drained only after that
turn's final delivery, so the barrier's reply orders after it. ``user_id`` picks that user's
private chat. ``attempts > 1`` is for streamed turns: a same-chat inbound landing while the
adapter closes such a turn can be dropped (#121393; no platform-visible signal marks that
window), so an unanswered barrier is re-sent under a new token after a short wait.
"""
rig.gw.wait_idle()
rig.director.script(token, answer(aid, "barrier"))
inbound = rig.drv.group(f"barrier [in:{token}]", mention=True) if group else rig.drv.dm(f"barrier [in:{token}]")
wait_reply(rig, inbound.chat_id, aid, f"barrier {token} answered")
for n in range(attempts):
tok = token if n == 0 else f"{token}r{n}"
aid = f"B-{tok}"
rig.director.script(tok, answer(aid, "barrier"))
text = f"barrier [in:{tok}]"
inbound = rig.drv.group(text, mention=True) if group else rig.drv.dm(text, user_id=user_id)
if n == attempts - 1:
wait_reply(rig, inbound.chat_id, aid, f"barrier {tok} answered")
return
try:
wait_until(lambda: complete(rig.drv.visible(inbound.chat_id), aid), f"barrier {tok}", timeout=8.0)
return
except AssertionError:
continue # dropped in the #121393 window: the next barrier still orders after the turn
def turns(rig: Rig, token: str) -> int:
@@ -87,14 +125,20 @@ def turns(rig: Rig, token: str) -> int:
# 1. inbound DM -> exactly one reply --------------------------------------------------------------
def dm_gets_exactly_one_reply(rig: Rig, tag: str) -> None:
token, aid = f"dm-{tag}", f"A-dm-{tag}"
rig.director.script(token, answer(aid, "hello back"))
# punctuation every markup dialect reserves: the adapter must escape it, not lean on a fallback
rig.director.script(token, answer(aid, "hello back (v1.2) - ok! #1 = {x}"))
inbound = rig.drv.dm(f"hello [in:{token}]")
wait_reply(rig, inbound.chat_id, aid)
barrier(rig, f"dmb-{tag}")
vis = copies(rig.drv.visible(inbound.chat_id), aid)
assert len(vis) == 1 and complete(vis, aid), f"expected one complete reply, got {vis}\n{rig.ctx()}"
assert turns(rig, token) == 1, f"inbound started {turns(rig, token)} model turns\n{rig.ctx()}"
assert len(rig.gw.user_rows(f"[in:{token}]")) == 1, "inbound persisted != once"
with rig.gate():
assert len(vis) == 1 and complete(vis, aid), f"expected one complete reply, got {vis}\n{rig.ctx()}"
assert "hello back (v1.2) - ok! #1 = {x}" in norm(vis[0].text), f"reply text mangled: {vis[0].text!r}"
assert turns(rig, token) == 1, f"inbound started {turns(rig, token)} model turns\n{rig.ctx()}"
assert len(rig.gw.user_rows(f"[in:{token}]")) == 1, "inbound persisted != once"
refused = rig.drv.format_rejections(inbound.chat_id)
assert not refused, (f"the platform refused the reply's formatting (a plain-text fallback hid it): "
f"{[c.response for c in refused]}\n{rig.ctx()}")
# 2. group without mention obeys require_mention --------------------------------------------------
@@ -168,43 +212,80 @@ def failed_continuation_is_retried_or_reported(rig: Rig, tag: str) -> None:
# 5. stream finalize rejected -> exactly one complete copy (#121108 pattern) -----------------------
def rejected_finalize_leaves_one_copy(rig: Rig, tag: str) -> None:
def rejected_finalize_leaves_one_copy(rig: Rig, tag: str, *, group: bool = False) -> None:
"""The platform refuses the call that would complete a streamed reply (``drv.fail_finalize``:
the finalize edit, the final send after draft previews, or the closing stream append)."""
token, aid = f"fin-{tag}", f"A-fin-{tag}"
body = " ".join(f"q{i:03d}" for i in range(120))
rig.director.script(token, answer(aid, body, chunk_chars=12, delay_per_chunk=0.03))
rig.drv.fail_edit(times=50, match=lambda text: foot(aid) in norm(text))
inbound = rig.drv.dm(f"stream it [in:{token}]")
faults = rig.drv.fail_finalize(lambda text: foot(aid) in norm(text), group=group)
text, user = f"stream it [in:{token}]", rig.fresh_user()
inbound = rig.drv.group(text, mention=True) if group else rig.drv.dm(text, user_id=user)
wait_reply(rig, inbound.chat_id, aid)
barrier(rig, f"finb-{tag}")
barrier(rig, f"finb-{tag}", group=group, user_id=user, attempts=3)
rig.drv.standin.clear_faults()
assert any(f.fired for f in faults), f"the finalize fault never fired: nothing was rejected\n{rig.ctx()}"
shown = _shown(rig, inbound.chat_id, aid, "q")
ctx = f"visible: {shown[:600]!r}\n{rig.ctx()}"
assert shown.count(head(aid)) == 1 and shown.count(foot(aid)) == 1, f"answer shown != once after a rejected finalize\n{ctx}"
assert _pwords(shown, "q") == _pwords(body, "q"), f"answer words lost/duplicated after a rejected finalize\n{ctx}"
ctx = f"visible: {shown[:300]!r} ... {shown[-300:]!r}\n{rig.ctx()}"
# the ending is shown exactly once: neither lost nor re-sent by a second final delivery
assert shown.count(foot(aid)) == 1, (f"answer ending shown {shown.count(foot(aid))}x after a rejected "
f"finalize (lost or re-sent)\n{ctx}")
assert turns(rig, token) == 1, rig.ctx()
with rig.gate():
# ... and so is everything before it, in order: no stale partial copy left beside it
assert shown.count(head(aid)) == 1 and _seamless(_pwords(shown, "q")) == _pwords(body, "q"), (
f"a partial copy of the answer was left visible after a rejected finalize "
f"(head shown {shown.count(head(aid))}x)\n{ctx}")
def _pwords(text: str, prefix: str) -> List[str]:
return re.findall(rf"\b{prefix}\d{{3}}\b", text)
def _seamless(words: List[str]) -> List[str]:
"""Collapse a word repeated across a message seam: the edit-fallback continuation backs its cut up
to the previous word boundary (``_continuation_text``), so the stuck preview's last word can
reappear at the start of the continuation. Tolerated here: it is not a second copy."""
return [w for i, w in enumerate(words) if i == 0 or w != words[i - 1]]
def _shown(rig: Rig, chat_id: str, aid: str, prefix: str) -> str:
"""Everything visible that carries a piece of THIS answer (the chat is shared across scenarios)."""
return " ".join(norm(v.text) for v in rig.drv.visible(chat_id)
if copies([v], aid) or _pwords(norm(v.text), prefix))
def streamed_reply_shown_once(rig: Rig, tag: str) -> None:
"""The plain streamed turn: previews then a final, and the user ends up with ONE copy (the
stream consumer's final-delivered bookkeeping is what stops the gateway's normal final send)."""
token, aid = f"st-{tag}", f"A-st-{tag}"
body = " ".join(f"s{i:03d}" for i in range(80))
rig.director.script(token, answer(aid, body, chunk_chars=12, delay_per_chunk=0.03))
user = rig.fresh_user()
inbound = rig.drv.dm(f"stream it [in:{token}]", user_id=user)
wait_reply(rig, inbound.chat_id, aid)
barrier(rig, f"stb-{tag}", user_id=user, attempts=3)
shown = _shown(rig, inbound.chat_id, aid, "s")
with rig.gate():
assert shown.count(head(aid)) == 1 and shown.count(foot(aid)) == 1, (
f"a streamed reply is shown != once: {shown[:300]!r}\n{rig.ctx()}")
assert _pwords(shown, "s") == _pwords(body, "s"), f"streamed words lost/duplicated: {shown[:300]!r}"
assert turns(rig, token) == 1, rig.ctx()
def streamed_reply_ending_in_whitespace_shown_once(rig: Rig, tag: str) -> None:
"""Models routinely end on a newline; the streamed transport must still finalize ONE message."""
token, aid = f"ws-{tag}", f"A-ws-{tag}"
body = " ".join(f"w{i:03d}" for i in range(80))
rig.director.script(token, Text(f"{head(aid)} {body} {foot(aid)}\n\n", chunk_chars=12, delay_per_chunk=0.03))
inbound = rig.drv.dm(f"stream with a trailing newline [in:{token}]")
user = rig.fresh_user()
inbound = rig.drv.dm(f"stream with a trailing newline [in:{token}]", user_id=user)
wait_reply(rig, inbound.chat_id, aid)
barrier(rig, f"wsb-{tag}")
barrier(rig, f"wsb-{tag}", user_id=user, attempts=3)
shown = _shown(rig, inbound.chat_id, aid, "w")
assert shown.count(head(aid)) == 1 and shown.count(foot(aid)) == 1, (
f"a streamed reply ending in whitespace is shown != once: {shown[:400]!r}\n{rig.ctx()}")
with rig.gate():
assert shown.count(head(aid)) == 1 and shown.count(foot(aid)) == 1, (
f"a streamed reply ending in whitespace is shown != once: {shown[:400]!r}\n{rig.ctx()}")
# 6. the platform redelivers the same inbound -> one reply (#119848 pattern) -----------------------
@@ -215,35 +296,47 @@ def redelivered_inbound_gets_one_reply(rig: Rig, tag: str) -> None:
wait_reply(rig, inbound.chat_id, aid)
rig.drv.redeliver(inbound)
barrier(rig, f"reb-{tag}")
assert turns(rig, token) == 1, f"redelivered inbound started {turns(rig, token)} turns\n{rig.ctx()}"
assert not copies(rig.drv.visible(inbound.chat_id), f"{aid}-dup"), "a duplicate reply was sent"
assert len(copies(rig.drv.visible(inbound.chat_id), aid)) == 1, rig.ctx()
with rig.gate():
assert turns(rig, token) == 1, f"redelivered inbound started {turns(rig, token)} turns\n{rig.ctx()}"
assert not copies(rig.drv.visible(inbound.chat_id), f"{aid}-dup"), "a duplicate reply was sent"
assert len(copies(rig.drv.visible(inbound.chat_id), aid)) == 1, rig.ctx()
# A replay that slips past dedup can ride INTO the barrier's turn (text batching merges the
# two), which the turn/reply counts above cannot see: the transcript can.
rows = rig.gw.user_rows(f"[in:{token}]")
assert len(rows) == 1, f"the redelivered inbound reached the transcript {len(rows)}x: {rows}\n{rig.ctx()}"
# 7. approval button click: allowlisted user accepted, stranger refused ----------------------------
def _approve_button(rig: Rig, chat_id: str) -> Dict[str, Any]:
def _button(rig: Rig, chat_id: str, label_re: str) -> Dict[str, Any]:
def pick() -> Any:
for b in rig.drv.buttons(chat_id):
label = str(b.get("text") or b.get("label") or "")
if re.search(r"once", label, re.I):
label = str(label.get("text", "")) if isinstance(label, dict) else label
if re.search(label_re, label, re.I):
return b
return None
return wait_until(pick, "an approve-once button under the approval prompt", timeout=TURN_TIMEOUT,
return wait_until(pick, f"a {label_re!r} button under the approval prompt", timeout=TURN_TIMEOUT,
on_timeout=rig.ctx)
def approval_click_by_allowlisted_user_runs_command(rig: Rig, tag: str, victim: Path) -> None:
"""A stranger's DENY must be refused (else it would deny the command); the owner's approve runs it."""
token, aid = f"ap-{tag}", f"A-ap-{tag}"
victim.mkdir(parents=True, exist_ok=True)
(victim / "f.txt").write_text("x")
rig.director.script(token, ToolCall("terminal", {"command": f"rm -rf {victim}"}), answer(aid, "cleaned"))
inbound = rig.drv.dm(f"clean up [in:{token}]")
button = _approve_button(rig, inbound.chat_id)
approve = _button(rig, inbound.chat_id, r"once")
deny = _button(rig, inbound.chat_id, r"deny")
assert victim.exists(), "dangerous command ran before approval"
rig.drv.click(inbound.chat_id, button, user_id=rig.drv.other_user_id) # a stranger first
rig.drv.click(inbound.chat_id, button)
stranger = rig.drv.click(inbound.chat_id, deny, user_id=rig.drv.other_user_id)
# the adapter has handled (and answered) the stranger's click before the owner clicks
wait_until(lambda: rig.drv.click_answered(stranger), "the stranger's click answered", timeout=TURN_TIMEOUT,
on_timeout=rig.ctx)
rig.drv.click(inbound.chat_id, approve)
wait_reply(rig, inbound.chat_id, aid, "turn completed after the approval click")
assert not victim.exists(), f"approved command did not run (click refused?)\n{rig.ctx()}"
assert not victim.exists(), (f"approved command did not run: the stranger's deny was obeyed or the "
f"owner's approval refused\n{rig.ctx()}")
# 8. platform_toolsets / disabled_toolsets honored (#121089) --------------------------------------
@@ -310,14 +403,15 @@ def heic_document_reaches_agent_as_image(rig: Rig, tag: str) -> None:
rig.director.script(token, answer(aid, "saw it"))
inbound = rig.drv.document(name, _image(fmt), mime, caption=f"look [in:{token}]")
try:
wait_reply(rig, inbound.chat_id, aid, timeout=30)
wait_reply(rig, inbound.chat_id, aid)
except AssertionError:
got[fmt] = False # the photo never started a turn at all
continue
got[fmt] = _model_saw_image(rig.director.requests[token][0])
log = rig.gw.grep(r"image|attach|heic|png|cache")
assert got["PNG"], f"control: a PNG sent as a file did not reach the model as an image\n{log}\n{rig.ctx()}"
assert got["HEIF"], f"a HEIC photo sent as a file did not reach the model as an image\n{log}"
with rig.gate():
assert got["HEIF"], f"a HEIC photo sent as a file did not reach the model as an image\n{log}"
# 10. planned restart: notice delivered once; a redelivered /restart does not loop ----------------
@@ -342,6 +436,11 @@ def planned_restart_notice_once(rig: Rig, tag: str, restart: Callable[[], None])
probe = rig.drv.dm(f"barrier [in:rsb-{tag}]")
wait_until(lambda: complete(rig.drv.visible(probe.chat_id), bid) or not rig.gw.alive(),
"barrier answered (or the gateway went down again)", timeout=TURN_TIMEOUT, on_timeout=rig.ctx)
assert rig.gw.alive(), f"a redelivered /restart restarted the gateway again\n{rig.ctx()}"
assert len(notices()) == 1, (f"expected exactly one restart notice from the new process, got "
f"{notices()} (a second restart ack means the replayed /restart was obeyed)\n{rig.ctx()}")
rc = rig.gw.proc.returncode if rig.gw.proc else None
# a crash is not the replay bug: only a second planned exit (75) is
assert rig.gw.alive() or rc == 75, f"gateway died (rc={rc}) after the replayed /restart\n{rig.ctx()}"
with rig.gate():
assert rig.gw.alive(), f"a redelivered /restart restarted the gateway again (exit 75)\n{rig.ctx()}"
assert len(notices()) == 1, (f"expected exactly one restart notice from the new process, got "
f"{notices()} (a second restart ack means the replayed /restart was obeyed)"
f"\n{rig.ctx()}")

View File

@@ -16,7 +16,7 @@ from dataclasses import dataclass
from pathlib import Path
from typing import Any, Callable, Dict, List, Optional
from tests.fakes.platforms._standin import Call, Visible
from tests.fakes.platforms._standin import Call, Fault, Visible
from tests.fakes.platforms.discord_standin import BOT_ID, MAX_TEXT, DiscordStandin
SHIM_DIR = Path(__file__).resolve().parents[3] / "fakes" / "platforms" / "discord_shim"
@@ -35,6 +35,7 @@ class DiscordDriver:
bot_id = BOT_ID
user_id = "200000000000000111"
other_user_id = "200000000000000222"
stream_users = tuple(f"2000000000000003{i:02d}" for i in range(1, 7))
group_chat_id = "1300000000000000001" # #general in the stand-in guild
home_channel = "1300000000000000002" # #home in the stand-in guild
@@ -42,7 +43,7 @@ class DiscordDriver:
self.standin = DiscordStandin()
self.standin.add_text_channel(self.group_chat_id, "general")
self.standin.add_text_channel(self.home_channel, "home")
for uid in (self.user_id, self.other_user_id):
for uid in (self.user_id, self.other_user_id, *self.stream_users):
self.standin.user(uid)["_in_guild"] = True
def start(self) -> None:
@@ -63,7 +64,7 @@ class DiscordDriver:
def gateway_env(self) -> Dict[str, str]:
return {
"DISCORD_BOT_TOKEN": self.standin.token, "DISCORD_ALLOWED_USERS": self.user_id,
"DISCORD_BOT_TOKEN": self.standin.token, "DISCORD_ALLOWED_USERS": ",".join((self.user_id, *self.stream_users)),
"DISCORD_HOME_CHANNEL": self.home_channel,
# one PUT instead of ~100 paced POSTs (safe policy sleeps 4.5 s per mutation)
"DISCORD_COMMAND_SYNC_POLICY": "bulk",
@@ -105,9 +106,13 @@ class DiscordDriver:
def buttons(self, chat_id: str) -> List[Dict[str, Any]]:
return self.standin.buttons(self._chat(chat_id))
def click(self, chat_id: str, button: Dict[str, Any], user_id: Optional[str] = None) -> None:
self.standin.click(str(user_id or self.user_id), self._chat(chat_id), str(button["message_id"]),
button["custom_id"], component_type=int(button.get("type") or 2))
def click(self, chat_id: str, button: Dict[str, Any], user_id: Optional[str] = None) -> str:
inter = self.standin.click(str(user_id or self.user_id), self._chat(chat_id), str(button["message_id"]),
button["custom_id"], component_type=int(button.get("type") or 2))
return str(inter["id"])
def click_answered(self, handle: str) -> bool:
return any(str(c.params.get("interaction_id")) == handle and not c.faulted for c in self.callback_answers())
def callback_answers(self) -> List[Call]:
return self.standin.calls_of("interaction_callback")
@@ -127,16 +132,24 @@ class DiscordDriver:
def edits(self, chat_id: str) -> List[Call]:
return self._ok(("edit_message",), chat_id)
def format_rejections(self, chat_id: str) -> List[Call]:
return [] # Discord renders markdown client-side: nothing to reject for formatting
def describe(self) -> str:
return self.standin.describe()
# faults ------------------------------------------------------------------------------------
def fail_send(self, *, times: int = 1, match: Optional[Callable[[str], bool]] = None) -> None:
def fail_send(self, *, times: int = 1, match: Optional[Callable[[str], bool]] = None) -> List[Fault]:
pred = (lambda p: match(str(p.get("content", "")))) if match else None
self.standin.fail("create_message", {"message": "Missing Access", "code": 50001},
status=403, times=times, match=pred)
return [self.standin.fail("create_message", {"message": "Missing Access", "code": 50001},
status=403, times=times, match=pred)]
def fail_edit(self, *, times: int = 1, match: Optional[Callable[[str], bool]] = None) -> None:
def fail_edit(self, *, times: int = 1, match: Optional[Callable[[str], bool]] = None) -> List[Fault]:
pred = (lambda p: match(str(p.get("content", "")))) if match else None
self.standin.fail("edit_message", {"message": "Unknown Message", "code": 10008}, status=404, times=times,
match=pred)
return [self.standin.fail("edit_message", {"message": "Unknown Message", "code": 10008}, status=404,
times=times, match=pred)]
def fail_finalize(self, has_footer: Callable[[str], bool], *, group: bool) -> List[Fault]:
"""Streaming edits one message in place (DMs and channels alike): every edit carrying the footer
is refused, so the finalize edit never lands."""
return self.fail_edit(times=50, match=has_footer)

View File

@@ -20,7 +20,7 @@ from dataclasses import dataclass
from pathlib import Path
from typing import Any, Callable, Dict, List, Optional
from tests.fakes.platforms._standin import Call, Visible
from tests.fakes.platforms._standin import Call, Fault, Visible
from tests.fakes.platforms.slack_standin import SlackStandin
_SHIM = Path(__file__).resolve().parents[3] / "fakes" / "platforms" / "slack_shim"
@@ -38,6 +38,7 @@ class SlackDriver:
limit = 39_000 # SlackAdapter.MAX_MESSAGE_LENGTH (Slack rejects > 40,000 with msg_too_long)
user_id = "U0000111"
other_user_id = "U0000222"
stream_users = tuple(f"U0000{i}" for i in range(301, 307))
group_chat_id = "C0000123"
home_channel = "C0000999"
@@ -55,7 +56,7 @@ class SlackDriver:
def gateway_env(self) -> Dict[str, str]:
return {"SLACK_BOT_TOKEN": self.standin.bot_token, "SLACK_APP_TOKEN": self.standin.app_token,
"SLACK_ALLOWED_USERS": self.user_id, "SLACK_HOME_CHANNEL": self.home_channel,
"SLACK_ALLOWED_USERS": ",".join((self.user_id, *self.stream_users)), "SLACK_HOME_CHANNEL": self.home_channel,
"HERMES_STANDIN_SLACK_API": self.standin.api_base, "PYTHONPATH_PREPEND": str(_SHIM)}
def connected(self) -> bool:
@@ -84,8 +85,13 @@ class SlackDriver:
def buttons(self, chat_id: str) -> List[Dict[str, Any]]:
return self.standin.buttons(chat_id)
def click(self, chat_id: str, button: Dict[str, Any], user_id: Optional[str] = None) -> None:
self.standin.block_action(user_id or self.user_id, chat_id, str(button["message_id"]), button)
def click(self, chat_id: str, button: Dict[str, Any], user_id: Optional[str] = None) -> Dict[str, Any]:
return self.standin.block_action(user_id or self.user_id, chat_id, str(button["message_id"]), button)
def click_answered(self, handle: Dict[str, Any]) -> bool:
"""Bolt acks an interactive envelope first, then authorizes the clicker synchronously; a
refused click is only logged, so the ack is the last platform-visible sign of it."""
return self.standin.acked(handle)
def slash(self, command: str, text: str = "", chat_id: Optional[str] = None) -> Dict[str, Any]:
return self.standin.slash(self.user_id, chat_id or self.standin.dm_channel(self.user_id), command, text)
@@ -113,15 +119,27 @@ class SlackDriver:
def edits(self, chat_id: str) -> List[Call]:
return self._ok(("chat.update", "chat.appendStream"), chat_id)
def format_rejections(self, chat_id: str) -> List[Call]:
return [] # mrkdwn never fails to parse; Slack renders what it cannot format as text
def describe(self) -> str:
return self.standin.describe()
# faults ------------------------------------------------------------------------------------
def fail_send(self, *, times: int = 1, match: Optional[Callable[[str], bool]] = None) -> None:
def fail_send(self, *, times: int = 1, match: Optional[Callable[[str], bool]] = None) -> List[Fault]:
pred = (lambda p: match(str(p.get("text", "")))) if match else None
self.standin.fail("chat.postMessage", {"ok": False, "error": "channel_not_found"}, times=times, match=pred)
return [self.standin.fail("chat.postMessage", {"ok": False, "error": "channel_not_found"}, times=times,
match=pred)]
def fail_edit(self, *, times: int = 1, match: Optional[Callable[[str], bool]] = None) -> None:
pred = (lambda p: match(str(p.get("text", "")) + str(p.get("markdown_text", "")))) if match else None
for method in ("chat.update", "chat.appendStream", "chat.stopStream"):
self.standin.fail(method, {"ok": False, "error": "cant_update_message"}, times=times, match=pred)
def fail_edit(self, *, times: int = 1, match: Optional[Callable[[str], bool]] = None) -> List[Fault]:
pred = (lambda p: match(str(p.get("text", "")))) if match else None
return [self.standin.fail("chat.update", {"ok": False, "error": "cant_update_message"}, times=times,
match=pred)]
def fail_finalize(self, has_footer: Callable[[str], bool], *, group: bool) -> List[Fault]:
"""Native streaming (startStream/appendStream/stopStream) is the transport in DMs and channels:
refuse the stream call whose resulting text completes the reply (``_stream_text`` = what the
message would read after it), however the deltas were cut."""
pred = lambda p: has_footer(str(p.get("_stream_text", ""))) # noqa: E731
return [self.standin.fail(m, {"ok": False, "error": "message_not_in_streaming_state"}, times=50, match=pred)
for m in ("chat.appendStream", "chat.stopStream")]

View File

@@ -2,14 +2,17 @@
Every adapter driver exposes the same surface (see ``_contract.py``):
* ``name``, ``limit`` (the platform text cap the adapter must split at), ``user_id`` (allowlisted)
* ``name``, ``limit`` (the platform text cap the adapter must split at), ``user_id`` (allowlisted),
``stream_users`` (more allowlisted users: a private chat of its own per streaming scenario)
* ``start()/stop()`` the stand-in, ``gateway_config()``/``gateway_env()`` for the child, ``connected()``
* inbound: ``dm(text)``, ``group(text, mention=)``, ``redeliver(inbound)``, ``document(...)``,
``click(inbound_chat, button)``
``click(inbound_chat, button)`` -> a handle, ``click_answered(handle)`` once the adapter answered it
* ground truth: ``visible(chat_id)`` (bot messages a human sees now), ``sends(chat_id)`` /
``edits(chat_id)`` (successful outbound create/edit calls), ``describe()``
* faults: ``fail_send(match=, times=)`` / ``fail_edit(times=)`` make the platform reject the call
with its documented error body.
* faults: ``fail_send(match=, times=)`` / ``fail_edit(times=)`` / ``fail_finalize(footer, group=)``
make the platform reject the call with its documented error body; each returns its ``Fault``s so a
scenario can assert the fault actually fired. ``format_rejections(chat_id)``: calls the platform
refused for their formatting (Telegram: MarkdownV2 ``can't parse entities``).
"""
from __future__ import annotations
@@ -17,7 +20,7 @@ from __future__ import annotations
from dataclasses import dataclass
from typing import Any, Callable, Dict, List, Optional
from tests.fakes.platforms._standin import Call, Visible
from tests.fakes.platforms._standin import Call, Fault, Visible
from tests.fakes.platforms.telegram_standin import MAX_TEXT, TelegramStandin
@@ -33,6 +36,7 @@ class TelegramDriver:
limit = MAX_TEXT
user_id = "111"
other_user_id = "222"
stream_users = ("112", "113", "114", "115", "116", "117")
group_chat_id = "-1001234567890"
home_channel = "999"
@@ -52,7 +56,8 @@ class TelegramDriver:
}}}}
def gateway_env(self) -> Dict[str, str]:
return {"TELEGRAM_BOT_TOKEN": self.standin.token, "TELEGRAM_ALLOWED_USERS": self.user_id,
return {"TELEGRAM_BOT_TOKEN": self.standin.token,
"TELEGRAM_ALLOWED_USERS": ",".join((self.user_id, *self.stream_users)),
"TELEGRAM_HOME_CHANNEL": self.home_channel, "HERMES_TELEGRAM_DISABLE_FALLBACK_IPS": "1",
# no text-batch debounce: one inbound is one turn, immediately
"HERMES_TELEGRAM_TEXT_BATCH_DELAY_SECONDS": "0", "HERMES_TELEGRAM_TEXT_BATCH_SPLIT_DELAY_SECONDS": "0"}
@@ -80,9 +85,13 @@ class TelegramDriver:
def buttons(self, chat_id: str) -> List[Dict[str, Any]]:
return self.standin.buttons(chat_id)
def click(self, chat_id: str, button: Dict[str, Any], user_id: Optional[str] = None) -> None:
self.standin.callback(int(user_id or self.user_id), int(chat_id), int(button["message_id"]),
button["callback_data"])
def click(self, chat_id: str, button: Dict[str, Any], user_id: Optional[str] = None) -> str:
update = self.standin.callback(int(user_id or self.user_id), int(chat_id), int(button["message_id"]),
button["callback_data"])
return update["callback_query"]["id"]
def click_answered(self, handle: str) -> bool:
return any(str(c.params.get("callback_query_id")) == handle and not c.faulted for c in self.callback_answers())
def callback_answers(self) -> List[Call]:
return self.standin.calls_of("answerCallbackQuery")
@@ -101,18 +110,34 @@ class TelegramDriver:
def edits(self, chat_id: str) -> List[Call]:
return self._ok(("editMessageText",), chat_id)
def format_rejections(self, chat_id: str) -> List[Call]:
return [c for c in self.standin.calls if c.faulted and str(c.params.get("chat_id")) == str(chat_id)
and "can't parse entities" in str(c.response)]
def describe(self) -> str:
return self.standin.describe()
# faults ------------------------------------------------------------------------------------
def fail_send(self, *, times: int = 1, match: Optional[Callable[[str], bool]] = None) -> None:
def fail_send(self, *, times: int = 1, match: Optional[Callable[[str], bool]] = None) -> List[Fault]:
pred = (lambda p: match(str(p.get("text", "")))) if match else None
self.standin.fail("sendMessage", {"ok": False, "error_code": 400,
"description": "Bad Request: chat not found"},
status=400, times=times, match=pred)
return [self.standin.fail("sendMessage", {"ok": False, "error_code": 400,
"description": "Bad Request: chat not found"},
status=400, times=times, match=pred)]
def fail_edit(self, *, times: int = 1, match: Optional[Callable[[str], bool]] = None) -> None:
def fail_edit(self, *, times: int = 1, match: Optional[Callable[[str], bool]] = None) -> List[Fault]:
pred = (lambda p: match(str(p.get("text", "")))) if match else None
self.standin.fail("editMessageText", {"ok": False, "error_code": 400,
"description": "Bad Request: message can't be edited"},
status=400, times=times, match=pred)
return [self.standin.fail("editMessageText", {"ok": False, "error_code": 400,
"description": "Bad Request: message can't be edited"},
status=400, times=times, match=pred)]
def fail_finalize(self, has_footer: Callable[[str], bool], *, group: bool) -> List[Fault]:
"""Reject the call that would complete the streamed reply.
Groups stream by editing one message (editMessageText): every edit carrying the footer is
refused, as for a message that can no longer be edited. Private chats stream through
``sendMessageDraft`` previews and the final text is a fresh ``sendMessage``: its first attempt
is refused, so the retry/fallback path has to produce the one copy.
"""
if group:
return self.fail_edit(times=50, match=has_footer)
return self.fail_send(times=1, match=has_footer)

View File

@@ -125,7 +125,7 @@ class GatewayUnderTest:
def db_path(self) -> Path:
return self.hermes_home / "state.db"
def start(self, timeout: float = 120.0) -> "GatewayUnderTest":
def start(self, timeout: float = 60.0) -> "GatewayUnderTest":
assert self.proc is None or self.proc.poll() is not None
log = open(self.log_path, "a", encoding="utf-8") # noqa: SIM115 - handed to the child
self.proc = subprocess.Popen(
@@ -184,8 +184,13 @@ class GatewayUnderTest:
except (OSError, ValueError, TypeError):
return None
def wait_idle(self, timeout: float = 90.0) -> None:
"""Every turn finished (not merely its reply visible): the next inbound starts a fresh turn."""
def wait_idle(self, timeout: float = 30.0) -> None:
"""No agent run is in flight (``active_agents == 0`` in gateway_state.json).
NOT delivery or adapter quiescence: the count drops when the agent run returns, before the
adapter sends the normal final reply and closes the turn, and a same-chat inbound in that
window can be dropped (#121393). Ordering after a turn's delivery comes from a same-chat
barrier (``_contract.barrier``); this only keeps a new scenario from racing the last one."""
wait_until(lambda: self.active_agents() == 0 or not self.alive(), "gateway idle (active_agents == 0)",
timeout=timeout, on_timeout=self.tail)

View File

@@ -1,19 +1,20 @@
"""Binds the shared contract (``_contract.py``) to one adapter driver for a test module.
A test module declares ``KNOWN`` (scenario -> "#<issue> <symptom>") and gets:
A test module declares ``KNOWN`` (scenario -> (pattern, "#<issue> <symptom>")) and gets:
* two module-scoped rigs: ``rig`` (the adapter's default delivery config, ``agent.disabled_toolsets:
[file]``, supervisor-owned so ``/restart`` exits 75) and ``rig_stream`` (edit-streaming on,
``platform_toolsets.<platform>: [file]``);
* one parametrized ``test_contract`` over every scenario, KNOWN ones as ``xfail(strict=True)`` so
a fix turns them red until the entry is removed.
[file]``, supervisor-owned so ``/restart`` exits 75) and ``rig_stream`` (streaming on with the
platform's default transport, ``platform_toolsets.<platform>: [file]``);
* one parametrized ``test_contract`` over every scenario. A KNOWN scenario runs its final assertions
under ``known_failure(pattern, reason)`` (``Rig.gate``): it xfails only while it fails with that
bug's own message, fails on anything else, and passes once the fix lands, in either merge order.
"""
from __future__ import annotations
import os
from pathlib import Path
from typing import Any, Callable, Dict, List, Tuple
from typing import Any, Callable, Dict, List, Optional, Tuple
import pytest
@@ -32,7 +33,10 @@ SCENARIOS: Dict[str, Tuple[str, Callable[..., None]]] = {
r, t, d / "victim")),
"disabled_toolsets": ("rig", lambda r, t, d: C.disabled_toolsets_are_honored(r, t)),
"heic_as_image": ("rig", lambda r, t, d: C.heic_document_reaches_agent_as_image(r, t)),
"stream_reply_once": ("rig_stream", lambda r, t, d: C.streamed_reply_shown_once(r, t)),
"stream_finalize_rejected": ("rig_stream", lambda r, t, d: C.rejected_finalize_leaves_one_copy(r, t)),
"stream_finalize_rejected_group": ("rig_stream",
lambda r, t, d: C.rejected_finalize_leaves_one_copy(r, t, group=True)),
"stream_trailing_whitespace": ("rig_stream",
lambda r, t, d: C.streamed_reply_ending_in_whitespace_shown_once(r, t)),
"platform_toolsets": ("rig_stream", lambda r, t, d: C.platform_toolsets_are_honored(r, t)),
@@ -41,24 +45,22 @@ SCENARIOS: Dict[str, Tuple[str, Callable[..., None]]] = {
}
def scenario_params(known: Dict[str, str], skip: Dict[str, str] | None = None) -> List[Any]:
out = []
for name in SCENARIOS:
marks = []
if name in known:
marks.append(pytest.mark.xfail(strict=True, reason=known[name]))
if skip and name in skip:
marks.append(pytest.mark.skip(reason=skip[name]))
out.append(pytest.param(name, id=name, marks=marks))
return out
def scenario_params(skip: Optional[Dict[str, str]] = None) -> List[Any]:
return [pytest.param(name, id=name, marks=[pytest.mark.skip(reason=skip[name])] if skip and name in skip else [])
for name in SCENARIOS]
def run_scenario(name: str, request: pytest.FixtureRequest, tmp_path: Path) -> None:
def run_scenario(name: str, request: pytest.FixtureRequest, tmp_path: Path,
known: Optional[Dict[str, Tuple[str, str]]] = None) -> None:
fixture, runner = SCENARIOS[name]
rig = request.getfixturevalue(fixture)
assert rig.gw.alive(), f"gateway died before {name}\n{rig.gw.tail()}"
rig.gw.wait_idle() # the previous scenario's last turn fully closed (see barrier())
runner(rig, name.replace("_", ""), tmp_path)
rig.gw.wait_idle() # no agent run of the previous scenario still in flight
rig.known = (known or {}).get(name)
try:
runner(rig, name.replace("_", ""), tmp_path)
finally:
rig.known = None
class _RestartableRig(C.Rig):
@@ -106,7 +108,9 @@ def rig_fixtures(driver_cls: type) -> Tuple[Any, Any]:
@pytest.fixture(scope="module")
def rig(tmp_path_factory: pytest.TempPathFactory):
r = _make(tmp_path_factory, "a", {"agent": {"disabled_toolsets": ["file"]}},
# model.supports_vision: the fake model has no vision metadata, so photos would be routed
# through vision pre-analysis and no image part could ever reach the model request
r = _make(tmp_path_factory, "a", {"agent": {"disabled_toolsets": ["file"]}, "model": {"supports_vision": True}},
{"HERMES_GATEWAY_EXTERNAL_SUPERVISOR": "1"})
yield r
_teardown(r)

View File

@@ -3,8 +3,8 @@
The child is ``hermes gateway run`` on a throwaway HOME; the adapter's own SDK (discord.py via the ``discord_shim`` sitecustomize) talks to
``tests/fakes/platforms/discord_standin.py``, a local stand-in shaped per the platform's published
API. Scenarios live in ``_contract.py`` and are identical for every adapter; this file only binds the
Discord driver and lists the scenarios that are red on main (``KNOWN``, strict xfail: a fix turns
the entry red until it is removed).
Discord driver and lists the scenarios that are red on main (``KNOWN``: scenario -> (the bug's failure-message
pattern, reason); see ``_suite.py``: xfail only on that message, pass once the fix lands).
"""
from __future__ import annotations
@@ -21,8 +21,10 @@ pytestmark = [
pytest.mark.skipif(sys.platform == "win32", reason="POSIX process-group gateway harness"),
]
KNOWN: dict[str, str] = {
"planned_restart_notice": "#121325 a replayed /restart restarts the gateway again (guard needs Telegram update ids)",
KNOWN: dict[str, tuple[str, str]] = {
"planned_restart_notice": (
r"a redelivered /restart restarted the gateway again|a second restart ack means the replayed /restart was obeyed",
"#121325 a replayed /restart restarts the gateway again (guard needs Telegram update ids)"),
}
SKIP: dict[str, str] = {
"heic_as_image": "image attachments are fetched by URL behind the SSRF guard, which refuses the loopback "
@@ -32,6 +34,6 @@ SKIP: dict[str, str] = {
rig, rig_stream = rig_fixtures(DiscordDriver)
@pytest.mark.parametrize("scenario", scenario_params(KNOWN, SKIP))
@pytest.mark.parametrize("scenario", scenario_params(SKIP))
def test_contract(scenario: str, request: pytest.FixtureRequest, tmp_path) -> None:
run_scenario(scenario, request, tmp_path)
run_scenario(scenario, request, tmp_path, KNOWN)

View File

@@ -3,8 +3,8 @@
The child is ``hermes gateway run`` on a throwaway HOME; the adapter's own SDK (slack_bolt Socket Mode + slack_sdk via the ``slack_shim`` sitecustomize) talks to
``tests/fakes/platforms/slack_standin.py``, a local stand-in shaped per the platform's published
API. Scenarios live in ``_contract.py`` and are identical for every adapter; this file only binds the
Slack driver and lists the scenarios that are red on main (``KNOWN``, strict xfail: a fix turns
the entry red until it is removed).
Slack driver and lists the scenarios that are red on main (``KNOWN``: scenario -> (the bug's failure-message
pattern, reason); see ``_suite.py``: xfail only on that message, pass once the fix lands).
"""
from __future__ import annotations
@@ -21,9 +21,17 @@ pytestmark = [
pytest.mark.skipif(sys.platform == "win32", reason="POSIX process-group gateway harness"),
]
KNOWN: dict[str, str] = {
"stream_trailing_whitespace": "#121326 native streaming re-posts the whole reply when it ends in whitespace",
"planned_restart_notice": "#121325 a replayed /restart restarts the gateway again (guard needs Telegram update ids)",
_PARTIAL = (r"a partial copy of the answer was left visible after a rejected finalize",
"#95430 a rejected closing appendStream/stopStream leaves the partial stream next to the re-posted answer")
KNOWN: dict[str, tuple[str, str]] = {
"stream_finalize_rejected": _PARTIAL,
"stream_finalize_rejected_group": _PARTIAL,
"stream_trailing_whitespace": (
r"a streamed reply ending in whitespace is shown != once",
"#121326 native streaming re-posts the whole reply when it ends in whitespace"),
"planned_restart_notice": (
r"a redelivered /restart restarted the gateway again|a second restart ack means the replayed /restart was obeyed",
"#121325 a replayed /restart restarts the gateway again (guard needs Telegram update ids)"),
}
SKIP: dict[str, str] = {
"heic_as_image": "the adapter only downloads https://*.slack.com file URLs (SSRF guard); a loopback "
@@ -33,6 +41,6 @@ SKIP: dict[str, str] = {
rig, rig_stream = rig_fixtures(SlackDriver)
@pytest.mark.parametrize("scenario", scenario_params(KNOWN, SKIP))
@pytest.mark.parametrize("scenario", scenario_params(SKIP))
def test_contract(scenario: str, request: pytest.FixtureRequest, tmp_path) -> None:
run_scenario(scenario, request, tmp_path)
run_scenario(scenario, request, tmp_path, KNOWN)

View File

@@ -3,8 +3,8 @@
The child is ``hermes gateway run`` on a throwaway HOME; the adapter's own SDK (python-telegram-bot via ``extra.base_url``) talks to
``tests/fakes/platforms/telegram_standin.py``, a local stand-in shaped per the platform's published
API. Scenarios live in ``_contract.py`` and are identical for every adapter; this file only binds the
Telegram driver and lists the scenarios that are red on main (``KNOWN``, strict xfail: a fix turns
the entry red until it is removed).
Telegram driver and lists the scenarios that are red on main (``KNOWN``: scenario -> (the bug's failure-message
pattern, reason); see ``_suite.py``: xfail only on that message, pass once the fix lands).
"""
from __future__ import annotations
@@ -21,14 +21,16 @@ pytestmark = [
pytest.mark.skipif(sys.platform == "win32", reason="POSIX process-group gateway harness"),
]
KNOWN: dict[str, str] = {
"heic_as_image": "#119593 HEIC photo sent as a file is refused by the image cache (never reaches the model)",
KNOWN: dict[str, tuple[str, str]] = {
"heic_as_image": (
r"a HEIC photo sent as a file did not reach the model as an image",
"#119593 HEIC photo sent as a file is refused by the image cache (never reaches the model)"),
}
SKIP: dict[str, str] = {}
rig, rig_stream = rig_fixtures(TelegramDriver)
@pytest.mark.parametrize("scenario", scenario_params(KNOWN, SKIP))
@pytest.mark.parametrize("scenario", scenario_params(SKIP))
def test_contract(scenario: str, request: pytest.FixtureRequest, tmp_path) -> None:
run_scenario(scenario, request, tmp_path)
run_scenario(scenario, request, tmp_path, KNOWN)

View File

@@ -6,7 +6,8 @@ discord.py has no base-URL knob: REST URLs are built from the class attribute
``GET /gateway/bot``; the resume URL comes from READY's ``resume_gateway_url``, which the stand-in
also points at itself). This module is put on a child process's ``PYTHONPATH`` by the Discord driver
and does nothing unless ``HERMES_STANDIN_DISCORD_API`` is set. It patches the SDK only; no Hermes
code is touched.
code is touched. Side effect: when set, discord.py (and yarl/aiohttp) are imported eagerly at
interpreter startup, before Hermes runs, rather than lazily by the adapter.
"""
import os

View File

@@ -238,7 +238,7 @@ class DiscordStandin(StandinServer):
return _json(body, status=401)
if name == "unknown":
body = {"message": "404: Not Found", "code": 0}
self.record(name, {**params, "_method": request.method, "_path": path}, body)
self.record(name, {**params, "_method": request.method, "_path": path}, body, faulted=True)
return _json(body, status=404)
fault = self.take_fault(name, params)
if fault is not None:
@@ -252,7 +252,7 @@ class DiscordStandin(StandinServer):
body = {"message": f"stand-in handler error: {exc!r}", "code": 0}
self.record(name, {**params, "_error": repr(exc)}, body, faulted=True)
return _json(body, status=500)
self.record(name, params, result)
self.record(name, params, result, faulted=status >= 400)
if status == 204:
return web.Response(status=204)
return _json(result, status=status)
@@ -344,6 +344,9 @@ class DiscordStandin(StandinServer):
if len(p.get("content") or "") > MAX_TEXT:
return 400, {"message": "Invalid Form Body", "code": 50035, "errors": {"content": {"_errors": [
{"code": "BASE_TYPE_MAX_LENGTH", "message": "Must be 2000 or fewer in length."}]}}}
if not str(p.get("content") or "").strip() and not any(
p.get(k) for k in ("embeds", "components", "_files", "sticker_ids", "poll", "attachments")):
return 400, {"message": "Cannot send an empty message", "code": 50006}
return 200, self._bot_message(cid, p)
def _edit(self, cid: str, mid: str, p: Dict[str, Any]) -> Tuple[int, Any]:

View File

@@ -6,6 +6,9 @@ when ``HERMES_STANDIN_SLACK_API`` is set (e.g. ``http://127.0.0.1:PORT/api/``) e
the Slack adapter builds defaults its ``base_url`` to the stand-in, which also moves
``apps.connections.open`` and therefore the Socket Mode websocket. Nothing in Hermes is touched: the
redirect lives at the SDK/HTTP boundary, exactly where DNS would send the real traffic.
Side effect: when the variable is set, slack_sdk (and aiohttp) are imported eagerly at interpreter
startup, before Hermes runs, so the child pays that import up front even if Slack never connects.
"""
import os

View File

@@ -10,8 +10,12 @@ Socket Mode client connects, receives ``hello``, and from then on the test pushe
Web API semantics kept faithful to Slack: every response is HTTP 200 JSON with ``ok``; failures are
``{"ok": false, "error": "<code>"}`` (the SDK raises ``SlackApiError`` on them); bodies arrive
form-encoded or JSON (structured fields like ``blocks`` may be JSON strings); ``chat.postMessage``
above 40,000 chars answers ``msg_too_long``. Unknown methods answer ``{"ok": true}`` and are still
recorded so a test can see them.
above 40,000 chars answers ``msg_too_long``. Unknown methods answer ``{"ok": false, "error":
"unknown_method"}`` like Slack (recorded as faulted), so an unmodelled call is never a silent success.
Native streams: a fault's ``match`` for ``chat.appendStream``/``chat.stopStream`` also sees
``_stream_text`` (the message text as it would read after this call), so a test can reject the call
that completes a given piece of the reply however the deltas were cut.
"""
from __future__ import annotations
@@ -126,13 +130,17 @@ class SlackStandin(StandinServer):
body = {"ok": False, "error": "invalid_auth"}
self.record(method, params, body, faulted=True)
return self._reply(body)
fault = self.take_fault(method, params)
fault = self.take_fault(method, self._fault_view(method, params))
if fault is not None:
self.record(method, params, fault.body, faulted=True)
return web.json_response(fault.body, status=fault.status)
handler = getattr(self, "_m_" + method.replace(".", "_"), None)
if handler is None:
body = {"ok": False, "error": "unknown_method"}
self.record(method, params, body, faulted=True)
return self._reply(body)
try:
result = handler(params) if handler else {}
result = handler(params)
except _ApiError as exc:
body = {"ok": False, "error": exc.code}
self.record(method, params, body, faulted=True)
@@ -141,6 +149,14 @@ class SlackStandin(StandinServer):
self.record(method, params, body)
return self._reply(body)
def _fault_view(self, method: str, params: Dict[str, Any]) -> Dict[str, Any]:
if method not in ("chat.appendStream", "chat.stopStream"):
return params
with self._lock:
vis = self._visible.get((str(params.get("channel", "")), str(params.get("ts", ""))))
before = vis.text if vis is not None else ""
return {**params, "_stream_text": before + str(params.get("markdown_text") or "")}
async def _upload(self, request: web.Request) -> web.Response:
fid = request.match_info["file_id"]
data = await request.read()
@@ -270,6 +286,11 @@ class SlackStandin(StandinServer):
def _m_conversations_info(self, p: Dict[str, Any]) -> Dict[str, Any]:
return {"channel": self._channel(str(p.get("channel", "")))}
def _m_users_conversations(self, _p: Dict[str, Any]) -> Dict[str, Any]:
with self._lock:
ids = sorted(c for c in self._history if not c.startswith("D"))
return {"channels": [self._channel(c) for c in ids], "response_metadata": {"next_cursor": ""}}
def _m_conversations_open(self, p: Dict[str, Any]) -> Dict[str, Any]:
users = str(p.get("users", "")).split(",")[0]
return {"channel": self._channel("D" + users[1:])}
@@ -381,6 +402,13 @@ class SlackStandin(StandinServer):
def _m_chat_stopStream(self, p: Dict[str, Any]) -> Dict[str, Any]:
return self._append(p, stop=True)
def _m_ok(self, _p: Dict[str, Any]) -> Dict[str, Any]:
return {}
# Side-effect-only methods the adapter calls (reactions, assistant thread status/title/prompts).
_m_reactions_add = _m_reactions_remove = _m_ok
_m_assistant_threads_setStatus = _m_assistant_threads_setTitle = _m_assistant_threads_setSuggestedPrompts = _m_ok
def _m_files_getUploadURLExternal(self, p: Dict[str, Any]) -> Dict[str, Any]:
fid = f"F{next(self._ids):08d}"
self.files[fid] = {"id": fid, "name": p.get("filename"), "title": p.get("filename"),

View File

@@ -4,14 +4,19 @@ The real ``plugins/platforms/telegram`` adapter reaches it through python-telegr
``base_url``/``base_file_url`` (``platforms.telegram.extra.base_url``), so every request crosses the
real PTB HTTP stack: ``POST {base}/bot<token>/<method>`` with form/JSON/multipart parameters, and
long-poll ``getUpdates`` that returns queued ``Update`` objects. Only the methods the adapter calls
are implemented; any other method answers ``ok: true, result: true`` and is still recorded, so a test
can see it.
are implemented; any other method answers ``404 Not Found`` like the real Bot API (recorded as
faulted), so a call the stand-in does not model can never be mistaken for a success.
Payload fidelity kept to the Bot API: message text must be 1..4096 chars after entity parsing
(``message text is empty`` / ``message is too long``), and with ``parse_mode=MarkdownV2`` every
reserved character outside an entity must be backslash-escaped, else ``400 can't parse entities``.
"""
from __future__ import annotations
import asyncio
import itertools
import re
import time
from typing import Any, Dict, List, Optional
@@ -22,6 +27,82 @@ from tests.fakes.platforms._standin import StandinServer, Visible, decode_value
BOT_ID = 7000000001
BOT_USERNAME = "hermes_standin_bot"
MAX_TEXT = 4096
# MarkdownV2 (https://core.telegram.org/bots/api#markdownv2-style): these must be escaped outside
# entities; inside ``code``/```pre``` only ` and \ are special.
_MDV2_RESERVED = set("_*[]()~`>#+-=|{}.!")
_MDV2_LINK = re.compile(r"\[((?:\\.|[^\]\\])*)\]\(((?:\\.|[^)\\])*)\)")
def mdv2_error(text: str) -> str | None:
"""The Bot API's ``can't parse entities`` reason for ``text`` in MarkdownV2, or None if it parses.
A faithful subset of the server's parser: escapes, code/pre spans, links, the paired style
markers (``*`` ``_`` ``__`` ``~`` ``||``) and line-leading ``>`` quotes. Any other reserved
character must carry a preceding backslash.
"""
i, n, open_marks = 0, len(text), []
while i < n:
ch = text[i]
if ch == "\\":
i += 2
continue
if text.startswith("```", i):
end = text.find("```", i + 3)
if end < 0:
return "Can't find end of Pre entity at byte offset %d" % i
i = end + 3
continue
if ch == "`":
end = text.find("`", i + 1)
if end < 0:
return "Can't find end of Code entity at byte offset %d" % i
i = end + 1
continue
if ch == "[":
m = _MDV2_LINK.match(text, i)
if not m:
return "Character '[' is reserved and must be escaped with the preceding '\\'"
i = m.end()
continue
if ch == ">" and (i == 0 or text[i - 1] == "\n"):
i += 1
continue
mark = next((mk for mk in ("||", "__", "*", "_", "~") if text.startswith(mk, i)), None)
if mark is not None:
if open_marks and open_marks[-1] == mark:
open_marks.pop()
else:
open_marks.append(mark)
i += len(mark)
continue
if ch in _MDV2_RESERVED:
return f"Character '{ch}' is reserved and must be escaped with the preceding '\\'"
i += 1
if open_marks:
return f"Can't find end of the entity starting with '{open_marks[-1]}'"
return None
class BotApiError(Exception):
"""A Bot API ``ok: false`` reply (``description`` + ``error_code``) from a method handler."""
def __init__(self, description: str, code: int = 400) -> None:
super().__init__(description)
self.description, self.code = description, code
def _check_text(p: Dict[str, Any], field: str = "text") -> str:
text = str(p.get(field) or "")
if p.get("parse_mode") == "MarkdownV2":
why = mdv2_error(text)
if why:
raise BotApiError(f"Bad Request: can't parse entities: {why}")
text = re.sub(r"\\(.)", r"\1", text) # length and emptiness count the parsed text
if not text.strip():
raise BotApiError("Bad Request: message text is empty")
if len(text) > MAX_TEXT:
raise BotApiError("Bad Request: message is too long")
return str(p.get(field) or "")
class TelegramStandin(StandinServer):
@@ -38,6 +119,8 @@ class TelegramStandin(StandinServer):
# (chat_id, message_id) -> Visible for BOT messages only
self._visible: Dict[tuple, Visible] = {}
self.chats: Dict[str, Dict[str, Any]] = {}
# (chat_id, draft_id) -> latest draft preview text (sendMessageDraft: ephemeral, not a message)
self.drafts: Dict[tuple, str] = {}
@property
def api_base(self) -> str:
@@ -86,7 +169,16 @@ class TelegramStandin(StandinServer):
self.record(method, params, fault.body, faulted=True)
return web.json_response(fault.body, status=fault.status)
handler = getattr(self, f"_m_{method}", None)
result = handler(params) if handler else True
if handler is None:
body = {"ok": False, "error_code": 404, "description": "Not Found"}
self.record(method, params, body, faulted=True)
return web.json_response(body, status=404)
try:
result = handler(params)
except BotApiError as exc:
body = {"ok": False, "error_code": exc.code, "description": exc.description}
self.record(method, params, body, faulted=True)
return web.json_response(body, status=exc.code)
self.record(method, params, result)
return web.json_response({"ok": True, "result": result})
@@ -160,7 +252,30 @@ class TelegramStandin(StandinServer):
return {**self._chat(p["chat_id"]), "accent_color_id": 0, "max_reaction_count": 11}
def _m_sendMessage(self, p: Dict[str, Any]) -> Dict[str, Any]:
return self._bot_message(p, text=p.get("text", ""))
return self._bot_message(p, text=_check_text(p))
def _m_sendMessageDraft(self, p: Dict[str, Any]) -> bool:
"""Animate a private-chat draft preview; not a message (no message_id)."""
if not int(p.get("draft_id") or 0):
raise BotApiError("Bad Request: draft_id must be non-zero")
text = _check_text(p)
with self._lock:
self.drafts[(str(p["chat_id"]), int(p["draft_id"]))] = text
return True
def _m_sendRichMessageDraft(self, p: Dict[str, Any]) -> bool:
if not int(p.get("draft_id") or 0):
raise BotApiError("Bad Request: draft_id must be non-zero")
with self._lock:
self.drafts[(str(p["chat_id"]), int(p["draft_id"]))] = str(p.get("text") or p.get("content") or "")
return True
def _ok_true(self, _p: Dict[str, Any]) -> bool:
return True
# Side-effect-only methods the adapter calls (reactions, typing, command menu, webhook reset).
_m_setMessageReaction = _m_sendChatAction = _m_setMyCommands = _m_deleteMyCommands = _ok_true
_m_setMyShortDescription = _m_setMyDescription = _m_deleteWebhook = _m_answerCallbackQuery = _ok_true
def _m_sendPhoto(self, p: Dict[str, Any]) -> Dict[str, Any]:
return self._bot_message(p, caption=p.get("caption", ""), _kind="photo",
@@ -171,6 +286,7 @@ class TelegramStandin(StandinServer):
document={"file_id": "out-doc", "file_unique_id": "od"})
def _m_editMessageText(self, p: Dict[str, Any]) -> Dict[str, Any]:
_check_text(p)
chat_id, mid = str(p["chat_id"]), str(p["message_id"])
with self._lock:
vis = self._visible.get((chat_id, mid))