Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
46 changes: 45 additions & 1 deletion gateway/run.py
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@
"""

import asyncio
import dataclasses
import json
import logging
import os
Expand Down Expand Up @@ -3145,7 +3146,50 @@ async def _handle_message(self, event: MessageEvent) -> Optional[str]:

# Internal events (e.g. background-process completion notifications)
# are system-generated and must skip user authorization.
if getattr(event, "internal", False):
is_internal = bool(getattr(event, "internal", False))

# Fire pre_gateway_dispatch plugin hook for user-originated messages.
# Plugins receive the MessageEvent and may return a dict influencing flow:
# {"action": "skip", "reason": ...} -> drop (no reply, plugin handled)
# {"action": "rewrite", "text": ...} -> replace event.text, continue
# {"action": "allow"} / None -> normal dispatch
# Hook runs BEFORE auth so plugins can handle unauthorized senders
# (e.g. customer handover ingest) without triggering the pairing flow.
if not is_internal:
try:
from hermes_cli.plugins import invoke_hook as _invoke_hook
_hook_results = _invoke_hook(
"pre_gateway_dispatch",
event=event,
gateway=self,
session_store=self.session_store,
)
except Exception as _hook_exc:
logger.warning("pre_gateway_dispatch invocation failed: %s", _hook_exc)
_hook_results = []

for _result in _hook_results:
if not isinstance(_result, dict):
continue
_action = _result.get("action")
if _action == "skip":
logger.info(
"pre_gateway_dispatch skip: reason=%s platform=%s chat=%s",
_result.get("reason"),
source.platform.value if source.platform else "unknown",
source.chat_id or "unknown",
)
return None
if _action == "rewrite":
_new_text = _result.get("text")
if isinstance(_new_text, str):
event = dataclasses.replace(event, text=_new_text)
source = event.source
break
if _action == "allow":
break

if is_internal:
pass
elif source.user_id is None:
# Messages with no user identity (Telegram service messages,
Expand Down
8 changes: 8 additions & 0 deletions hermes_cli/plugins.py
Original file line number Diff line number Diff line change
Expand Up @@ -71,6 +71,14 @@
"on_session_finalize",
"on_session_reset",
"subagent_stop",
# Gateway pre-dispatch hook. Fired once per incoming MessageEvent
# after the internal-event guard but BEFORE auth/pairing and agent
# dispatch. Plugins may return a dict to influence flow:
# {"action": "skip", "reason": "..."} -> drop message (no reply)
# {"action": "rewrite", "text": "..."} -> replace event.text, continue
# {"action": "allow"} / None -> normal dispatch
# Kwargs: event: MessageEvent, gateway: GatewayRunner, session_store.
"pre_gateway_dispatch",
}

ENTRY_POINTS_GROUP = "hermes_agent.plugins"
Expand Down
1 change: 1 addition & 0 deletions scripts/release.py
Original file line number Diff line number Diff line change
Expand Up @@ -60,6 +60,7 @@
"abner.the.foreman@agentmail.to": "Abnertheforeman",
"harryykyle1@gmail.com": "hharry11",
"kshitijk4poor@gmail.com": "kshitijk4poor",
"keira.voss94@gmail.com": "keiravoss94",
"16443023+stablegenius49@users.noreply.github.com": "stablegenius49",
"185121704+stablegenius49@users.noreply.github.com": "stablegenius49",
"101283333+batuhankocyigit@users.noreply.github.com": "batuhankocyigit",
Expand Down
179 changes: 179 additions & 0 deletions tests/gateway/test_pre_gateway_dispatch.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,179 @@
"""Tests for the pre_gateway_dispatch plugin hook.

