refactor(web_server helpers): second pass — helper dedupe in windows_ssh_runtime, webhook via atomic_json_write, table packing
- windows_ssh_runtime: _open_existing shared by _read_shared/remove_artifact; _win32 builds its namespace via importlib; reparse/path checks folded - webhook: _save_subscriptions -> utils.atomic_json_write(mode=0o600) (same fchmod-before- rename + post-replace chmod semantics); base URL builder tightened - web_server_messaging: override table tuples single-line, catalog entry builder flattened, WhatsApp payload from a field tuple; channel keys as one comprehension - web_server_oauth: poller bodies read sess fields inline; status dicts packed - web_server_gateway: health URL normalisation via one regex; Popen detach kwargs inline; topology cache getter collapsed - web_server_cron/xai_retirement/win_pty_bridge/worktree_gc: small collapses Verification: routes identical, --help byte-identical (webhook/worktree + subcommands), golden corpus identical, 189 test files / 2925 passed / 0 failed.
This commit is contained in:
@@ -118,7 +118,6 @@ def _cron_profile_home(profile: Optional[str]) -> Tuple[str, Path]:
|
||||
"""Resolve a profile query value to (profile_name, HERMES_HOME)."""
|
||||
from hermes_cli.web_server import _cron_default_profile
|
||||
from hermes_cli import profiles as profiles_mod
|
||||
|
||||
raw = (profile or _cron_default_profile()).strip() or "default"
|
||||
try:
|
||||
canon = profiles_mod.normalize_profile_name(raw)
|
||||
@@ -148,7 +147,6 @@ def _cron_store_scope(home: Path):
|
||||
"""
|
||||
from cron import jobs as cron_jobs
|
||||
from hermes_constants import reset_hermes_home_override, set_hermes_home_override
|
||||
|
||||
token = set_hermes_home_override(str(home))
|
||||
try:
|
||||
with cron_jobs.use_cron_store(home):
|
||||
@@ -190,18 +188,16 @@ def _notify_cron_provider_for_profile(target_profile: Optional[str]) -> None:
|
||||
from cron.scheduler_provider import InProcessCronScheduler, resolve_cron_scheduler
|
||||
with _cron_store_scope(home):
|
||||
provider = resolve_cron_scheduler()
|
||||
if not isinstance(provider, InProcessCronScheduler):
|
||||
profile_names = [str(p.get("name") or "") for p in _cron_profile_dicts()]
|
||||
if len([n for n in profile_names if n]) > 1:
|
||||
_log.warning(
|
||||
"Skipping cron provider reconcile for profile %s: "
|
||||
"external provider '%s' reconcile is not "
|
||||
"profile-scoped and would disarm other profiles' "
|
||||
"armed one-shots. The mutated profile re-arms "
|
||||
"idempotently on its next fire/start.",
|
||||
target_profile,
|
||||
provider.name)
|
||||
return
|
||||
external = not isinstance(provider, InProcessCronScheduler)
|
||||
if external and sum(1 for p in _cron_profile_dicts() if p.get("name")) > 1:
|
||||
_log.warning(
|
||||
"Skipping cron provider reconcile for profile %s: "
|
||||
"external provider '%s' reconcile is not "
|
||||
"profile-scoped and would disarm other profiles' "
|
||||
"armed one-shots. The mutated profile re-arms "
|
||||
"idempotently on its next fire/start.", target_profile, provider.name,
|
||||
)
|
||||
return
|
||||
provider.on_jobs_changed()
|
||||
except Exception:
|
||||
_log.debug("Cron provider reconciliation failed for profile %s", target_profile, exc_info=True)
|
||||
@@ -293,7 +289,6 @@ def _fire_cron_job_for_profile(profile: str, job_id: str, *, force: bool = False
|
||||
from hermes_cli.web_server import _cron_profile_home
|
||||
_profile_name, home = _cron_profile_home(profile)
|
||||
from cron.scheduler_provider import provider_supports_force_fire, resolve_cron_scheduler
|
||||
|
||||
with _cron_store_scope(home):
|
||||
provider = resolve_cron_scheduler()
|
||||
if force:
|
||||
@@ -339,7 +334,6 @@ def _gateway_fire_endpoint(profile: str, home: Path) -> str:
|
||||
"""
|
||||
from hermes_cli.web_server import _cron_default_profile, load_config
|
||||
import os as _os
|
||||
|
||||
multiplex = False
|
||||
try:
|
||||
from gateway.config import _env_multiplex_profiles_override
|
||||
@@ -356,9 +350,8 @@ def _gateway_fire_endpoint(profile: str, home: Path) -> str:
|
||||
listener_profile, listener_home = "default", get_default_hermes_root()
|
||||
_log.info(
|
||||
"cron fire: multiplex gateway — resolving api_server port for %s "
|
||||
"from the default profile's listener (%s)",
|
||||
profile,
|
||||
listener_home)
|
||||
"from the default profile's listener (%s)", profile, listener_home,
|
||||
)
|
||||
|
||||
port = 0
|
||||
try:
|
||||
@@ -407,14 +400,11 @@ async def _forward_cron_fire_to_gateway(
|
||||
_profile_name, home = _cron_profile_home(profile)
|
||||
url = _gateway_fire_endpoint(_profile_name, home)
|
||||
import httpx
|
||||
|
||||
try:
|
||||
async with httpx.AsyncClient(timeout=10.0) as client:
|
||||
resp = await client.post(url, json={"job_id": job_id}, headers={"Authorization": authorization})
|
||||
except Exception as exc:
|
||||
_log.warning(
|
||||
"cron fire forward to %s failed (%s: %s); returning 503 for NAS retry",
|
||||
url, type(exc).__name__, exc)
|
||||
_log.warning("cron fire forward to %s failed (%s: %s); returning 503 for NAS retry", url, type(exc).__name__, exc)
|
||||
return None
|
||||
try:
|
||||
body = resp.json()
|
||||
@@ -437,13 +427,8 @@ def _gateway_intentionally_stopped(profile: Optional[str]) -> bool:
|
||||
"""
|
||||
from hermes_cli.web_server import _cron_profile_home
|
||||
import json as _json
|
||||
|
||||
try:
|
||||
_name, home = _cron_profile_home(profile)
|
||||
state_file = home / "gateway_state.json"
|
||||
if not state_file.exists():
|
||||
return False
|
||||
data = _json.loads(state_file.read_text(encoding="utf-8"))
|
||||
data = _json.loads((_cron_profile_home(profile)[1] / "gateway_state.json").read_text(encoding="utf-8"))
|
||||
return isinstance(data, dict) and data.get("desired_state") == "stopped"
|
||||
except Exception:
|
||||
return False
|
||||
|
||||
@@ -34,11 +34,7 @@ def _probe_gateway_health() -> tuple[bool, dict | None]:
|
||||
from hermes_cli.web_server import _GATEWAY_HEALTH_TIMEOUT, _GATEWAY_HEALTH_URL
|
||||
if not _GATEWAY_HEALTH_URL:
|
||||
return False, None
|
||||
base = _GATEWAY_HEALTH_URL.rstrip("/")
|
||||
if base.endswith("/health/detailed"):
|
||||
base = base[: -len("/health/detailed")]
|
||||
elif base.endswith("/health"):
|
||||
base = base[: -len("/health")]
|
||||
base = re.sub(r"/health(/detailed)?$", "", _GATEWAY_HEALTH_URL.rstrip("/"))
|
||||
for path in (f"{base}/health/detailed", f"{base}/health"):
|
||||
try:
|
||||
req = urllib.request.Request(path, method="GET")
|
||||
@@ -54,16 +50,11 @@ def _probe_gateway_health() -> tuple[bool, dict | None]:
|
||||
# Mirrors PORT_BINDING_PLATFORM_VALUES (gateway/config.py) and each adapter's DEFAULT_PORT /
|
||||
# DEFAULT_WEBHOOK_PORT. Display-only data for the topology readout, not a bind source.
|
||||
_PORT_BINDING_PLATFORM_PORTS: Dict[str, Tuple[str, int]] = {
|
||||
"webhook": ("port", 8644),
|
||||
"api_server": ("port", 8642),
|
||||
"msgraph_webhook": ("port", 8646),
|
||||
"feishu": ("webhook_port", 8765),
|
||||
"wecom_callback": ("port", 8645),
|
||||
"bluebubbles": ("webhook_port", 8645),
|
||||
"sms": ("webhook_port", 8080),
|
||||
"whatsapp_cloud": ("webhook_port", 8090),
|
||||
"line": ("port", 8646),
|
||||
"teams": ("port", 3978)}
|
||||
"webhook": ("port", 8644), "api_server": ("port", 8642), "msgraph_webhook": ("port", 8646),
|
||||
"feishu": ("webhook_port", 8765), "wecom_callback": ("port", 8645), "bluebubbles": ("webhook_port", 8645),
|
||||
"sms": ("webhook_port", 8080), "whatsapp_cloud": ("webhook_port", 8090), "line": ("port", 8646),
|
||||
"teams": ("port", 3978),
|
||||
}
|
||||
|
||||
# Platform states that mean the adapter is NOT serving its port right now.
|
||||
_PLATFORM_DEAD_STATES = frozenset({"fatal", "disconnected", "stopped"})
|
||||
@@ -216,12 +207,9 @@ _TOPOLOGY_CACHE_TTL = 10.0
|
||||
|
||||
|
||||
def _topology_cache_get(fn: Any) -> Optional[Dict[str, Any]]:
|
||||
if (
|
||||
_TOPOLOGY_CACHE["data"] is not None
|
||||
and _TOPOLOGY_CACHE["fn"] is fn
|
||||
and time.monotonic() - _TOPOLOGY_CACHE["ts"] < _TOPOLOGY_CACHE_TTL):
|
||||
return _TOPOLOGY_CACHE["data"]
|
||||
return None
|
||||
c = _TOPOLOGY_CACHE
|
||||
fresh = c["fn"] is fn and time.monotonic() - c["ts"] < _TOPOLOGY_CACHE_TTL
|
||||
return c["data"] if fresh and c["data"] is not None else None
|
||||
|
||||
|
||||
def _collect_profile_gateway_topology_cached() -> Dict[str, Any]:
|
||||
@@ -269,11 +257,9 @@ def _display_system_platform(*, system: str, release: str, version: str, platfor
|
||||
return {"os": system, "os_release": release, "os_version": version, "platform": platform_label}
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Gateway + update actions (invoked from the Status page). Spawned detached so the request
|
||||
# returns immediately; stdin is DEVNULL so stray input() fails fast; stdout/stderr stream to
|
||||
# ~/.hermes/logs/<action>.log which the dashboard tails.
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
_ACTION_LOG_DIR: Path = get_hermes_home() / "logs"
|
||||
|
||||
@@ -356,18 +342,11 @@ def _spawn_hermes_action(
|
||||
# the gateway's own restart watcher does.
|
||||
action_env = {**os.environ, "HERMES_NONINTERACTIVE": "1"}
|
||||
action_env.pop("_HERMES_GATEWAY", None)
|
||||
popen_kwargs: Dict[str, Any] = {
|
||||
"cwd": str(PROJECT_ROOT),
|
||||
"stdin": subprocess.DEVNULL,
|
||||
"stdout": log_file,
|
||||
"stderr": subprocess.STDOUT,
|
||||
"env": {**action_env, **(env_overrides or {})}}
|
||||
if sys.platform == "win32":
|
||||
popen_kwargs["creationflags"] = windows_detach_flags()
|
||||
else:
|
||||
popen_kwargs["start_new_session"] = True
|
||||
|
||||
proc = subprocess.Popen(cmd, **popen_kwargs)
|
||||
detach = {"creationflags": windows_detach_flags()} if sys.platform == "win32" else {"start_new_session": True}
|
||||
proc = subprocess.Popen(
|
||||
cmd, cwd=str(PROJECT_ROOT), stdin=subprocess.DEVNULL, stdout=log_file, stderr=subprocess.STDOUT,
|
||||
env={**action_env, **(env_overrides or {})}, **detach,
|
||||
)
|
||||
log_file.close() # child holds its own dup'd fd; keeping ours leaks one per action
|
||||
_ACTION_RESULTS.pop(name, None)
|
||||
_ACTION_COMMANDS[name] = tuple(subcommand)
|
||||
@@ -407,7 +386,6 @@ def _split_text_for_speak_stream(text: str, cap: int) -> list:
|
||||
reflows whitespace (sentences re-joined with single spaces) and has no fence semantics.
|
||||
"""
|
||||
from tools.tts_streaming import SENTENCE_BOUNDARY_RE as _SENTENCE_BOUNDARY_RE
|
||||
|
||||
cap = cap if cap and cap > 0 else 4000
|
||||
pieces, buf = [], ""
|
||||
for sentence in filter(str.strip, _SENTENCE_BOUNDARY_RE.split(text)):
|
||||
|
||||
@@ -27,15 +27,13 @@ _log = logging.getLogger("hermes_cli.web_server")
|
||||
# and pulls required_env from a plugin's PlatformEntry when available.
|
||||
_PLATFORM_OVERRIDES: dict[str, dict[str, Any]] = {
|
||||
"telegram": {
|
||||
"name": "Telegram",
|
||||
"description": "Run Hermes from Telegram DMs, groups, and topics.",
|
||||
"name": "Telegram", "description": "Run Hermes from Telegram DMs, groups, and topics.",
|
||||
"docs_url": "https://core.telegram.org/bots/features#botfather",
|
||||
"env_vars": ("TELEGRAM_BOT_TOKEN", "TELEGRAM_ALLOWED_USERS", "TELEGRAM_PROXY"),
|
||||
"required_env": ("TELEGRAM_BOT_TOKEN",),
|
||||
},
|
||||
"discord": {
|
||||
"name": "Discord",
|
||||
"description": "Connect Hermes to Discord DMs, channels, and threads.",
|
||||
"name": "Discord", "description": "Connect Hermes to Discord DMs, channels, and threads.",
|
||||
"docs_url": "https://discord.com/developers/applications",
|
||||
"env_vars": ("DISCORD_BOT_TOKEN", "DISCORD_ALLOWED_USERS"),
|
||||
"required_env": ("DISCORD_BOT_TOKEN",),
|
||||
@@ -55,20 +53,15 @@ _PLATFORM_OVERRIDES: dict[str, dict[str, Any]] = {
|
||||
"required_env": ("MATTERMOST_URL", "MATTERMOST_TOKEN"),
|
||||
},
|
||||
"matrix": {
|
||||
"name": "Matrix",
|
||||
"description": "Use Hermes in Matrix rooms and direct messages.",
|
||||
"name": "Matrix", "description": "Use Hermes in Matrix rooms and direct messages.",
|
||||
"docs_url": "https://matrix.org/ecosystem/servers/",
|
||||
"env_vars": (
|
||||
"MATRIX_HOMESERVER",
|
||||
"MATRIX_ACCESS_TOKEN",
|
||||
"MATRIX_USER_ID",
|
||||
"MATRIX_ALLOWED_USERS",
|
||||
"MATRIX_HOMESERVER", "MATRIX_ACCESS_TOKEN", "MATRIX_USER_ID", "MATRIX_ALLOWED_USERS",
|
||||
),
|
||||
"required_env": ("MATRIX_HOMESERVER", "MATRIX_ACCESS_TOKEN", "MATRIX_USER_ID"),
|
||||
},
|
||||
"signal": {
|
||||
"name": "Signal",
|
||||
"description": "Connect through a signal-cli REST bridge.",
|
||||
"name": "Signal", "description": "Connect through a signal-cli REST bridge.",
|
||||
"docs_url": "https://github.com/bbernhard/signal-cli-rest-api",
|
||||
"env_vars": ("SIGNAL_HTTP_URL", "SIGNAL_ACCOUNT", "SIGNAL_ALLOWED_USERS"),
|
||||
"required_env": ("SIGNAL_HTTP_URL", "SIGNAL_ACCOUNT"),
|
||||
@@ -78,10 +71,7 @@ _PLATFORM_OVERRIDES: dict[str, dict[str, Any]] = {
|
||||
"description": "Use Hermes through the bundled WhatsApp bridge with QR-based auth.",
|
||||
"docs_url": "https://github.com/tulir/whatsmeow",
|
||||
"env_vars": (
|
||||
"WHATSAPP_ENABLED",
|
||||
"WHATSAPP_MODE",
|
||||
"WHATSAPP_DM_POLICY",
|
||||
"WHATSAPP_ALLOWED_USERS",
|
||||
"WHATSAPP_ENABLED", "WHATSAPP_MODE", "WHATSAPP_DM_POLICY", "WHATSAPP_ALLOWED_USERS",
|
||||
),
|
||||
"required_env": (),
|
||||
},
|
||||
@@ -89,69 +79,52 @@ _PLATFORM_OVERRIDES: dict[str, dict[str, Any]] = {
|
||||
"name": "Home Assistant",
|
||||
"description": "Control your smart home from Hermes via Home Assistant.",
|
||||
"docs_url": "https://www.home-assistant.io/docs/authentication/",
|
||||
"env_vars": ("HASS_URL", "HASS_TOKEN"),
|
||||
"required_env": ("HASS_URL", "HASS_TOKEN"),
|
||||
"env_vars": ("HASS_URL", "HASS_TOKEN"), "required_env": ("HASS_URL", "HASS_TOKEN"),
|
||||
},
|
||||
"email": {
|
||||
"name": "Email",
|
||||
"description": "Talk to Hermes through an IMAP/SMTP mailbox.",
|
||||
"name": "Email", "description": "Talk to Hermes through an IMAP/SMTP mailbox.",
|
||||
"docs_url": "https://hermes-agent.nousresearch.com/docs/user-guide/messaging/",
|
||||
"env_vars": ("EMAIL_ADDRESS", "EMAIL_PASSWORD", "EMAIL_IMAP_HOST", "EMAIL_SMTP_HOST"),
|
||||
"required_env": ("EMAIL_ADDRESS", "EMAIL_PASSWORD", "EMAIL_IMAP_HOST", "EMAIL_SMTP_HOST"),
|
||||
},
|
||||
"sms": {
|
||||
"name": "SMS (Twilio)",
|
||||
"description": "Send and receive text messages via Twilio.",
|
||||
"name": "SMS (Twilio)", "description": "Send and receive text messages via Twilio.",
|
||||
"docs_url": "https://www.twilio.com/console",
|
||||
"env_vars": ("TWILIO_ACCOUNT_SID", "TWILIO_AUTH_TOKEN"),
|
||||
"required_env": ("TWILIO_ACCOUNT_SID", "TWILIO_AUTH_TOKEN"),
|
||||
},
|
||||
"dingtalk": {
|
||||
"name": "DingTalk",
|
||||
"description": "Connect Hermes to DingTalk groups (钉钉).",
|
||||
"name": "DingTalk", "description": "Connect Hermes to DingTalk groups (钉钉).",
|
||||
"docs_url": "https://open.dingtalk.com/document/orgapp/the-robot-development-process",
|
||||
"env_vars": ("DINGTALK_CLIENT_ID", "DINGTALK_CLIENT_SECRET"),
|
||||
"required_env": ("DINGTALK_CLIENT_ID", "DINGTALK_CLIENT_SECRET"),
|
||||
},
|
||||
"feishu": {
|
||||
"name": "Feishu / Lark",
|
||||
"description": "Use Hermes inside Feishu / Lark.",
|
||||
"name": "Feishu / Lark", "description": "Use Hermes inside Feishu / Lark.",
|
||||
"docs_url": "https://open.feishu.cn/document/uAjLw4CM/ukTMukTMukTM/reference/im-v1/intro",
|
||||
"env_vars": (
|
||||
"FEISHU_APP_ID",
|
||||
"FEISHU_APP_SECRET",
|
||||
"FEISHU_ENCRYPT_KEY",
|
||||
"FEISHU_VERIFICATION_TOKEN",
|
||||
"FEISHU_APP_ID", "FEISHU_APP_SECRET", "FEISHU_ENCRYPT_KEY", "FEISHU_VERIFICATION_TOKEN",
|
||||
),
|
||||
"required_env": ("FEISHU_APP_ID", "FEISHU_APP_SECRET"),
|
||||
},
|
||||
"google_chat": {
|
||||
"name": "Google Chat",
|
||||
"description": "Connect Hermes to Google Chat via Cloud Pub/Sub.",
|
||||
"name": "Google Chat", "description": "Connect Hermes to Google Chat via Cloud Pub/Sub.",
|
||||
"docs_url": "https://hermes-agent.nousresearch.com/docs/user-guide/messaging/google_chat",
|
||||
},
|
||||
"wecom": {
|
||||
"name": "WeCom (group bot)",
|
||||
"description": "Send-only WeCom group bot via webhook.",
|
||||
"name": "WeCom (group bot)", "description": "Send-only WeCom group bot via webhook.",
|
||||
"docs_url": "https://developer.work.weixin.qq.com/document/path/91770",
|
||||
"env_vars": ("WECOM_BOT_ID", "WECOM_SECRET"),
|
||||
"required_env": ("WECOM_BOT_ID",),
|
||||
"env_vars": ("WECOM_BOT_ID", "WECOM_SECRET"), "required_env": ("WECOM_BOT_ID",),
|
||||
},
|
||||
"wecom_callback": {
|
||||
"name": "WeCom (app)",
|
||||
"description": "Two-way WeCom integration via callback app.",
|
||||
"name": "WeCom (app)", "description": "Two-way WeCom integration via callback app.",
|
||||
"docs_url": "https://developer.work.weixin.qq.com/document/path/90930",
|
||||
"env_vars": (
|
||||
"WECOM_CALLBACK_CORP_ID",
|
||||
"WECOM_CALLBACK_CORP_SECRET",
|
||||
"WECOM_CALLBACK_AGENT_ID",
|
||||
"WECOM_CALLBACK_TOKEN",
|
||||
"WECOM_CALLBACK_ENCODING_AES_KEY",
|
||||
"WECOM_CALLBACK_CORP_ID", "WECOM_CALLBACK_CORP_SECRET", "WECOM_CALLBACK_AGENT_ID",
|
||||
"WECOM_CALLBACK_TOKEN", "WECOM_CALLBACK_ENCODING_AES_KEY",
|
||||
),
|
||||
"required_env": (
|
||||
"WECOM_CALLBACK_CORP_ID",
|
||||
"WECOM_CALLBACK_CORP_SECRET",
|
||||
"WECOM_CALLBACK_AGENT_ID",
|
||||
"WECOM_CALLBACK_CORP_ID", "WECOM_CALLBACK_CORP_SECRET", "WECOM_CALLBACK_AGENT_ID",
|
||||
),
|
||||
},
|
||||
"weixin": {
|
||||
@@ -166,15 +139,12 @@ _PLATFORM_OVERRIDES: dict[str, dict[str, Any]] = {
|
||||
"description": "Use Hermes through iMessage via a BlueBubbles server.",
|
||||
"docs_url": "https://bluebubbles.app/",
|
||||
"env_vars": (
|
||||
"BLUEBUBBLES_SERVER_URL",
|
||||
"BLUEBUBBLES_PASSWORD",
|
||||
"BLUEBUBBLES_ALLOWED_USERS",
|
||||
"BLUEBUBBLES_SERVER_URL", "BLUEBUBBLES_PASSWORD", "BLUEBUBBLES_ALLOWED_USERS",
|
||||
),
|
||||
"required_env": ("BLUEBUBBLES_SERVER_URL", "BLUEBUBBLES_PASSWORD"),
|
||||
},
|
||||
"qqbot": {
|
||||
"name": "QQ Bot",
|
||||
"description": "Connect Hermes to a QQ Bot from the QQ Open Platform.",
|
||||
"name": "QQ Bot", "description": "Connect Hermes to a QQ Bot from the QQ Open Platform.",
|
||||
"docs_url": "https://q.qq.com",
|
||||
"env_vars": ("QQ_APP_ID", "QQ_CLIENT_SECRET", "QQ_ALLOWED_USERS"),
|
||||
"required_env": ("QQ_APP_ID", "QQ_CLIENT_SECRET"),
|
||||
@@ -214,9 +184,7 @@ _PLATFORM_OVERRIDES: dict[str, dict[str, Any]] = {
|
||||
"docs_url": "https://hermes-agent.nousresearch.com/docs/user-guide/messaging/simplex",
|
||||
},
|
||||
"yuanbao": {
|
||||
"name": "Yuanbao (元宝)",
|
||||
"description": "Connect Hermes to Tencent Yuanbao.",
|
||||
"docs_url": "",
|
||||
"name": "Yuanbao (元宝)", "description": "Connect Hermes to Tencent Yuanbao.", "docs_url": "",
|
||||
"required_env": (),
|
||||
},
|
||||
"api_server": {
|
||||
@@ -224,10 +192,7 @@ _PLATFORM_OVERRIDES: dict[str, dict[str, Any]] = {
|
||||
"description": "Expose Hermes as an OpenAI-compatible HTTP API for tools like Open WebUI.",
|
||||
"docs_url": "https://hermes-agent.nousresearch.com/docs/user-guide/messaging/",
|
||||
"env_vars": (
|
||||
"API_SERVER_ENABLED",
|
||||
"API_SERVER_KEY",
|
||||
"API_SERVER_PORT",
|
||||
"API_SERVER_HOST",
|
||||
"API_SERVER_ENABLED", "API_SERVER_KEY", "API_SERVER_PORT", "API_SERVER_HOST",
|
||||
"API_SERVER_MODEL_NAME",
|
||||
),
|
||||
"required_env": (),
|
||||
@@ -236,8 +201,7 @@ _PLATFORM_OVERRIDES: dict[str, dict[str, Any]] = {
|
||||
"name": "Webhooks",
|
||||
"description": "Receive events from GitHub, GitLab, and other webhook sources.",
|
||||
"docs_url": "https://hermes-agent.nousresearch.com/docs/user-guide/messaging/webhooks/",
|
||||
"env_vars": ("WEBHOOK_ENABLED", "WEBHOOK_PORT", "WEBHOOK_SECRET"),
|
||||
"required_env": (),
|
||||
"env_vars": ("WEBHOOK_ENABLED", "WEBHOOK_PORT", "WEBHOOK_SECRET"), "required_env": (),
|
||||
},
|
||||
"msgraph_webhook": {
|
||||
"name": "Microsoft Graph Webhook",
|
||||
@@ -253,8 +217,7 @@ _PLATFORM_OVERRIDES: dict[str, dict[str, Any]] = {
|
||||
"relay": {
|
||||
"name": "Relay (experimental)",
|
||||
"description": "Generic relay adapter fronted by the Hermes Relay connector.",
|
||||
"docs_url": "",
|
||||
"required_env": (),
|
||||
"docs_url": "", "required_env": (),
|
||||
},
|
||||
}
|
||||
|
||||
@@ -309,10 +272,7 @@ def _channel_managed_env_keys() -> frozenset[str]:
|
||||
"""Env-var keys owned by a Channels page platform card; the Keys/Env page hides them so the
|
||||
same fields aren't duplicated. Best-effort: if the catalog can't be built, nothing is hidden."""
|
||||
try:
|
||||
keys: set[str] = set()
|
||||
for entry in _messaging_platform_catalog():
|
||||
keys.update(entry.get("env_vars", ()))
|
||||
return frozenset(keys)
|
||||
return frozenset(k for entry in _messaging_platform_catalog() for k in entry.get("env_vars", ()))
|
||||
except Exception:
|
||||
_log.debug("could not build channel-managed env key set", exc_info=True)
|
||||
return frozenset()
|
||||
@@ -366,26 +326,18 @@ def _build_catalog_entry(platform_id: str, plugin_entry: Any | None = None) -> d
|
||||
override = _PLATFORM_OVERRIDES.get(platform_id, {})
|
||||
if "required_env" in override:
|
||||
required_env = tuple(override["required_env"])
|
||||
elif plugin_entry is not None:
|
||||
required_env = tuple(plugin_entry.required_env or ())
|
||||
else:
|
||||
required_env = ()
|
||||
if override.get("name"):
|
||||
name = override["name"]
|
||||
elif plugin_entry is not None and plugin_entry.label:
|
||||
name = plugin_entry.label
|
||||
else:
|
||||
name = platform_id.replace("_", " ").title()
|
||||
description = override.get("description")
|
||||
if not description and plugin_entry is not None:
|
||||
description = plugin_entry.install_hint or ""
|
||||
required_env = tuple(plugin_entry.required_env or ()) if plugin_entry is not None else ()
|
||||
plugin_label = plugin_entry.label if plugin_entry is not None else None
|
||||
plugin_hint = (plugin_entry.install_hint or "") if plugin_entry is not None else None
|
||||
return {
|
||||
"id": platform_id,
|
||||
"name": name,
|
||||
"description": description or "",
|
||||
"name": override.get("name") or plugin_label or platform_id.replace("_", " ").title(),
|
||||
"description": override.get("description") or plugin_hint or "",
|
||||
"docs_url": override.get("docs_url", ""),
|
||||
"env_vars": _merge_platform_env_vars(platform_id, override, plugin_entry),
|
||||
"required_env": required_env}
|
||||
"required_env": required_env,
|
||||
}
|
||||
|
||||
|
||||
def _write_platform_enabled(platform_id: str, enabled: bool) -> None:
|
||||
@@ -414,22 +366,17 @@ _whatsapp_onboarding_sessions: dict[str, _WhatsAppOnboardingSession] = {}
|
||||
|
||||
def _whatsapp_session_path() -> Path:
|
||||
from hermes_constants import get_hermes_dir
|
||||
|
||||
return get_hermes_dir("platforms/whatsapp/session", "whatsapp/session")
|
||||
|
||||
|
||||
_WHATSAPP_PAYLOAD_FIELDS = (
|
||||
"status", "qr_payload", "expires_at", "mode", "allowed_users", "account_id", "account_name",
|
||||
"account_phone", "error",
|
||||
)
|
||||
|
||||
|
||||
def _whatsapp_onboarding_payload(pairing_id: str, record: _WhatsAppOnboardingSession) -> dict[str, Any]:
|
||||
return {
|
||||
"pairing_id": pairing_id,
|
||||
"status": record.status,
|
||||
"qr_payload": record.qr_payload,
|
||||
"expires_at": record.expires_at,
|
||||
"mode": record.mode,
|
||||
"allowed_users": record.allowed_users,
|
||||
"account_id": record.account_id,
|
||||
"account_name": record.account_name,
|
||||
"account_phone": record.account_phone,
|
||||
"error": record.error}
|
||||
return {"pairing_id": pairing_id, **{f: getattr(record, f) for f in _WHATSAPP_PAYLOAD_FIELDS}}
|
||||
|
||||
|
||||
def _restart_gateway_after_whatsapp_onboarding(profile: Optional[str] = None) -> dict[str, Any]:
|
||||
@@ -478,7 +425,6 @@ def _telegram_onboarding_request_sync(
|
||||
method: str, path: str, *, body: dict[str, Any] | None = None, bearer_token: str | None = None
|
||||
) -> dict[str, Any]:
|
||||
import httpx
|
||||
|
||||
headers = {"Accept": "application/json", "User-Agent": _TELEGRAM_ONBOARDING_USER_AGENT}
|
||||
request_kwargs: dict[str, Any] = {}
|
||||
if body is not None:
|
||||
@@ -486,11 +432,9 @@ def _telegram_onboarding_request_sync(
|
||||
request_kwargs["json"] = body
|
||||
if bearer_token:
|
||||
headers["Authorization"] = f"Bearer {bearer_token}"
|
||||
|
||||
url = f"{_telegram_onboarding_base_url()}{path}"
|
||||
try:
|
||||
with httpx.Client(timeout=httpx.Timeout(10.0)) as client:
|
||||
response = client.request(method, url, headers=headers, **request_kwargs)
|
||||
response = client.request(method, f"{_telegram_onboarding_base_url()}{path}", headers=headers, **request_kwargs)
|
||||
response.raise_for_status()
|
||||
except httpx.HTTPStatusError as exc:
|
||||
try:
|
||||
|
||||
@@ -37,12 +37,10 @@ def _truncate_token(value: Optional[str], visible: int = 6) -> str:
|
||||
|
||||
def _token_status(source: str, source_label: str, creds: Dict[str, Any]) -> Dict[str, Any]:
|
||||
return {
|
||||
"logged_in": True,
|
||||
"source": source,
|
||||
"source_label": source_label,
|
||||
"logged_in": True, "source": source, "source_label": source_label,
|
||||
"token_preview": _truncate_token(creds.get("accessToken")),
|
||||
"expires_at": creds.get("expiresAt"),
|
||||
"has_refresh_token": bool(creds.get("refreshToken"))}
|
||||
"expires_at": creds.get("expiresAt"), "has_refresh_token": bool(creds.get("refreshToken")),
|
||||
}
|
||||
|
||||
|
||||
def _anthropic_oauth_status() -> Dict[str, Any]:
|
||||
@@ -68,17 +66,13 @@ def _anthropic_oauth_status() -> Dict[str, Any]:
|
||||
pass
|
||||
from hermes_cli.config import get_env_value
|
||||
from hermes_cli.env_loader import format_secret_source_suffix
|
||||
|
||||
for var in env_var_order:
|
||||
value = get_env_value(var) or os.getenv(var)
|
||||
if value:
|
||||
return {
|
||||
"logged_in": True,
|
||||
"source": "env_var",
|
||||
"source_label": f"{var}{format_secret_source_suffix(var)}",
|
||||
"token_preview": _truncate_token(value),
|
||||
"expires_at": None,
|
||||
"has_refresh_token": False}
|
||||
"logged_in": True, "source": "env_var", "source_label": f"{var}{format_secret_source_suffix(var)}",
|
||||
"token_preview": _truncate_token(value), "expires_at": None, "has_refresh_token": False,
|
||||
}
|
||||
return dict(_LOGGED_OUT)
|
||||
|
||||
|
||||
@@ -113,13 +107,9 @@ def _copilot_acp_status() -> Dict[str, Any]:
|
||||
else:
|
||||
source_label = "GitHub Copilot CLI not found on PATH"
|
||||
return {
|
||||
"logged_in": verified,
|
||||
"source": "copilot_cli",
|
||||
"source_label": source_label,
|
||||
"token_preview": None,
|
||||
"expires_at": None,
|
||||
"has_refresh_token": False,
|
||||
"configured": configured}
|
||||
"logged_in": verified, "source": "copilot_cli", "source_label": source_label, "token_preview": None,
|
||||
"expires_at": None, "has_refresh_token": False, "configured": configured,
|
||||
}
|
||||
|
||||
|
||||
def _external_process_cli_command(provider_id: str, default: str) -> str:
|
||||
@@ -229,20 +219,13 @@ def _nous_poller(session_id: str, sess: Dict[str, Any]) -> None:
|
||||
from hermes_cli.web_server import _profile_scope
|
||||
from hermes_cli.auth import _poll_for_token, persist_nous_credentials, refresh_nous_oauth_from_state
|
||||
import httpx
|
||||
portal_base_url = sess["portal_base_url"]
|
||||
client_id = sess["client_id"]
|
||||
device_code = sess["device_code"]
|
||||
interval = sess["interval"]
|
||||
scope = sess.get("scope")
|
||||
expires_in = max(60, int(sess["expires_at"] - time.time()))
|
||||
portal_base_url, client_id = sess["portal_base_url"], sess["client_id"]
|
||||
with httpx.Client(timeout=httpx.Timeout(15.0), headers={"Accept": "application/json"}) as client:
|
||||
token_data = _poll_for_token(
|
||||
client=client,
|
||||
portal_base_url=portal_base_url,
|
||||
client_id=client_id,
|
||||
device_code=device_code,
|
||||
expires_in=expires_in,
|
||||
poll_interval=interval)
|
||||
client=client, portal_base_url=portal_base_url, client_id=client_id,
|
||||
device_code=sess["device_code"], expires_in=max(60, int(sess["expires_at"] - time.time())),
|
||||
poll_interval=sess["interval"],
|
||||
)
|
||||
# Same post-processing as _nous_device_code_login (validate/refresh JWT)
|
||||
now = datetime.now(timezone.utc)
|
||||
token_ttl = int(token_data.get("expires_in") or 0)
|
||||
@@ -250,15 +233,17 @@ def _nous_poller(session_id: str, sess: Dict[str, Any]) -> None:
|
||||
"portal_base_url": portal_base_url,
|
||||
"inference_base_url": token_data.get("inference_base_url"),
|
||||
"client_id": client_id,
|
||||
"scope": token_data.get("scope") or scope,
|
||||
"scope": token_data.get("scope") or sess.get("scope"),
|
||||
"token_type": token_data.get("token_type", "Bearer"),
|
||||
"access_token": token_data["access_token"],
|
||||
"refresh_token": token_data.get("refresh_token"),
|
||||
"obtained_at": now.isoformat(),
|
||||
"expires_at": (
|
||||
datetime.fromtimestamp(now.timestamp() + token_ttl, tz=timezone.utc).isoformat()
|
||||
if token_ttl else None),
|
||||
"expires_in": token_ttl}
|
||||
if token_ttl else None
|
||||
),
|
||||
"expires_in": token_ttl,
|
||||
}
|
||||
with _profile_scope(_oauth_session_profile(session_id)):
|
||||
full_state = refresh_nous_oauth_from_state(auth_state, timeout_seconds=15.0, force_refresh=False)
|
||||
persist_nous_credentials(full_state)
|
||||
@@ -272,33 +257,21 @@ def _minimax_poller(session_id: str, sess: Dict[str, Any]) -> None:
|
||||
Region is fixed to "global" here; cn-region operators use the CLI's ``--region cn``."""
|
||||
from hermes_cli.web_server import _profile_scope
|
||||
from hermes_cli.auth import (
|
||||
_minimax_poll_token,
|
||||
_minimax_resolve_token_expiry_unix,
|
||||
_minimax_save_auth_state,
|
||||
MINIMAX_OAUTH_GLOBAL_INFERENCE,
|
||||
MINIMAX_OAUTH_SCOPE)
|
||||
_minimax_poll_token, _minimax_resolve_token_expiry_unix, _minimax_save_auth_state,
|
||||
MINIMAX_OAUTH_GLOBAL_INFERENCE, MINIMAX_OAUTH_SCOPE,
|
||||
)
|
||||
import httpx
|
||||
portal_base_url = sess["portal_base_url"]
|
||||
client_id = sess["client_id"]
|
||||
user_code = sess["user_code"]
|
||||
code_verifier = sess["code_verifier"]
|
||||
interval_ms = sess.get("interval_ms")
|
||||
expired_in_raw = sess["expired_in_raw"]
|
||||
portal_base_url, client_id = sess["portal_base_url"], sess["client_id"]
|
||||
with httpx.Client(
|
||||
timeout=httpx.Timeout(15.0), headers={"Accept": "application/json"}, follow_redirects=True
|
||||
) as client:
|
||||
token_data = _minimax_poll_token(
|
||||
client=client,
|
||||
portal_base_url=portal_base_url,
|
||||
client_id=client_id,
|
||||
user_code=user_code,
|
||||
code_verifier=code_verifier,
|
||||
expired_in=expired_in_raw,
|
||||
interval_ms=interval_ms)
|
||||
client=client, portal_base_url=portal_base_url, client_id=client_id,
|
||||
user_code=sess["user_code"], code_verifier=sess["code_verifier"],
|
||||
expired_in=sess["expired_in_raw"], interval_ms=sess.get("interval_ms"),
|
||||
)
|
||||
now = datetime.now(timezone.utc)
|
||||
expires_at_ts = _minimax_resolve_token_expiry_unix(
|
||||
int(token_data["expired_in"]), now=now)
|
||||
expires_in_s = max(0, int(expires_at_ts - now.timestamp()))
|
||||
expires_at_ts = _minimax_resolve_token_expiry_unix(int(token_data["expired_in"]), now=now)
|
||||
auth_state = {
|
||||
"provider": "minimax-oauth",
|
||||
"region": sess.get("region", "global"),
|
||||
@@ -311,9 +284,9 @@ def _minimax_poller(session_id: str, sess: Dict[str, Any]) -> None:
|
||||
"refresh_token": token_data["refresh_token"],
|
||||
"resource_url": token_data.get("resource_url"),
|
||||
"obtained_at": now.isoformat(),
|
||||
"expires_at": datetime.fromtimestamp(
|
||||
expires_at_ts, tz=timezone.utc).isoformat(),
|
||||
"expires_in": expires_in_s}
|
||||
"expires_at": datetime.fromtimestamp(expires_at_ts, tz=timezone.utc).isoformat(),
|
||||
"expires_in": max(0, int(expires_at_ts - now.timestamp())),
|
||||
}
|
||||
with _profile_scope(_oauth_session_profile(session_id)):
|
||||
_minimax_save_auth_state(auth_state)
|
||||
|
||||
@@ -324,39 +297,29 @@ def _xai_device_poller(session_id: str, sess: Dict[str, Any]) -> None:
|
||||
from hermes_cli.web_server import _profile_scope
|
||||
import httpx
|
||||
from hermes_cli.auth import (
|
||||
_save_xai_oauth_tokens,
|
||||
_xai_oauth_discovery,
|
||||
_xai_oauth_poll_device_token,
|
||||
mark_provider_active_if_unset,
|
||||
unsuppress_credential_source)
|
||||
_save_xai_oauth_tokens, _xai_oauth_discovery, _xai_oauth_poll_device_token,
|
||||
mark_provider_active_if_unset, unsuppress_credential_source,
|
||||
)
|
||||
|
||||
device_code = sess["device_code"]
|
||||
interval = int(sess["interval"])
|
||||
expires_in = max(60, int(sess["expires_at"] - time.time()))
|
||||
discovery = _xai_oauth_discovery(20.0)
|
||||
with httpx.Client(
|
||||
timeout=httpx.Timeout(20.0), headers={"Accept": "application/json"}
|
||||
) as client:
|
||||
with httpx.Client(timeout=httpx.Timeout(20.0), headers={"Accept": "application/json"}) as client:
|
||||
token_data = _xai_oauth_poll_device_token(
|
||||
client,
|
||||
token_endpoint=discovery["token_endpoint"],
|
||||
device_code=device_code,
|
||||
expires_in=expires_in,
|
||||
poll_interval=interval)
|
||||
client, token_endpoint=discovery["token_endpoint"], device_code=sess["device_code"],
|
||||
expires_in=max(60, int(sess["expires_at"] - time.time())), poll_interval=int(sess["interval"]),
|
||||
)
|
||||
tokens = {
|
||||
"access_token": str(token_data.get("access_token", "") or "").strip(),
|
||||
"refresh_token": str(token_data.get("refresh_token", "") or "").strip(),
|
||||
"id_token": str(token_data.get("id_token", "") or "").strip(),
|
||||
"expires_in": token_data.get("expires_in"),
|
||||
"token_type": str(token_data.get("token_type") or "Bearer").strip() or "Bearer"}
|
||||
"token_type": str(token_data.get("token_type") or "Bearer").strip() or "Bearer",
|
||||
}
|
||||
with _profile_scope(_oauth_session_profile(session_id)):
|
||||
# set_active=False: persist without hijacking an existing active chat provider.
|
||||
_save_xai_oauth_tokens(
|
||||
tokens,
|
||||
discovery=discovery,
|
||||
tokens, discovery=discovery, auth_mode="oauth_device_code", set_active=False,
|
||||
last_refresh=datetime.now(timezone.utc).isoformat().replace("+00:00", "Z"),
|
||||
auth_mode="oauth_device_code",
|
||||
set_active=False)
|
||||
)
|
||||
# Mirror `hermes auth add xai-oauth`: first credential may become active; never overwrite.
|
||||
mark_provider_active_if_unset("xai-oauth")
|
||||
# The singleton write is the source of truth (the pool load seeds it as the canonical
|
||||
|
||||
@@ -3,17 +3,15 @@
|
||||
import hashlib
|
||||
import hmac
|
||||
import json
|
||||
import os
|
||||
import re
|
||||
import secrets
|
||||
import tempfile
|
||||
import time
|
||||
import urllib.request
|
||||
from pathlib import Path
|
||||
from typing import Dict
|
||||
|
||||
from hermes_constants import display_hermes_home
|
||||
from utils import atomic_replace
|
||||
from utils import atomic_json_write
|
||||
from hermes_cli.config import cfg_get
|
||||
|
||||
|
||||
@@ -21,13 +19,9 @@ _SUBSCRIPTIONS_FILENAME = "webhook_subscriptions.json"
|
||||
_SUBSCRIPTIONS_FILE_MODE = 0o600
|
||||
|
||||
|
||||
def _hermes_home() -> Path:
|
||||
from hermes_constants import get_hermes_home
|
||||
return get_hermes_home()
|
||||
|
||||
|
||||
def _subscriptions_path() -> Path:
|
||||
return _hermes_home() / _SUBSCRIPTIONS_FILENAME
|
||||
from hermes_constants import get_hermes_home
|
||||
return get_hermes_home() / _SUBSCRIPTIONS_FILENAME
|
||||
|
||||
|
||||
def _load_subscriptions() -> Dict[str, dict]:
|
||||
@@ -42,27 +36,9 @@ def _load_subscriptions() -> Dict[str, dict]:
|
||||
|
||||
|
||||
def _save_subscriptions(subs: Dict[str, dict]) -> None:
|
||||
path = _subscriptions_path()
|
||||
path.parent.mkdir(parents=True, exist_ok=True)
|
||||
# The file holds per-route HMAC secrets: chmod 0o600 the temp file BEFORE the atomic rename so
|
||||
# a permissive umask can't expose them in the create→rename window.
|
||||
fd, tmp_name = tempfile.mkstemp(prefix=f".{path.name}.", suffix=".tmp", dir=path.parent, text=True)
|
||||
tmp_path = Path(tmp_name)
|
||||
try:
|
||||
with os.fdopen(fd, "w", encoding="utf-8") as fh:
|
||||
json.dump(subs, fh, indent=2, ensure_ascii=False)
|
||||
fh.flush()
|
||||
os.fsync(fh.fileno())
|
||||
os.chmod(tmp_path, _SUBSCRIPTIONS_FILE_MODE)
|
||||
atomic_replace(tmp_path, path)
|
||||
# Re-assert: the destination may have existed with a broader mode that the replace kept.
|
||||
os.chmod(path, _SUBSCRIPTIONS_FILE_MODE)
|
||||
except Exception:
|
||||
try:
|
||||
tmp_path.unlink(missing_ok=True)
|
||||
except OSError:
|
||||
pass
|
||||
raise
|
||||
# The file holds per-route HMAC secrets: atomic_json_write fchmods the temp file 0o600 BEFORE the
|
||||
# rename (no umask window) and re-asserts the mode on the destination afterwards.
|
||||
atomic_json_write(_subscriptions_path(), subs, mode=_SUBSCRIPTIONS_FILE_MODE)
|
||||
|
||||
|
||||
def _get_webhook_config() -> dict:
|
||||
@@ -82,11 +58,10 @@ def _is_webhook_enabled() -> bool:
|
||||
def _get_webhook_base_url() -> str:
|
||||
wh = _get_webhook_config().get("extra", {})
|
||||
host = wh.get("host")
|
||||
port = wh.get("port", 8644)
|
||||
display_host = "localhost" if not host or host in {"0.0.0.0", "::"} else host
|
||||
if ":" in display_host and not display_host.startswith("["):
|
||||
display_host = f"[{display_host}]"
|
||||
return f"http://{display_host}:{port}"
|
||||
return f"http://{display_host}:{wh.get('port', 8644)}"
|
||||
|
||||
|
||||
def _setup_hint() -> str:
|
||||
@@ -172,8 +147,7 @@ def _cmd_subscribe(args):
|
||||
print(" Mode: direct delivery (no agent, zero LLM cost)")
|
||||
if route.get("prompt"):
|
||||
prompt_preview = route["prompt"][:80] + ("..." if len(route["prompt"]) > 80 else "")
|
||||
label = "Message" if route.get("deliver_only") else "Prompt"
|
||||
print(f" {label}: {prompt_preview}")
|
||||
print(f" {'Message' if route.get('deliver_only') else 'Prompt'}: {prompt_preview}")
|
||||
if route.get("script"):
|
||||
print(f" Script: {route['script']}")
|
||||
print("\n Configure your service to POST to the URL above.")
|
||||
|
||||
@@ -60,21 +60,17 @@ class WinPtyBridge:
|
||||
spawn_env = (
|
||||
build_subprocess_env(scrub_secrets=False, inherit_profile_home=False)
|
||||
if env is None else dict(env))
|
||||
if not spawn_env.get("TERM"):
|
||||
spawn_env["TERM"] = "xterm-256color"
|
||||
spawn_env["TERM"] = spawn_env.get("TERM") or "xterm-256color"
|
||||
# pywinpty mirrors ptyprocess: dimensions=(rows, cols).
|
||||
proc = PtyProcess.spawn(list(argv), cwd=cwd, env=spawn_env, dimensions=(rows, cols)) # type: ignore[union-attr]
|
||||
return cls(proc)
|
||||
return cls(PtyProcess.spawn(list(argv), cwd=cwd, env=spawn_env, dimensions=(rows, cols))) # type: ignore[union-attr]
|
||||
|
||||
@property
|
||||
def pid(self) -> int:
|
||||
return int(self._proc.pid)
|
||||
|
||||
def is_alive(self) -> bool:
|
||||
if self._closed:
|
||||
return False
|
||||
try:
|
||||
return bool(self._proc.isalive())
|
||||
return not self._closed and bool(self._proc.isalive())
|
||||
except Exception:
|
||||
return False
|
||||
|
||||
|
||||
@@ -2,6 +2,7 @@
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import importlib
|
||||
import json
|
||||
import os
|
||||
import re
|
||||
@@ -29,14 +30,8 @@ def _win32() -> Any:
|
||||
win32security); import is deferred so the module imports on non-Windows hosts."""
|
||||
if sys.platform != "win32":
|
||||
raise RuntimeError("Windows SSH runtime is only available on Windows")
|
||||
import ntsecuritycon
|
||||
import pywintypes
|
||||
import win32api
|
||||
import win32con
|
||||
import win32file
|
||||
import win32security
|
||||
return SimpleNamespace(ntsecuritycon=ntsecuritycon, pywintypes=pywintypes, win32api=win32api,
|
||||
win32con=win32con, win32file=win32file, win32security=win32security)
|
||||
names = ("ntsecuritycon", "pywintypes", "win32api", "win32con", "win32file", "win32security")
|
||||
return SimpleNamespace(**{name: importlib.import_module(name) for name in names})
|
||||
|
||||
|
||||
def _check(pattern: re.Pattern, value: str, message: str) -> str:
|
||||
@@ -136,14 +131,10 @@ def _open(path: Path, access: int, creation: int, flags: int, share: int = 0):
|
||||
win32file = _win32().win32file
|
||||
handle = win32file.CreateFile(str(path), access, share, _security_attributes(), creation, flags, None)
|
||||
try:
|
||||
actual = win32file.GetFinalPathNameByHandle(handle, 0)
|
||||
if actual.startswith("\\\\?\\"):
|
||||
actual = actual[4:]
|
||||
expected = os.path.abspath(str(path))
|
||||
if os.path.normcase(actual) != os.path.normcase(expected):
|
||||
actual = win32file.GetFinalPathNameByHandle(handle, 0).removeprefix("\\\\?\\")
|
||||
if os.path.normcase(actual) != os.path.normcase(os.path.abspath(str(path))):
|
||||
raise OSError("Windows SSH runtime handle escaped its expected path")
|
||||
attributes = win32file.GetFileInformationByHandle(handle)[0]
|
||||
if attributes & 0x400:
|
||||
if win32file.GetFileInformationByHandle(handle)[0] & 0x400: # FILE_ATTRIBUTE_REPARSE_POINT
|
||||
raise OSError("Windows SSH runtime path contains a reparse point")
|
||||
_verify_security(handle)
|
||||
return handle
|
||||
@@ -152,21 +143,28 @@ def _open(path: Path, access: int, creation: int, flags: int, share: int = 0):
|
||||
raise
|
||||
|
||||
|
||||
def _read_shared(path: Path, limit: int, share: int) -> bytes | None:
|
||||
"""Open ``path`` read-only with ``share`` and read up to ``limit`` bytes; None when missing."""
|
||||
def _open_existing(path: Path, access: int, extra_flags: int = 0, share: int = 0):
|
||||
"""``_open`` an existing file (FILE_ATTRIBUTE_NORMAL | reparse guard); None when it is missing."""
|
||||
w = _win32()
|
||||
pywintypes, win32con, win32file = w.pywintypes, w.win32con, w.win32file
|
||||
try:
|
||||
handle = _open(path, win32con.GENERIC_READ | win32con.READ_CONTROL, win32con.OPEN_EXISTING,
|
||||
win32con.FILE_ATTRIBUTE_NORMAL | _OPEN_REPARSE_POINT, share)
|
||||
except pywintypes.error as exc:
|
||||
return _open(path, access, w.win32con.OPEN_EXISTING,
|
||||
w.win32con.FILE_ATTRIBUTE_NORMAL | _OPEN_REPARSE_POINT | extra_flags, share)
|
||||
except w.pywintypes.error as exc:
|
||||
if exc.winerror in (2, 3):
|
||||
return None
|
||||
raise
|
||||
|
||||
|
||||
def _read_shared(path: Path, limit: int, share: int) -> bytes | None:
|
||||
"""Read up to ``limit`` bytes of ``path`` opened read-only with ``share``; None when missing."""
|
||||
w = _win32()
|
||||
handle = _open_existing(path, w.win32con.GENERIC_READ | w.win32con.READ_CONTROL, share=share)
|
||||
if handle is None:
|
||||
return None
|
||||
try:
|
||||
return win32file.ReadFile(handle, limit)[1]
|
||||
return w.win32file.ReadFile(handle, limit)[1]
|
||||
finally:
|
||||
win32file.CloseHandle(handle)
|
||||
w.win32file.CloseHandle(handle)
|
||||
|
||||
|
||||
def _write_new(path: Path, data: bytes, share: int = 0) -> None:
|
||||
@@ -200,8 +198,7 @@ def _ensure_directory(path: Path) -> None:
|
||||
|
||||
|
||||
def _ensure_scope(ownership_id: str) -> Path:
|
||||
root = _root()
|
||||
_ensure_directory(root)
|
||||
_ensure_directory(_root())
|
||||
directory = _directory(ownership_id)
|
||||
_ensure_directory(directory)
|
||||
return directory
|
||||
@@ -224,9 +221,8 @@ def read_token(path_value: str) -> str:
|
||||
w = _win32()
|
||||
win32con, win32file = w.win32con, w.win32file
|
||||
path = Path(path_value)
|
||||
root = _root()
|
||||
try:
|
||||
relative = path.relative_to(root)
|
||||
relative = path.relative_to(_root())
|
||||
except ValueError as exc:
|
||||
raise SystemExit("--ssh-session-token-file must be under the desktop-ssh directory") from exc
|
||||
if len(relative.parts) != 2 or not _HEX32.fullmatch(relative.parts[0]) or not re.fullmatch(r"[0-9a-f]{16}\.token", relative.parts[1]):
|
||||
@@ -284,15 +280,10 @@ def write_lock(ownership_id: str, payload: dict[str, Any]) -> None:
|
||||
|
||||
def remove_artifact(path: Path) -> bool:
|
||||
w = _win32()
|
||||
pywintypes, win32con, win32file = w.pywintypes, w.win32con, w.win32file
|
||||
try:
|
||||
handle = _open(path, win32con.DELETE | win32con.READ_CONTROL, win32con.OPEN_EXISTING,
|
||||
win32con.FILE_ATTRIBUTE_NORMAL | _OPEN_REPARSE_POINT | _DELETE_ON_CLOSE)
|
||||
except pywintypes.error as exc:
|
||||
if exc.winerror in (2, 3):
|
||||
return False
|
||||
raise
|
||||
win32file.CloseHandle(handle)
|
||||
handle = _open_existing(path, w.win32con.DELETE | w.win32con.READ_CONTROL, _DELETE_ON_CLOSE)
|
||||
if handle is None:
|
||||
return False
|
||||
w.win32file.CloseHandle(handle)
|
||||
return True
|
||||
|
||||
|
||||
@@ -388,7 +379,6 @@ def spawn_backend(payload: dict[str, Any]) -> dict[str, Any]:
|
||||
raise ValueError("Hermes path must be absolute")
|
||||
hermes_path = os.path.abspath(configured_path)
|
||||
token_path = str(_token_path(ownership_id, spawn_nonce))
|
||||
log_path = _log_path(ownership_id, spawn_nonce)
|
||||
profile = str(payload.get("profile") or "")
|
||||
if len(profile) > 256 or any(ch in profile for ch in "\x00\r\n"):
|
||||
raise ValueError("invalid profile")
|
||||
@@ -412,7 +402,7 @@ def spawn_backend(payload: dict[str, Any]) -> dict[str, Any]:
|
||||
env["VIRTUAL_ENV"] = os.path.dirname(venv_dir)
|
||||
env.pop("PYTHONPATH", None)
|
||||
_ensure_scope(ownership_id)
|
||||
creationflags = 0x00000008 | 0x00000200 | 0x01000000
|
||||
log_path = _log_path(ownership_id, spawn_nonce)
|
||||
win32con = _win32().win32con
|
||||
log_handle = _open(log_path, win32con.GENERIC_WRITE | win32con.READ_CONTROL,
|
||||
win32con.CREATE_NEW, win32con.FILE_ATTRIBUTE_NORMAL | _OPEN_REPARSE_POINT,
|
||||
@@ -420,8 +410,9 @@ def spawn_backend(payload: dict[str, Any]) -> dict[str, Any]:
|
||||
import msvcrt
|
||||
log_fd = msvcrt.open_osfhandle(int(log_handle), os.O_WRONLY)
|
||||
with os.fdopen(log_fd, "wb", buffering=0) as log_stream:
|
||||
# DETACHED_PROCESS | CREATE_NEW_PROCESS_GROUP | CREATE_BREAKAWAY_FROM_JOB
|
||||
process = subprocess.Popen(args, stdin=subprocess.DEVNULL, stdout=log_stream, stderr=log_stream,
|
||||
close_fds=True, creationflags=creationflags, env=env)
|
||||
close_fds=True, creationflags=0x00000008 | 0x00000200 | 0x01000000, env=env)
|
||||
creation_time_ns = int(__import__("psutil").Process(process.pid).create_time() * 1_000_000_000)
|
||||
return {"pid": process.pid, "creationTimeNs": str(creation_time_ns),
|
||||
"logPath": str(log_path), "tokenPath": token_path}
|
||||
|
||||
@@ -60,7 +60,6 @@ _ACTIONS = {"list": _list, "prune": _prune}
|
||||
|
||||
def cmd_worktree(args) -> int:
|
||||
from hermes_cli import worktree_gc
|
||||
|
||||
repo_root = getattr(args, "repo", None)
|
||||
if not repo_root:
|
||||
import cli as _cli
|
||||
|
||||
@@ -97,12 +97,11 @@ def _archive_untracked(tree: Path, untracked: List[str]) -> Optional[Path]:
|
||||
src = tree / rel
|
||||
if not src.exists() or src.is_symlink():
|
||||
continue
|
||||
target = dest / rel
|
||||
target.parent.mkdir(parents=True, exist_ok=True)
|
||||
(dest / rel).parent.mkdir(parents=True, exist_ok=True)
|
||||
if src.is_dir():
|
||||
shutil.copytree(src, target, dirs_exist_ok=True)
|
||||
shutil.copytree(src, dest / rel, dirs_exist_ok=True)
|
||||
else:
|
||||
shutil.copy2(src, target)
|
||||
shutil.copy2(src, dest / rel)
|
||||
return dest if dest.exists() else None
|
||||
except Exception as exc:
|
||||
logger.warning("Could not archive untracked files from %s: %s", tree, exc)
|
||||
@@ -138,7 +137,6 @@ def _classify_tree(_cli, repo_root: str, entry: Path, merge_cache, remote_heads)
|
||||
def audit_worktrees(repo_root: str, *, with_sizes: bool = True) -> List[TreeRecord]:
|
||||
"""Classify every tree under ``.worktrees/`` without mutating anything."""
|
||||
import cli as _cli # lazy: cli.py is heavy
|
||||
|
||||
worktrees_dir = Path(repo_root) / ".worktrees"
|
||||
if not worktrees_dir.exists():
|
||||
return []
|
||||
@@ -228,7 +226,6 @@ def audit_branches(repo_root: str) -> List[BranchRecord]:
|
||||
"""Classify EVERY local branch: deletable when fully merged OR every commit is patch-equivalent
|
||||
upstream (``git cherry``) and not checked out. The gate is content reachability, not name."""
|
||||
import cli as _cli
|
||||
|
||||
if _cli._repo_is_shallow(repo_root):
|
||||
_cli._deepen_shallow_repo(repo_root)
|
||||
|
||||
@@ -280,7 +277,6 @@ def audit_branches(repo_root: str) -> List[BranchRecord]:
|
||||
|
||||
# Read-only, so parallel: hundreds of local branches × ~0.2-1s cherry probes would be minutes.
|
||||
import concurrent.futures
|
||||
|
||||
workers = max(1, min(8, (os.cpu_count() or 4), len(branches)))
|
||||
if workers > 1:
|
||||
try:
|
||||
|
||||
@@ -115,15 +115,15 @@ class ApplyResult:
|
||||
|
||||
def _walk_to_parent(yaml_doc: Any, dotted_path: str) -> "tuple[Any, str]":
|
||||
"""Resolve a dotted slot path to (parent_mapping, leaf_key)."""
|
||||
parts = dotted_path.split(".")
|
||||
if len(parts) < 2:
|
||||
*parents, leaf = dotted_path.split(".")
|
||||
if not parents:
|
||||
raise ValueError(f"Path must have at least one parent: {dotted_path!r}")
|
||||
node = yaml_doc
|
||||
for segment in parts[:-1]:
|
||||
for segment in parents:
|
||||
if not isinstance(node, dict) or segment not in node:
|
||||
raise KeyError(f"Path segment {segment!r} missing in {dotted_path!r}")
|
||||
node = node[segment]
|
||||
return node, parts[-1]
|
||||
return node, leaf
|
||||
|
||||
|
||||
def apply_migration(
|
||||
@@ -133,7 +133,6 @@ def apply_migration(
|
||||
Unless ``backup=False`` a copy goes to ``<config_path>.bak-pre-migrate-xai-YYYYMMDD-HHMMSS``.
|
||||
"""
|
||||
from ruamel.yaml import YAML # local import — avoid hard dep at module load
|
||||
|
||||
config_path = Path(config_path)
|
||||
if not config_path.exists():
|
||||
raise FileNotFoundError(config_path)
|
||||
@@ -169,7 +168,6 @@ def apply_migration(
|
||||
|
||||
from hermes_cli.config import require_readable_config_before_write
|
||||
from utils import atomic_write_text
|
||||
|
||||
require_readable_config_before_write(config_path)
|
||||
# Dump to a buffer, then atomic-write: ``open(path, "w")`` truncates before the dump runs, so a
|
||||
# crash mid-write would leave config.yaml empty (and with ``--no-backup`` that is the only
|
||||
|
||||
Reference in New Issue
Block a user