From 8f6a568646bebfe5f71c529defb3ec156df3cad2 Mon Sep 17 00:00:00 2001 From: m4 Date: Wed, 12 Aug 2026 16:42:10 +0800 Subject: [PATCH] feat(update): system update/status/rollback HTTP endpoints with system:write scope --- EvoScientist/langgraph_dev/http.py | 273 ++++++++++++++++++++++++++++- tests/test_system_update_http.py | 196 +++++++++++++++++++++ 2 files changed, 468 insertions(+), 1 deletion(-) create mode 100644 tests/test_system_update_http.py diff --git a/EvoScientist/langgraph_dev/http.py b/EvoScientist/langgraph_dev/http.py index ad89b9c..4a658b5 100644 --- a/EvoScientist/langgraph_dev/http.py +++ b/EvoScientist/langgraph_dev/http.py @@ -42,6 +42,16 @@ from starlette.routing import Route from EvoScientist.config.legacy_artifacts import assert_no_legacy_artifacts from EvoScientist.model_registry.http_api import model_registry_routes +from EvoScientist.updater import ( + DEPLOY_DOCKER, + UpdateInProgressError, + build_plan, + copy_updater, + detect_deployment, + probe_system, + spawn_updater, + write_plan_locked, +) from EvoScientist.sessions import ( MAIN_THREAD_FILTER_PARAMS, MAIN_THREAD_FILTER_SQL, @@ -53,6 +63,237 @@ from EvoScientist.sessions import ( _logger = logging.getLogger(__name__) +async def _authenticate_system_scoped(request: Request, scope: str) -> JSONResponse | None: + """Shared auth for /internal/system/* routes; None means authorized.""" + import uuid + + from EvoScientist.model_registry.errors import ( + PLATFORM_CONFIG_MISSING, + ModelRegistryError, + ) + from EvoScientist.model_registry.http_api import get_default_services + from EvoScientist.model_registry.platform import PlatformConfigError + + request_id = uuid.uuid4().hex + try: + try: + services = await asyncio.to_thread(get_default_services) + except PlatformConfigError: + raise ModelRegistryError( + PLATFORM_CONFIG_MISSING, + "The platform security configuration is missing or invalid; " + "the system endpoints are unavailable.", + ) from None + await asyncio.to_thread( + services.authenticator.authenticate, + request.headers, + required_scope=scope, + require_thread_id=False, + ) + except ModelRegistryError as exc: + return JSONResponse( + exc.payload(request_id=request_id), status_code=exc.http_status + ) + return None + + +async def _authenticate_system(request: Request) -> JSONResponse | None: + return await _authenticate_system_scoped(request, "system:read") + + +async def _authenticate_system_write(request: Request) -> JSONResponse | None: + """system:write variant for mutating routes (update / rollback).""" + return await _authenticate_system_scoped(request, "system:write") + + +async def get_system_version(request: Request) -> JSONResponse: + """Report installed vs latest published version (update checker).""" + from EvoScientist.update_check import get_update_info + + denied = await _authenticate_system(request) + if denied is not None: + return denied + info = await asyncio.to_thread( + get_update_info, force=request.query_params.get("force") == "true" + ) + return JSONResponse(info) + + +async def post_system_version_download(request: Request) -> JSONResponse: + """Stage the latest (or requested) release artifact under the workspace.""" + from EvoScientist.update_check import UpdateDownloadError, download_update + + denied = await _authenticate_system(request) + if denied is not None: + return denied + try: + payload = await request.json() + except json.JSONDecodeError: + payload = None + version = None + if isinstance(payload, dict) and payload.get("version") is not None: + if not isinstance(payload["version"], str): + return JSONResponse({"error": "version must be a string"}, status_code=400) + version = payload["version"] + try: + result = await asyncio.to_thread(download_update, version) + except UpdateDownloadError as exc: + return JSONResponse({"error": str(exc)}, status_code=400) + return JSONResponse(result) + + +def _schedule_self_exit(delay: float = 0.5) -> None: + """SIGTERM this process after the response has been flushed to the client.""" + import signal + import threading + + def _exit() -> None: + os.kill(os.getpid(), signal.SIGTERM) + + threading.Timer(delay, _exit).start() + + +def _updates_staging() -> Path: + from EvoScientist.update_check import _staging_dir + + return _staging_dir() + + +async def _start_apply(request: Request, version: str) -> JSONResponse: + """Shared update/rollback pipeline after the target version is resolved.""" + from EvoScientist import update_check + + probe = probe_system() + kind = detect_deployment(probe) + if kind == DEPLOY_DOCKER: + return JSONResponse( + { + "status": "manual", + "guidance": "In-place update is not supported inside Docker; " + "run: docker compose pull && docker compose up -d", + } + ) + try: + staged = await asyncio.to_thread(update_check.download_update, version) + except update_check.UpdateDownloadError as exc: + return JSONResponse( + {"code": "DOWNLOAD_FAILED", "message": str(exc)}, status_code=502 + ) + + import sys as _sys + + staging = _updates_staging() + plan = build_plan( + version=staged["version"], + artifact=Path(staged["path"]), + kind=kind, + python=_sys.executable, + argv=_sys.argv, + cwd=Path.cwd(), + deploy_mode=os.environ.get("EVOSCIENTIST_DEPLOY_MODE", ""), + staging=staging, + systemd_unit=os.environ.get("EVOSCIENTIST_SYSTEMD_UNIT") or None, + ) + plan_path = staging / f"v{staged['version']}" / "plan.json" + try: + write_plan_locked(plan, plan_path) + except UpdateInProgressError: + return JSONResponse( + {"code": "UPDATE_IN_PROGRESS", "message": "an update is already in progress"}, + status_code=409, + ) + updater_copy = copy_updater(plan_path.parent) + spawn_updater(updater_copy, plan_path, os.getpid()) + _schedule_self_exit() + return JSONResponse( + { + "operation_id": f"upd-{staged['version']}", + "status": "applying", + "need_restart": True, + }, + status_code=202, + ) + + +async def post_system_update(request: Request) -> JSONResponse: + """One-click update to the latest published version.""" + from EvoScientist.update_check import get_update_info + + denied = await _authenticate_system_write(request) + if denied is not None: + return denied + info = await asyncio.to_thread(get_update_info, force=True) + if not info["has_update"]: + return JSONResponse( + { + "code": "ALREADY_UP_TO_DATE", + "current_version": info["current_version"], + "latest_version": info["latest_version"], + }, + status_code=409, + ) + if info.get("breaking_db") and request.query_params.get("confirm_breaking") != "true": + return JSONResponse( + { + "code": "BREAKING_DB_CONFIRM_REQUIRED", + "message": "target version contains breaking DB changes; " + "retry with ?confirm_breaking=true", + }, + status_code=409, + ) + return await _start_apply(request, info["latest_version"]) + + +async def get_system_update_status(request: Request) -> JSONResponse: + """Report the result of the most recent update attempt.""" + denied = await _authenticate_system(request) + if denied is not None: + return denied + result_file = _updates_staging() / "last-result.json" + if not result_file.exists(): + return JSONResponse({"status": "none"}) + return JSONResponse(json.loads(result_file.read_text(encoding="utf-8"))) + + +async def get_system_rollback_versions(request: Request) -> JSONResponse: + """List versions the installation may be rolled back to.""" + from EvoScientist.update_check import list_rollback_versions + + denied = await _authenticate_system(request) + if denied is not None: + return denied + versions = await asyncio.to_thread(list_rollback_versions) + return JSONResponse({"versions": versions}) + + +async def post_system_update_rollback(request: Request) -> JSONResponse: + """Roll back to an allowed older version via the same updater pipeline.""" + from EvoScientist.update_check import is_allowed_rollback + + denied = await _authenticate_system_write(request) + if denied is not None: + return denied + try: + payload = await request.json() + except json.JSONDecodeError: + payload = None + version = payload.get("version") if isinstance(payload, dict) else None + if not isinstance(version, str) or not version.strip(): + return JSONResponse( + {"code": "VERSION_REQUIRED", "message": 'body must be {"version": "x.y.z"}'}, + status_code=400, + ) + if not await asyncio.to_thread(is_allowed_rollback, version): + return JSONResponse( + { + "code": "ROLLBACK_VERSION_NOT_ALLOWED", + "message": f"{version} is not in the rollback list", + }, + status_code=400, + ) + return await _start_apply(request, version.strip().lstrip("vV")) + + def _message_type(message: Any) -> str | None: if isinstance(message, dict): role = message.get("role") @@ -276,7 +517,7 @@ async def _read_thread_runtime_state( async def get_final_answer(request: Request) -> JSONResponse: """Return the latest checkpointed assistant answer for a WebUI thread. - This is a recovery surface for the WebUI stream consumer: when browser-side + This is a recovery surface for the WebUI stream consumer. When browser-side SSE is interrupted but the langgraph run continues server-side, the final answer is already persisted in ``sessions.db``. The route centralizes the non-trivial "latest AIMessage text only" extraction so the browser does not @@ -875,6 +1116,36 @@ app = Starlette( # /api/config, /api/default-model) and the admin-token check were # removed with the unified model configuration refactor. *model_registry_routes(), + Route( + "/internal/system/version", + get_system_version, + methods=["GET"], + ), + Route( + "/internal/system/version/download", + post_system_version_download, + methods=["POST"], + ), + Route( + "/internal/system/update", + post_system_update, + methods=["POST"], + ), + Route( + "/internal/system/update/status", + get_system_update_status, + methods=["GET"], + ), + Route( + "/internal/system/update/rollback", + post_system_update_rollback, + methods=["POST"], + ), + Route( + "/internal/system/rollback-versions", + get_system_rollback_versions, + methods=["GET"], + ), Route( "/api/threads/{thread_id}/final-answer", get_final_answer, diff --git a/tests/test_system_update_http.py b/tests/test_system_update_http.py new file mode 100644 index 0000000..2924f19 --- /dev/null +++ b/tests/test_system_update_http.py @@ -0,0 +1,196 @@ +"""HTTP-layer tests for the system update routes.""" + +import json +from pathlib import Path + +import pytest +from starlette.applications import Starlette +from starlette.routing import Route +from starlette.testclient import TestClient + +from EvoScientist import update_check +from EvoScientist.langgraph_dev import http as http_mod + + +async def _allow(request): + return None + + +@pytest.fixture +def client(monkeypatch, tmp_path): + monkeypatch.setattr(http_mod, "_authenticate_system", _allow, raising=False) + monkeypatch.setattr(http_mod, "_authenticate_system_write", _allow, raising=False) + monkeypatch.setattr(http_mod, "_schedule_self_exit", lambda delay=0.5: None) + monkeypatch.setenv("EVOSCIENTIST_UPDATE_STAGING_DIR", str(tmp_path)) + app = Starlette( + routes=[ + Route("/internal/system/update", http_mod.post_system_update, methods=["POST"]), + Route( + "/internal/system/update/status", + http_mod.get_system_update_status, + methods=["GET"], + ), + Route( + "/internal/system/rollback-versions", + http_mod.get_system_rollback_versions, + methods=["GET"], + ), + Route( + "/internal/system/update/rollback", + http_mod.post_system_update_rollback, + methods=["POST"], + ), + ] + ) + return TestClient(app) + + +def _info(has_update, breaking_db=False): + return { + "current_version": "0.2.8", + "latest_version": "0.3.0", + "has_update": has_update, + "release_url": None, + "release_notes": None, + "published_at": None, + "cached": False, + "warning": None, + "breaking_db": breaking_db, + } + + +def _probe_docker(): + from EvoScientist.updater import SystemProbe + + return SystemProbe(in_container=True, invocation_id=None, uv_tool_pkg=False, pipx_pkg=False) + + +def _probe_uvtool(): + from EvoScientist.updater import SystemProbe + + return SystemProbe(in_container=False, invocation_id=None, uv_tool_pkg=True, pipx_pkg=False) + + +def test_update_409_when_no_update(client, monkeypatch): + monkeypatch.setattr(update_check, "get_update_info", lambda force=False: _info(has_update=False)) + res = client.post("/internal/system/update") + assert res.status_code == 409 + assert res.json()["code"] == "ALREADY_UP_TO_DATE" + + +def test_update_docker_returns_guidance(client, monkeypatch): + monkeypatch.setattr(update_check, "get_update_info", lambda force=False: _info(has_update=True)) + monkeypatch.setattr(http_mod, "probe_system", lambda environ=None: _probe_docker()) + res = client.post("/internal/system/update") + assert res.status_code == 200 + assert res.json()["status"] == "manual" + assert "docker compose" in res.json()["guidance"] + + +def test_update_202_spawns_updater(client, monkeypatch, tmp_path): + spawned = {} + monkeypatch.setattr(update_check, "get_update_info", lambda force=False: _info(has_update=True)) + monkeypatch.setattr(http_mod, "probe_system", lambda environ=None: _probe_uvtool()) + monkeypatch.setattr( + update_check, + "download_update", + lambda version=None: { + "version": "0.3.0", + "file": "EvoScientist-0.3.0-py3-none-any.whl", + "path": str(tmp_path / "v0.3.0" / "EvoScientist-0.3.0-py3-none-any.whl"), + "suggested_command": "...", + }, + ) + (tmp_path / "v0.3.0").mkdir(parents=True) + (tmp_path / "v0.3.0" / "EvoScientist-0.3.0-py3-none-any.whl").write_bytes(b"x") + monkeypatch.setattr( + http_mod, + "spawn_updater", + lambda u, p, pid: spawned.setdefault("ok", (str(u), str(p), pid)), + ) + res = client.post("/internal/system/update") + assert res.status_code == 202 + body = res.json() + assert body["status"] == "applying" + assert body["need_restart"] is True + assert spawned["ok"][2] > 0 # parent pid + plan = json.loads((tmp_path / "v0.3.0" / "plan.json").read_text()) + assert plan["version"] == "0.3.0" + assert plan["respawn_command"] + + +def test_update_409_when_update_in_progress(client, monkeypatch, tmp_path): + monkeypatch.setattr(update_check, "get_update_info", lambda force=False: _info(has_update=True)) + monkeypatch.setattr(http_mod, "probe_system", lambda environ=None: _probe_uvtool()) + monkeypatch.setattr( + update_check, + "download_update", + lambda version=None: { + "version": "0.3.0", + "file": "f.whl", + "path": str(tmp_path / "v0.3.0" / "f.whl"), + "suggested_command": "...", + }, + ) + (tmp_path / "v0.3.0").mkdir(parents=True) + (tmp_path / "v0.3.0" / "f.whl").write_bytes(b"x") + (tmp_path / "v0.3.0" / "plan.lock").write_text("") # lock held + monkeypatch.setattr(http_mod, "spawn_updater", lambda u, p, pid: None) + res = client.post("/internal/system/update") + assert res.status_code == 409 + assert res.json()["code"] == "UPDATE_IN_PROGRESS" + + +def test_update_breaking_requires_confirm(client, monkeypatch): + monkeypatch.setattr( + update_check, + "get_update_info", + lambda force=False: _info(has_update=True, breaking_db=True), + ) + res = client.post("/internal/system/update") + assert res.status_code == 409 + assert res.json()["code"] == "BREAKING_DB_CONFIRM_REQUIRED" + + +def test_status_none_then_result(client, tmp_path): + assert client.get("/internal/system/update/status").json()["status"] == "none" + (tmp_path / "last-result.json").write_text(json.dumps({"status": "success", "version": "0.3.0"})) + assert client.get("/internal/system/update/status").json()["status"] == "success" + + +def test_rollback_versions_passthrough(client, monkeypatch): + monkeypatch.setattr( + update_check, + "list_rollback_versions", + lambda limit=3: [{"version": "0.2.9", "published_at": "p", "release_url": "u"}], + ) + res = client.get("/internal/system/rollback-versions") + assert res.json()["versions"][0]["version"] == "0.2.9" + + +def test_rollback_rejects_disallowed_version(client, monkeypatch): + monkeypatch.setattr(update_check, "is_allowed_rollback", lambda v: False) + res = client.post("/internal/system/update/rollback", json={"version": "0.0.1"}) + assert res.status_code == 400 + + +def test_rollback_202_uses_same_pipeline(client, monkeypatch, tmp_path): + monkeypatch.setattr(update_check, "is_allowed_rollback", lambda v: True) + monkeypatch.setattr(http_mod, "probe_system", lambda environ=None: _probe_uvtool()) + monkeypatch.setattr( + update_check, + "download_update", + lambda version=None: { + "version": version, + "file": "f.whl", + "path": str(tmp_path / f"v{version}" / "f.whl"), + "suggested_command": "...", + }, + ) + (tmp_path / "v0.2.9").mkdir(parents=True) + (tmp_path / "v0.2.9" / "f.whl").write_bytes(b"x") + monkeypatch.setattr(http_mod, "spawn_updater", lambda u, p, pid: None) + res = client.post("/internal/system/update/rollback", json={"version": "0.2.9"}) + assert res.status_code == 202 + plan = json.loads((tmp_path / "v0.2.9" / "plan.json").read_text()) + assert plan["version"] == "0.2.9"