The hook allows plugins to intercept incoming messages before auth and
agent dispatch. It runs in _handle_message and acts on returned action
dicts: {"action": "skip"|"rewrite"|"allow"}.
"""

from types import SimpleNamespace
from unittest.mock import AsyncMock, MagicMock

import pytest

from gateway.config import GatewayConfig, Platform, PlatformConfig
from gateway.platforms.base import MessageEvent
from gateway.session import SessionSource


def _clear_auth_env(monkeypatch) -> None:
for key in (
"TELEGRAM_ALLOWED_USERS",
"WHATSAPP_ALLOWED_USERS",
"GATEWAY_ALLOWED_USERS",
"TELEGRAM_ALLOW_ALL_USERS",
"WHATSAPP_ALLOW_ALL_USERS",
"GATEWAY_ALLOW_ALL_USERS",
):
monkeypatch.delenv(key, raising=False)


def _make_event(text: str = "hello", platform: Platform = Platform.WHATSAPP) -> MessageEvent:
return MessageEvent(
text=text,
message_id="m1",
source=SessionSource(
platform=platform,
user_id="15551234567@s.whatsapp.net",
chat_id="15551234567@s.whatsapp.net",
user_name="tester",
chat_type="dm",
),
)


def _make_runner(platform: Platform):
from gateway.run import GatewayRunner

config = GatewayConfig(
platforms={platform: PlatformConfig(enabled=True)},
)
runner = object.__new__(GatewayRunner)
runner.config = config
adapter = SimpleNamespace(send=AsyncMock())
runner.adapters = {platform: adapter}
runner.pairing_store = MagicMock()
runner.pairing_store.is_approved.return_value = False
runner.pairing_store._is_rate_limited.return_value = False
runner.session_store = MagicMock()
runner._running_agents = {}
runner._update_prompt_pending = {}
return runner, adapter


@pytest.mark.asyncio
async def test_hook_skip_short_circuits_dispatch(monkeypatch):
"""A plugin returning {'action': 'skip'} drops the message before auth."""
_clear_auth_env(monkeypatch)

def _fake_hook(name, **kwargs):
if name == "pre_gateway_dispatch":
return [{"action": "skip", "reason": "plugin-handled"}]
return []

monkeypatch.setattr("hermes_cli.plugins.invoke_hook", _fake_hook)

runner, adapter = _make_runner(Platform.WHATSAPP)

result = await runner._handle_message(_make_event("hi"))

assert result is None
adapter.send.assert_not_awaited()
runner.pairing_store.generate_code.assert_not_called()


@pytest.mark.asyncio
async def test_hook_rewrite_replaces_event_text(monkeypatch):
"""A plugin returning {'action': 'rewrite', 'text': ...} mutates event.text."""
_clear_auth_env(monkeypatch)
monkeypatch.setenv("WHATSAPP_ALLOWED_USERS", "*")

seen_text = {}

def _fake_hook(name, **kwargs):
if name == "pre_gateway_dispatch":
return [{"action": "rewrite", "text": "REWRITTEN"}]
return []

async def _capture(event, source, _quick_key, _run_generation):
seen_text["value"] = event.text
return "ok"

monkeypatch.setattr("hermes_cli.plugins.invoke_hook", _fake_hook)

runner, _adapter = _make_runner(Platform.WHATSAPP)
runner._handle_message_with_agent = _capture # noqa: SLF001

await runner._handle_message(_make_event("original"))

assert seen_text.get("value") == "REWRITTEN"


@pytest.mark.asyncio
async def test_hook_allow_falls_through_to_auth(monkeypatch):
"""A plugin returning {'action': 'allow'} continues to normal dispatch."""
_clear_auth_env(monkeypatch)
# No allowed users set → auth fails → pairing flow triggers.
monkeypatch.delenv("WHATSAPP_ALLOWED_USERS", raising=False)

def _fake_hook(name, **kwargs):
if name == "pre_gateway_dispatch":
return [{"action": "allow"}]
return []

monkeypatch.setattr("hermes_cli.plugins.invoke_hook", _fake_hook)

runner, adapter = _make_runner(Platform.WHATSAPP)
runner.pairing_store.generate_code.return_value = "12345"

result = await runner._handle_message(_make_event("hi"))

# auth chain ran → pairing code was generated
assert result is None
runner.pairing_store.generate_code.assert_called_once()


@pytest.mark.asyncio
async def test_hook_exception_does_not_break_dispatch(monkeypatch):
"""A raising plugin hook does not break the gateway."""
_clear_auth_env(monkeypatch)
monkeypatch.delenv("WHATSAPP_ALLOWED_USERS", raising=False)

def _fake_hook(name, **kwargs):
raise RuntimeError("plugin blew up")

monkeypatch.setattr("hermes_cli.plugins.invoke_hook", _fake_hook)

runner, _adapter = _make_runner(Platform.WHATSAPP)
runner.pairing_store.generate_code.return_value = None

# Should not raise; falls through to auth chain.
result = await runner._handle_message(_make_event("hi"))
assert result is None


@pytest.mark.asyncio
async def test_internal_events_bypass_hook(monkeypatch):
"""Internal events (event.internal=True) skip the plugin hook entirely."""
_clear_auth_env(monkeypatch)
monkeypatch.setenv("WHATSAPP_ALLOWED_USERS", "*")

called = {"count": 0}

def _fake_hook(name, **kwargs):
called["count"] += 1
return [{"action": "skip"}]

async def _capture(event, source, _quick_key, _run_generation):
return "ok"

monkeypatch.setattr("hermes_cli.plugins.invoke_hook", _fake_hook)

runner, _adapter = _make_runner(Platform.WHATSAPP)
runner._handle_message_with_agent = _capture # noqa: SLF001

event = _make_event("hi")
event.internal = True

# Even though the hook would say skip, internal events bypass it.
await runner._handle_message(event)
assert called["count"] == 0
27 changes: 27 additions & 0 deletions tests/hermes_cli/test_plugins.py
Original file line number Diff line number Diff line change
Expand Up @@ -330,6 +330,33 @@ def test_valid_hooks_include_request_scoped_api_hooks(self):
assert "transform_terminal_output" in VALID_HOOKS
assert "transform_tool_result" in VALID_HOOKS

def test_valid_hooks_include_pre_gateway_dispatch(self):
assert "pre_gateway_dispatch" in VALID_HOOKS

def test_pre_gateway_dispatch_collects_action_dicts(self, tmp_path, monkeypatch):
"""pre_gateway_dispatch callbacks return action dicts (skip/rewrite/allow)."""
plugins_dir = tmp_path / "hermes_test" / "plugins"
_make_plugin_dir(
plugins_dir, "predispatch_plugin",
register_body=(
'ctx.register_hook("pre_gateway_dispatch", '
'lambda **kw: {"action": "skip", "reason": "test"})'
),
)
monkeypatch.setenv("HERMES_HOME", str(tmp_path / "hermes_test"))

mgr = PluginManager()
mgr.discover_and_load()

results = mgr.invoke_hook(
"pre_gateway_dispatch",
event=object(),
gateway=object(),
session_store=object(),
)
assert len(results) == 1
assert results[0] == {"action": "skip", "reason": "test"}

def test_register_and_invoke_hook(self, tmp_path, monkeypatch):
"""Registered hooks are called on invoke_hook()."""
plugins_dir = tmp_path / "hermes_test" / "plugins"
Expand Down
63 changes: 63 additions & 0 deletions website/docs/user-guide/features/hooks.md
Original file line number Diff line number Diff line change
Expand Up @@ -247,6 +247,7 @@ def register(ctx):
| [`on_session_finalize`](#on_session_finalize) | CLI/gateway tears down an active session (flush, save, stats) | ignored |
| [`on_session_reset`](#on_session_reset) | Gateway swaps in a fresh session key (e.g. `/new`, `/reset`) | ignored |
| [`subagent_stop`](#subagent_stop) | A `delegate_task` child has exited | ignored |
| [`pre_gateway_dispatch`](#pre_gateway_dispatch) | Gateway received a user message, before auth + dispatch | `{"action": "skip" \| "rewrite" \| "allow", ...}` to influence flow |

---

Expand Down Expand Up @@ -708,6 +709,68 @@ With heavy delegation (e.g. orchestrator roles × 5 leaves × nested depth), `su

---

### `pre_gateway_dispatch`

Fires **once per incoming `MessageEvent`** in the gateway, after the internal-event guard but **before** auth/pairing and agent dispatch. This is the interception point for gateway-level message-flow policies (listen-only windows, human handover, per-chat routing, etc.) that don't fit cleanly into any single platform adapter.

**Callback signature:**

```python
def my_callback(event, gateway, session_store, **kwargs):
```

| Parameter | Type | Description |
|-----------|------|-------------|
| `event` | `MessageEvent` | The normalized inbound message (has `.text`, `.source`, `.message_id`, `.internal`, etc.). |
| `gateway` | `GatewayRunner` | The active gateway runner, so plugins can call `gateway.adapters[platform].send(...)` for side-channel replies (owner notifications, etc.). |
| `session_store` | `SessionStore` | For silent transcript ingestion via `session_store.append_to_transcript(...)`. |

**Fires:** In `gateway/run.py`, inside `GatewayRunner._handle_message()`, immediately after `is_internal` is computed. **Internal events skip the hook entirely** (they are system-generated — background-process completions, etc. — and must not be gate-kept by user-facing policy).

**Return value:** `None` or a dict. The first recognized action dict wins; remaining plugin results are ignored. Exceptions in plugin callbacks are caught and logged; the gateway always falls through to normal dispatch on error.

| Return | Effect |
|--------|--------|
| `{"action": "skip", "reason": "..."}` | Drop the message — no agent reply, no pairing flow, no auth. Plugin is assumed to have handled it (e.g. silent-ingested into the transcript). |
| `{"action": "rewrite", "text": "new text"}` | Replace `event.text`, then continue normal dispatch with the modified event. Useful for collapsing buffered ambient messages into a single prompt. |
| `{"action": "allow"}` / `None` | Normal dispatch — runs the full auth / pairing / agent-loop chain. |

**Use cases:** Listen-only group chats (only respond when tagged; buffer ambient messages into context); human handover (silent-ingest customer messages while owner handles the chat manually); per-profile rate limiting; policy-driven routing.

**Example — drop unauthorized DMs silently without triggering the pairing code:**

```python
def deny_unauthorized_dms(event, **kwargs):
src = event.source
if src.chat_type == "dm" and not _is_approved_user(src.user_id):
return {"action": "skip", "reason": "unauthorized-dm"}
return None

def register(ctx):
ctx.register_hook("pre_gateway_dispatch", deny_unauthorized_dms)
```

**Example — rewrite an ambient-message buffer into a single prompt on mention:**

```python
_buffers = {}

def buffer_or_rewrite(event, **kwargs):
key = (event.source.platform, event.source.chat_id)
buf = _buffers.setdefault(key, [])
if _bot_mentioned(event.text):
combined = "\n".join(buf + [event.text])
buf.clear()
return {"action": "rewrite", "text": combined}
buf.append(event.text)
return {"action": "skip", "reason": "ambient-buffered"}

def register(ctx):
ctx.register_hook("pre_gateway_dispatch", buffer_or_rewrite)
```

---

## Shell Hooks

Declare shell-script hooks in your `cli-config.yaml` and Hermes will run them as subprocesses whenever the corresponding plugin-hook event fires — in both CLI and gateway sessions. No Python plugin authoring required.
Expand Down
1 change: 1 addition & 0 deletions website/docs/user-guide/features/plugins.md
Original file line number Diff line number Diff line change
Expand Up @@ -141,6 +141,7 @@ Plugins can register callbacks for these lifecycle events. See the **[Event Hook
| [`post_llm_call`](/docs/user-guide/features/hooks#post_llm_call) | Once per turn, after the LLM loop (successful turns only) |
| [`on_session_start`](/docs/user-guide/features/hooks#on_session_start) | New session created (first turn only) |
| [`on_session_end`](/docs/user-guide/features/hooks#on_session_end) | End of every `run_conversation` call + CLI exit handler |
| [`pre_gateway_dispatch`](/docs/user-guide/features/hooks#pre_gateway_dispatch) | Gateway received a user message, before auth + dispatch. Return `{"action": "skip" \| "rewrite" \| "allow", ...}` to influence flow. |

## Plugin types

Expand Down
Loading