Skip to content
Open
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
442 changes: 442 additions & 0 deletions api/agent_update.py

Large diffs are not rendered by default.

27 changes: 20 additions & 7 deletions api/gateway_chat.py
Original file line number Diff line number Diff line change
Expand Up @@ -278,15 +278,28 @@ def webui_gateway_chat_enabled(config_data=None, environ: dict[str, str] | None
return webui_chat_backend_mode(config_data, environ) == "gateway"


def _gateway_base_url(config_data=None, environ: dict[str, str] | None = None) -> str:
def resolve_effective_gateway_target(config_data=None, environ: dict[str, str] | None = None) -> dict[str, object]:
"""Resolve the single URL that receives browser chat traffic.

Health URLs are deliberately absent here; they are observation inputs and
cannot change lifecycle ownership.
"""
source = os.environ if environ is None else environ
cfg = config_data if isinstance(config_data, dict) else {}
raw = str(
source.get(_WEBUI_GATEWAY_BASE_URL_ENV)
or cfg.get("webui_gateway_base_url")
or "http://127.0.0.1:8642"
).strip()
return raw.rstrip("/") or "http://127.0.0.1:8642"
raw = str(source.get(_WEBUI_GATEWAY_BASE_URL_ENV) or cfg.get("webui_gateway_base_url") or "http://127.0.0.1:8642").strip()
url = raw.rstrip("/") or "http://127.0.0.1:8642"
try:
import ipaddress
host = (urllib.parse.urlparse(url).hostname or "").strip().lower().rstrip(".")
local = host in {"localhost", "localhost.localdomain"} or ipaddress.ip_address(host).is_loopback
except ValueError:
local = False
source_name = "environment" if source.get(_WEBUI_GATEWAY_BASE_URL_ENV) else "config" if cfg.get("webui_gateway_base_url") else "default"
return {"url": url, "source": source_name, "local": local}


def _gateway_base_url(config_data=None, environ: dict[str, str] | None = None) -> str:
return str(resolve_effective_gateway_target(config_data, environ)["url"])


def _gateway_api_key(environ: dict[str, str] | None = None) -> str:
Expand Down
163 changes: 32 additions & 131 deletions api/updates.py
Original file line number Diff line number Diff line change
Expand Up @@ -24,9 +24,6 @@
from pathlib import Path
from urllib.parse import urlparse

from api.agent_health import get_active_profile_gateway_running_pid
from api.gateway_restart import restart_active_profile_gateway
from api.profiles import get_active_profile_name
from api.config import REPO_ROOT, STREAMS, STREAMS_LOCK

logger = logging.getLogger(__name__)
Expand All @@ -44,7 +41,6 @@
_check_in_progress = False
_apply_lock = threading.Lock() # prevents concurrent stash/pull/pop on same repo
CACHE_TTL = 1800 # 30 minutes
_AGENT_GATEWAY_RESTART_RETRY_DELAY_S = 1.0
_FORCE_DIRTY_PROBE_TIMEOUT = 5
_GIT_DIAGNOSTIC_MAX_CHARS = 300
_CREDENTIAL_IN_URL_RE = re.compile(r"([a-zA-Z][a-zA-Z0-9+.-]*://)([^/@\s'\"]+)@")
Expand Down Expand Up @@ -316,14 +312,16 @@ def apply_clear_lock(target: str) -> dict:
return {'ok': False, 'message': 'Update already in progress'}

try:
if target == 'agent':
return _apply_agent_transaction()
if target == 'webui':
path = REPO_ROOT
elif target == 'agent':
path = _AGENT_DIR
else:
return {'ok': False, 'message': f'Unknown target: {target}'}

if path is None or not (path / '.git').exists():
if path is None:
return {'ok': False, 'message': 'Not a git repository'}
if not (path / '.git').exists():
return {'ok': False, 'message': 'Not a git repository'}

inv = _inventory_locks(path)
Expand Down Expand Up @@ -1833,84 +1831,6 @@ def _do():
threading.Thread(target=_do, daemon=True).start()


def _ensure_gateway_restart_for_agent_update() -> tuple[bool, dict]:
"""Run the active-profile gateway restart when agent checkout changed.

Returns:
(ok, restart_payload) where:
- ok is False when restart did not complete and callers must abort success.
- restart_payload contains helper status fields for response shaping.
"""
target_profile = str(get_active_profile_name() or "default").strip() or "default"
gateway_pid_before_restart = get_active_profile_gateway_running_pid(profile=target_profile)
restart_result = restart_active_profile_gateway(profile=target_profile)
status = str(restart_result.get("status") or "")
if status in {"completed", "in_progress"}:
return True, restart_result
if status != "failed":
return False, restart_result

# launchd can briefly fail to spawn the replacement gateway while it is
# rotating the supervised process (#6045). Retry exactly once after a
# bounded delay so an already-applied Agent update is not reported as a
# complete failure because of that transient process handoff.
time.sleep(_AGENT_GATEWAY_RESTART_RETRY_DELAY_S)
retry_result = restart_active_profile_gateway(profile=target_profile)
retry_status = str(retry_result.get("status") or "")
if retry_status in {"completed", "in_progress"}:
return True, {
**retry_result,
"retry_attempted": True,
"initial_failure": restart_result.get("message"),
}
if retry_status != "failed":
return False, {
**retry_result,
"retry_attempted": True,
"initial_failure": restart_result.get("message"),
}

# A restart command can still exit non-zero after launchd has recovered the
# service. Only accept that recovery when the confirmed local PID changed;
# a merely-alive old gateway has not loaded the updated Agent checkout.
time.sleep(_AGENT_GATEWAY_RESTART_RETRY_DELAY_S)
gateway_pid_after_retry = get_active_profile_gateway_running_pid(profile=target_profile)
if (
gateway_pid_before_restart is not None
and gateway_pid_after_retry is not None
and gateway_pid_after_retry != gateway_pid_before_restart
):
return True, {
"status": "completed",
"message": "Gateway service recovered after a transient restart failure",
"retry_attempted": True,
"process_replaced": True,
"initial_failure": restart_result.get("message"),
"retry_failure": retry_result.get("message"),
}

initial_message = str(restart_result.get("message") or "Restart failed")
retry_message = str(retry_result.get("message") or "retry did not complete")
return False, {
**retry_result,
"message": f"{initial_message}; recovery retry did not complete: {retry_message}",
"retry_attempted": True,
"initial_failure": restart_result.get("message"),
}


def _agent_gateway_restart_failure_message(target: str, restart_result: dict) -> str:
if restart_result.get("message"):
return (
f'{target} updated, but gateway restart did not complete: '
f'{restart_result["message"]}. Run `hermes gateway restart` manually.'
)
return (
f'{target} updated, but gateway restart did not complete. '
'Run `hermes gateway restart` manually.'
)


def _discard_local_changes(path: Path, reset_ref: str) -> bool:
"""Discard local changes and reset *path* to *reset_ref*."""
# Do not use -x: ignored build/cache artifacts should survive force update.
Expand Down Expand Up @@ -1961,15 +1881,19 @@ def apply_force_update(target: str, channel=None) -> dict:
if blocker_snapshot.get('restart_blocked'):
return _restart_blocked_response(target, blocker_snapshot)

if target == 'agent':
if not _apply_lock.acquire(blocking=False):
return {'ok': False, 'message': 'Update already in progress'}
try:
return _apply_agent_transaction(force=True)
finally:
_apply_lock.release()

if not _apply_lock.acquire(blocking=False):
return {'ok': False, 'message': 'Update already in progress'}
try:
if target == 'webui':
path = REPO_ROOT
elif target == 'agent':
path = _AGENT_DIR
# Channel is WebUI-only — the Agent always uses the default channel.
channel = DEFAULT_UPDATE_CHANNEL
else:
return {'ok': False, 'message': f'Unknown target: {target}'}

Expand Down Expand Up @@ -2017,8 +1941,8 @@ def apply_force_update(target: str, channel=None) -> dict:
'channel': channel,
}

# Rewind guard (Codex CORE #3): refuse to reset --hard onto a ref that
# is an ANCESTOR of HEAD — that would downgrade the checkout. This is the
# Rewind guard: refuse to reset --hard onto a ref that
# is an ANCESTOR of HEAD, since that would downgrade the checkout. This is the
# switch-back-to-stable-while-ahead case. A ref that is a descendant of
# HEAD (normal update / opt-in to experimental) fast-forwards fine and is
# allowed. Refs on a divergent line (neither ancestor nor descendant) are
Expand All @@ -2043,16 +1967,6 @@ def apply_force_update(target: str, channel=None) -> dict:
with _cache_lock:
_update_cache['checked_at'] = 0

if target == 'agent':
gateway_ok, gateway_result = _ensure_gateway_restart_for_agent_update()
if not gateway_ok:
return {
'ok': False,
'message': _agent_gateway_restart_failure_message(target, gateway_result),
'target': target,
'gateway_restart': gateway_result.get('status'),
}

_schedule_restart()

response = {
Expand All @@ -2061,8 +1975,6 @@ def apply_force_update(target: str, channel=None) -> dict:
'target': target,
'restart_scheduled': True,
}
if target == 'agent':
response['gateway_restart'] = gateway_result.get('status')
return response
finally:
_apply_lock.release()
Expand All @@ -2085,6 +1997,21 @@ def apply_update(target, channel=None):
_apply_lock.release()


