simplify(compat): gateway — re-land 92d0bd0d73 (reverted by stale-index commit b818085298)

Re-applies the gateway compat removal byte-for-byte; see 92d0bd0d73 for the
full inventory (30 re-exports/aliases + 2 shim modules dropped, 3 shim-only
names re-removed, 24 callers + 34 test files repointed). No new changes.
This commit is contained in:
Teknium
2026-09-03 13:12:50 -07:00
parent 89fbd5d4d3
commit 2031c819fe
67 changed files with 160 additions and 290 deletions
+1 -2
View File
@@ -1,11 +1,10 @@
"""Hermes Gateway - multi-platform messaging integration (sessions, context
injection, delivery routing, platform-specific toolsets)."""
from .config import GatewayConfig, PlatformConfig, HomeChannel, load_gateway_config
from .config import GatewayConfig, PlatformConfig, HomeChannel, SessionResetPolicy, load_gateway_config
from .session import (
SessionContext,
SessionStore,
SessionResetPolicy,
build_session_context_prompt,
)
from .delivery import DeliveryRouter, DeliveryTarget
+2 -6
View File
@@ -15,7 +15,7 @@ import time
from pathlib import Path
from typing import Any, Optional
from gateway.kanban_watchers_common import ( # noqa: F401 (tests import via origin)
from gateway.kanban_watchers_common import (
_acquire_singleton_lock,
_kanban_dispatch_allowed,
_release_singleton_lock,
@@ -24,11 +24,7 @@ from gateway.kanban_watchers_common import ( # noqa: F401 (tests import via or
_to_thread_process_service,
logger,
)
from gateway.kanban_watchers_notifier import ( # noqa: F401 (_wake_scope_id: tests import via origin)
_KanbanNotification,
_notifier_collect,
_wake_scope_id,
)
from gateway.kanban_watchers_notifier import _KanbanNotification, _notifier_collect
from gateway.kanban_watchers_dispatcher import (
_KanbanDispatcher,
_log_spawn_results,
+1 -18
View File
@@ -2,21 +2,4 @@
from .base import BasePlatformAdapter, MessageEvent, SendResult
# QQAdapter / YuanbaoAdapter are exposed lazily (PEP 562 ``__getattr__``): eager
# imports cost ~48 ms / ~8 MB RSS on every CLI invocation and nothing in-tree
# imports them from the package root.
__all__ = ["BasePlatformAdapter", "MessageEvent", "SendResult", "QQAdapter", "YuanbaoAdapter"]
_LAZY_ADAPTERS = {"QQAdapter": ".qqbot", "YuanbaoAdapter": ".yuanbao"}
def __getattr__(name):
module = _LAZY_ADAPTERS.get(name)
if module is None:
raise AttributeError(f"module {__name__!r} has no attribute {name!r}")
from importlib import import_module
return getattr(import_module(module, __name__), name)
def __dir__():
return sorted(__all__)
__all__ = ["BasePlatformAdapter", "MessageEvent", "SendResult"]
-1
View File
@@ -122,7 +122,6 @@ from gateway.platforms import api_server_runs as _api_runs
from gateway.platforms.api_server_openai_routes import OpenAICompatRoutesMixin
from gateway.platforms.base import (
MEDIA_TAG_CLEANUP_RE, BasePlatformAdapter, SendResult, is_network_accessible, validate_media_delivery_path)
# Re-exported here for existing imports and constructor monkeypatches.
from gateway.platforms.api_server_run_idempotency import RunIdempotencyStore
from agent.redact import redact_sensitive_text
from agent.interrupt_compat import request_hard_interrupt
+2 -1
View File
@@ -432,7 +432,8 @@ sys.path.insert(0, str(Path(__file__).resolve().parents[2]))
from gateway.config import Platform, PlatformConfig
from gateway.platforms.helpers import fence_state_after
from gateway.session import SessionSource, TranscriptReadError, build_session_key
from gateway.session import SessionSource, build_session_key
from gateway.session_transcript import TranscriptReadError
from hermes_constants import get_default_hermes_root, get_hermes_dir, get_hermes_home
if TYPE_CHECKING:
+1 -4
View File
@@ -632,9 +632,6 @@ class SignalAdapter(BasePlatformAdapter):
logger.log(fail_level, "Signal RPC %s failed: %s", method, e)
return None
# Backward-compatible alias for the shared formatting helper.
_markdown_to_signal = staticmethod(markdown_to_signal)
def format_message(self, content: str) -> str:
"""Plain-text fallback for the base-class send path; send() applies rich styles itself."""
return content
@@ -714,7 +711,7 @@ class SignalAdapter(BasePlatformAdapter):
if not content or not content.strip():
return SendResult(success=True, message_id=None)
base_params = await self._with_target({"account": self.account}, chat_id)
chunks = self._split_signal_formatted_message(*self._markdown_to_signal(content), self.MAX_MESSAGE_LENGTH)
chunks = self._split_signal_formatted_message(*markdown_to_signal(content), self.MAX_MESSAGE_LENGTH)
last_result = None
for idx, (plain_text, text_styles) in enumerate(chunks, start=1):
params: Dict[str, Any] = dict(base_params, message=plain_text)
+7 -38
View File
@@ -62,7 +62,8 @@ from gateway.platforms.yuanbao_proto import (
encode_send_private_heartbeat, encode_send_group_heartbeat, encode_query_group_info,
encode_get_group_member_list, next_seq_no,
)
from gateway.session import TranscriptReadError, build_session_key
from gateway.session import build_session_key
from gateway.session_transcript import TranscriptReadError
logger = logging.getLogger(__name__)
@@ -144,20 +145,7 @@ async def _cancel_task(task: asyncio.Task) -> None:
class MarkdownProcessor:
"""Thin delegates to the shared fence-aware chunker in gateway.platforms.helpers; method names
kept for existing call sites and tests."""
@staticmethod
def has_unclosed_fence(text: str) -> bool:
return _mdchunk.text_has_unclosed_fence(text)
@staticmethod
def ends_with_table_row(text: str) -> bool:
return _mdchunk.text_ends_with_table_row(text)
@staticmethod
def split_at_paragraph_boundary(text: str, max_chars: int, len_fn: Optional[Callable[[str], int]] = None) -> tuple[str, str]:
return _mdchunk.split_at_paragraph_boundary(text, max_chars, len_fn=len_fn)
"""Yuanbao's fence/table-aware chunking policy over the shared chunker in gateway.platforms.helpers."""
@classmethod
def chunk_markdown_text(cls, text: str, max_chars: int = 4000, len_fn: Optional[Callable[[str], int]] = None) -> list[str]:
"""<= max_chars chunks at paragraph boundaries, never inside a fence or table (an oversized
@@ -2199,7 +2187,7 @@ class MediaSendHandler(ABC):
caption: Optional[str] = None, **kwargs: Any) -> "SendResult":
if adapter._connection.ws is None:
return SendResult(success=False, error="Not connected", retryable=True)
adapter._outbound.cancel_slow_notifier(chat_id)
adapter._outbound.slow_notifier.cancel(chat_id)
try:
file_bytes, filename, content_type = await self.acquire_file(adapter, **kwargs)
if self.needs_cos_upload():
@@ -2420,7 +2408,7 @@ class MessageSender:
adapter = self._adapter
if adapter._connection.ws is None:
return SendResult(success=False, error="Not connected", retryable=True)
adapter._outbound.cancel_slow_notifier(chat_id)
adapter._outbound.slow_notifier.cancel(chat_id)
async with self.get_chat_lock(chat_id):
content_to_send = self.strip_cron_wrapper(content)
chunks = self.truncate_message(content_to_send, adapter.MAX_TEXT_CHUNK)
@@ -2594,21 +2582,11 @@ class MessageSender:
class OutboundManager:
"""Composes MessageSender, HeartbeatManager and SlowResponseNotifier (sender cancels the notifier
before a send and emits the FINISH heartbeat after)."""
CHAT_DICT_MAX_SIZE: ClassVar[int] = MessageSender.CHAT_DICT_MAX_SIZE
def __init__(self, adapter: "YuanbaoAdapter") -> None:
self._adapter = adapter
self.sender: MessageSender = MessageSender(adapter)
self.heartbeat: HeartbeatManager = HeartbeatManager(adapter)
self.slow_notifier: SlowResponseNotifier = SlowResponseNotifier(adapter, self.sender)
# Delegates kept for callers/tests that address the outbound facade.
self.start_slow_notifier = self.slow_notifier.start
self.cancel_slow_notifier = self.slow_notifier.cancel
self.get_chat_lock = self.sender.get_chat_lock
@property
def _chat_locks(self) -> collections.OrderedDict:
return self.sender._chat_locks
async def close(self) -> None:
await self.sender.close()
@@ -2748,11 +2726,11 @@ class YuanbaoAdapter(BasePlatformAdapter):
async def _process_message_background(self, event, session_key: str) -> None:
"""Wrap base class processing with a slow-response notifier."""
chat_id = event.source.chat_id
await self._outbound.start_slow_notifier(chat_id)
await self._outbound.slow_notifier.start(chat_id)
try:
await super()._process_message_background(event, session_key)
finally:
self._outbound.cancel_slow_notifier(chat_id)
self._outbound.slow_notifier.cancel(chat_id)
# Clear RecallGuard tracking only if our msg_id is still current: a concurrent message may have
# overwritten it (the drain task then owns it); id-less events never wrote one and must not pop.
msg_id = event.message_id
@@ -2824,12 +2802,3 @@ class YuanbaoAdapter(BasePlatformAdapter):
"""Current valid sign token (module-level cache)."""
return await SignManager.get_token(self._app_key, self._app_secret, self._api_domain, route_env=self._route_env)
# Module-level delegates kept for external importers (tools/send_message_tool, tools/yuanbao_tools).
def get_active_adapter() -> Optional["YuanbaoAdapter"]:
return YuanbaoAdapter.get_active()
async def send_yuanbao_direct(adapter: "YuanbaoAdapter", chat_id: str, message: str,
media_files: Optional[List[Tuple[str, bool]]] = None) -> Dict[str, Any]:
return await adapter._outbound.sender.send_direct(chat_id, message, media_files)
-41
View File
@@ -93,44 +93,3 @@ class CapabilityDescriptor:
else ()
)
return cls(**filtered)
@classmethod
def from_platform_entry(
cls,
entry,
*,
len_unit: str = "chars",
supports_draft_streaming: bool = False,
supports_edit: bool = True,
supports_threads: bool = False,
markdown_dialect: str = "plain",
) -> "CapabilityDescriptor":
"""Project a ``gateway.platform_registry.PlatformEntry`` into a descriptor.
Demonstrates the descriptor is a *subset/projection* of what
``PlatformEntry`` already encodes, not a parallel concept: ``label``,
``max_message_length``, ``emoji``, ``platform_hint``, ``pii_safe`` and
the platform name come straight off the entry. The runtime capability
bits that ``PlatformEntry`` does NOT encode (length unit, draft/edit/
thread/markdown behavior) are supplied by the caller — in production
the connector fills these from the live adapter's capability methods.
``max_message_length`` of 0 on a ``PlatformEntry`` means "no limit";
we map that to the stream_consumer default of 4096 so the descriptor
always carries a concrete chunking bound.
"""
max_len = getattr(entry, "max_message_length", 0) or 4096
return cls(
contract_version=CONTRACT_VERSION,
platform=entry.name,
label=entry.label,
max_message_length=max_len,
supports_draft_streaming=supports_draft_streaming,
supports_edit=supports_edit,
supports_threads=supports_threads,
markdown_dialect=markdown_dialect,
len_unit=len_unit,
emoji=getattr(entry, "emoji", "\U0001f50c"),
platform_hint=getattr(entry, "platform_hint", ""),
pii_safe=getattr(entry, "pii_safe", False),
)
+2 -1
View File
@@ -28,7 +28,8 @@ from pathlib import Path
from typing import Any, Awaitable, Callable, Dict, Optional
if TYPE_CHECKING: # string annotations only; never imported at runtime (cycle)
from gateway.run import GatewayRunner, TurnRunner # noqa: F401
from gateway.run import GatewayRunner # noqa: F401
from gateway.run_turn_runner import TurnRunner # noqa: F401
# Log-record parity with the origin module.
logger = logging.getLogger("gateway.run")
+2 -1
View File
@@ -17,7 +17,8 @@ from gateway.session import SessionSource, build_session_context_prompt
from hermes_cli.config import cfg_get
if TYPE_CHECKING: # string annotations only; never imported at runtime (cycle)
from gateway.run import GatewayRunner, TurnRunner # noqa: F401
from gateway.run import GatewayRunner # noqa: F401
from gateway.run_turn_runner import TurnRunner # noqa: F401
# Log-record parity with the origin module.
logger = logging.getLogger("gateway.run")
+3 -2
View File
@@ -21,7 +21,8 @@ from gateway.session import SessionSource
from typing import Any, Dict, Optional, Union
if TYPE_CHECKING: # string annotations only; never imported at runtime (cycle)
from gateway.run import GatewayRunner, TurnRunner # noqa: F401
from gateway.run import GatewayRunner # noqa: F401
from gateway.run_turn_runner import TurnRunner # noqa: F401
# Log-record parity with the origin module.
logger = logging.getLogger("gateway.run")
@@ -275,7 +276,7 @@ class GatewayBusySessionMixin:
)
def _queue_or_replace_pending_event(self, session_key: str, event: MessageEvent) -> None:
from gateway.run import merge_pending_message_event
from gateway.platforms.base import merge_pending_message_event
adapter = self._adapter_for_source(event.source)
if not adapter:
return
+2 -1
View File
@@ -30,7 +30,8 @@ from hermes_cli.fallback_config import get_fallback_chain
from utils import is_truthy_value
if TYPE_CHECKING: # string annotations only; never imported at runtime (cycle)
from gateway.run import GatewayRunner, TurnRunner # noqa: F401
from gateway.run import GatewayRunner # noqa: F401
from gateway.run_turn_runner import TurnRunner # noqa: F401
# Log-record parity with the origin module.
logger = logging.getLogger("gateway.run")
+2 -1
View File
@@ -16,7 +16,8 @@ from typing import TYPE_CHECKING, Any
from gateway.platforms.base import MessageEvent, MessageType
if TYPE_CHECKING: # string annotations only; never imported at runtime (cycle)
from gateway.run import GatewayRunner, TurnRunner # noqa: F401
from gateway.run import GatewayRunner # noqa: F401
from gateway.run_turn_runner import TurnRunner # noqa: F401
# Log-record parity with the origin module.
logger = logging.getLogger("gateway.run")
+3 -2
View File
@@ -27,7 +27,8 @@ from gateway.turn_lease import TurnLeaseTimeoutError
from typing import Any, Dict, List, Optional, Tuple
if TYPE_CHECKING: # string annotations only; never imported at runtime (cycle)
from gateway.run import GatewayRunner, TurnRunner # noqa: F401
from gateway.run import GatewayRunner # noqa: F401
from gateway.run_turn_runner import TurnRunner # noqa: F401
# Log-record parity with the origin module.
logger = logging.getLogger("gateway.run")
@@ -485,7 +486,7 @@ class GatewayInboundMixin:
self, source: SessionSource, _quick_key: str, event: "MessageEvent", *, merge_text: bool = False
) -> None:
"""Merge *event* into the source adapter's pending slot (no-op without an adapter)."""
from gateway.run import merge_pending_message_event
from gateway.platforms.base import merge_pending_message_event
adapter = self._adapter_for_source(source)
if adapter:
merge_pending_message_event(adapter._pending_messages, _quick_key, event, merge_text=merge_text)
+4 -3
View File
@@ -598,7 +598,8 @@ class GatewayNotificationsMixin:
async def _send_restart_notification(self) -> Optional[tuple[str, str, Optional[str]]]:
"""Notify the chat that initiated /restart that the gateway is back."""
from gateway.run import _hermes_home, _non_conversational_metadata, resolve_delivery_transport
from gateway.delivery import resolve_delivery_transport
from gateway.run import _hermes_home, _non_conversational_metadata
notify_path = _hermes_home / ".restart_notify.json"
if not notify_path.exists():
return None
@@ -647,7 +648,7 @@ class GatewayNotificationsMixin:
def _home_channel_transports(self):
"""Yield ``(platform, platform_cfg, home, transport)`` for every home channel with a live transport."""
from gateway.run import resolve_delivery_transport
from gateway.delivery import resolve_delivery_transport
for platform, platform_cfg in self.config.platforms.items():
home = platform_cfg.home_channel
if not home or not home.chat_id:
@@ -866,7 +867,7 @@ class GatewayNotificationsMixin:
"""Adapter for a synthetic-event platform: alias-aware transport resolver first (one
Platform.RELAY adapter fronts N logical platforms; native wins), literal ``p.value`` scan as
fallback for minimal runner stubs / exotic platform strings when the resolver can't run."""
from gateway.run import resolve_delivery_transport
from gateway.delivery import resolve_delivery_transport
try:
_transport = resolve_delivery_transport(Platform(platform_name), self.config, self.adapters)
except Exception:
+4 -6
View File
@@ -1056,7 +1056,7 @@ class GatewayShutdownMixin:
def _increment_restart_failure_counts(self, active_session_keys: set) -> None:
"""Increment persisted restart-failure counters for active sessions; drop the rest (loop broken)."""
from gateway.run import atomic_json_write
from utils import atomic_json_write
path = self._stuck_loop_counts_path()
counts = self._read_json_counts(path) or {}
with suppress(Exception):
@@ -1091,7 +1091,7 @@ class GatewayShutdownMixin:
async def _clear_restart_failure_count(self, session_key: str) -> None:
"""Clear a completed session's restart-failure counter off-loop (atomic_json_write fsyncs)."""
from gateway.run import atomic_json_write
from utils import atomic_json_write
path = self._stuck_loop_counts_path()
if not path.exists():
return
@@ -1666,10 +1666,8 @@ class GatewayShutdownMixin:
def _stop_persist_exit_state(self, ctx: "GatewayShutdownMixin._StopContext") -> None:
"""PID/lock release, clean-shutdown marker, restart markers, terminal runtime status."""
from gateway.run import (
_hermes_home, _planned_restart_notification_path, _shutdown_gateway_health_export,
atomic_json_write,
)
from gateway.run import _hermes_home, _planned_restart_notification_path, _shutdown_gateway_health_export
from utils import atomic_json_write
from gateway.status import remove_pid_file, release_gateway_runtime_lock
remove_pid_file()
release_gateway_runtime_lock()
+3 -3
View File
@@ -523,7 +523,7 @@ class GatewayStartupMixin:
See #69089.
"""
from gateway.run import _arm_loop_floor_timer, start_loop_liveness_watchdog
from gateway.shutdown_watchdog import _arm_loop_floor_timer, start_loop_liveness_watchdog
config = getattr(self, "config", None)
if config is not None and not getattr(config, "loop_watchdog", True):
return
@@ -673,7 +673,7 @@ class GatewayStartupMixin:
# Loop live: the loop-liveness watchdog takes over from the startup watchdog. Disarm even
# when loop guards are config-disabled; only inside this branch (no live loop = stay armed).
with _log_suppressed(logging.DEBUG, "Startup watchdog disarm failed", exc_info=True):
from gateway.startup_watchdog import disarm_startup_watchdog
from hermes_startup_watchdog import disarm_startup_watchdog
disarm_startup_watchdog()
logger.info("Session storage: %s", self.config.sessions_dir)
self._start_log_systemd_timing_alignment()
@@ -1326,7 +1326,7 @@ class GatewayStartupMixin:
self, row: Dict[str, Any], profile_name: Optional[str]
) -> "GatewayStartupMixin._HandoffDestination":
"""Resolve platform, transport, home channel, thread and destination source for a row."""
from gateway.run import resolve_delivery_transport
from gateway.delivery import resolve_delivery_transport
cli_session_id = row["id"]
platform_name = (row.get("handoff_platform") or "").strip().lower()
if not platform_name:
+2 -1
View File
@@ -21,7 +21,8 @@ from gateway.session import SessionSource
from utils import is_truthy_value
if TYPE_CHECKING: # string annotations only; never imported at runtime (cycle)
from gateway.run import GatewayRunner, TurnRunner # noqa: F401
from gateway.run import GatewayRunner # noqa: F401
from gateway.run_turn_runner import TurnRunner # noqa: F401
# Log-record parity with the origin module.
logger = logging.getLogger("gateway.run")
+8 -5
View File
@@ -22,9 +22,10 @@ from gateway.config import Platform
from gateway.media_repair import repair_explicit_computer_use_media_paths
from gateway.platforms.base import BasePlatformAdapter, MessageEvent
from gateway.session import (
SessionSource, TranscriptReadError, _session_key_namespace, build_channel_continuity_note,
SessionSource, _session_key_namespace, build_channel_continuity_note,
build_session_context,
)
from gateway.session_transcript import TranscriptReadError
from gateway.turn_context import TurnContext
from gateway.turn_lease import DEFAULT_LEASE_WAIT, TurnLeaseTimeoutError
from hermes_constants import get_hermes_home_override
@@ -33,7 +34,8 @@ from typing import Any, Callable, Dict, List, Optional, Tuple
from utils import base_url_hostname
if TYPE_CHECKING: # string annotations only; never imported at runtime (cycle)
from gateway.run import GatewayRunner, TurnRunner # noqa: F401
from gateway.run import GatewayRunner # noqa: F401
from gateway.run_turn_runner import TurnRunner # noqa: F401
# Log-record parity with the origin module.
logger = logging.getLogger("gateway.run")
@@ -1428,7 +1430,7 @@ class GatewayTurnMixin:
last_reasoning = agent_result.get("last_reasoning")
if not (_show_reasoning_effective and response and not _intentional_silence and last_reasoning):
return response
from gateway.stream_consumer import escape_code_fences_for_display
from gateway.stream_consumer_fences import escape_code_fences_for_display
# Collapse long reasoning to keep messages readable
lines = last_reasoning.strip().splitlines()
if len(lines) > 15:
@@ -2692,7 +2694,7 @@ class GatewayTurnMixin:
"""Build the ``TurnContext`` and its ``TurnRunner``; ``turn_params`` (history, context_prompt,
session_id, persist_user_*, …) are stored verbatim. Returns ``(turn_ctx, turn_runner,
cleanup_adapter)``."""
from gateway.run import TurnRunner
from gateway.run_turn_runner import TurnRunner
# Discord voice "verbal ack" on the FIRST tool call (discord.voice_fx.enabled): resolve the
# guild whose voice connection is bound to this text channel (mirrors DiscordAdapter.play_tts).
_voice_ack_guild: List[Optional[int]] = [None]
@@ -3390,7 +3392,8 @@ class GatewayTurnMixin:
response: Any, result: Any, stream_task: Any,
) -> Any:
"""Run the queued / interrupting follow-up as the next turn (recursive ``_run_agent``)."""
from gateway.run import _preserve_queued_followup_history_offset, merge_pending_message_event
from gateway.platforms.base import merge_pending_message_event
from gateway.run import _preserve_queued_followup_history_offset
source, session_id, session_key, run_generation = (
turn_ctx.source, turn_ctx.session_id, turn_ctx.session_key, turn_ctx.run_generation,
)
+1 -1
View File
@@ -29,7 +29,7 @@ from hermes_cli.config import cfg_get
from utils import is_truthy_value
if TYPE_CHECKING: # string annotations only; never imported at runtime (cycle)
from gateway.run import GatewayRunner, TurnRunner # noqa: F401
from gateway.run import GatewayRunner # noqa: F401
# Log-record parity with the origin module.
logger = logging.getLogger("gateway.run")
+3 -7
View File
@@ -13,15 +13,11 @@ from dataclasses import dataclass, field, fields
from typing import Dict, List, Optional, Any
from .config import Platform, GatewayConfig, HomeChannel
from .config import SessionResetPolicy # noqa: F401 — re-exported via gateway/__init__.py
from .whatsapp_identity import canonical_whatsapp_identifier
from .whatsapp_identity import normalize_whatsapp_identifier # noqa: F401 — re-exported
from gateway.session_persistence import SessionPersistenceMixin, _DB_UNPINNED # noqa: F401
from gateway.session_persistence import SessionPersistenceMixin, _DB_UNPINNED
from gateway.session_recovery import SessionRecoveryMixin
from gateway.session_lifecycle import ( # noqa: F401 — _now & co. re-exported for callers/tests
SessionLifecycleMixin, _iso, _new_session_id, _now, _parse_iso, auto_continue_freshness_window,
)
from gateway.session_transcript import SessionTranscriptMixin, TranscriptReadError # noqa: F401
from gateway.session_lifecycle import SessionLifecycleMixin, _iso, _new_session_id, _now, _parse_iso
from gateway.session_transcript import SessionTranscriptMixin
logger = logging.getLogger(__name__)
+1 -1
View File
@@ -341,7 +341,7 @@ class SessionPersistenceMixin:
def _stale_entry_verdict(self, key: str, entry, row):
"""For a routing entry whose row has ended: ``"prune"``, a replacement entry (repoint), or
None (keep as-is)."""
from gateway.session import _now
from gateway.session_lifecycle import _now
recovered_entry = None
if entry.origin is not None:
try:
-3
View File
@@ -267,6 +267,3 @@ def parse_systemd_duration_to_us(raw: str) -> Optional[int]:
return None
return total_us if total_us > 0 else None
# Backward-compat private alias (pre-promotion name).
_parse_systemd_duration_to_us = parse_systemd_duration_to_us
+4 -5
View File
@@ -22,13 +22,12 @@ from typing import Optional, Union
from agent.i18n import t
from gateway.config import HomeChannel, Platform, PlatformConfig, persist_home_channel
from gateway.platforms.base import EphemeralReply, MessageEvent
from gateway.session import AsyncSessionStore, TranscriptReadError
from gateway.session import AsyncSessionStore
from gateway.session_transcript import TranscriptReadError
from gateway.slash_commands_goals import GatewayGoalCommandsMixin
from gateway.slash_commands_model import ( # noqa: F401 — _model_switch_skew_guard re-exported for tests
GatewayModelCommandsMixin,
_model_switch_skew_guard)
from gateway.slash_commands_model import GatewayModelCommandsMixin
from gateway.slash_commands_session import GatewaySessionCommandsMixin
from gateway.slash_commands_status import HISTORY_UNREADABLE, GatewayStatusCommandsMixin # noqa: F401 — re-exported
from gateway.slash_commands_status import GatewayStatusCommandsMixin
from hermes_cli.config import atomic_config_write, cfg_get
from utils import atomic_json_write, is_truthy_value
+2 -1
View File
@@ -18,7 +18,8 @@ from agent.i18n import t
from agent.turn_context import extract_api_content_sidecar
from gateway.config import Platform
from gateway.platforms.base import EphemeralReply, MessageEvent, MessageType
from gateway.session import SessionSource, TranscriptReadError, build_session_key, is_shared_multi_user_session
from gateway.session import SessionSource, build_session_key, is_shared_multi_user_session
from gateway.session_transcript import TranscriptReadError
from gateway.slash_commands_status import HISTORY_UNREADABLE
logger = logging.getLogger("gateway.run") # log-record parity with gateway/run.py
+1 -1
View File
@@ -15,7 +15,7 @@ from agent.account_usage import fetch_account_usage, render_account_usage_lines
from agent.i18n import t
from gateway.config import Platform
from gateway.platforms.base import MessageEvent
from gateway.session import TranscriptReadError
from gateway.session_transcript import TranscriptReadError
# Log-record parity with gateway/run.py and the origin module.
logger = logging.getLogger("gateway.run")
+1 -3
View File
@@ -30,9 +30,7 @@ from gateway.config import (
from gateway.response_filters import (
is_intentional_silence_response as _is_intentional_silence_response,
is_partial_silence_marker as _is_partial_silence_marker)
from gateway.stream_consumer_fences import ( # noqa: F401 (re-exported)
ensure_closed_code_fences,
escape_code_fences_for_display)
from gateway.stream_consumer_fences import ensure_closed_code_fences
from gateway.stream_consumer_transport import StreamTransportMixin
from gateway.stream_consumer_fallback import StreamFallbackMixin
from gateway.stream_consumer_think import StreamThinkFilterMixin
-3
View File
@@ -76,9 +76,6 @@ class SessionTurnLeaseRegistry:
self._leases: Dict[str, _SessionLease] = {}
self._max_entries = max(1, int(max_entries))
def __len__(self) -> int:
return len(self._leases)
def _get_or_create(self, session_id: str) -> _SessionLease:
if (lease := self._leases.get(session_id)) is None:
self._evict_idle()
+1 -1
View File
@@ -4500,7 +4500,7 @@ def _respawn_storm_backoff() -> None:
)
# Tell the startup watchdog the backoff sleep is intentional, not a parked deadlock.
try:
from gateway.startup_watchdog import kick_startup_watchdog
from hermes_startup_watchdog import kick_startup_watchdog
kick_startup_watchdog(extra_s=_storm.backoff_s)
except Exception:
pass
@@ -25,14 +25,18 @@ from gateway.platforms.yuanbao import MessageSender, YuanbaoAdapter
# Helpers
# ---------------------------------------------------------------------------
class _OutboundStub:
async def start_slow_notifier(self, chat_id): # noqa: ANN001
class _SlowNotifierStub:
async def start(self, chat_id): # noqa: ANN001
pass
def cancel_slow_notifier(self, chat_id): # noqa: ANN001
def cancel(self, chat_id): # noqa: ANN001
pass
class _OutboundStub:
slow_notifier = _SlowNotifierStub()
def _bare_adapter():
"""YuanbaoAdapter instance without running its heavy __init__."""
adapter = object.__new__(YuanbaoAdapter)
@@ -24,7 +24,8 @@ import pytest
import gateway.run as gateway_run
from gateway.config import GatewayConfig, Platform
from gateway.platforms.base import MessageEvent
from gateway.session import SessionEntry, SessionSource, TranscriptReadError
from gateway.session import SessionEntry, SessionSource
from gateway.session_transcript import TranscriptReadError
def _bootstrap(monkeypatch, tmp_path):
@@ -6,17 +6,10 @@ import pytest
from gateway.config import Platform, PlatformConfig
from gateway.platforms import api_server
from gateway.platforms.api_server_run_idempotency import (
RunIdempotencyStore as ExtractedRunIdempotencyStore,
)
def test_run_idempotency_store_remains_reexported_from_api_server():
assert api_server.RunIdempotencyStore is ExtractedRunIdempotencyStore
@pytest.mark.asyncio
async def test_api_server_constructor_uses_legacy_run_store_monkeypatch(monkeypatch):
async def test_api_server_constructor_uses_module_run_store_binding(monkeypatch):
store = MagicMock()
store_factory = MagicMock(return_value=store)
monkeypatch.setattr(api_server, "RunIdempotencyStore", store_factory)
+5 -5
View File
@@ -959,7 +959,7 @@ class TestRunsProviderAuthFailure:
def _use_idempotency_db(adapter, path):
from gateway.platforms.api_server import RunIdempotencyStore
from gateway.platforms.api_server_run_idempotency import RunIdempotencyStore
adapter._run_idempotency_store.close()
adapter._run_idempotency_store = RunIdempotencyStore(str(path))
@@ -1114,7 +1114,7 @@ class TestRunIdempotency:
assert calls == 1
def test_restart_durability_and_terminal_semantics(self, tmp_path):
from gateway.platforms.api_server import RunIdempotencyStore
from gateway.platforms.api_server_run_idempotency import RunIdempotencyStore
path = tmp_path / "idem.db"
for terminal in ("completed", "failed", "cancelled"):
@@ -1141,7 +1141,7 @@ class TestRunIdempotency:
restarted.close()
def test_tenant_isolation_and_retention(self, tmp_path):
from gateway.platforms.api_server import RunIdempotencyStore
from gateway.platforms.api_server_run_idempotency import RunIdempotencyStore
store = RunIdempotencyStore(str(tmp_path / "idem.db"))
assert (
@@ -1157,7 +1157,7 @@ class TestRunIdempotency:
def test_retention_never_releases_an_active_idempotency_reservation(
self, tmp_path
):
from gateway.platforms.api_server import RunIdempotencyStore
from gateway.platforms.api_server_run_idempotency import RunIdempotencyStore
store = RunIdempotencyStore(str(tmp_path / "idem.db"))
with patch("gateway.platforms.api_server.time.time", return_value=100):
@@ -1375,7 +1375,7 @@ class TestRunIdempotency:
async def test_dead_owner_nonterminal_status_becomes_interrupted(
self, tmp_path
):
from gateway.platforms.api_server import RunIdempotencyStore
from gateway.platforms.api_server_run_idempotency import RunIdempotencyStore
path = tmp_path / "idem.db"
scope = hashlib.sha256(
+1 -1
View File
@@ -208,7 +208,7 @@ class TestBusySessionAck:
agent.steer = MagicMock(return_value=True)
runner._running_agents[sk] = agent
with patch("gateway.run.merge_pending_message_event") as mock_merge:
with patch("gateway.platforms.base.merge_pending_message_event") as mock_merge:
await runner._handle_active_session_busy_message(event, sk)
# VERIFY: Agent was steered, NOT interrupted
@@ -144,7 +144,7 @@ class TestBusyHandlerDemotesInterruptForCompression:
runner.adapters[event.source.platform] = adapter
runner._session_db._db.get_compression_lock_holder.return_value = "compressing"
with patch("gateway.run.merge_pending_message_event"):
with patch("gateway.platforms.base.merge_pending_message_event"):
await runner._handle_active_session_busy_message(event, sk)
adapter._send_with_retry.assert_called_once()
@@ -172,7 +172,6 @@ def test_turn_route_merges_fast_mode_with_provider_request_overrides():
@pytest.mark.asyncio
async def test_run_agent_preserves_provider_request_overrides_on_gateway_path(monkeypatch):
monkeypatch.setattr(gateway_run, "_load_gateway_config", lambda: {})
monkeypatch.setattr(gateway_run, "load_dotenv", lambda *args, **kwargs: None)
monkeypatch.setattr(gateway_run, "_load_gateway_runtime_config", lambda: {})
monkeypatch.setattr(gateway_run, "_resolve_gateway_model", lambda config=None: "gpt-5.4")
monkeypatch.setattr(
@@ -229,7 +228,6 @@ async def test_reused_agent_turn_merges_request_overrides_not_overwrite(monkeypa
fast-mode key while the provider extra_body survives.
"""
monkeypatch.setattr(gateway_run, "_load_gateway_config", lambda: {})
monkeypatch.setattr(gateway_run, "load_dotenv", lambda *args, **kwargs: None)
monkeypatch.setattr(gateway_run, "_load_gateway_runtime_config", lambda: {})
monkeypatch.setattr(gateway_run, "_resolve_gateway_model", lambda config=None: "gpt-5.4")
monkeypatch.setattr(
@@ -6,7 +6,7 @@ B1: Escape triple-backtick markers inside reasoning text before wrapping
"""
import pytest
from gateway.stream_consumer import escape_code_fences_for_display
from gateway.stream_consumer_fences import escape_code_fences_for_display
class TestEscapeCodeFencesForDisplay:
+2 -2
View File
@@ -7,7 +7,7 @@ import pytest
import gateway.run as gateway_run
from gateway.config import HomeChannel, Platform
from gateway.platforms.base import MessageEvent
from gateway.restart import GATEWAY_SERVICE_RESTART_EXIT_CODE
from gateway.restart import DEFAULT_GATEWAY_POST_INTERRUPT_GRACE_TIMEOUT, GATEWAY_SERVICE_RESTART_EXIT_CODE
from gateway.session import build_session_key
from tests.gateway.restart_test_helpers import make_restart_runner, make_restart_source
@@ -217,7 +217,7 @@ def test_post_interrupt_grace_tolerates_duck_typed_runner():
assert (
gateway_run.GatewayRunner._post_interrupt_grace_timeout(runner)
== gateway_run.DEFAULT_GATEWAY_POST_INTERRUPT_GRACE_TIMEOUT
== DEFAULT_GATEWAY_POST_INTERRUPT_GRACE_TIMEOUT
)
@pytest.mark.asyncio
@@ -124,7 +124,7 @@ async def test_secondary_profile_handoff_uses_its_own_adapter(monkeypatch):
used = {}
monkeypatch.setattr(
"gateway.run.resolve_delivery_transport", _spy_transport_factory(used),
"gateway.delivery.resolve_delivery_transport", _spy_transport_factory(used),
)
# The watcher would already be inside _profile_runtime_scope here, so a
# fresh load resolves the secondary's config.
@@ -155,7 +155,7 @@ async def test_default_profile_handoff_keeps_primary_adapter(monkeypatch):
used = {}
monkeypatch.setattr(
"gateway.run.resolve_delivery_transport", _spy_transport_factory(used),
"gateway.delivery.resolve_delivery_transport", _spy_transport_factory(used),
)
await runner._process_handoff(
@@ -181,7 +181,7 @@ async def test_secondary_profile_config_load_failure_fails_closed(monkeypatch):
raise RuntimeError("config.yaml exploded")
monkeypatch.setattr(
"gateway.run.resolve_delivery_transport", _spy_transport_factory(used),
"gateway.delivery.resolve_delivery_transport", _spy_transport_factory(used),
)
monkeypatch.setattr("gateway.run.load_gateway_config", _boom)
@@ -206,7 +206,7 @@ async def test_secondary_profile_without_live_adapters_fails_loudly(monkeypatch)
runner._profile_adapters = {}
monkeypatch.setattr(
"gateway.run.resolve_delivery_transport", _spy_transport_factory({}),
"gateway.delivery.resolve_delivery_transport", _spy_transport_factory({}),
)
with pytest.raises(RuntimeError, match="no live adapters"):
@@ -11,7 +11,7 @@ from __future__ import annotations
import pytest
from gateway.kanban_watchers import _resolve_auto_decompose_settings
from gateway.kanban_watchers_common import _resolve_auto_decompose_settings
def test_enabled_by_default_when_key_absent():
+1 -1
View File
@@ -4,7 +4,7 @@ from pathlib import Path
from gateway.config import Platform
from gateway.kanban_watchers import (
from gateway.kanban_watchers_common import (
_acquire_singleton_lock,
_release_singleton_lock,
)
+1 -1
View File
@@ -12,7 +12,7 @@ from dataclasses import replace
from unittest.mock import AsyncMock, MagicMock
from gateway.config import Platform, PlatformConfig
from gateway.kanban_watchers import _wake_scope_id
from gateway.kanban_watchers_notifier import _wake_scope_id
from gateway.run import GatewayRunner
from gateway.session import build_session_key
from hermes_cli import kanban_db as kb
+2 -2
View File
@@ -292,10 +292,10 @@ def test_gateway_runner_liveness_guards_start_and_stop():
with (
patch(
"gateway.run._arm_loop_floor_timer", return_value=floor_timer
"gateway.shutdown_watchdog._arm_loop_floor_timer", return_value=floor_timer
) as arm_floor,
patch(
"gateway.run.start_loop_liveness_watchdog", return_value=watchdog
"gateway.shutdown_watchdog.start_loop_liveness_watchdog", return_value=watchdog
) as start_watchdog,
):
runner._start_loop_liveness_guards(loop)
+1 -1
View File
@@ -534,7 +534,7 @@ class TestSpawnSupervised:
delegated_child_context,
is_delegated_child_context,
)
from gateway.kanban_watchers import _to_thread_process_service
from gateway.kanban_watchers_common import _to_thread_process_service
from hermes_cli.kanban_db import _assert_not_delegated_child_mutation
with delegated_child_context():
-1
View File
@@ -149,7 +149,6 @@ class TestReasoningCommand:
monkeypatch.setattr(gateway_run, "_hermes_home", hermes_home)
monkeypatch.setattr(gateway_run, "_env_path", hermes_home / ".env")
monkeypatch.setattr(gateway_run, "load_dotenv", lambda *args, **kwargs: None)
monkeypatch.setattr(
gateway_run,
"_resolve_runtime_agent_kwargs",
@@ -77,7 +77,6 @@ async def test_restart_command_uses_atomic_json_writes_for_marker_files(tmp_path
# run.py); it uses that module's top-level atomic_json_write import.
import gateway.slash_commands as gateway_slash
monkeypatch.setattr(gateway_slash, "atomic_json_write", _fake_atomic_json_write)
monkeypatch.setattr(gateway_run, "atomic_json_write", _fake_atomic_json_write)
runner, _adapter = make_restart_runner()
runner.request_restart = MagicMock(return_value=True)
@@ -202,7 +202,8 @@ class TestPeerResolutionRecency:
class TestLoadTranscriptReroutes:
def test_load_transcript_raises_when_message_read_fails(self, tmp_path, monkeypatch):
from gateway.session import SessionStore, TranscriptReadError
from gateway.session import SessionStore
from gateway.session_transcript import TranscriptReadError
from gateway.config import GatewayConfig
+4 -2
View File
@@ -109,7 +109,8 @@ def test_runtime_health_is_sanitized_and_recovers() -> None:
def test_session_store_and_runner_reopen_after_failed_construction(monkeypatch, tmp_path) -> None:
import hermes_state
from gateway.run import GatewayRunner, _SESSION_DB_UNPINNED
from gateway.session import SessionStore, _DB_UNPINNED
from gateway.session import SessionStore
from gateway.session_persistence import _DB_UNPINNED
db_path = tmp_path / "state.db"
clock = _Clock()
@@ -260,7 +261,8 @@ def test_close_all_preserves_inflight_failure() -> None:
def test_recovered_db_rows_survive_fallback_structural_save(monkeypatch, tmp_path) -> None:
import hermes_state
from gateway.config import GatewayConfig, Platform
from gateway.session import SessionEntry, SessionSource, SessionStore, _now
from gateway.session import SessionEntry, SessionSource, SessionStore
from gateway.session_lifecycle import _now
db_path = tmp_path / "state.db"
sessions_dir = tmp_path / "sessions"
@@ -14,7 +14,8 @@ import types
import pytest
from gateway.config import Platform
from gateway.run import GatewayRunner, TurnRunner
from gateway.run import GatewayRunner
from gateway.run_turn_runner import TurnRunner
def _attach(lane):
+3 -3
View File
@@ -117,15 +117,15 @@ class TestSpawnAsyncDiagnostic:
# ---------------------------------------------------------------------------
# _parse_systemd_duration_to_us
# parse_systemd_duration_to_us
# ---------------------------------------------------------------------------
class TestParseSystemdDuration:
def test_seconds(self):
assert sf._parse_systemd_duration_to_us("90s") == 90 * 1_000_000
assert sf.parse_systemd_duration_to_us("90s") == 90 * 1_000_000
def test_minutes(self):
assert sf._parse_systemd_duration_to_us("3min") == 180 * 1_000_000
assert sf.parse_systemd_duration_to_us("3min") == 180 * 1_000_000
# ---------------------------------------------------------------------------
+3 -8
View File
@@ -1,4 +1,4 @@
"""Tests for Signal _markdown_to_signal() formatting.
"""Tests for Signal markdown_to_signal() formatting.
Covers the markdown-to-bodyRanges conversion pipeline: bold, italic,
strikethrough, monospace, code blocks, headings, and — critically — the
@@ -17,13 +17,8 @@ from gateway.platforms.signal_format import markdown_to_signal
# ---------------------------------------------------------------------------
def _m2s(text: str):
"""Shorthand: call the static method and return (plain_text, styles)."""
return SignalAdapter._markdown_to_signal(text)
def test_shared_helper_matches_signal_adapter_wrapper():
text = "🙂 **bold** and `code`"
assert markdown_to_signal(text) == SignalAdapter._markdown_to_signal(text)
"""Shorthand: return (plain_text, styles)."""
return markdown_to_signal(text)
def _style_types(styles: list[str]) -> list[str]:
+1 -9
View File
@@ -2,8 +2,7 @@
The watchdog covers the pre-event-loop window: armed at process entry
(before the gateway package imports — the implementation is the stdlib-only
top-level module ``hermes_startup_watchdog``; ``gateway.startup_watchdog``
is a re-export shim), disarmed once the gateway's asyncio loop is confirmed
top-level module ``hermes_startup_watchdog``), disarmed once the gateway's asyncio loop is confirmed
live. If neither happens within the deadline — and the process shows no CPU
progress, so slow-but-alive schema migrations are exempt — it must dump
diagnostics, record a lifecycle exit, and hard-exit with the service-restart
@@ -84,13 +83,6 @@ class TestContracts:
assert SERVICE_RESTART_EXIT_CODE == GATEWAY_SERVICE_RESTART_EXIT_CODE
def test_gateway_shim_reexports_same_objects(self):
import gateway.startup_watchdog as shim
assert shim.arm_startup_watchdog is arm_startup_watchdog
assert shim.disarm_startup_watchdog is disarm_startup_watchdog
assert shim.kick_startup_watchdog is kick_startup_watchdog
def test_implementation_module_is_stdlib_only(self):
"""Import-lightness is a correctness property (arm-before-imports,
no import-lock dependence at fire time): the implementation module
@@ -97,7 +97,6 @@ def _setup_monkeypatches(monkeypatch, tmp_path):
(tmp_path / "config.yaml").write_text("agent:\n model: test-model\n", encoding="utf-8")
monkeypatch.setattr(gateway_run, "_hermes_home", tmp_path)
monkeypatch.setattr(gateway_run, "_env_path", tmp_path / ".env")
monkeypatch.setattr(gateway_run, "load_dotenv", lambda *args, **kwargs: None)
monkeypatch.setattr(gateway_run, "_load_gateway_config", lambda: {})
monkeypatch.setattr(
gateway_run,
@@ -112,7 +112,8 @@ def _make_runner_and_captured(monkeypatch, run_still_current=True):
progress_queue=None,
_loop_for_step=None,
)
return run_mod.TurnRunner(_StubGatewayRunner(), ctx), captured
from gateway.run_turn_runner import TurnRunner
return TurnRunner(_StubGatewayRunner(), ctx), captured
class TestGatewayFailureNotice:
@@ -187,7 +187,7 @@ class TestBusyHandlerDemotesInterruptForSubagents:
runner._running_agents[sk] = parent
runner.adapters[event.source.platform] = adapter
with patch("gateway.run.merge_pending_message_event"):
with patch("gateway.platforms.base.merge_pending_message_event"):
await runner._handle_active_session_busy_message(event, sk)
parent.interrupt.assert_called_once_with("please stop")
@@ -208,7 +208,7 @@ class TestBusyHandlerDemotesInterruptForSubagents:
runner._running_agents[sk] = parent
runner.adapters[event.source.platform] = adapter
with patch("gateway.run.merge_pending_message_event"):
with patch("gateway.platforms.base.merge_pending_message_event"):
await runner._handle_active_session_busy_message(event, sk)
parent.interrupt.assert_not_called()
@@ -236,7 +236,7 @@ class TestBusyHandlerDemotesInterruptForSubagents:
runner._running_agents[sk] = parent
runner.adapters[event.source.platform] = adapter
with patch("gateway.run.merge_pending_message_event"):
with patch("gateway.platforms.base.merge_pending_message_event"):
await runner._handle_active_session_busy_message(event, sk)
parent.steer.assert_called_once_with("course-correct")
@@ -26,7 +26,8 @@ import sqlite3
import pytest
from gateway.config import GatewayConfig
from gateway.session import SessionStore, TranscriptReadError
from gateway.session import SessionStore
from gateway.session_transcript import TranscriptReadError
@pytest.fixture
@@ -83,7 +84,7 @@ class TestLoadTranscriptReadFailure:
class TestSlashCommandsOnUnreadableTranscript:
def test_history_unreadable_text_is_explicit(self):
from gateway.slash_commands import HISTORY_UNREADABLE
from gateway.slash_commands_status import HISTORY_UNREADABLE
assert "unreadable" in HISTORY_UNREADABLE
assert "not a new conversation" in HISTORY_UNREADABLE
+3 -3
View File
@@ -21,7 +21,7 @@ from gateway.turn_context import TurnContext
def _make_runner(ctx):
from gateway.run import TurnRunner
from gateway.run_turn_runner import TurnRunner
class _StubGatewayRunner:
def _adapter_for_source(self, source):
@@ -51,7 +51,7 @@ class TestTurnContext:
class TestTurnRunner:
def test_methods_exist_and_bind(self):
from gateway.run import TurnRunner
from gateway.run_turn_runner import TurnRunner
ctx = TurnContext()
runner = _make_runner(ctx)
@@ -137,7 +137,7 @@ class TestTurnRunner:
_hooks_ref=SimpleNamespace(loaded_hooks=False),
)
from gateway.run import TurnRunner
from gateway.run_turn_runner import TurnRunner
result = TurnRunner(gateway_runner, ctx).run_sync()
-16
View File
@@ -494,19 +494,3 @@ def test_runner_release_turn_lease_is_token_scoped_and_bare_safe():
_run(scenario())
def test_registry_len_reports_tracked_sessions():
"""``len(registry)`` is public API (plugins/diagnostics size the registry
with it): 0 when empty, one per tracked session_id, and it follows eviction."""
async def scenario():
registry = SessionTurnLeaseRegistry(max_entries=8)
assert len(registry) == 0
a = await registry.acquire("s1", owner_key="k1", generation=1, timeout=1)
assert len(registry) == 1
b = await registry.acquire("s2", owner_key="k2", generation=1, timeout=1)
assert len(registry) == 2
registry.release(a)
registry.release(b)
# Released (idle) entries stay tracked until eviction, matching _leases.
assert len(registry) == len(registry._leases) == 2
_run(scenario())
+2 -2
View File
@@ -53,13 +53,13 @@ class TestShort:
class TestModelSwitchSkewGuard:
def test_guard_returns_none_without_skew(self, monkeypatch):
from gateway import slash_commands
from gateway import slash_commands_model as slash_commands
monkeypatch.setattr(code_skew, "detect_code_skew", lambda: None)
assert slash_commands._model_switch_skew_guard() is None
def test_guard_message_names_revs_and_restart(self, monkeypatch):
from gateway import slash_commands
from gateway import slash_commands_model as slash_commands
monkeypatch.setattr(code_skew, "detect_code_skew", lambda: ("abc1234567", "def4567890"))
msg = slash_commands._model_switch_skew_guard()
+1 -1
View File
@@ -164,7 +164,7 @@ def test_cron_tick_resumes_after_disengage(hermes_home, monkeypatch):
def test_kanban_dispatch_blocked_when_engaged(hermes_home):
from gateway.kanban_watchers import _kanban_dispatch_allowed
from gateway.kanban_watchers_common import _kanban_dispatch_allowed
assert _kanban_dispatch_allowed() is True
estop.engage(reason="test")
+13 -13
View File
@@ -276,35 +276,35 @@ class TestP0ChatLockEviction:
def test_eviction_skips_locked(self):
"""When eviction is needed, locked entries are skipped."""
adapter = YuanbaoAdapter(make_config())
from gateway.platforms.yuanbao import OutboundManager
from gateway.platforms.yuanbao import MessageSender
# Fill to capacity with unlocked locks
for i in range(OutboundManager.CHAT_DICT_MAX_SIZE):
adapter._outbound._chat_locks[f"chat_{i}"] = asyncio.Lock()
for i in range(MessageSender.CHAT_DICT_MAX_SIZE):
adapter._outbound.sender._chat_locks[f"chat_{i}"] = asyncio.Lock()
# Lock the oldest entry
oldest_key = next(iter(adapter._outbound._chat_locks))
oldest_lock = adapter._outbound._chat_locks[oldest_key]
oldest_key = next(iter(adapter._outbound.sender._chat_locks))
oldest_lock = adapter._outbound.sender._chat_locks[oldest_key]
# Simulate a held lock by acquiring it in a non-async way (set _locked)
# asyncio.Lock is not held until actually acquired; so we test the
# method logic by acquiring the first lock manually.
# For a sync test, we check that get_chat_lock doesn't crash.
new_lock = adapter._outbound.get_chat_lock("new_chat")
assert "new_chat" in adapter._outbound._chat_locks
new_lock = adapter._outbound.sender.get_chat_lock("new_chat")
assert "new_chat" in adapter._outbound.sender._chat_locks
assert isinstance(new_lock, asyncio.Lock)
# The oldest unlocked entry should have been evicted
assert len(adapter._outbound._chat_locks) == OutboundManager.CHAT_DICT_MAX_SIZE
assert len(adapter._outbound.sender._chat_locks) == MessageSender.CHAT_DICT_MAX_SIZE
def test_move_to_end_on_access(self):
"""Accessing an existing key moves it to the end (MRU)."""
adapter = YuanbaoAdapter(make_config())
adapter._outbound._chat_locks["a"] = asyncio.Lock()
adapter._outbound._chat_locks["b"] = asyncio.Lock()
adapter._outbound._chat_locks["c"] = asyncio.Lock()
adapter._outbound.sender._chat_locks["a"] = asyncio.Lock()
adapter._outbound.sender._chat_locks["b"] = asyncio.Lock()
adapter._outbound.sender._chat_locks["c"] = asyncio.Lock()
# Access "a" — should move to end
adapter._outbound.get_chat_lock("a")
keys = list(adapter._outbound._chat_locks.keys())
adapter._outbound.sender.get_chat_lock("a")
keys = list(adapter._outbound.sender._chat_locks.keys())
assert keys[-1] == "a"
assert keys[0] == "b"
+12 -11
View File
@@ -16,6 +16,7 @@ import unittest
# Ensure project root is on the path
sys.path.insert(0, os.path.join(os.path.dirname(__file__), '..'))
from gateway.platforms import helpers as _mdchunk
from gateway.platforms.yuanbao import MarkdownProcessor
@@ -23,10 +24,10 @@ from gateway.platforms.yuanbao import MarkdownProcessor
class TestHasUnclosedFence(unittest.TestCase):
def test_unclosed_fence(self):
self.assertTrue(MarkdownProcessor.has_unclosed_fence("```python\ncode"))
self.assertTrue(_mdchunk.text_has_unclosed_fence("```python\ncode"))
def test_closed_fence(self):
self.assertFalse(MarkdownProcessor.has_unclosed_fence("```python\ncode\n```"))
self.assertFalse(_mdchunk.text_has_unclosed_fence("```python\ncode\n```"))
@@ -35,25 +36,25 @@ class TestHasUnclosedFence(unittest.TestCase):
def test_inline_backtick_ignored(self):
text = "`inline code` is fine"
self.assertFalse(MarkdownProcessor.has_unclosed_fence(text))
self.assertFalse(_mdchunk.text_has_unclosed_fence(text))
# ============ ends_with_table_row ============
class TestEndsWithTableRow(unittest.TestCase):
def test_simple_table_row(self):
self.assertTrue(MarkdownProcessor.ends_with_table_row("| col1 | col2 |"))
self.assertTrue(_mdchunk.text_ends_with_table_row("| col1 | col2 |"))
def test_table_row_in_middle(self):
text = "| col1 | col2 |\nsome other text"
self.assertFalse(MarkdownProcessor.ends_with_table_row(text))
self.assertFalse(_mdchunk.text_ends_with_table_row(text))
def test_table_separator_row(self):
self.assertTrue(MarkdownProcessor.ends_with_table_row("| --- | --- |"))
self.assertTrue(_mdchunk.text_ends_with_table_row("| --- | --- |"))
@@ -62,13 +63,13 @@ class TestEndsWithTableRow(unittest.TestCase):
class TestSplitAtParagraphBoundary(unittest.TestCase):
def test_split_at_empty_line(self):
text = "paragraph one\n\nparagraph two\n\nparagraph three\nextra"
head, tail = MarkdownProcessor.split_at_paragraph_boundary(text, 30)
head, tail = _mdchunk.split_at_paragraph_boundary(text, 30)
self.assertLessEqual(len(head), 30)
self.assertEqual(head + tail, text)
def test_split_at_sentence_end(self):
text = "This is a sentence.\nNext line.\nAnother line."
head, tail = MarkdownProcessor.split_at_paragraph_boundary(text, 25)
head, tail = _mdchunk.split_at_paragraph_boundary(text, 25)
self.assertLessEqual(len(head), 25)
self.assertEqual(head + tail, text)
@@ -76,7 +77,7 @@ class TestSplitAtParagraphBoundary(unittest.TestCase):
def test_chinese_sentence_boundary(self):
text = "这是第一句话。\n这是第二句话。\n这是第三句话。"
head, tail = MarkdownProcessor.split_at_paragraph_boundary(text, 15)
head, tail = _mdchunk.split_at_paragraph_boundary(text, 15)
self.assertLessEqual(len(head), 15)
self.assertEqual(head + tail, text)
@@ -106,7 +107,7 @@ class TestChunkMarkdownText(unittest.TestCase):
text = "Some intro text.\n\n" + table + "\n\nSome outro text."
result = MarkdownProcessor.chunk_markdown_text(text, 3000)
for chunk in result:
self.assertFalse(MarkdownProcessor.has_unclosed_fence(chunk))
self.assertFalse(_mdchunk.text_has_unclosed_fence(chunk))
def test_multiple_paragraphs(self):
@@ -165,7 +166,7 @@ def test_large_fence_kept_whole():
# 代码块应在同一个 chunk 中(允许超出 max_chars)
fence_chunks = [c for c in chunks if "```python" in c]
for c in fence_chunks:
assert not MarkdownProcessor.has_unclosed_fence(c)
assert not _mdchunk.text_has_unclosed_fence(c)
+5 -6
View File
@@ -1,7 +1,7 @@
"""test_yuanbao_reconnect_set_active.py - Verify _do_reconnect restores the active singleton.
Regression test for #58363: after a WS disconnect/reconnect cycle,
``get_active_adapter()`` must return the live adapter (not ``None``).
``YuanbaoAdapter.get_active()`` must return the live adapter (not ``None``).
The original ``_do_reconnect()`` succeeded but never called
``YuanbaoAdapter.set_active()``, leaving the singleton permanently
``None`` until a full gateway restart.
@@ -20,7 +20,6 @@ import pytest
from gateway.platforms.yuanbao import (
YuanbaoAdapter,
ConnectionManager,
get_active_adapter,
)
@@ -65,7 +64,7 @@ async def test_do_reconnect_calls_set_active_on_success():
):
# Clear any existing active instance
YuanbaoAdapter.set_active(None)
assert get_active_adapter() is None
assert YuanbaoAdapter.get_active() is None
# Run reconnect
result = await cm._do_reconnect()
@@ -74,7 +73,7 @@ async def test_do_reconnect_calls_set_active_on_success():
assert result is True
# After successful reconnect, get_active() must return the adapter
assert get_active_adapter() is adapter
assert YuanbaoAdapter.get_active() is adapter
@pytest.mark.asyncio
@@ -94,7 +93,7 @@ async def test_do_reconnect_does_not_set_active_on_failure():
):
# Clear any existing active instance
YuanbaoAdapter.set_active(None)
assert get_active_adapter() is None
assert YuanbaoAdapter.get_active() is None
# Run reconnect - should fail
result = await cm._do_reconnect()
@@ -103,4 +102,4 @@ async def test_do_reconnect_does_not_set_active_on_failure():
assert result is False
# get_active() should still be None
assert get_active_adapter() is None
assert YuanbaoAdapter.get_active() is None
@@ -72,7 +72,7 @@ async def test_in_process_scoped_transport_contract_finishes_headlessly(
PlatformConfig(enabled=True, extra={"key": "target-peer-key-1234567890"})
)
target._run_idempotency_store.close()
from gateway.platforms.api_server import RunIdempotencyStore
from gateway.platforms.api_server_run_idempotency import RunIdempotencyStore
target._run_idempotency_store = RunIdempotencyStore(
str(tmp_path / "target-runs.db")
+3 -3
View File
@@ -628,13 +628,13 @@ async def _send_qqbot(pconfig, chat_id, message):
async def _send_yuanbao(chat_id, message, media_files=None):
"""Send via the running Yuanbao adapter's persistent WebSocket (no throwaway client possible)."""
try:
from gateway.platforms.yuanbao import get_active_adapter, send_yuanbao_direct
from gateway.platforms.yuanbao import YuanbaoAdapter
except ImportError:
return _error("Yuanbao adapter module not available.")
adapter = get_active_adapter()
adapter = YuanbaoAdapter.get_active()
if adapter is None:
return _error("Yuanbao adapter is not running. Start the gateway with yuanbao platform enabled first.")
try:
return await send_yuanbao_direct(adapter, chat_id, message, media_files=media_files)
return await adapter._outbound.sender.send_direct(chat_id, message, media_files)
except Exception as e:
return _error(f"Yuanbao send failed: {e}")
+3 -3
View File
@@ -3,7 +3,7 @@
get_group_info / query_group_members / search_sticker / send_sticker / send_dm. Sticker
flow mirrors chatbot-web's sticker-search/sticker-send: the LLM should search_sticker for
a sticker_id (or pass the Chinese name), then send_sticker — never bare Unicode emoji.
The active adapter singleton lives in ``gateway.platforms.yuanbao.get_active_adapter``.
The active adapter singleton lives in ``gateway.platforms.yuanbao.YuanbaoAdapter.get_active``.
"""
from __future__ import annotations
@@ -56,8 +56,8 @@ def _yb_tool(label: str):
def _get_active_adapter():
"""Lazy import to avoid ImportError when gateway.platforms.yuanbao is unavailable."""
with suppress(ImportError):
from gateway.platforms.yuanbao import get_active_adapter
return get_active_adapter()
from gateway.platforms.yuanbao import YuanbaoAdapter
return YuanbaoAdapter.get_active()
return None
+1 -1
View File
@@ -148,7 +148,7 @@ def _room_link_run_storage_durable() -> bool:
if store is None:
# This process does not construct the API adapter that owns the store; open the
# same shared SQLite store lazily so negotiation reflects the real replay boundary.
from gateway.platforms.api_server import RunIdempotencyStore
from gateway.platforms.api_server_run_idempotency import RunIdempotencyStore
with _run_store_lock:
store = getattr(_bound_server, "_run_idempotency_store", None)
if store is None: