diff --git a/flocks/console/login.py b/flocks/console/login.py index ad79843c7..0fb2d179a 100644 --- a/flocks/console/login.py +++ b/flocks/console/login.py @@ -23,6 +23,80 @@ def _shared_console_session_path() -> Path: return Path(raw).expanduser() / "run" / "console-session.json" +def _flocks_root() -> Path: + return Path(os.getenv("FLOCKS_ROOT", str(Path.home() / ".flocks"))).expanduser() + + +def _read_json_file(path: Path) -> dict[str, Any]: + try: + payload = json.loads(path.read_text(encoding="utf-8")) + except Exception: + return {} + return payload if isinstance(payload, dict) else {} + + +def _read_pro_bundle_marker() -> dict[str, Any]: + return _read_json_file(_flocks_root() / "run" / "pro-bundle-installed.json") + + +def _local_pro_license_path() -> Path: + return _flocks_root() / "flockspro" / "license.json" + + +def _read_local_pro_license_state() -> dict[str, Any]: + return _read_json_file(_local_pro_license_path()) + + +def _read_local_pro_license_id() -> str: + state = _read_local_pro_license_state() + payload = state.get("payload") if isinstance(state.get("payload"), dict) else {} + return str(state.get("license_id") or payload.get("license_id") or "").strip() + + +def _read_local_pro_license_status() -> str: + state = _read_local_pro_license_state() + payload = state.get("payload") if isinstance(state.get("payload"), dict) else {} + return str( + state.get("license_status") + or state.get("status") + or payload.get("license_status") + or payload.get("status") + or "" + ).strip() + + +def _pending_pro_bundle_install_receipt_path() -> Path: + return _flocks_root() / "run" / "pro-bundle-install-receipt-pending.json" + + +def _sync_local_pro_license_from_heartbeat_response(data: dict[str, Any]) -> None: + license_path = _local_pro_license_path() + state = _read_json_file(license_path) + now_ts = int(datetime.now(UTC).timestamp()) + if state: + changed = False + patch_token = data.get("license_patch") or data.get("latest_patch") + if isinstance(patch_token, str) and patch_token: + patches = state.get("patches") if isinstance(state.get("patches"), list) else [] + if patch_token not in patches: + state["patches"] = [*patches, patch_token] + changed = True + if state.get("last_sync_at") != now_ts: + state["last_sync_at"] = now_ts + state["last_heartbeat_ok_at"] = now_ts + changed = True + if changed: + license_path.parent.mkdir(parents=True, exist_ok=True) + license_path.write_text(json.dumps(state, ensure_ascii=False, indent=2), encoding="utf-8") + + revoked_license_ids = data.get("revoked_license_ids") + if isinstance(revoked_license_ids, list): + revocation_path = _flocks_root() / "flockspro" / "revocation.json" + revocation_path.parent.mkdir(parents=True, exist_ok=True) + payload = {"revoked_license_ids": sorted({str(item) for item in revoked_license_ids})} + revocation_path.write_text(json.dumps(payload, ensure_ascii=False, indent=2), encoding="utf-8") + + def _write_shared_console_session(session: dict[str, Any]) -> None: path = _shared_console_session_path() path.parent.mkdir(parents=True, exist_ok=True) @@ -42,6 +116,29 @@ def _write_shared_console_session(session: dict[str, Any]) -> None: pass +def read_shared_console_session() -> dict[str, Any] | None: + path = _shared_console_session_path() + try: + payload = json.loads(path.read_text(encoding="utf-8")) + except (FileNotFoundError, OSError, json.JSONDecodeError): + return None + if not isinstance(payload, dict): + return None + token = str(payload.get("console_session_token") or "").strip() + fingerprint = str(payload.get("fingerprint") or "").strip() + install_id = str(payload.get("install_id") or "").strip() + if not token or not fingerprint or not install_id: + return None + expires_at = str(payload.get("expires_at") or "").strip() + if expires_at: + try: + if _parse_iso(expires_at) <= datetime.now(UTC): + return None + except ValueError: + return None + return payload + + def _delete_shared_console_session() -> None: path = _shared_console_session_path() try: @@ -284,39 +381,162 @@ def _runtime_version() -> str: return str(__version__).lstrip("v") @classmethod - async def send_heartbeat(cls) -> dict[str, Any]: - session = await cls._require_session() - console_base = cls.console_base_url() + def runtime_version_payload(cls, *, pro_component_version: str | None = None) -> dict[str, str]: + marker = _read_pro_bundle_marker() + core_version = str( + marker.get("core_version") + or cls._runtime_version() + ).strip() + bundle_version = str( + marker.get("bundle_version") + or "" + ).strip() + pro_component_version = str(marker.get("flockspro_component_version") or pro_component_version or "").strip() + has_pro_bundle = bool(bundle_version or pro_component_version) + edition = "flockspro" if has_pro_bundle else "oss" payload = { - "fingerprint": session["fingerprint"], - "install_id": session["install_id"], + "edition": edition, + } + if core_version: + payload["core_version"] = core_version + if edition == "flockspro" and bundle_version: + payload["bundle_version"] = bundle_version + if edition == "flockspro" and pro_component_version: + payload["flockspro_component_version"] = pro_component_version + return {key: value for key, value in payload.items() if value} + + @classmethod + def heartbeat_payload( + cls, + session: dict[str, Any], + *, + status: str = "ok", + license_id: str | None = None, + pro_component_version: str | None = None, + ) -> dict[str, Any]: + version_payload = cls.runtime_version_payload(pro_component_version=pro_component_version) + return { + "fingerprint": session.get("fingerprint"), + "install_id": session.get("install_id"), "console_login_id": session.get("console_login_id"), "sent_at": _now_iso(), - "status": "ok", + "status": status, + "license_id": license_id or None, + **version_payload, } - if not console_base: + + @classmethod + async def send_heartbeat_for_session( + cls, + *, + session: dict[str, Any], + status: str = "ok", + license_id: str | None = None, + heartbeat_url: str | None = None, + report_install_receipt: bool = False, + pro_component_version: str | None = None, + ) -> dict[str, Any]: + console_base = cls.console_base_url() + payload = cls.heartbeat_payload( + session, + status=status, + license_id=license_id, + pro_component_version=pro_component_version, + ) + target_url = heartbeat_url or (f"{console_base}/v1/heartbeats" if console_base else "") + if not target_url: return {"ok": True, "mode": "mock", "node": payload} + token = str(session.get("console_session_token") or "").strip() + if not token: + raise ValueError("console_session_token 缺失,无法发送心跳") async with httpx.AsyncClient(timeout=10) as client: resp = await client.post( - f"{console_base}/v1/heartbeats", + target_url, json=payload, - headers={"Authorization": f"Bearer {session['console_session_token']}"}, + headers={"Authorization": f"Bearer {token}"}, ) if resp.status_code in {401, 403}: raise ValueError("console 会话已失效,请重新登录") resp.raise_for_status() - return resp.json() + data = resp.json() + _sync_local_pro_license_from_heartbeat_response(data) + if report_install_receipt: + await cls._report_pending_pro_bundle_install_receipt(client=client, session=session) + return data + + @classmethod + async def send_heartbeat(cls) -> dict[str, Any]: + session = await cls._require_session() + return await cls.send_heartbeat_for_session( + session=session, + status=_read_local_pro_license_status() or "ok", + license_id=_read_local_pro_license_id() or None, + report_install_receipt=True, + ) + + @classmethod + async def report_pending_pro_bundle_install_receipt(cls) -> bool: + try: + session = await cls._require_session() + except Exception: + session = read_shared_console_session() + if not session: + return False + async with httpx.AsyncClient(timeout=10) as client: + return await cls._report_pending_pro_bundle_install_receipt(client=client, session=session) + + @classmethod + def _console_base_url_for_session(cls, session: dict[str, Any]) -> str: + console_base = cls.console_base_url() or str(session.get("console_base_url") or "").strip().rstrip("/") + if console_base: + return console_base + shared_session = read_shared_console_session() + return str((shared_session or {}).get("console_base_url") or "").strip().rstrip("/") + + @classmethod + async def _report_pending_pro_bundle_install_receipt( + cls, + *, + client: httpx.AsyncClient, + session: dict[str, Any], + ) -> bool: + console_base = cls._console_base_url_for_session(session) + token = str(session.get("console_session_token") or "").strip() + if not console_base or not token: + return False + path = _pending_pro_bundle_install_receipt_path() + payload = _read_json_file(path) + if not payload: + return False + payload = { + **payload, + "fingerprint": session.get("fingerprint"), + "install_id": session.get("install_id"), + "license_id": payload.get("license_id") or _read_local_pro_license_id() or None, + } + try: + resp = await client.post( + f"{console_base}/v1/pro-bundles/installations", + json=payload, + headers={"Authorization": f"Bearer {token}"}, + ) + if resp.status_code in {200, 201, 202}: + path.unlink(missing_ok=True) + return True + except Exception: + return False + return False @classmethod async def sync_node_profile(cls, *, force: bool = False, source: str = "scheduled") -> dict[str, Any]: _ = force session = await cls._require_session() console_base = cls.console_base_url() + version_payload = cls.runtime_version_payload() payload = { "fingerprint": session["fingerprint"], "install_id": session["install_id"], - "edition": cls._edition(), - "version": cls._runtime_version(), + **version_payload, "source": source, "sent_at": _now_iso(), } diff --git a/flocks/server/routes/console_upgrade.py b/flocks/server/routes/console_upgrade.py index 9539bf71a..6fb12cd26 100644 --- a/flocks/server/routes/console_upgrade.py +++ b/flocks/server/routes/console_upgrade.py @@ -317,8 +317,8 @@ def _enrich_record_from_install_marker(record: dict[str, Any]) -> dict[str, Any] marker = _read_pro_bundle_install_marker() if marker: details.setdefault("auto_install_release_id", marker.get("release_id") or marker.get("bundle_release_id")) - details.setdefault("auto_install_version", marker.get("installed_version")) - details.setdefault("auto_install_pro_version", marker.get("flockspro_component_version")) + details.setdefault("auto_install_bundle_version", marker.get("bundle_version")) + details.setdefault("auto_install_pro_component_version", marker.get("flockspro_component_version")) details.setdefault("flockspro_component_version", marker.get("flockspro_component_version")) details.setdefault("auto_install_build_id", marker.get("build_id")) @@ -442,19 +442,16 @@ def _record_target_bundle(record: dict[str, Any]) -> dict[str, str]: "release_id": release_id, "bundle_release_id": _clean_bundle_value(details.get("bundle_release_id") or release_id), "build_id": _clean_bundle_value(details.get("target_build_id") or latest_bundle.get("build_id")), - "display_version": _clean_bundle_value( - details.get("target_display_version") - or details.get("auto_install_target") - or latest_bundle.get("display_version") + "bundle_version_update_to": _clean_bundle_value( + details.get("bundle_version_update_to") + or latest_bundle.get("bundle_version") ), - "core_version": _clean_bundle_value( - details.get("target_core_version") - or details.get("target_oss_version") + "core_version_update_to": _clean_bundle_value( + details.get("core_version_update_to") or latest_bundle.get("core_version") - or latest_bundle.get("oss_version") ), - "flockspro_component_version": _clean_bundle_value( - details.get("target_flockspro_component_version") + "flockspro_component_version_update_to": _clean_bundle_value( + details.get("flockspro_component_version_update_to") or latest_bundle.get("flockspro_component_version") ), } @@ -467,18 +464,18 @@ def _target_bundle_fingerprint_matches(target: dict[str, str], marker: dict[str, if build_id and marker_build_id: return marker_build_id == build_id - pro_version = target.get("flockspro_component_version") + pro_version = target.get("flockspro_component_version_update_to") marker_pro_version = _clean_bundle_value(marker.get("flockspro_component_version")) if pro_version and marker_pro_version: return marker_pro_version == pro_version - display_version = target.get("display_version") - marker_display_version = _clean_bundle_value(marker.get("installed_version") or marker.get("display_version")) - if display_version and marker_display_version: - return _clean_version_value(marker_display_version) == _clean_version_value(display_version) + bundle_version = target.get("bundle_version_update_to") + marker_bundle_version = _clean_bundle_value(marker.get("bundle_version")) + if bundle_version and marker_bundle_version: + return _clean_version_value(marker_bundle_version) == _clean_version_value(bundle_version) - core_version = target.get("core_version") or target.get("oss_version") - marker_core_version = _clean_bundle_value(marker.get("core_version") or marker.get("oss_version")) + core_version = target.get("core_version_update_to") + marker_core_version = _clean_bundle_value(marker.get("core_version")) if core_version and marker_core_version: return _clean_version_value(marker_core_version) == _clean_version_value(core_version) @@ -510,7 +507,7 @@ async def _run_auto_upgrade_install(record: dict[str, Any]) -> dict[str, Any]: marker = _read_pro_bundle_install_marker() if _is_pro_component_installed() and _marker_matches_target_bundle(marker, record): details["auto_install_release_id"] = marker.get("release_id") or marker.get("bundle_release_id") - details["auto_install_version"] = marker.get("installed_version") + details["auto_install_bundle_version"] = marker.get("bundle_version") await _maybe_activate_pro_license(record, allow_fallback=False) await _maybe_refresh_pro_license(record) capability = _record_pro_capability(details) @@ -538,8 +535,8 @@ async def _run_auto_upgrade_install(record: dict[str, Any]) -> dict[str, Any]: "done" if final_stage == "done" and capability.get("pro_enabled") else "license_inactive" ) details["auto_install_release_id"] = marker.get("release_id") or marker.get("bundle_release_id") - details["auto_install_version"] = marker.get("installed_version") - details["auto_install_pro_version"] = marker.get("flockspro_component_version") + details["auto_install_bundle_version"] = marker.get("bundle_version") + details["auto_install_pro_component_version"] = marker.get("flockspro_component_version") details["auto_install_completed_at"] = datetime.now(UTC).isoformat() details["auto_install_message"] = final_message _enrich_record_from_install_marker(record) @@ -565,7 +562,7 @@ def _marker_indicates_pro_bundle_installed(marker: dict[str, Any]) -> bool: return False return any( str(marker.get(key) or "").strip() - for key in ("installed_at", "installed_version", "bundle_version", "flockspro_component_version", "build_id") + for key in ("installed_at", "bundle_version", "flockspro_component_version", "build_id") ) @@ -600,15 +597,12 @@ async def _report_pro_bundle_installation( "license_id": _record_license_id(record), "fingerprint": console_session.get("fingerprint"), "install_id": console_session.get("install_id"), - "installed_version": source.get("installed_version") - or source.get("display_version") - or target.get("display_version") - or details.get("auto_install_target") - or details.get("auto_install_version") - or "", - "core_version": source.get("core_version") or source.get("oss_version") or target.get("core_version"), - "oss_version": source.get("core_version") or source.get("oss_version") or target.get("core_version"), - "flockspro_component_version": source.get("flockspro_component_version") or target.get("flockspro_component_version"), + "bundle_version": source.get("bundle_version") or target.get("bundle_version_update_to") or "", + "core_version": source.get("core_version") or target.get("core_version_update_to"), + "flockspro_component_version": ( + source.get("flockspro_component_version") + or target.get("flockspro_component_version_update_to") + ), "build_id": source.get("build_id") or target.get("build_id"), "install_result": install_result, "error_message": error_message, @@ -688,8 +682,8 @@ async def _finalize_restarting_upgrade_if_installed(record: dict[str, Any]) -> d await _maybe_refresh_pro_license(record) capability = _record_pro_capability(details) details["auto_install_result"] = "done" if capability.get("pro_enabled") else "license_inactive" - details["auto_install_version"] = marker.get("installed_version") or marker.get("display_version") - details["auto_install_pro_version"] = marker.get("flockspro_component_version") + details["auto_install_bundle_version"] = marker.get("bundle_version") + details["auto_install_pro_component_version"] = marker.get("flockspro_component_version") details["auto_install_completed_at"] = datetime.now(UTC).isoformat() details["auto_install_message"] = "Upgrade completed after service restart" _enrich_record_from_install_marker(record) @@ -877,7 +871,8 @@ async def get_pro_package_status(request: Request) -> dict[str, Any]: "installed": installed, "runtime_importable": runtime_importable, "install_marker_present": install_marker_present, - "installed_version": marker.get("installed_version"), + "bundle_version": marker.get("bundle_version"), + "core_version": marker.get("core_version"), "flockspro_component_version": marker.get("flockspro_component_version"), "build_id": marker.get("build_id"), "installed_at": marker.get("installed_at"), @@ -996,8 +991,8 @@ async def _stream(): details["auto_install_result"] = "done" else: details["auto_install_result"] = "license_inactive" - details["auto_install_version"] = marker.get("installed_version") - details["auto_install_pro_version"] = marker.get("flockspro_component_version") + details["auto_install_bundle_version"] = marker.get("bundle_version") + details["auto_install_pro_component_version"] = marker.get("flockspro_component_version") details["auto_install_completed_at"] = datetime.now(UTC).isoformat() details["auto_install_message"] = progress.message _enrich_record_from_install_marker(raw) diff --git a/flocks/server/routes/flockspro_license.py b/flocks/server/routes/flockspro_license.py index 10b132273..b7f041b89 100644 --- a/flocks/server/routes/flockspro_license.py +++ b/flocks/server/routes/flockspro_license.py @@ -11,6 +11,7 @@ from fastapi import APIRouter, Request +from flocks.console.login import ConsoleLoginService from flocks.server.auth import require_user from flocks.server.routes.console_upgrade import _get_pro_capability_status, _is_pro_component_installed @@ -50,6 +51,8 @@ async def refresh_flockspro_license_status(request: Request) -> dict[str, Any]: return _inactive_status("flockspro_not_installed") try: + await ConsoleLoginService.send_heartbeat() + from flockspro.license.runtime import get_license_checker # type: ignore[import-not-found] checker = get_license_checker() diff --git a/flocks/updater/restart_handoff.py b/flocks/updater/restart_handoff.py index f6350b400..dcfd1505c 100644 --- a/flocks/updater/restart_handoff.py +++ b/flocks/updater/restart_handoff.py @@ -161,6 +161,22 @@ def _run_upgrade_tasks(args: argparse.Namespace) -> str | None: ) +def _report_pending_pro_bundle_install_receipt(args: argparse.Namespace) -> None: + if not args.pro_bundle_manifest_path: + return + try: + from flocks.console.login import ConsoleLoginService + + reported = asyncio.run(ConsoleLoginService.report_pending_pro_bundle_install_receipt()) + except Exception as exc: + _record_handoff_log(f"install_receipt_report_failed error={exc}") + return + if reported: + _record_handoff_log("install_receipt_reported") + else: + _record_handoff_log("install_receipt_report_skipped") + + def _rollback_failed_upgrade(args: argparse.Namespace, error: str) -> None: from flocks.updater import updater @@ -214,6 +230,7 @@ def run(argv: Sequence[str] | None = None) -> int: _rollback_failed_upgrade(args, task_error) _cleanup_dir(args.cleanup_dir) return 1 + _report_pending_pro_bundle_install_receipt(args) try: process = subprocess.Popen( diff --git a/flocks/updater/updater.py b/flocks/updater/updater.py index d076af229..7337d6dbe 100644 --- a/flocks/updater/updater.py +++ b/flocks/updater/updater.py @@ -31,7 +31,7 @@ from datetime import datetime, timezone from pathlib import Path, PureWindowsPath from typing import Any, AsyncGenerator, Awaitable, Callable -from urllib.parse import quote, urlparse +from urllib.parse import quote, urlencode, urlparse import httpx @@ -1082,31 +1082,15 @@ def _archive_format_for_url(url: str, manifest_format: str | None = None) -> str return "tar.gz" -def _console_manifest_display_version(data: dict[str, Any]) -> str: - display_version = str(data.get("display_version") or data.get("version") or data.get("latest_version") or "").strip() - if display_version: - return display_version - core_version = str(data.get("core_version") or "").strip() - if core_version: - return core_version - oss_version = str(data.get("oss_version") or "").strip() - if oss_version: - return oss_version - compare_version = str(data.get("compare_version") or "").strip() - return f"v{compare_version}" if compare_version and not compare_version.startswith(("v", "V")) else compare_version +def _console_manifest_bundle_version(data: dict[str, Any]) -> str: + return str(data.get("bundle_version") or "").strip() def _pro_bundle_core_version(data: dict[str, Any]) -> str: - core_version = str(data.get("core_version") or "").strip() - if core_version: - return core_version - oss_version = str(data.get("oss_version") or "").strip() - if oss_version: - return oss_version - return _console_manifest_display_version(data) + return str(data.get("core_version") or "").strip() -def _pro_bundle_oss_version(data: dict[str, Any]) -> str: +def _pro_bundle_core_version_for_compare(data: dict[str, Any]) -> str: return _pro_bundle_core_version(data) @@ -1117,22 +1101,21 @@ def _version_label(version: str | None) -> str: return normalized if normalized.startswith(("v", "V")) else f"v{normalized}" -def _is_pro_bundle_oss_older_than_local(manifest: dict[str, Any], current_version: str | None = None) -> bool: - bundle_oss_version = _pro_bundle_oss_version(manifest) - if not bundle_oss_version: +def _is_pro_bundle_core_older_than_local(manifest: dict[str, Any], current_version: str | None = None) -> bool: + bundle_core_version = _pro_bundle_core_version_for_compare(manifest) + if not bundle_core_version: return False local_version = str(current_version or get_current_version() or "").strip() if not local_version: return False - return _parse_version(bundle_oss_version) < _parse_version(local_version) + return _parse_version(bundle_core_version) < _parse_version(local_version) -def _effective_pro_bundle_manifest(manifest: dict[str, Any], effective_oss_version: str) -> dict[str, Any]: +def _effective_pro_bundle_manifest(manifest: dict[str, Any], effective_core_version: str) -> dict[str, Any]: payload = dict(manifest) - effective_label = _version_label(effective_oss_version) + effective_label = _version_label(effective_core_version) if effective_label: payload["core_version"] = effective_label - payload["oss_version"] = effective_label return payload @@ -1222,6 +1205,17 @@ async def _fetch_gitlab_release( ) +def _read_local_pro_license_id() -> str: + license_path = _flocks_root() / "flockspro" / "license.json" + try: + payload = json.loads(license_path.read_text(encoding="utf-8")) + except Exception: + return "" + if not isinstance(payload, dict): + return "" + return str(payload.get("license_id") or "").strip() + + async def _load_console_session_token() -> str | None: def _token_from_payload(payload: Any) -> str | None: if not isinstance(payload, dict): @@ -1273,8 +1267,14 @@ async def _fetch_console_manifest_release_info(console_session_token: str | None raise ValueError("FLOCKS_CONSOLE_BASE_URL 未配置,无法使用 console-manifest 源") channel = (os.getenv("FLOCKS_UPDATE_CHANNEL") or "flockspro").strip() or "flockspro" - url = f"{manifest_base}/v1/manifest/latest?channel={channel}" + license_id = (os.getenv("FLOCKSPRO_LICENSE_ID") or _read_local_pro_license_id()).strip() + query = {"channel": channel} + if license_id: + query["license_id"] = license_id + url = f"{manifest_base}/v1/manifest/latest?{urlencode(query)}" headers: dict[str, str] = {} + if license_id: + headers["x-license-id"] = license_id token = str(console_session_token or "").strip() or await _load_console_session_token() if token: headers["Authorization"] = f"Bearer {token}" @@ -1297,9 +1297,9 @@ async def _fetch_console_manifest_release_info(console_session_token: str | None if datetime.now(timezone.utc) < frozen_until: raise ValueError("console manifest channel frozen_until not reached") - latest = _console_manifest_display_version(data) + latest = _console_manifest_bundle_version(data) if not latest: - raise ValueError("manifest 响应缺少 compare_version/display_version") + raise ValueError("manifest 响应缺少 bundle_version") bundle_url = str( data.get("bundle_url") or data.get("url") @@ -1848,13 +1848,12 @@ def _merge_console_manifest_release_identity( merged["release_id"] = release_id if bundle_release_id and not merged.get("bundle_release_id"): merged["bundle_release_id"] = bundle_release_id - for key in ("display_version", "version", "latest_version", "compare_version", "flockspro_component_version", "build_id"): + for key in ("bundle_version", "compare_version", "flockspro_component_version", "build_id"): if console_manifest.get(key): merged[key] = console_manifest.get(key) - core_version = console_manifest.get("core_version") or console_manifest.get("oss_version") or merged.get("core_version") or merged.get("oss_version") + core_version = console_manifest.get("core_version") or merged.get("core_version") if core_version: merged["core_version"] = core_version - merged["oss_version"] = core_version return merged @@ -1863,22 +1862,47 @@ def _write_pro_bundle_install_marker(manifest: dict[str, Any], *, bundle_sha256: marker.parent.mkdir(parents=True, exist_ok=True) release_id = manifest.get("release_id") or manifest.get("bundle_release_id") bundle_release_id = manifest.get("bundle_release_id") or manifest.get("release_id") - display_version = _console_manifest_display_version(manifest) - core_version = manifest.get("core_version") or manifest.get("oss_version") + bundle_version = _console_manifest_bundle_version(manifest) + core_version = manifest.get("core_version") payload = { "release_id": release_id, "bundle_release_id": bundle_release_id, - "bundle_version": display_version, - "display_version": display_version, - "installed_version": display_version, + "bundle_version": bundle_version, "core_version": core_version, - "oss_version": core_version, "flockspro_component_version": manifest.get("flockspro_component_version"), "build_id": manifest.get("build_id"), "bundle_sha256": bundle_sha256 or manifest.get("bundle_sha256"), "installed_at": datetime.now(timezone.utc).isoformat(), } marker.write_text(json.dumps(payload, ensure_ascii=True, sort_keys=True), encoding="utf-8") + _write_pending_pro_bundle_install_receipt(payload) + + +def _write_pending_pro_bundle_install_receipt(marker_payload: dict[str, Any]) -> None: + receipt_path = _flocks_root() / "run" / "pro-bundle-install-receipt-pending.json" + receipt_path.parent.mkdir(parents=True, exist_ok=True) + bundle_version = str( + marker_payload.get("bundle_version") + or "" + ).strip() + core_version = str(marker_payload.get("core_version") or "").strip() + pro_component_version = str(marker_payload.get("flockspro_component_version") or "").strip() + receipt = { + "release_id": marker_payload.get("release_id") or marker_payload.get("bundle_release_id"), + "bundle_release_id": marker_payload.get("bundle_release_id") or marker_payload.get("release_id"), + "license_id": _read_local_pro_license_id() or None, + "bundle_version": bundle_version, + "core_version": core_version, + "flockspro_component_version": pro_component_version, + "build_id": marker_payload.get("build_id"), + "install_result": "success", + "reported_at": datetime.now(timezone.utc).isoformat(), + } + receipt_path.write_text(json.dumps(receipt, ensure_ascii=True, sort_keys=True), encoding="utf-8") + try: + os.chmod(receipt_path, 0o600) + except OSError: + pass class _NullConsole: @@ -2533,20 +2557,12 @@ def _read_pro_bundle_install_marker() -> dict[str, Any]: def _read_pro_bundle_installed_bundle_version() -> str: payload = _read_pro_bundle_install_marker() - for key in ("bundle_version", "installed_version", "display_version"): - version = str(payload.get(key) or "").strip() - if version: - return version - return "" + return str(payload.get("bundle_version") or "").strip() def _read_pro_bundle_installed_core_version() -> str: payload = _read_pro_bundle_install_marker() - for key in ("core_version", "oss_version"): - version = str(payload.get(key) or "").strip() - if version: - return version - return "" + return str(payload.get("core_version") or "").strip() def _read_pro_bundle_installed_component_version() -> str: @@ -3069,13 +3085,13 @@ async def _queue_download_progress(progress: UpdateProgress) -> None: yield UpdateProgress(stage="error", message=msg, success=False) return if profile.sources == ["console-manifest"]: - skip_core_replace = _is_pro_bundle_oss_older_than_local(pro_bundle_manifest, current_version) + skip_core_replace = _is_pro_bundle_core_older_than_local(pro_bundle_manifest, current_version) if skip_core_replace: - bundle_oss_version = _pro_bundle_oss_version(pro_bundle_manifest) + bundle_core_version = _pro_bundle_core_version_for_compare(pro_bundle_manifest) pro_bundle_manifest = _effective_pro_bundle_manifest(pro_bundle_manifest, current_version) log.info( "updater.pro_bundle.keep_local_core", - {"local_version": current_version, "bundle_oss_version": bundle_oss_version}, + {"local_version": current_version, "bundle_core_version": bundle_core_version}, ) else: effective_update_version = latest_tag diff --git a/tests/console/test_console_login_heartbeat.py b/tests/console/test_console_login_heartbeat.py new file mode 100644 index 000000000..5103b58af --- /dev/null +++ b/tests/console/test_console_login_heartbeat.py @@ -0,0 +1,224 @@ +import asyncio +import json + +from flocks.console import login as login_mod +from flocks.console.login import ConsoleLoginService + + +def test_heartbeat_payload_reports_oss_for_core_only_install(tmp_path, monkeypatch): + monkeypatch.setenv("FLOCKS_ROOT", str(tmp_path)) + monkeypatch.setenv("FLOCKS_EDITION", "flockspro") + monkeypatch.setattr(ConsoleLoginService, "_runtime_version", staticmethod(lambda: "2026.7.3.3")) + + payload = ConsoleLoginService.heartbeat_payload( + { + "console_session_token": "cs_heartbeat", + "fingerprint": "fp_heartbeat", + "install_id": "inst_heartbeat", + }, + ) + + assert payload["edition"] == "oss" + assert payload["core_version"] == "2026.7.3.3" + assert "version" not in payload + assert "bundle_version" not in payload + assert "flockspro_component_version" not in payload + assert "version_info" not in payload + + +def test_heartbeat_payload_includes_pro_runtime_versions(tmp_path, monkeypatch): + monkeypatch.setenv("FLOCKS_ROOT", str(tmp_path)) + monkeypatch.setattr(ConsoleLoginService, "_runtime_version", staticmethod(lambda: "2026.7.3")) + + marker = tmp_path / "run" / "pro-bundle-installed.json" + marker.parent.mkdir(parents=True) + marker.write_text( + json.dumps( + { + "bundle_version": "v2026.7.3", + "core_version": "v2026.7.3", + } + ), + encoding="utf-8", + ) + + payload = ConsoleLoginService.heartbeat_payload( + { + "console_session_token": "cs_heartbeat", + "fingerprint": "fp_heartbeat", + "install_id": "inst_heartbeat", + }, + status="poc", + license_id="lic_heartbeat", + pro_component_version="2026.7.3.1", + ) + + assert payload["fingerprint"] == "fp_heartbeat" + assert payload["install_id"] == "inst_heartbeat" + assert payload["status"] == "poc" + assert payload["license_id"] == "lic_heartbeat" + assert payload["edition"] == "flockspro" + assert payload["bundle_version"] == "v2026.7.3" + assert payload["core_version"] == "v2026.7.3" + assert payload["flockspro_component_version"] == "2026.7.3.1" + assert "version" not in payload + assert "version_info" not in payload + + +def test_send_heartbeat_uses_local_pro_license_and_applies_response(tmp_path, monkeypatch): + monkeypatch.setenv("FLOCKS_ROOT", str(tmp_path)) + monkeypatch.setenv("FLOCKS_CONSOLE_BASE_URL", "https://console.example.com") + monkeypatch.setattr(ConsoleLoginService, "_runtime_version", staticmethod(lambda: "2026.7.3")) + + license_path = tmp_path / "flockspro" / "license.json" + license_path.parent.mkdir(parents=True) + license_path.write_text( + json.dumps( + { + "license_id": "lic_core", + "payload": {"license_id": "lic_core", "status": "poc"}, + "patches": [], + } + ), + encoding="utf-8", + ) + marker = tmp_path / "run" / "pro-bundle-installed.json" + marker.parent.mkdir(parents=True) + marker.write_text( + json.dumps( + { + "bundle_version": "v2026.7.3", + "core_version": "v2026.7.3", + "flockspro_component_version": "2026.7.3.1", + } + ), + encoding="utf-8", + ) + + async def _require_session(cls): + return { + "console_session_token": "cs_core", + "fingerprint": "fp_core", + "install_id": "inst_core", + } + + captured: dict[str, object] = {} + + class _Response: + status_code = 200 + + def raise_for_status(self): + return None + + def json(self): + return { + "license_patch": "patch_token_1", + "revoked_license_ids": ["lic_revoked"], + } + + class _Client: + def __init__(self, *_args, **_kwargs): + pass + + async def __aenter__(self): + return self + + async def __aexit__(self, exc_type, exc, tb): + return False + + async def post(self, url, json=None, headers=None): + captured["url"] = url + captured["json"] = json + captured["headers"] = headers + return _Response() + + monkeypatch.setattr(ConsoleLoginService, "_require_session", classmethod(_require_session)) + monkeypatch.setattr(login_mod.httpx, "AsyncClient", _Client) + + asyncio.run(ConsoleLoginService.send_heartbeat()) + + assert captured["url"] == "https://console.example.com/v1/heartbeats" + assert captured["headers"] == {"Authorization": "Bearer cs_core"} + payload = captured["json"] + assert payload["status"] == "poc" + assert payload["license_id"] == "lic_core" + assert payload["bundle_version"] == "v2026.7.3" + assert payload["core_version"] == "v2026.7.3" + assert payload["flockspro_component_version"] == "2026.7.3.1" + assert "version" not in payload + assert "version_info" not in payload + + updated = json.loads(license_path.read_text(encoding="utf-8")) + assert updated["patches"] == ["patch_token_1"] + assert updated["last_sync_at"] + revocation = json.loads((tmp_path / "flockspro" / "revocation.json").read_text(encoding="utf-8")) + assert revocation == {"revoked_license_ids": ["lic_revoked"]} + + +def test_report_pending_pro_bundle_install_receipt_posts_and_deletes(tmp_path, monkeypatch): + monkeypatch.setenv("FLOCKS_ROOT", str(tmp_path)) + monkeypatch.delenv("FLOCKS_CONSOLE_BASE_URL", raising=False) + + license_path = tmp_path / "flockspro" / "license.json" + license_path.parent.mkdir(parents=True) + license_path.write_text(json.dumps({"license_id": "lic_pending"}), encoding="utf-8") + pending_path = tmp_path / "run" / "pro-bundle-install-receipt-pending.json" + pending_path.parent.mkdir(parents=True) + pending_path.write_text( + json.dumps( + { + "release_id": "rel_pending", + "bundle_release_id": "rel_pending", + "bundle_version": "2026.7.3.5", + "core_version": "2026.7.3.5", + "flockspro_component_version": "2026.7.3.3", + "build_id": "job_pending", + "install_result": "success", + } + ), + encoding="utf-8", + ) + + async def _require_session(cls): + return { + "console_session_token": "cs_pending", + "fingerprint": "fp_pending", + "install_id": "inst_pending", + "console_base_url": "http://127.0.0.1:18001", + } + + captured: dict[str, object] = {} + + class _Response: + status_code = 200 + + class _Client: + def __init__(self, *_args, **_kwargs): + pass + + async def __aenter__(self): + return self + + async def __aexit__(self, exc_type, exc, tb): + return False + + async def post(self, url, json=None, headers=None): + captured["url"] = url + captured["json"] = json + captured["headers"] = headers + return _Response() + + monkeypatch.setattr(ConsoleLoginService, "_require_session", classmethod(_require_session)) + monkeypatch.setattr(login_mod.httpx, "AsyncClient", _Client) + + reported = asyncio.run(ConsoleLoginService.report_pending_pro_bundle_install_receipt()) + + assert reported is True + assert not pending_path.exists() + assert captured["url"] == "http://127.0.0.1:18001/v1/pro-bundles/installations" + assert captured["headers"] == {"Authorization": "Bearer cs_pending"} + payload = captured["json"] + assert payload["fingerprint"] == "fp_pending" + assert payload["install_id"] == "inst_pending" + assert payload["license_id"] == "lic_pending" + assert payload["bundle_version"] == "2026.7.3.5" diff --git a/tests/server/routes/test_console_upgrade_routes.py b/tests/server/routes/test_console_upgrade_routes.py index 575e1e10e..6a638d3bc 100644 --- a/tests/server/routes/test_console_upgrade_routes.py +++ b/tests/server/routes/test_console_upgrade_routes.py @@ -248,7 +248,7 @@ async def test_pro_package_status_reports_installed_marker( console_routes, "_read_pro_bundle_install_marker", lambda: { - "installed_version": "pro-v2026-05-13-3", + "bundle_version": "pro-v2026-05-13-3", "flockspro_component_version": "1.2.3", "build_id": "build_1", "installed_at": "2026-05-15T12:00:00+00:00", @@ -276,7 +276,7 @@ async def test_pro_package_status_treats_install_marker_as_installed( console_routes, "_read_pro_bundle_install_marker", lambda: { - "installed_version": "pro-v2026.6.23", + "bundle_version": "pro-v2026.6.23", "flockspro_component_version": "2026.6.23", "installed_at": "2026-06-29T04:00:00+00:00", }, @@ -337,6 +337,50 @@ async def test_flockspro_license_status_delegates_to_pro_runtime(monkeypatch: py assert payload["license_id"] == "lic_1" +async def test_flockspro_license_refresh_sends_heartbeat_from_core(monkeypatch: pytest.MonkeyPatch): + from flocks.server.routes import flockspro_license as license_routes + + app = FastAPI() + app.include_router(license_routes.router, prefix="/api/flockspro/license") + monkeypatch.setattr(license_routes, "_is_pro_component_installed", lambda: True) + monkeypatch.setattr( + license_routes, + "_get_pro_capability_status", + lambda: {"active": True, "pro_enabled": True, "license_status": "poc", "license_id": "lic_1"}, + ) + monkeypatch.setattr(license_routes, "require_user", lambda _req: _mock_admin()) + + heartbeat_calls: list[str] = [] + refresh_calls: list[str] = [] + + async def _send_heartbeat(): + heartbeat_calls.append("sent") + return {"ok": True} + + class _Checker: + async def refresh(self): + refresh_calls.append("refreshed") + return {"active": True} + + runtime_module = ModuleType("flockspro.license.runtime") + runtime_module.get_license_checker = lambda: _Checker() + license_module = ModuleType("flockspro.license") + flockspro_module = ModuleType("flockspro") + monkeypatch.setitem(__import__("sys").modules, "flockspro", flockspro_module) + monkeypatch.setitem(__import__("sys").modules, "flockspro.license", license_module) + monkeypatch.setitem(__import__("sys").modules, "flockspro.license.runtime", runtime_module) + monkeypatch.setattr(license_routes.ConsoleLoginService, "send_heartbeat", _send_heartbeat) + + transport = httpx.ASGITransport(app=app) + async with AsyncClient(transport=transport, base_url="http://test") as local_client: + resp = await local_client.post("/api/flockspro/license/refresh") + + assert resp.status_code == status.HTTP_200_OK + assert heartbeat_calls == ["sent"] + assert refresh_calls == ["refreshed"] + assert resp.json()["license_id"] == "lic_1" + + async def test_create_upgrade_request_does_not_link_previous_request_when_omitted( client: AsyncClient, monkeypatch: pytest.MonkeyPatch, @@ -768,7 +812,7 @@ async def _fake_report(record: dict, *, install_result: str, error_message: str console_routes, "_read_pro_bundle_install_marker", lambda: { - "installed_version": "v2026.6.5", + "bundle_version": "v2026.6.5", "flockspro_component_version": "v2026.6.5", }, ) @@ -812,7 +856,7 @@ async def test_restarting_request_reports_receipt_after_service_restart( "approved_bundle_release_id": "rel_restart", "latest_pro_bundle": { "release_id": "rel_restart", - "display_version": "v2026.6.24", + "bundle_version": "v2026.6.24", "core_version": "v2026.6.21", "flockspro_component_version": "v2026.6.24", "build_id": "job_restart", @@ -844,7 +888,7 @@ async def _fake_report(record: dict, *, install_result: str, error_message: str lambda: { "release_id": "rel_restart", "bundle_release_id": "rel_restart", - "installed_version": "v2026.6.24", + "bundle_version": "v2026.6.24", "core_version": "v2026.6.21", "flockspro_component_version": "v2026.6.24", "build_id": "job_restart", @@ -857,7 +901,7 @@ async def _fake_report(record: dict, *, install_result: str, error_message: str payload = resp.json() assert payload["status"] == "activated" assert payload["details"]["auto_install_result"] == "done" - assert payload["details"]["auto_install_version"] == "v2026.6.24" + assert payload["details"]["auto_install_bundle_version"] == "v2026.6.24" assert reported == [("success", None)] @@ -977,7 +1021,7 @@ async def _fake_report(record: dict, *, install_result: str, error_message: str monkeypatch.setattr( console_routes, "_read_pro_bundle_install_marker", - lambda: {"installed_version": "v2026.5.9"} if installed else {}, + lambda: {"bundle_version": "v2026.5.9"} if installed else {}, ) resp = await client.post(f"/api/console/upgrade-requests/{request_id}/start") @@ -987,7 +1031,7 @@ async def _fake_report(record: dict, *, install_result: str, error_message: str stored = await Storage.get(f"console:upgrade_request:{request_id}") assert stored["status"] == "activated" assert stored["details"]["auto_install_result"] == "done" - assert stored["details"]["auto_install_version"] == "v2026.5.9" + assert stored["details"]["auto_install_bundle_version"] == "v2026.5.9" async def test_start_revoked_request_does_not_reinstall( @@ -1042,7 +1086,7 @@ async def _noop(_record: dict, **_kwargs): monkeypatch.setattr( console_routes, "_read_pro_bundle_install_marker", - lambda: {"installed_version": "v2026.5.9"}, + lambda: {"bundle_version": "v2026.5.9"}, ) record = { @@ -1057,7 +1101,7 @@ async def _noop(_record: dict, **_kwargs): payload = await console_routes._maybe_auto_activate_upgrade(record) assert payload["status"] == "activated" assert payload["details"]["auto_install_result"] == "already_latest" - assert payload["details"]["auto_install_version"] == "v2026.5.9" + assert payload["details"]["auto_install_bundle_version"] == "v2026.5.9" assert reported == [("success", None)] @@ -1071,7 +1115,7 @@ async def test_auto_activate_reinstalls_when_existing_pro_marker_is_not_target_b "payload": { "release_id": "rel_20260601", "bundle_release_id": "rel_20260601", - "installed_version": "v2026.6.1", + "bundle_version": "v2026.6.1", "flockspro_component_version": "v2026.6.1", "build_id": "job_20260601", } @@ -1084,7 +1128,7 @@ async def _fake_perform_pro_bundle_install(*args, **kwargs): marker_state["payload"] = { "release_id": "rel_20260605", "bundle_release_id": "rel_20260605", - "installed_version": "v2026.6.5", + "bundle_version": "v2026.6.5", "flockspro_component_version": "v2026.6.5", "build_id": "job_20260605", } @@ -1113,7 +1157,7 @@ async def _noop(_record: dict, **_kwargs): "approved_bundle_release_id": "rel_20260605", "latest_pro_bundle": { "release_id": "rel_20260605", - "display_version": "v2026.6.5", + "bundle_version": "v2026.6.5", "flockspro_component_version": "v2026.6.5", "build_id": "job_20260605", }, @@ -1127,7 +1171,7 @@ async def _noop(_record: dict, **_kwargs): assert payload["status"] == "activated" assert payload["details"]["auto_install_result"] == "done" assert payload["details"]["auto_install_release_id"] == "rel_20260605" - assert payload["details"]["auto_install_version"] == "v2026.6.5" + assert payload["details"]["auto_install_bundle_version"] == "v2026.6.5" assert reported == [("success", None)] @@ -1153,7 +1197,7 @@ async def _noop(_record: dict, **_kwargs): "_get_pro_capability_status", lambda: {"pro_enabled": False, "active": False, "license_status": "expired", "inactive_reason": "expired"}, ) - monkeypatch.setattr(console_routes, "_read_pro_bundle_install_marker", lambda: {"installed_version": "v2026.5.9"}) + monkeypatch.setattr(console_routes, "_read_pro_bundle_install_marker", lambda: {"bundle_version": "v2026.5.9"}) record = { "request_id": "req_auto_inactive", @@ -1205,7 +1249,7 @@ async def _noop(_record: dict, **_kwargs): monkeypatch.setattr( console_routes, "_read_pro_bundle_install_marker", - lambda: {"installed_version": "v2026.5.9"} if installed else {}, + lambda: {"bundle_version": "v2026.5.9"} if installed else {}, ) record = { @@ -1220,7 +1264,7 @@ async def _noop(_record: dict, **_kwargs): payload = await console_routes._maybe_auto_activate_upgrade(record) assert payload["status"] == "activated" assert payload["details"]["auto_install_result"] == "done" - assert payload["details"]["auto_install_version"] == "v2026.5.9" + assert payload["details"]["auto_install_bundle_version"] == "v2026.5.9" async def test_report_pro_bundle_installation_uses_license_id( @@ -1253,7 +1297,7 @@ async def post(self, url, json=None, headers=None): monkeypatch.setattr( console_routes, "_read_pro_bundle_install_marker", - lambda: {"installed_version": "v2026.5.9"}, + lambda: {"bundle_version": "v2026.5.9"}, ) record = { @@ -1266,7 +1310,7 @@ async def post(self, url, json=None, headers=None): "approved_bundle_release_id": "rel_receipt", "latest_pro_bundle": { "release_id": "rel_receipt", - "display_version": "v2026.6.5", + "bundle_version": "v2026.6.5", "core_version": "v2026.6.1", "flockspro_component_version": "v2026.6.5", "build_id": "job_receipt", @@ -1281,7 +1325,7 @@ async def post(self, url, json=None, headers=None): assert posted_payloads[0]["release_id"] == "rel_receipt" assert posted_payloads[0]["bundle_release_id"] == "rel_receipt" assert posted_payloads[0]["core_version"] == "v2026.6.1" - assert posted_payloads[0]["oss_version"] == "v2026.6.1" + assert "oss_version" not in posted_payloads[0] assert posted_payloads[0]["build_id"] == "job_receipt" @@ -1317,7 +1361,7 @@ async def post(self, url, json=None, headers=None): lambda: { "release_id": "rel_old", "bundle_release_id": "rel_old", - "installed_version": "v2026.6.1", + "bundle_version": "v2026.6.1", "flockspro_component_version": "v2026.6.1", "build_id": "job_old", }, @@ -1333,8 +1377,8 @@ async def post(self, url, json=None, headers=None): "approved_bundle_release_id": "rel_new", "latest_pro_bundle": { "release_id": "rel_new", - "display_version": "v2026.6.5", - "oss_version": "v2026.6.5", + "bundle_version": "v2026.6.5", + "core_version": "v2026.6.5", "flockspro_component_version": "v2026.6.5", "build_id": "job_new", }, @@ -1349,6 +1393,6 @@ async def post(self, url, json=None, headers=None): assert posted_payloads[0]["release_id"] == "rel_new" assert posted_payloads[0]["bundle_release_id"] == "rel_new" - assert posted_payloads[0]["installed_version"] == "v2026.6.5" + assert posted_payloads[0]["bundle_version"] == "v2026.6.5" assert posted_payloads[0]["build_id"] == "job_new" assert posted_payloads[0]["install_result"] == "failed" diff --git a/tests/updater/test_restart_handoff.py b/tests/updater/test_restart_handoff.py index 6f9ceaab2..b9e5e1e5d 100644 --- a/tests/updater/test_restart_handoff.py +++ b/tests/updater/test_restart_handoff.py @@ -77,6 +77,40 @@ def test_run_waits_for_parent_and_backend_port_before_spawning( ] +def test_run_reports_pending_install_receipt_after_pro_bundle_tasks( + monkeypatch, + tmp_path: Path, +) -> None: + events: list[str] = [] + restart_argv = ["python.exe", "-m", "flocks.cli.main", "serve"] + manifest = tmp_path / "manifest.json" + manifest.write_text("{}", encoding="utf-8") + args = _handoff_args(tmp_path, restart_argv) + separator_index = args.index("--") + args[separator_index:separator_index] = ["--pro-bundle-manifest-path", str(manifest)] + + monkeypatch.setattr(restart_handoff, "_record_handoff_log", lambda message: events.append(f"log:{message}")) + monkeypatch.setattr(restart_handoff, "_wait_for_parent_exit", lambda parent_pid: True) + monkeypatch.setattr(restart_handoff, "_ensure_backend_port_free", lambda backend_port, backend_pid_file: True) + monkeypatch.setattr(restart_handoff, "_run_upgrade_tasks", lambda args: events.append("tasks") or None) + monkeypatch.setattr( + restart_handoff, + "_report_pending_pro_bundle_install_receipt", + lambda args: events.append("receipt"), + ) + monkeypatch.setattr( + restart_handoff.subprocess, + "Popen", + lambda argv, cwd=None, close_fds=False: events.append(f"spawn:{list(argv)}") or SimpleNamespace(pid=4321), + ) + monkeypatch.setattr(restart_handoff, "_record_backend_runtime_if_direct_serve", lambda *_args, **_kwargs: None) + + code = restart_handoff.run(args) + + assert code == 0 + assert events[1:4] == ["tasks", "receipt", f"spawn:{restart_argv}"] + + def test_run_does_not_spawn_when_parent_exit_times_out(monkeypatch, tmp_path: Path) -> None: events: list[str] = [] diff --git a/tests/updater/test_updater_console_manifest_bundle.py b/tests/updater/test_updater_console_manifest_bundle.py index 518f3b862..2389f618c 100644 --- a/tests/updater/test_updater_console_manifest_bundle.py +++ b/tests/updater/test_updater_console_manifest_bundle.py @@ -15,6 +15,9 @@ async def test_fetch_console_manifest_release_uses_bundle_url(monkeypatch: pytes from flocks.storage.storage import Storage monkeypatch.setenv("FLOCKS_ROOT", str(tmp_path)) + license_path = tmp_path / "flockspro" / "license.json" + license_path.parent.mkdir(parents=True) + license_path.write_text('{"license_id": "lic_manifest"}', encoding="utf-8") await Storage.set("console:session", {"console_session_token": "cs_manifest"}, "json") class _Resp: @@ -23,11 +26,11 @@ def raise_for_status(self) -> None: def json(self) -> dict: return { - "display_version": "v2026.5.10", + "bundle_version": "v2026.5.10", "compare_version": "2026.5.10", "bundle_url": "https://cdn.example.com/flockspro-bundle-v2026.5.10.tar.gz", "bundle_sha256": "abc123", - "oss_version": "v2026.5.10", + "core_version": "v2026.5.10", "flockspro_component_version": "pro-v2026-5-10", "release_notes": "bundle release", } @@ -41,7 +44,11 @@ async def __aexit__(self, exc_type, exc, tb): async def get(self, url, headers=None, follow_redirects=True): assert "channel=flockspro" in url - assert headers == {"Authorization": "Bearer cs_manifest"} + assert "license_id=lic_manifest" in url + assert headers == { + "x-license-id": "lic_manifest", + "Authorization": "Bearer cs_manifest", + } return _Resp() monkeypatch.setenv("FLOCKS_CONSOLE_BASE_URL", "https://console.example.com") @@ -68,7 +75,8 @@ async def test_check_update_uses_pro_marker_bundle_version_and_component_metadat marker.parent.mkdir(parents=True) marker.write_text( """{ - "installed_version": "v2026.5.23", + "bundle_version": "v2026.5.23", + "core_version": "v2026.5.23", "flockspro_component_version": "pro-v2026-05-23" }""", encoding="utf-8", @@ -86,7 +94,8 @@ async def _fake_manifest_info(): bundle_sha256=None, bundle_format="zip", manifest={ - "display_version": "v2026.5.23", + "bundle_version": "v2026.5.23", + "core_version": "v2026.5.23", "flockspro_component_version": "pro-v2026-05-23", }, ) @@ -119,7 +128,8 @@ async def test_check_update_force_console_manifest_uses_bundle_versions(monkeypa marker.parent.mkdir(parents=True) marker.write_text( """{ - "installed_version": "v2026.5.23", + "bundle_version": "v2026.5.23", + "core_version": "v2026.5.23", "flockspro_component_version": "pro-v2026-05-23" }""", encoding="utf-8", @@ -137,7 +147,7 @@ async def _fake_manifest_info(): bundle_sha256="abc123", bundle_format="zip", manifest={ - "display_version": "v2026.5.24", + "bundle_version": "v2026.5.24", "core_version": "v2026.5.23", "flockspro_component_version": "pro-v2026-05-24", }, @@ -173,7 +183,8 @@ async def test_check_update_force_console_manifest_detects_component_only_update marker.parent.mkdir(parents=True) marker.write_text( """{ - "installed_version": "v2026.6.18", + "bundle_version": "v2026.6.18", + "core_version": "v2026.6.18", "flockspro_component_version": "v2026.6.1" }""", encoding="utf-8", @@ -191,8 +202,8 @@ async def _fake_manifest_info(): bundle_sha256="def456", bundle_format="zip", manifest={ - "display_version": "v2026.6.18", - "oss_version": "v2026.6.18", + "bundle_version": "v2026.6.18", + "core_version": "v2026.6.18", "flockspro_component_version": "v2026.6.2", }, ) @@ -220,7 +231,7 @@ async def test_check_update_force_console_manifest_reports_stale_product_marker_ marker.parent.mkdir(parents=True) marker.write_text( """{ - "installed_version": "v2026.6.22", + "bundle_version": "v2026.6.22", "core_version": "v2026.6.21", "flockspro_component_version": "v2026.6.23" }""", @@ -239,7 +250,7 @@ async def _fake_manifest_info(): bundle_sha256="ghi789", bundle_format="zip", manifest={ - "display_version": "v2026.6.23", + "bundle_version": "v2026.6.23", "core_version": "v2026.6.21", "flockspro_component_version": "v2026.6.23", }, @@ -271,29 +282,34 @@ def test_console_manifest_release_identity_writes_product_and_core_versions( monkeypatch.setenv("FLOCKS_ROOT", str(tmp_path)) merged = updater._merge_console_manifest_release_identity( { - "display_version": "v2026.6.21", + "bundle_version": "v2026.6.21", "core_version": "v2026.6.21", "flockspro_component_version": "v2026.6.23", }, { "release_id": "rel_623", - "display_version": "v2026.6.23", + "bundle_version": "v2026.6.23", "core_version": "v2026.6.21", "flockspro_component_version": "v2026.6.23", "build_id": "job_623", }, ) - assert merged["display_version"] == "v2026.6.23" + assert merged["bundle_version"] == "v2026.6.23" assert merged["core_version"] == "v2026.6.21" updater._write_pro_bundle_install_marker(merged, bundle_sha256="sha623") marker = json.loads((tmp_path / "run" / "pro-bundle-installed.json").read_text(encoding="utf-8")) - assert marker["installed_version"] == "v2026.6.23" + assert marker["bundle_version"] == "v2026.6.23" assert marker["core_version"] == "v2026.6.21" - assert marker["oss_version"] == "v2026.6.21" assert marker["flockspro_component_version"] == "v2026.6.23" assert marker["build_id"] == "job_623" + pending = json.loads((tmp_path / "run" / "pro-bundle-install-receipt-pending.json").read_text(encoding="utf-8")) + assert pending["install_result"] == "success" + assert pending["bundle_version"] == "v2026.6.23" + assert pending["core_version"] == "v2026.6.21" + assert pending["flockspro_component_version"] == "v2026.6.23" + assert "version_info" not in pending @pytest.mark.asyncio @@ -358,7 +374,7 @@ def raise_for_status(self) -> None: def json(self) -> dict: return { - "display_version": "v2026.5.10", + "bundle_version": "v2026.5.10", "bundle_url": "https://cdn.example.com/flockspro-bundle-v2026.5.10.tar.gz", "frozen": True, } @@ -577,8 +593,8 @@ async def test_perform_pro_bundle_install_replaces_core_and_installs_wheel( wheel.write_bytes(b"fake-wheel") (bundle_root / "manifest.json").write_text( """{ - "display_version": "v2026.5.10", - "oss_version": "v2026.5.10", + "bundle_version": "v2026.5.10", + "core_version": "v2026.5.10", "flockspro_component_version": "pro-v2026-5-10", "flockspro_wheel": "wheels/flockspro-0.1.0-py3-none-any.whl", "build_id": "job_test" @@ -628,12 +644,12 @@ async def _fake_run_async(cmd, **_kwargs): marker = tmp_path / "flocks-root" / "run" / "pro-bundle-installed.json" assert marker.is_file() marker_payload = __import__("json").loads(marker.read_text(encoding="utf-8")) - assert marker_payload["display_version"] == "v2026.5.10" - assert marker_payload["oss_version"] == "v2026.5.10" + assert marker_payload["bundle_version"] == "v2026.5.10" + assert marker_payload["core_version"] == "v2026.5.10" @pytest.mark.asyncio -async def test_perform_pro_bundle_install_keeps_newer_local_core_when_bundle_oss_is_older( +async def test_perform_pro_bundle_install_keeps_newer_local_core_when_bundle_core_is_older( monkeypatch: pytest.MonkeyPatch, tmp_path, ) -> None: @@ -648,8 +664,8 @@ async def test_perform_pro_bundle_install_keeps_newer_local_core_when_bundle_oss wheel.write_bytes(b"fake-wheel") (bundle_root / "manifest.json").write_text( """{ - "display_version": "v2026.6.13", - "oss_version": "v2026.6.13", + "bundle_version": "v2026.6.13", + "core_version": "v2026.6.13", "flockspro_component_version": "v2026.6.2", "flockspro_wheel": "wheels/flockspro-0.2.0-py3-none-any.whl", "build_id": "job_new_pro_old_core" @@ -683,8 +699,8 @@ async def _fake_manifest_info(): bundle_format="zip", manifest={ "release_id": "rel_new_pro_old_core", - "display_version": "v2026.6.13", - "oss_version": "v2026.6.13", + "bundle_version": "v2026.6.13", + "core_version": "v2026.6.13", "flockspro_component_version": "v2026.6.2", "build_id": "job_new_pro_old_core", }, @@ -718,10 +734,8 @@ async def _fake_run_async(cmd, **_kwargs): marker = tmp_path / "flocks-root" / "run" / "pro-bundle-installed.json" marker_payload = __import__("json").loads(marker.read_text(encoding="utf-8")) assert marker_payload["release_id"] == "rel_new_pro_old_core" - assert marker_payload["display_version"] == "v2026.6.13" - assert marker_payload["installed_version"] == "v2026.6.13" + assert marker_payload["bundle_version"] == "v2026.6.13" assert marker_payload["core_version"] == "v2026.6.18" - assert marker_payload["oss_version"] == "v2026.6.18" assert marker_payload["flockspro_component_version"] == "v2026.6.2" @@ -740,8 +754,8 @@ async def test_perform_pro_bundle_install_schedules_restart_before_stream_can_cl wheel.write_bytes(b"fake-wheel") (bundle_root / "manifest.json").write_text( """{ - "display_version": "v2026.5.10", - "oss_version": "v2026.5.10", + "bundle_version": "v2026.5.10", + "core_version": "v2026.5.10", "flockspro_component_version": "pro-v2026-5-10", "flockspro_wheel": "wheels/flockspro-0.1.0-py3-none-any.whl", "build_id": "job_test" @@ -792,8 +806,8 @@ async def _async_manifest_info(bundle): bundle_sha256=None, bundle_format="zip", manifest={ - "display_version": "v2026.5.10", - "oss_version": "v2026.5.10", + "bundle_version": "v2026.5.10", + "core_version": "v2026.5.10", "flockspro_component_version": "pro-v2026-5-10", "build_id": "job_test", }, diff --git a/tests/updater/test_updater_edition_sources.py b/tests/updater/test_updater_edition_sources.py index 2798ed541..a58b33ac9 100644 --- a/tests/updater/test_updater_edition_sources.py +++ b/tests/updater/test_updater_edition_sources.py @@ -10,7 +10,8 @@ async def test_installed_pro_bundle_marker_without_active_license_keeps_oss_sour marker.parent.mkdir(parents=True) marker.write_text( """{ - "installed_version": "v2026.5.23", + "bundle_version": "v2026.5.23", + "core_version": "v2026.5.23", "flockspro_component_version": "pro-v2026-05-23" }""", encoding="utf-8", diff --git a/webui/src/api/consoleUpgrade.ts b/webui/src/api/consoleUpgrade.ts index e211efcb1..d49d8d2a9 100644 --- a/webui/src/api/consoleUpgrade.ts +++ b/webui/src/api/consoleUpgrade.ts @@ -37,8 +37,11 @@ export interface UpgradeRequestDetails { notes?: string | null; auto_install_target?: string; auto_install_version?: string; + auto_install_bundle_version?: string; auto_install_pro_version?: string; + auto_install_pro_component_version?: string; flockspro_component_version?: string; + bundle_version_update_to?: string; auto_install_result?: string; auto_install_completed_at?: string; license_refreshed_at?: string; @@ -66,7 +69,9 @@ export interface ProPackageStatus { installed: boolean; runtime_importable?: boolean | null; install_marker_present?: boolean | null; + bundle_version?: string | null; installed_version?: string | null; + core_version?: string | null; flockspro_component_version?: string | null; build_id?: string | null; installed_at?: string | null; diff --git a/webui/src/components/common/UpdateModal.tsx b/webui/src/components/common/UpdateModal.tsx index fbfbdb12b..ab312d31f 100644 --- a/webui/src/components/common/UpdateModal.tsx +++ b/webui/src/components/common/UpdateModal.tsx @@ -29,6 +29,28 @@ function formatUpdateVersion(version?: string | null): string { return /^(pro-)?v/i.test(raw) ? raw : `v${raw}`; } +function clampPercent(value?: number | null): number | null { + if (typeof value !== 'number' || Number.isNaN(value)) { + return null; + } + return Math.max(0, Math.min(100, Math.round(value))); +} + +function formatBytes(value?: number | null): string { + if (typeof value !== 'number' || !Number.isFinite(value) || value < 0) { + return '—'; + } + const units = ['B', 'KB', 'MB', 'GB']; + let size = value; + let unitIndex = 0; + while (size >= 1024 && unitIndex < units.length - 1) { + size /= 1024; + unitIndex += 1; + } + const precision = unitIndex === 0 || size >= 10 ? 0 : 1; + return `${size.toFixed(precision)} ${units[unitIndex]}`; +} + interface UpdateModalProps { initialInfo?: VersionInfo | null; edition?: UpdateEdition; @@ -152,22 +174,53 @@ export default function UpdateModal({ initialInfo, edition = 'flocks', canUpgrad }; const renderStep = (step: UpdateProgress, index: number) => { - const label = t(`stageLabels.${step.stage}`, { defaultValue: step.stage }); + const label = edition === 'flockspro' && step.stage === 'fetching' + ? t('stageLabels.fetchingPro') + : t(`stageLabels.${step.stage}`, { defaultValue: step.stage }); const isError = step.stage === 'error'; const isSpinning = step.stage === 'restarting'; + const downloadPercent = clampPercent(step.percent); + const hasDownloadProgress = step.stage === 'fetching' && typeof step.downloaded_bytes === 'number'; const detail = step.pro_component_filename || step.bundle_filename || step.message; return ( -