def _apply_agent_transaction(*, force: bool = False) -> dict:
"""Map the Agent-owned transaction into the WebUI response contract."""
from api import agent_update

result = agent_update.apply_agent_update(force=force)
if result.get('reload_eligible'):
with _cache_lock:
_update_cache['checked_at'] = 0
_schedule_restart()
result['restart_scheduled'] = True
else:
result.pop('restart_scheduled', None)
return result


def _restore_stash_after_pull_failure(
target: str,
path: Path,
Expand Down Expand Up @@ -2123,13 +2050,10 @@ def _restore_stash_after_pull_failure(
def _apply_update_inner(target, channel=DEFAULT_UPDATE_CHANNEL):
"""Inner implementation of apply_update, called under _apply_lock."""
channel = _normalize_channel(channel)
if target == 'agent':
return _apply_agent_transaction()
if target == 'webui':
path = REPO_ROOT
elif target == 'agent':
path = _AGENT_DIR
# Channel is WebUI-only — the Agent always uses the default channel
# regardless of the user's WebUI selection (see check_for_updates).
channel = DEFAULT_UPDATE_CHANNEL
else:
return {'ok': False, 'message': f'Unknown target: {target}'}

Expand Down Expand Up @@ -2357,15 +2281,6 @@ def _apply_update_inner(target, channel=DEFAULT_UPDATE_CHANNEL):
with _cache_lock:
_update_cache['checked_at'] = 0

if target == 'agent':
gateway_ok, gateway_result = _ensure_gateway_restart_for_agent_update()
if not gateway_ok:
return {
'ok': False,
'message': _agent_gateway_restart_failure_message(target, gateway_result),
'target': target,
'gateway_restart': gateway_result.get('status'),
}
_schedule_restart()
response = {
'ok': True,
Expand All @@ -2381,24 +2296,12 @@ def _apply_update_inner(target, channel=DEFAULT_UPDATE_CHANNEL):
'restart_scheduled': True,
'stash_conflict': True,
}
if target == 'agent':
response['gateway_restart'] = gateway_result.get('status')
return response

# Invalidate cache
with _cache_lock:
_update_cache['checked_at'] = 0

if target == 'agent':
gateway_ok, gateway_result = _ensure_gateway_restart_for_agent_update()
if not gateway_ok:
return {
'ok': False,
'message': _agent_gateway_restart_failure_message(target, gateway_result),
'target': target,
'gateway_restart': gateway_result.get('status'),
}

# Schedule a self-restart so the updated code is loaded fresh. A plain
# git pull leaves stale Python modules in sys.modules — agent imports that
# reference new symbols (functions, classes) added in the update will fail
Expand All @@ -2424,6 +2327,4 @@ def _apply_update_inner(target, channel=DEFAULT_UPDATE_CHANNEL):
'target': target,
'restart_scheduled': True,
}
if target == 'agent':
response['gateway_restart'] = gateway_result.get('status')
return response
Loading
Loading