Merge branch 'simp/r3-29-F' into simp/integration3

This commit is contained in:
Teknium
2026-09-02 22:41:06 -07:00
5 changed files with 148 additions and 185 deletions
+19 -50
View File
@@ -1,22 +1,13 @@
"""User OAuth helper for the Google Chat gateway adapter.
Google Chat's ``media.upload`` rejects service-account auth, so for native file
attachments each user grants the bot ``chat.messages.create`` ONCE in their own
DM; the bot stores per-user refresh tokens and uploads *as the user*
(https://developers.google.com/chat/api/guides/auth/users).
Library: load_user_credentials(email=None), refresh_or_none(creds, email=None),
build_user_chat_service(creds), list_authorized_emails().
CLI (driven by ``/setup-files``): --check | --client-secret PATH | --auth-url |
--auth-code CODE | --revoke | --install-deps [--email EMAIL] (legacy single-user
mode when --email is omitted).
Token storage layout (all under ``${HERMES_HOME}``):
- Per-user tokens: ``google_chat_user_tokens/<sanitized_email>.json``
- Legacy single-user: ``google_chat_user_token.json``
- Per-user pending PKCE: ``google_chat_user_oauth_pending/<sanitized_email>.json``
- Legacy pending state: ``google_chat_user_oauth_pending.json``
- OAuth client secret: ``google_chat_user_client_secret.json``
attachments each user grants the bot ``chat.messages.create`` ONCE in their own DM;
the bot stores per-user refresh tokens and uploads *as the user*
(https://developers.google.com/chat/api/guides/auth/users). Library API for the
adapter plus a CLI driven by ``/setup-files`` (``--help``; ``--email`` omitted ==
legacy single-user mode). Files under ``${HERMES_HOME}``: ``google_chat_user_tokens/
<email>.json`` (per-user) / ``google_chat_user_token.json`` (legacy); pending PKCE state
in ``google_chat_user_oauth_pending[/<email>].json``; ``google_chat_user_client_secret.json``.
"""
from __future__ import annotations
@@ -26,7 +17,6 @@ import json
import logging
import os
import re
import secrets
import stat
import sys
from importlib.metadata import version as _distribution_version
@@ -36,7 +26,7 @@ from typing import Any, List, NoReturn, Optional, Tuple
from packaging.requirements import Requirement
from hermes_constants import display_hermes_home, get_hermes_home
from utils import atomic_replace
from utils import atomic_write_text
# Pinned legacy logger name so operator log filters keep matching (see adapter.py).
logger = logging.getLogger("gateway.platforms.google_chat_user_oauth")
@@ -65,11 +55,6 @@ _REQUIRED_PACKAGES = [
_REDIRECT_URI = "http://localhost:1"
def _hermes_home() -> Path:
"""Resolve HERMES_HOME at call time (late-binding for tests / profile switches)."""
return get_hermes_home()
def _sanitize_email(email: str) -> str:
cleaned = _EMAIL_FS_RE.sub("_", (email or "").strip().lower())
return cleaned or "_unknown_"
@@ -81,27 +66,25 @@ def _token_rel(email: Optional[str]) -> str:
def _user_tokens_dir() -> Path:
return _hermes_home() / "google_chat_user_tokens"
return get_hermes_home() / "google_chat_user_tokens"
def _token_path(email: Optional[str] = None) -> Path:
"""Per-user token path for ``email``, or the legacy single-user path."""
return _hermes_home() / _token_rel(email)
return get_hermes_home() / _token_rel(email)
def _client_secret_path() -> Path:
return _hermes_home() / "google_chat_user_client_secret.json"
return get_hermes_home() / "google_chat_user_client_secret.json"
def _pending_auth_path(email: Optional[str] = None) -> Path:
if email:
return _hermes_home() / "google_chat_user_oauth_pending" / f"{_sanitize_email(email)}.json"
return _hermes_home() / "google_chat_user_oauth_pending.json"
return get_hermes_home() / "google_chat_user_oauth_pending" / f"{_sanitize_email(email)}.json"
return get_hermes_home() / "google_chat_user_oauth_pending.json"
# =============================================================================
# Library API — called from the adapter at runtime
# =============================================================================
# -- Library API — called from the adapter at runtime -------------------------
def _refresh_and_persist(creds: Any, token_path: Path, request_cls: Any, *, failure_msg: str) -> Optional[Any]:
@@ -192,9 +175,7 @@ def _persist_credentials(creds: Any, token_path: Path) -> None:
logger.debug("[google_chat_user_oauth] failed to persist credentials at %s", token_path, exc_info=True)
# =============================================================================
# CLI commands — driven by the agent via /setup-files
# =============================================================================
# -- CLI commands — driven by the agent via /setup-files ----------------------
def _normalize_authorized_user_payload(payload: dict) -> dict:
@@ -213,24 +194,12 @@ def _chmod_quiet(path: Path, mode: int) -> None:
def _write_private_json(path: Path, data: Any) -> None:
"""Atomically write JSON with 0o600 permissions where supported."""
"""Atomically write JSON with 0o600 permissions (0o700 parent) where supported."""
path.parent.mkdir(parents=True, exist_ok=True)
_chmod_quiet(path.parent, 0o700)
tmp_path = path.with_suffix(f".tmp.{os.getpid()}.{secrets.token_hex(4)}")
try:
fd = os.open(str(tmp_path), os.O_WRONLY | os.O_CREAT | os.O_EXCL, stat.S_IRUSR | stat.S_IWUSR)
with os.fdopen(fd, "w", encoding="utf-8") as fh:
json.dump(data, fh, indent=2, ensure_ascii=False)
fh.flush()
os.fsync(fh.fileno())
atomic_replace(tmp_path, path)
_chmod_quiet(path, stat.S_IRUSR | stat.S_IWUSR)
finally:
try:
if tmp_path.exists():
tmp_path.unlink()
except OSError:
pass
# mkstemp's 0o600 temp + atomic rename never exposes the token at process umask.
atomic_write_text(path, json.dumps(data, indent=2, ensure_ascii=False), create_mode=0o600)
_chmod_quiet(path, stat.S_IRUSR | stat.S_IWUSR)
def _fail(*lines: str) -> NoReturn:
+47 -49
View File
@@ -30,6 +30,15 @@ _START_INSTRUCTIONS = (
"`/setup-files <PASTE_URL>` (or just the `code=...` value).\n\n"
"Tip: the URL contains your access grant — keep it private."
)
_START_EXIT_TEXT = (
"❌ Couldn't generate the OAuth URL. Check the gateway logs and verify the client_secret.json is valid."
)
_EXCHANGE_EXIT_TEXT = (
"❌ Token exchange failed. The code may have expired or the URL is malformed. "
"Send `/setup-files start` to get a fresh OAuth URL."
)
_REVOKE_EXIT_OUTPUT = "Revoke completed (some steps may have been skipped)."
_EXITED = object() # _run_helper marker: helper called sys.exit but the step tolerates it
async def _run_captured(fn: Callable[..., Any], *args: Any) -> str:
@@ -70,6 +79,33 @@ async def handle_setup_files_command(
except Exception:
logger.debug("[GoogleChat] /setup-files reply send failed", exc_info=True)
async def _run_helper(step: str, exit_text: Optional[str], fn: Callable[..., Any], *args: Any):
"""Captured helper output; ``None`` after replying on failure. ``exit_text``
is the reply on ``SystemExit`` (the helpers' failure signal); ``None``
tolerates the exit and returns ``_EXITED``."""
try:
return await _run_captured(fn, *args)
except SystemExit:
if exit_text is None:
return _EXITED
await _reply(exit_text)
except Exception as exc:
logger.warning("[GoogleChat] /setup-files %s failed: %s", step, exc)
await _reply(f"❌ Error{' revoking' if step == 'revoke' else ''}: {exc}")
return None
def _set_user_creds(creds: Any, api: Any) -> None:
"""Set (or evict, with ``None``) only the sender's slot: Bob revoking must not
break Alice's per-user token nor the shared legacy fallback."""
if not sender_key:
adapter._user_credentials, adapter._user_chat_api = creds, api
elif creds is None:
adapter._user_creds_by_email.pop(sender_key, None)
adapter._user_chat_api_by_email.pop(sender_key, None)
else:
adapter._user_creds_by_email[sender_key] = creds
adapter._user_chat_api_by_email[sender_key] = api
if not arg:
client_secret_present = oauth_helper._client_secret_path().exists()
token_path = oauth_helper._token_path(sender_key)
@@ -96,75 +132,37 @@ async def handle_setup_files_command(
"`/setup-files` (no args) for setup instructions."
)
return True
try:
output = await _run_captured(oauth_helper.get_auth_url, sender_key)
auth_url = output.strip().splitlines()[-1]
except SystemExit:
await _reply(
"❌ Couldn't generate the OAuth URL. Check the gateway logs and verify "
"the client_secret.json is valid."
)
return True
except Exception as exc:
logger.warning("[GoogleChat] /setup-files start failed: %s", exc)
await _reply(f"❌ Error: {exc}")
return True
await _reply(_START_INSTRUCTIONS.format(auth_url=auth_url))
output = await _run_helper("start", _START_EXIT_TEXT, oauth_helper.get_auth_url, sender_key)
if output is not None:
await _reply(_START_INSTRUCTIONS.format(auth_url=output.strip().splitlines()[-1]))
return True
if arg == "revoke":
try:
output = (await _run_captured(oauth_helper.revoke, sender_key)).strip() or "Revoked."
except SystemExit:
output = "Revoke completed (some steps may have been skipped)."
except Exception as exc:
logger.warning("[GoogleChat] /setup-files revoke failed: %s", exc)
await _reply(f"❌ Error revoking: {exc}")
output = await _run_helper("revoke", None, oauth_helper.revoke, sender_key)
if output is None:
return True
# Evict only the sender's slot: Bob revoking must not break Alice's
# per-user token nor the shared legacy fallback.
if sender_key:
adapter._user_creds_by_email.pop(sender_key, None)
adapter._user_chat_api_by_email.pop(sender_key, None)
else:
adapter._user_credentials = None
adapter._user_chat_api = None
output = _REVOKE_EXIT_OUTPUT if output is _EXITED else (output.strip() or "Revoked.")
_set_user_creds(None, None)
await _reply(f"✅ Done.\n```\n{output}\n```")
return True
# Anything else is the auth code or the pasted failed-redirect URL.
try:
output = (await _run_captured(oauth_helper.exchange_auth_code, arg, sender_key)).strip()
except SystemExit:
await _reply(
"❌ Token exchange failed. The code may have expired or the URL is malformed. "
"Send `/setup-files start` to get a fresh OAuth URL."
)
output = await _run_helper("exchange", _EXCHANGE_EXIT_TEXT, oauth_helper.exchange_auth_code, arg, sender_key)
if output is None:
return True
except Exception as exc:
logger.warning("[GoogleChat] /setup-files exchange failed: %s", exc)
await _reply(f"❌ Error: {exc}")
return True
# Re-load credentials so the next file send uses them without a gateway restart.
try:
new_creds = await asyncio.to_thread(oauth_helper.load_user_credentials, sender_key)
if new_creds is not None:
new_api = await asyncio.to_thread(lambda: oauth_helper.build_user_chat_service(new_creds))
if sender_key:
adapter._user_creds_by_email[sender_key] = new_creds
adapter._user_chat_api_by_email[sender_key] = new_api
else:
adapter._user_credentials = new_creds
adapter._user_chat_api = new_api
_set_user_creds(new_creds, new_api)
await _reply("✅ Authorized! Native attachment delivery is now active. Try asking me to send you a PDF.")
return True
except Exception as exc:
logger.warning("[GoogleChat] post-exchange creds load failed: %s", exc)
await _reply(
"⚠️ Token exchanged but the gateway couldn't load the new credentials in-memory. "
f"Restart the gateway and the token at `{oauth_helper._token_path(sender_key)}` will be picked up.\n"
f"Helper output:\n```\n{output}\n```"
f"Helper output:\n```\n{output.strip()}\n```"
)
return True
+29 -42
View File
@@ -14,7 +14,7 @@ import os
import time
import uuid
from datetime import datetime
from typing import Any, Callable, Dict, Optional, Set
from typing import Any, Dict, Optional, Set
try:
import aiohttp
@@ -52,38 +52,21 @@ def _triggered(val: str) -> str:
return "triggered" if val == "on" else "cleared"
def _describe_climate(name, old_val, new_val, attrs) -> str:
temp = attrs.get("current_temperature", "?")
target = attrs.get("temperature", "?")
return (
f"[Home Assistant] {name}: HVAC mode changed from "
f"'{old_val}' to '{new_val}' (current: {temp}, target: {target})"
)
def _describe_sensor(name, old_val, new_val, attrs) -> str:
unit = attrs.get("unit_of_measurement", "")
return f"[Home Assistant] {name}: changed from {old_val}{unit} to {new_val}{unit}"
def _describe_on_off(name, old_val, new_val, attrs) -> str:
return f"[Home Assistant] {name}: turned {_on_off(new_val)}"
# domain -> (name, old_val, new_val, attrs) -> human-readable description
_DOMAIN_DESCRIBERS: Dict[str, Callable[..., str]] = {
"climate": _describe_climate,
"sensor": _describe_sensor,
"binary_sensor": lambda name, old_val, new_val, attrs: (
f"[Home Assistant] {name}: {_triggered(new_val)} (was {_triggered(old_val)})"
),
"light": _describe_on_off,
"switch": _describe_on_off,
"fan": _describe_on_off,
"alarm_control_panel": lambda name, old_val, new_val, attrs: (
f"[Home Assistant] {name}: alarm state changed from '{old_val}' to '{new_val}'"
# domain -> description template; see ``_format_state_change`` for the fields.
_TURNED = "[Home Assistant] {name}: turned {on_off}"
_DOMAIN_TEMPLATES = {
"climate": (
"[Home Assistant] {name}: HVAC mode changed from "
"'{old}' to '{new}' (current: {temp}, target: {target})"
),
"sensor": "[Home Assistant] {name}: changed from {old}{unit} to {new}{unit}",
"binary_sensor": "[Home Assistant] {name}: {triggered} (was {was_triggered})",
"light": _TURNED,
"switch": _TURNED,
"fan": _TURNED,
"alarm_control_panel": "[Home Assistant] {name}: alarm state changed from '{old}' to '{new}'",
}
_DEFAULT_TEMPLATE = "[Home Assistant] {name} ({entity_id}): changed from '{old}' to '{new}'"
class HomeAssistantAdapter(BasePlatformAdapter):
@@ -174,12 +157,15 @@ class HomeAssistantAdapter(BasePlatformAdapter):
return False
return True
@staticmethod
async def _close(obj) -> None:
if obj and not obj.closed:
await obj.close()
async def _cleanup_ws(self) -> None:
if self._ws and not self._ws.closed:
await self._ws.close()
await self._close(self._ws)
self._ws = None
if self._session and not self._session.closed:
await self._session.close()
await self._close(self._session)
self._session = None
async def disconnect(self) -> None:
@@ -192,8 +178,7 @@ class HomeAssistantAdapter(BasePlatformAdapter):
pass
self._listen_task = None
await self._cleanup_ws()
if self._rest_session and not self._rest_session.closed:
await self._rest_session.close()
await self._close(self._rest_session)
self._rest_session = None
logger.info("[%s] Disconnected", self.name)
@@ -280,11 +265,13 @@ class HomeAssistantAdapter(BasePlatformAdapter):
if old_val == new_val:
return None
attrs = new_state.get("attributes", {})
name = attrs.get("friendly_name", entity_id)
describe = _DOMAIN_DESCRIBERS.get(_domain_of(entity_id))
if describe is not None:
return describe(name, old_val, new_val, attrs)
return f"[Home Assistant] {name} ({entity_id}): changed from '{old_val}' to '{new_val}'"
template = _DOMAIN_TEMPLATES.get(_domain_of(entity_id), _DEFAULT_TEMPLATE)
return template.format(
name=attrs.get("friendly_name", entity_id), entity_id=entity_id, old=old_val, new=new_val,
temp=attrs.get("current_temperature", "?"), target=attrs.get("temperature", "?"),
unit=attrs.get("unit_of_measurement", ""), on_off=_on_off(new_val),
triggered=_triggered(new_val), was_triggered=_triggered(old_val),
)
# -- Outbound messaging -------------------------------------------------
+44 -23
View File
@@ -3,17 +3,16 @@
Subscribes to a topic on ntfy.sh or a self-hosted server via HTTP streaming
(``/json`` with ``poll=false``) and publishes replies via HTTP POST (httpx only).
config.yaml ``platforms.ntfy.extra``: ``server`` (default https://ntfy.sh),
``topic`` (subscribe, required), ``publish_topic`` (defaults to topic), ``token``
(Bearer or ``user:pass`` Basic), ``markdown`` (default false).
Env (read at construct time; env wins over ``extra``): NTFY_TOPIC, NTFY_SERVER_URL,
NTFY_TOKEN, NTFY_PUBLISH_TOPIC, NTFY_MARKDOWN ("true"/"1"/"yes"), NTFY_ALLOWED_USERS
(topic names), NTFY_ALLOW_ALL_USERS (dev only), NTFY_HOME_CHANNEL, NTFY_HOME_CHANNEL_NAME.
config.yaml ``platforms.ntfy.extra``: ``server`` (default https://ntfy.sh), ``topic``
(subscribe, required), ``publish_topic`` (defaults to topic), ``token`` (Bearer or
``user:pass`` Basic), ``markdown`` (default false). Env (read at construct time; env
wins over ``extra``): NTFY_TOPIC, NTFY_SERVER_URL, NTFY_TOKEN, NTFY_PUBLISH_TOPIC,
NTFY_MARKDOWN ("true"/"1"/"yes"), NTFY_ALLOWED_USERS (topic names), NTFY_ALLOW_ALL_USERS
(dev only), NTFY_HOME_CHANNEL, NTFY_HOME_CHANNEL_NAME.
Identity model: ntfy has no authenticated user identity; ``title`` is
publisher-controlled and NOT used for authorization. Each topic is a single
trusted channel (``user_id`` == topic name). Protect the topic with a read
token for any real trust boundary.
Identity model: ntfy has no authenticated user identity; ``title`` is publisher-controlled
and NOT used for authorization. Each topic is a single trusted channel (``user_id`` ==
topic name). Protect the topic with a read token for any real trust boundary.
"""
import asyncio
@@ -50,7 +49,6 @@ RECONNECT_BACKOFF = [2, 5, 10, 30, 60]
STREAM_TIMEOUT_SECONDS = 90 # ntfy keepalive default is 55s; give margin
_ECHO_TAG = "hermes-agent" # tag added to outgoing messages for echo-loop prevention
_MARKDOWN_TRUTHY = ("1", "true", "yes")
_BASE_HEADERS = {"Content-Type": "text/plain; charset=utf-8", "X-Tags": _ECHO_TAG}
def _build_auth_header(token: str) -> Dict[str, str]:
@@ -68,6 +66,19 @@ def _build_auth_header(token: str) -> Dict[str, str]:
return {"Authorization": f"Bearer {token}"}
def _publish_headers(token: str, markdown: bool, *, auth_first: bool = True) -> Dict[str, str]:
"""Headers for a publish POST: auth (if any), plain-text body, echo tag, optional X-Markdown.
``auth_first`` pins the header order each call site has always sent on the wire.
"""
auth = _build_auth_header(token)
base = {"Content-Type": "text/plain; charset=utf-8", "X-Tags": _ECHO_TAG}
headers = {**auth, **base} if auth_first else {**base, **auth}
if markdown:
headers["X-Markdown"] = "true"
return headers
def _truncate_body(message: str, *, context: str) -> bytes:
"""Apply the ntfy 4096-char limit, logging a warning (tagged ``context``) on truncation."""
if len(message) > MAX_MESSAGE_LENGTH:
@@ -111,7 +122,9 @@ class NtfyAdapter(BasePlatformAdapter):
def __init__(self, config: PlatformConfig):
super().__init__(config=config, platform=Platform("ntfy"))
extra = config.extra or {}
self._server: str = (extra.get("server") or _get_scoped_secret("NTFY_SERVER_URL", DEFAULT_SERVER)).rstrip("/")
self._server: str = (
extra.get("server") or _get_scoped_secret("NTFY_SERVER_URL", DEFAULT_SERVER)
).rstrip("/")
self._topic: str = extra.get("topic") or _get_scoped_secret("NTFY_TOPIC", "")
self._publish_topic: str = (
extra.get("publish_topic") or _get_scoped_secret("NTFY_PUBLISH_TOPIC", "") or self._topic
@@ -248,11 +261,14 @@ class NtfyAdapter(BasePlatformAdapter):
# No native user identity on ntfy: the publisher-controlled title must
# NOT drive authorization, so user_id is fixed to the topic name.
topic = event.get("topic") or self._topic
source = self.build_source(chat_id=topic, chat_name=topic, chat_type="dm", user_id=topic, user_name=topic)
source = self.build_source(
chat_id=topic, chat_name=topic, chat_type="dm", user_id=topic, user_name=topic,
)
unix_ts = event.get("time")
try:
timestamp = (
datetime.fromtimestamp(int(unix_ts), tz=timezone.utc) if unix_ts else datetime.now(tz=timezone.utc)
datetime.fromtimestamp(int(unix_ts), tz=timezone.utc)
if unix_ts else datetime.now(tz=timezone.utc)
)
except (ValueError, OSError, TypeError):
timestamp = datetime.now(tz=timezone.utc)
@@ -285,9 +301,7 @@ class NtfyAdapter(BasePlatformAdapter):
if not self._http_client:
return SendResult(success=False, error="HTTP client not initialized")
url = f"{self._server}/{publish_topic}"
headers = {**self._auth_headers(), **_BASE_HEADERS}
if (self.config.extra or {}).get("markdown", False):
headers["X-Markdown"] = "true"
headers = _publish_headers(self._token, bool((self.config.extra or {}).get("markdown", False)))
if len(content) > self.MAX_MESSAGE_LENGTH:
logger.warning(
"[%s] Message truncated from %d to %d chars (ntfy limit)",
@@ -295,7 +309,9 @@ class NtfyAdapter(BasePlatformAdapter):
)
body = content[:self.MAX_MESSAGE_LENGTH]
try:
resp = await self._http_client.post(url, content=body.encode("utf-8"), headers=headers, timeout=15.0)
resp = await self._http_client.post(
url, content=body.encode("utf-8"), headers=headers, timeout=15.0,
)
if resp.status_code < 300:
return SendResult(success=True, message_id=_response_message_id(resp))
body_text = resp.text
@@ -328,7 +344,10 @@ def _env_enablement() -> dict | None:
topic = _get_scoped_secret("NTFY_TOPIC", "").strip()
if not topic:
return None
seed: dict = {"topic": topic, "server": _get_scoped_secret("NTFY_SERVER_URL", DEFAULT_SERVER).rstrip("/")}
seed: dict = {
"topic": topic,
"server": _get_scoped_secret("NTFY_SERVER_URL", DEFAULT_SERVER).rstrip("/"),
}
for key, env in (("publish_topic", "NTFY_PUBLISH_TOPIC"), ("token", "NTFY_TOKEN")):
value = _get_scoped_secret(env, "").strip()
if value:
@@ -367,16 +386,18 @@ async def _standalone_send(
return {"error": "ntfy standalone send: NTFY_TOPIC not configured"}
token = extra.get("token") or _get_scoped_secret("NTFY_TOKEN", "")
markdown_env = _get_scoped_secret("NTFY_MARKDOWN", "").strip().lower()
headers = {**_BASE_HEADERS, **_build_auth_header(token)}
if bool(extra.get("markdown")) or markdown_env in _MARKDOWN_TRUTHY:
headers["X-Markdown"] = "true"
headers = _publish_headers(
token, bool(extra.get("markdown")) or markdown_env in _MARKDOWN_TRUTHY, auth_first=False,
)
body = _truncate_body(message, context="ntfy standalone")
try:
async with httpx.AsyncClient(timeout=15.0) as client:
resp = await client.post(f"{server}/{publish_topic}", content=body, headers=headers)
if resp.status_code >= 300:
return {"error": f"ntfy HTTP {resp.status_code}: {resp.text[:200]}"}
return {"success": True, "platform": "ntfy", "chat_id": publish_topic, "message_id": _response_message_id(resp)}
return {
"success": True, "platform": "ntfy", "chat_id": publish_topic, "message_id": _response_message_id(resp),
}
except Exception as e:
return {"error": f"ntfy standalone send failed: {e}"}
+9 -21
View File
@@ -2,19 +2,11 @@
Outbound SMS via the Twilio REST API; inbound via an aiohttp webhook server.
Shares credentials with the optional telephony skill — same env vars:
- TWILIO_ACCOUNT_SID
- TWILIO_AUTH_TOKEN
- TWILIO_PHONE_NUMBER (E.164 from-number, e.g. +15551234567)
Gateway-specific env vars:
- SMS_WEBHOOK_PORT (default 8080)
- SMS_WEBHOOK_HOST (default 127.0.0.1)
- SMS_WEBHOOK_URL (public URL for Twilio signature validation — required)
- SMS_INSECURE_NO_SIGNATURE (true to disable signature validation — dev only)
- SMS_ALLOWED_USERS (comma-separated E.164 phone numbers)
- SMS_ALLOW_ALL_USERS (true/false)
- SMS_HOME_CHANNEL (phone number for cron delivery)
Env vars — shared with the telephony skill: TWILIO_ACCOUNT_SID, TWILIO_AUTH_TOKEN,
TWILIO_PHONE_NUMBER (E.164 from-number). Gateway-specific: SMS_WEBHOOK_PORT (8080),
SMS_WEBHOOK_HOST (127.0.0.1), SMS_WEBHOOK_URL (public URL for Twilio signature
validation — required), SMS_INSECURE_NO_SIGNATURE (true disables validation — dev only),
SMS_ALLOWED_USERS (comma-separated E.164), SMS_ALLOW_ALL_USERS, SMS_HOME_CHANNEL (cron).
"""
import asyncio
@@ -101,9 +93,6 @@ class SmsAdapter(BasePlatformAdapter):
self._runner = None
self._http_session: Optional["aiohttp.ClientSession"] = None
def _basic_auth_header(self) -> str:
return _basic_auth(self._account_sid, self._auth_token)
# -- Lifecycle -----------------------------------------------------------
async def connect(self, *, is_reconnect: bool = False) -> bool:
@@ -169,7 +158,7 @@ class SmsAdapter(BasePlatformAdapter):
chunks = self.truncate_message(self.format_message(content))
last_result = SendResult(success=True)
url = f"{TWILIO_API_BASE}/{self._account_sid}/Messages.json"
headers = {"Authorization": self._basic_auth_header()}
headers = {"Authorization": _basic_auth(self._account_sid, self._auth_token)}
session = self._http_session or _new_session(trust_env=gateway_trust_env())
try:
for chunk in chunks:
@@ -287,11 +276,10 @@ class SmsAdapter(BasePlatformAdapter):
return _twiml_response()
# -- Plugin registration -----------------------------------------------------
# TWILIO_* env→PlatformConfig seeding stays in core (gateway/config.py).
# -- Plugin registration (TWILIO_* env→PlatformConfig seeding stays in gateway/config.py)
# Standalone-send markdown stripping (looser than helpers.strip_markdown: no
# word-boundary guards on underscores, ``[a-z]*`` fence tags — kept for parity).
# Standalone-send markdown stripping: looser than helpers.strip_markdown (no
# word-boundary guards on underscores, ``[a-z]*`` fence tags) — kept for parity.
_SMS_MARKDOWN_SUBS = (
(re.compile(r"\*\*(.+?)\*\*", re.DOTALL), r"\1"),
(re.compile(r"\*(.+?)\*", re.DOTALL), r"\1"),