fix(gateway): skill-slash fallthrough and the Telegram inline picker run off the event loop
The gateway's idle-command path resolved skill slash commands inline on the loop: a cold skill scan, skill file loads and the unavailable-skill rglob over every skills dir. On a 1.5k-skill install that held the loop ~2 minutes, the loop-liveness watchdog fired and the gateway exited mid-session (#111091). _hm_skill_slash_rewrite now runs through _run_in_executor_with_context so the profile contextvars the scan is scoped to survive the hop. Known commands still short-circuit before any I/O (previous commit). Same class in the Telegram inline picker: build_inline_results rebuilds the command/skill catalog per keystroke via _collect_gateway_skill_entries, the same path-resolution pass #110707 traced in the command menu. One invariant test: a known command triggers no scan; an unknown command's scan runs while the loop keeps ticking. Red on origin/main and on the reorder-only tree, green here.
This commit is contained in:
@@ -1206,7 +1206,12 @@ class GatewayInboundMixin:
|
||||
if not _handled:
|
||||
_handled, _result, command = await self._hm_dispatch_quick_and_plugin_commands(event, source, command)
|
||||
if not _handled:
|
||||
_result = self._hm_skill_slash_rewrite(event, source, _quick_key, command)
|
||||
# Skill-slash resolution is disk-bound (cold skill scan, skill file loads, the
|
||||
# unavailable-skill rglob over every skills dir) and uncached on a first hit; on a
|
||||
# large install it held the loop past the liveness watchdog (#111091). The executor
|
||||
# hop carries the profile contextvars the scan is scoped to.
|
||||
_result = await self._run_in_executor_with_context(
|
||||
self._hm_skill_slash_rewrite, event, source, _quick_key, command)
|
||||
_handled = _result is not None
|
||||
return _handled, _result
|
||||
|
||||
|
||||
@@ -4267,7 +4267,9 @@ class TelegramAdapter(BasePlatformAdapter):
|
||||
try:
|
||||
from telegram import InlineQueryResultArticle, InputTextMessageContent
|
||||
from plugins.platforms.telegram.inline_picker import CACHE_TIME_SECONDS as _CACHE, build_inline_results
|
||||
results, next_offset = build_inline_results(
|
||||
# Per-keystroke catalog build resolves every skill path; keep it off the loop (#110707).
|
||||
results, next_offset = await asyncio.to_thread(
|
||||
build_inline_results,
|
||||
getattr(inline_query, "query", "") or "", offset=getattr(inline_query, "offset", "") or "")
|
||||
articles = [
|
||||
InlineQueryResultArticle(
|
||||
|
||||
@@ -0,0 +1,64 @@
|
||||
"""Skill-slash fallthrough must not hold the gateway event loop (#111091)."""
|
||||
|
||||
import asyncio
|
||||
import threading
|
||||
from types import SimpleNamespace
|
||||
from unittest.mock import patch
|
||||
|
||||
import pytest
|
||||
|
||||
from gateway.config import Platform
|
||||
from gateway.session import SessionSource
|
||||
|
||||
|
||||
def _runner(command):
|
||||
from gateway.run import GatewayRunner
|
||||
|
||||
runner = object.__new__(GatewayRunner)
|
||||
|
||||
async def _resolve(event, source, qk):
|
||||
return False, None, command, None
|
||||
|
||||
async def _canonical(event, source, qk, canonical):
|
||||
return False, None
|
||||
|
||||
async def _quick(event, source, cmd):
|
||||
return False, None, cmd
|
||||
|
||||
runner._hm_resolve_command = _resolve
|
||||
runner._hm_dispatch_canonical_command = _canonical
|
||||
runner._hm_dispatch_quick_and_plugin_commands = _quick
|
||||
return runner
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_unavailable_skill_scan_skips_known_commands_and_runs_off_loop():
|
||||
import gateway.run as gateway_run
|
||||
|
||||
source = SessionSource(platform=Platform.DISCORD, chat_id="c1")
|
||||
event = SimpleNamespace(text="/x", get_command_args=lambda: "")
|
||||
scanned = []
|
||||
loop_was_free = []
|
||||
scan_started = threading.Event()
|
||||
loop_ticked = threading.Event()
|
||||
|
||||
def _slow_scan(command):
|
||||
scanned.append(command)
|
||||
scan_started.set()
|
||||
# True only if the loop ran while this scan was in flight, i.e. the scan is off-loop.
|
||||
loop_was_free.append(loop_ticked.wait(timeout=1))
|
||||
return None
|
||||
|
||||
with patch.object(gateway_run, "_check_unavailable_skill", _slow_scan):
|
||||
# A registered gateway command returns before any filesystem walk.
|
||||
handled, reply = await _runner("steer")._hm_dispatch_idle_commands(event, source, "qk")
|
||||
assert (handled, reply, scanned) == (False, None, [])
|
||||
|
||||
# An unknown command still consults the hint, but the scan runs off the loop.
|
||||
task = asyncio.create_task(_runner("no-such-skill")._hm_dispatch_idle_commands(event, source, "qk"))
|
||||
await asyncio.to_thread(scan_started.wait, 1)
|
||||
loop_ticked.set() # only reachable mid-scan when the loop is free
|
||||
handled, reply = await task
|
||||
assert scanned == ["no-such-skill"]
|
||||
assert loop_was_free == [True]
|
||||
assert handled and "Unknown command" in reply
|
||||
Reference in New Issue
Block a user