feat(update): system update/status/rollback HTTP endpoints with system:write scope

This commit is contained in:
m4
2026-08-12 16:42:10 +08:00
parent 91e2a87be5
commit 8f6a568646
2 changed files with 468 additions and 1 deletions
+272 -1
View File
@@ -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,
+196
View File
@@ -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"