From 79f4f01954bc07fd2a6c4e9b46d6f4283068610d Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Thu, 8 Oct 2026 17:45:47 +0000 Subject: [PATCH] Add a drain-and-restore path for a live serve upgrade (MN-REQ-06.14). Admin upgrade-prepare snapshots every loaded session, refuses new work with a retryable serve_draining error, and blocks ready-to-stop when a session cannot be saved. Startup restores from the manifest and checks counts and checksums. MCP and the gateway retry across the gap. Co-authored-by: chouswei --- docs/README.md | 1 + docs/operations/README.md | 1 + docs/operations/safe-upgrade.md | 112 +++ parts/common/memnet/memnet/cli.py | 65 ++ .../common/memnet/memnet/local_ipc_gateway.py | 3 + parts/common/memnet/memnet/serve.py | 52 +- parts/common/memnet/memnet/session.py | 3 + parts/common/memnet/memnet/snapshot.py | 14 +- parts/common/memnet/memnet/upgrade.py | 646 ++++++++++++++++++ parts/common/memnet/memnet/upgrade_retry.py | 67 ++ parts/common/memnet/memnet/upgrade_run.py | 258 +++++++ .../memnet-mcp/software/memnet_mcp/client.py | 16 +- .../software/memnet_mcp/product_gateway.py | 44 +- .../software/memnet_mcp/serve_bridge.py | 9 +- pyproject.toml | 1 + sysml-models/models/deploy.sysml | 69 ++ sysml-models/models/implementation.sysml | 38 ++ sysml-models/models/requirements.sysml | 58 ++ sysml-models/models/verify.sysml | 70 ++ sysml-models/outputs/README.md | 1 + sysml-models/outputs/product-nest-one-page.md | 1 + .../outputs/safe-upgrade-case-study.md | 20 + .../outputs/ssot-to-code-allocate-map.md | 4 + sysml-models/outputs/system-design-notes.md | 1 + tests/conftest.py | 10 + tests/test_safe_upgrade.py | 556 +++++++++++++++ tests/test_sysml_safe_upgrade.py | 69 ++ 27 files changed, 2148 insertions(+), 41 deletions(-) create mode 100644 docs/operations/safe-upgrade.md create mode 100644 parts/common/memnet/memnet/upgrade.py create mode 100644 parts/common/memnet/memnet/upgrade_retry.py create mode 100644 parts/common/memnet/memnet/upgrade_run.py create mode 100644 sysml-models/outputs/safe-upgrade-case-study.md create mode 100644 tests/test_safe_upgrade.py create mode 100644 tests/test_sysml_safe_upgrade.py diff --git a/docs/README.md b/docs/README.md index 90363baf..0a3452ee 100644 --- a/docs/README.md +++ b/docs/README.md @@ -68,6 +68,7 @@ Multitask MUST for this product. Index: [`operations/README.md`](operations/READ | [`operations/product-gateway-contract.md`](operations/product-gateway-contract.md) | Product gateway on memnet-mcp (MN-REQ-06.12): route by product, house, or session | | [`operations/cluster-route-vs-slice-hand-carry.md`](operations/cluster-route-vs-slice-hand-carry.md) | Two named moves: ClusterRoute vs SliceHandCarry (#191 / #47 cousin; inventOnly) | | [`operations/admin-usage-report.md`](operations/admin-usage-report.md) | Admin-only serve usage JSON for a product-gate admin MCP (opaque alias; not agent MCP) | +| [`operations/safe-upgrade.md`](operations/safe-upgrade.md) | Drain and restore a live serve without dropping sessions (MN-REQ-06.14) | | [`operations/one-session-per-document.md`](operations/one-session-per-document.md) | One MemNet session per document over loopback serve (no MCP front) | Product skill: [`.cursor/skills/memnet-reference/`](../.cursor/skills/memnet-reference/). SysML trail: MN-REQ-12 → [`sysml-models/outputs/multitask-case-study.md`](../sysml-models/outputs/multitask-case-study.md). diff --git a/docs/operations/README.md b/docs/operations/README.md index 4bf2e9c0..b5ae5eab 100644 --- a/docs/operations/README.md +++ b/docs/operations/README.md @@ -12,6 +12,7 @@ Agent operating doctrine for this product (not domain recipes). | [`product-gateway-contract.md`](product-gateway-contract.md) | Product gateway on memnet-mcp (MN-REQ-06.12): per-product route to one owning serve | | [`cluster-route-vs-slice-hand-carry.md`](cluster-route-vs-slice-hand-carry.md) | Two named moves: ClusterRoute vs SliceHandCarry (#191 / #47 cousin; inventOnly) | | [`admin-usage-report.md`](admin-usage-report.md) | Admin-only serve usage JSON (opaque alias; not agent MCP; MN-REQ-06.11) | +| [`safe-upgrade.md`](safe-upgrade.md) | Drain, restore, and `memnet-upgrade` so a serve swap keeps sessions (MN-REQ-06.14) | | [`one-session-per-document.md`](one-session-per-document.md) | Product gate: one serve session per document over loopback (no MCP); 0.19.18 probe | Application pattern for `modelbasedPrj-*` / `SysMLEdgePrj-*`: [`../application-notes/system/llm-system-dev-multitask.md`](../application-notes/system/llm-system-dev-multitask.md). Index: [`../README.md`](../README.md). diff --git a/docs/operations/safe-upgrade.md b/docs/operations/safe-upgrade.md new file mode 100644 index 00000000..6b96c50e --- /dev/null +++ b/docs/operations/safe-upgrade.md @@ -0,0 +1,112 @@ +# Safe serve upgrade + +Upgrade `memnet-serve` without silently dropping a loaded session (MN-REQ-06.14). The agent loop does not change: cue, then `pin_map`, then `mutate`. This is a patch-level cut on 0.19. Do not bump the package version from this procedure. + +The manual procedure this replaces is a side-by-side virtualenv: install the new version, save every session, stop the old serve, start the new serve on the same port, reload, and check counts. Anything that failed to save was easy to miss. The drain below names every session it could not snapshot and refuses a ready-to-stop result unless you pass an explicit override. + +## Order + +Deploy client tolerance **before** the serve swap. + +1. Install the new `memnet-mcp` (and the droplet gateway, if this host is the gateway) so `serve_draining` and a brief connection refusal are retried. Leave the gateway's `pinned_version` on the version that is still running. +2. Confirm `MEMNET_UPGRADE_RETRY_S` is unset or about `30` on those client processes. `0` disables retry. +3. Only then run `memnet-upgrade` against the serve. + +During the gap, new `session_open` calls receive `@ERR: serve_draining|retry_after_s=` and clients back off. After the old process has exited and before the new one listens, connections are refused. That refusal is retried for the same window. A command the serve already accepted is not retried on timeout. + +Downtime that remains: a few seconds while the port is closed, covered by that retry. Sessions keep their ids, CapsPolicy ACL bindings, TTL expiry clock, and house (the session tag map). + +## What the serve does + +Admin credential: `MEMNET_ADMIN_TOKEN` (same rule as the usage report). Unset is `@ERR: admin_unconfigured`. A mismatch is `@ERR: admin_denied`. This command is not an agent MCP tool. + +```bash +memnet admin upgrade-prepare --state-dir "$MEMNET_STATE_DIR" +``` + +Envelope form: `{"upgrade_prepare": true, "admin_token": ""}`. + +Drain behaviour: + +- Refuse new `session_open` with `serve_draining`. +- Finish commands already in flight, then refuse further commands. +- Snapshot every loaded session with the lossless snapshot writer. +- Write `upgrade-manifest.json` and `upgrade-snapshots/` under `MEMNET_STATE_DIR` (default `~/.local/state/memnet`). The manifest records session ids, row and edge counts, sha256 checksums, the serve version, and snapshot format `1`. +- An ACL and clock passport is an extra `# upgrade-passport` line. Older v1 loaders skip `#` lines, so a rollback can still read the graph. +- If any session raises `snapshot_unsaveable` (or any other save error), the command exits non-zero, lists that session in `upgrade-blocked.json`, and does **not** write a ready manifest. The serve keeps accepting work. Pass `--allow-unsaved` only when you accept dropping those named sessions. +- Do not stop the process unless stderr contains `@STAT: upgrade_prepare|ready|`. + +The new process reads the manifest on startup, reloads each snapshot, and checks counts and checksums. It prints: + +```text +@STAT: upgrade_restore|ok||failed| +``` + +Snapshot files stay on disk until `memnet admin upgrade-retire` after a clean restore report. A format this process does not support, or a checksum or parse failure, exits `3` and does not delete or rewrite the files. Restarting before retire replays the upgrade snapshots (writes made after the restore are not in those files). Retire as soon as the stat line is clean. + +## Helper + +`memnet-upgrade` encodes the side-by-side venv procedure. It refuses to start unless `--clients-ready` is set. + +```bash +memnet-upgrade \ + --clients-ready \ + --new-python /opt/memnet/venv-new/bin/python \ + --new-exec "/opt/memnet/venv-new/bin/memnet serve --host 127.0.0.1 --port 18765" \ + --state-dir /var/lib/memnet \ + --unit /etc/systemd/system/memnet-serve.service \ + --restart-unit memnet-serve +``` + +Steps: + +1. Preflight: import `memnet` with the new interpreter and check `SUPPORTED_SNAPSHOT_FORMATS` against the current manifest (or format `1` when no manifest exists yet). +2. `admin upgrade-prepare` on the running serve. +3. Copy the unit to `memnet-serve.service.bak` and replace `ExecStart=`. +4. If `--gateway-config` and `--product` are set, copy the config to `*.bak` and set that product's `pinned_version` to the new version **after** the old serve is drained and **before** the new process listens. +5. Restart the unit. The new process restores and writes `upgrade-restore.json`. +6. If `failed` is not `0`, write the unit backup and the config backup back and restart again. + +`--state-dir` must be the same directory as the unit's `Environment=MEMNET_STATE_DIR`. The serve reads that variable at startup; the prepare command also receives `--state-dir`. + +Build the new venv beside the old one. Do not overwrite it until the restore has verified. + +```bash +python -m venv /opt/memnet/venv-new +/opt/memnet/venv-new/bin/pip install "memnet-llm==" +``` + +Use the index you trust. Do not put tokens or session ids in the unit file. Keep `MEMNET_ADMIN_TOKEN` in a root-only `EnvironmentFile`. + +## rpi5-syson + +Local engine on the Pi: + +| Process | Port | Unit (adjust to the host) | +|---------|------|---------------------------| +| `memnet serve` | `18765` | `memnet-serve.service` | +| `memnet-mcp` streamable-http | `18766` | `memnet-mcp-http.service` | + +Upgrade the MCP unit first (retry code, same serve version pin). Then `memnet-upgrade` the serve unit on `18765` with `MEMNET_STATE_DIR` on the Pi disk. Leave `18766` up so clients retry against the serve gap. + +## Endleaf engine serve + +The Endleaf product backend is the serve named in the droplet gateway registry (host and port live in that config, not in this repo). Use that unit's port and state directory. The drain and restore are the same command. Do not point the gateway at a second backend to "move" live sessions. + +## Droplet gateway + +The gateway process is `memnet-mcp --transport gateway` with `MEMNET_GATEWAY_CONFIG`. Ship the retry-capable build first, with `pinned_version` still equal to the running engine. Then, in the same `memnet-upgrade` invocation that swaps the engine, pass: + +```bash +--gateway-config /etc/memnet/gateway.json \ +--product endleaf \ +--restart-unit memnet-gateway +``` + +The pin changes only after drain succeeds, and it is restored from `gateway.json.bak` if the engine restore fails. Restart the gateway after the pin write so it stops caching the previous version (`version_cache_s`). No plaintext credentials in git. Hashes stay in the config file on the droplet. + +## Rollback + +A failed verification restores the previous `ExecStart=` and the previous pin, then restarts. Snapshot files are still in the state directory. The old venv loads format v1 snapshots (`#` passport lines are ignored). That reload keeps the session id and the graph. The old loader rebases the TTL clock from `ttl_minutes` and does not read ACL bindings from the passport; re-grant those on the old serve if you stay there. The new serve's restore path keeps the id, the ACL bindings, the expiry instant, and the house. + +Neighbourhood reserves are not part of the snapshot. Take a fresh reserve after the new serve is up. diff --git a/parts/common/memnet/memnet/cli.py b/parts/common/memnet/memnet/cli.py index 995955bb..8652aabd 100644 --- a/parts/common/memnet/memnet/cli.py +++ b/parts/common/memnet/memnet/cli.py @@ -528,6 +528,71 @@ def admin_usage_report( raise typer.Exit(code) +@admin_app.command("upgrade-prepare") +def admin_upgrade_prepare( + token: Annotated[ + str | None, + typer.Option( + "--token", + help="Admin token presented by the caller (do not log). Not MEMNET_ADMIN_TOKEN.", + ), + ] = None, + allow_unsaved: Annotated[ + bool, + typer.Option( + "--allow-unsaved", + help="Reach ready-to-stop even when a session could not be snapshotted.", + ), + ] = False, + state_dir: Annotated[ + str | None, + typer.Option("--state-dir", help="Manifest directory (default MEMNET_STATE_DIR)."), + ] = None, +) -> None: + """Snapshot every loaded session and write an upgrade manifest. Not agent MCP.""" + from pathlib import Path + + from memnet.admin_usage import caller_token + from memnet.upgrade import prepare_upgrade + + presented = token if token else caller_token() + directory = Path(state_dir) if state_dir else None + try: + result = prepare_upgrade(presented, allow_unsaved=allow_unsaved, directory=directory) + except MemNetError as exc: + _handle_error(exc) + return + sys.stdout.write(result.stdout_json()) + sys.stderr.write(result.stat + "\n") + if result.exit_code: + sys.stderr.write(f"@ERR: upgrade_blocked|unsaved|{len(result.unsaved)}\n") + raise typer.Exit(result.exit_code) + + +@admin_app.command("upgrade-retire") +def admin_upgrade_retire( + token: Annotated[ + str | None, + typer.Option("--token", help="Admin token presented by the caller (do not log)."), + ] = None, + state_dir: Annotated[str | None, typer.Option("--state-dir")] = None, +) -> None: + """Stop replaying an upgrade manifest after a verified restore. Keeps the files.""" + from pathlib import Path + + from memnet.admin_usage import authenticate, caller_token + from memnet.upgrade import retire_manifest + + presented = token if token else caller_token() + try: + authenticate(presented) + retire_manifest(Path(state_dir) if state_dir else None) + except MemNetError as exc: + _handle_error(exc) + emit_stdout('{"ok": true, "retired": true}') + emit_stderr("@STAT: upgrade_retire|1|-") + + @session_app.command("expire-status") def session_expire_status() -> None: """Booleans for expire-save (serve_status). Path redacted; no sids.""" diff --git a/parts/common/memnet/memnet/local_ipc_gateway.py b/parts/common/memnet/memnet/local_ipc_gateway.py index bbc182c5..b703dbe3 100644 --- a/parts/common/memnet/memnet/local_ipc_gateway.py +++ b/parts/common/memnet/memnet/local_ipc_gateway.py @@ -78,6 +78,9 @@ def run_ipc_serve(path: str | None = None) -> None: sock_path = resolve_ipc_path(path) os.environ["MEMNET_IPC_SOCKET"] = sock_path os.environ["MEMNET_SERVE_INTERNAL"] = "1" + from memnet.upgrade import startup_restore_or_exit + + startup_restore_or_exit() try: from memnet.cheap_llm_import_guard import maybe_install_cheap_llm_import_guard diff --git a/parts/common/memnet/memnet/serve.py b/parts/common/memnet/memnet/serve.py index ea9672d3..48582f10 100644 --- a/parts/common/memnet/memnet/serve.py +++ b/parts/common/memnet/memnet/serve.py @@ -88,6 +88,14 @@ def _handle_request(payload: dict[str, Any]) -> dict[str, Any]: token_s = admin_token if isinstance(admin_token, str) else None if payload.get("admin_usage") is True: return usage_report_envelope(token_s) + if payload.get("upgrade_prepare") is True: + from memnet.upgrade import upgrade_prepare_envelope + + return upgrade_prepare_envelope(token_s, allow_unsaved=bool(payload.get("allow_unsaved"))) + if payload.get("upgrade_retire") is True: + from memnet.upgrade import upgrade_retire_envelope + + return upgrade_retire_envelope(token_s) argv = payload.get("args", []) if not isinstance(argv, list): @@ -99,26 +107,39 @@ def _handle_request(payload: dict[str, Any]) -> dict[str, Any]: "stdout": "", "stderr": "@ERR: bad_request|stdin must be a string\n", } + from memnet.upgrade import drain_gate, draining_envelope + + counted, refuse = drain_gate.begin(argv) + if refuse: + return draining_envelope() os.environ["MEMNET_SERVE_INTERNAL"] = "1" from memnet.cli import app from memnet.output import capture_request_stdio code = 0 + stdout = "" + stderr = "" token_ctx = set_caller_token(token_s) - with capture_request_stdio(stdin_text if isinstance(stdin_text, str) else None) as (out, err): - try: - result = app(argv, prog_name="memnet", standalone_mode=False) - if isinstance(result, int) and result != 0: - code = result - except SystemExit as exc: - code = int(exc.code) if isinstance(exc.code, int) else 1 - except Exception as exc: - code = 1 - err.write(f"@ERR: internal|{type(exc).__name__}: {exc}\n") - finally: - reset_caller_token(token_ctx) - stdout = out.getvalue() - stderr = err.getvalue() + try: + with capture_request_stdio(stdin_text if isinstance(stdin_text, str) else None) as ( + out, + err, + ): + try: + result = app(argv, prog_name="memnet", standalone_mode=False) + if isinstance(result, int) and result != 0: + code = result + except SystemExit as exc: + code = int(exc.code) if isinstance(exc.code, int) else 1 + except Exception as exc: + code = 1 + err.write(f"@ERR: internal|{type(exc).__name__}: {exc}\n") + finally: + reset_caller_token(token_ctx) + stdout = out.getvalue() + stderr = err.getvalue() + finally: + drain_gate.end(counted) return {"exit_code": code, "stdout": stdout, "stderr": stderr} @@ -232,6 +253,9 @@ def run_serve(host: str | None = None, port: int | None = None) -> None: port = port or serve_port() validate_serve_bind_host(host) os.environ["MEMNET_SERVE_INTERNAL"] = "1" + from memnet.upgrade import startup_restore_or_exit + + startup_restore_or_exit() try: from memnet.admin_usage import mark_serve_start diff --git a/parts/common/memnet/memnet/session.py b/parts/common/memnet/memnet/session.py index a7df6a6d..1694d800 100644 --- a/parts/common/memnet/memnet/session.py +++ b/parts/common/memnet/memnet/session.py @@ -318,6 +318,9 @@ def open_session( caps: Caps | None = None, product: str | None = None, ) -> SessionStore: + from memnet.upgrade import refuse_new_session + + refuse_new_session() caps = caps or Caps() purge_expired(caps) if count_sessions(caps) >= caps.max_sessions: diff --git a/parts/common/memnet/memnet/snapshot.py b/parts/common/memnet/memnet/snapshot.py index 8eedc960..0a6b80ea 100644 --- a/parts/common/memnet/memnet/snapshot.py +++ b/parts/common/memnet/memnet/snapshot.py @@ -336,6 +336,7 @@ def load_snapshot( ttl_minutes: int | None = None, keep_id: bool = False, hide_path: bool = False, + preserve_clocks: bool = False, ) -> SessionStore: caps = caps or Caps() purge_expired(caps) @@ -371,6 +372,7 @@ def load_snapshot( caps=caps, ttl_minutes=ttl_minutes, keep_id=keep_id, + preserve_clocks=preserve_clocks, ) @@ -380,6 +382,7 @@ def load_snapshot_text( caps: Caps | None = None, ttl_minutes: int | None = None, keep_id: bool = False, + preserve_clocks: bool = False, ) -> SessionStore: caps = caps or Caps() meta, map_lines, rel_lines, rec_lines = _parse_sections(split_snapshot_lines(text)) @@ -395,11 +398,16 @@ def load_snapshot_text( ttl = ttl_minutes if ttl_minutes is not None else meta.ttl_minutes if ttl < 1 or ttl > 1440: raise MemNetError("bad_ttl", "ttl must be 1..1440") - expires = now + timedelta(minutes=ttl) + if preserve_clocks: + created_at = meta.created_at + expires_at = meta.expires_at + else: + created_at = now.isoformat().replace("+00:00", "Z") + expires_at = (now + timedelta(minutes=ttl)).isoformat().replace("+00:00", "Z") new_meta = SessionMeta( session_id=session_id, - created_at=now.isoformat().replace("+00:00", "Z"), - expires_at=expires.isoformat().replace("+00:00", "Z"), + created_at=created_at, + expires_at=expires_at, ttl_minutes=ttl, has_writes=meta.has_writes, modified_at=meta.modified_at, diff --git a/parts/common/memnet/memnet/upgrade.py b/parts/common/memnet/memnet/upgrade.py new file mode 100644 index 00000000..3f54ac9e --- /dev/null +++ b/parts/common/memnet/memnet/upgrade.py @@ -0,0 +1,646 @@ +"""Safe serve upgrade: drain, manifest, and startup restore (MN-REQ-06.14). + +Admin-only. Not an agent MCP tool. A session that cannot be snapshotted +blocks ready-to-stop unless the operator passes allow-unsaved. Snapshot +files are not deleted when restore fails. +""" + +from __future__ import annotations + +import hashlib +import json +import os +import threading +import time +from dataclasses import dataclass, field +from pathlib import Path +from typing import Any + +from memnet import __version__ +from memnet.acl import SessionAcl, parse_write_scope +from memnet.admin_usage import authenticate +from memnet.exceptions import MemNetError +from memnet.registry import count, get_entry, list_entries +from memnet.snapshot import SNAPSHOT_MAGIC, load_snapshot_text, write_snapshot + +SUPPORTED_SNAPSHOT_FORMATS = (1,) +MANIFEST_FORMAT = 1 +MANIFEST_NAME = "upgrade-manifest.json" +RESTORE_NAME = "upgrade-restore.json" +BLOCKED_NAME = "upgrade-blocked.json" +SNAPSHOT_DIRNAME = "upgrade-snapshots" +PASSPORT_PREFIX = "# upgrade-passport " +ENV_STATE_DIR = "MEMNET_STATE_DIR" +ENV_RETRY_AFTER = "MEMNET_UPGRADE_RETRY_AFTER_S" +ENV_DRAIN_WAIT = "MEMNET_UPGRADE_DRAIN_WAIT_S" +DEFAULT_RETRY_AFTER_S = 5 +DEFAULT_DRAIN_WAIT_S = 20.0 + +_PHASE_OPEN = "open" +_PHASE_QUIESCE = "quiesce" +_PHASE_READY = "ready" + + +def snapshot_format_version() -> int: + """Format written by this patch. Older v1 snapshots still load.""" + if not SNAPSHOT_MAGIC.endswith("v1"): + raise MemNetError("bad_snapshot", "snapshot magic is not v1") + return 1 + + +def state_dir() -> Path: + raw = (os.environ.get(ENV_STATE_DIR) or "").strip() + if raw: + return Path(raw) + return Path.home() / ".local" / "state" / "memnet" + + +def _retry_after_s() -> int: + raw = (os.environ.get(ENV_RETRY_AFTER) or "").strip() + if not raw: + return DEFAULT_RETRY_AFTER_S + try: + value = int(raw) + except ValueError: + return DEFAULT_RETRY_AFTER_S + return value if value > 0 else DEFAULT_RETRY_AFTER_S + + +def _drain_wait_s() -> float: + raw = (os.environ.get(ENV_DRAIN_WAIT) or "").strip() + if not raw: + return DEFAULT_DRAIN_WAIT_S + try: + value = float(raw) + except ValueError: + return DEFAULT_DRAIN_WAIT_S + return value if value > 0 else DEFAULT_DRAIN_WAIT_S + + +def _is_upgrade_argv(argv: list[Any]) -> bool: + if len(argv) < 2: + return False + return str(argv[0]) == "admin" and str(argv[1]) in { + "upgrade-prepare", + "upgrade-retire", + } + + +class DrainGate: + """Process-wide drain. New work is refused once quiesce begins.""" + + def __init__(self) -> None: + self._cv = threading.Condition(threading.Lock()) + self._phase = _PHASE_OPEN + self._inflight = 0 + self.retry_after_s = DEFAULT_RETRY_AFTER_S + + def reset(self) -> None: + with self._cv: + self._phase = _PHASE_OPEN + self._inflight = 0 + self._cv.notify_all() + + @property + def phase(self) -> str: + with self._cv: + return self._phase + + def set_phase_for_test(self, phase: str) -> None: + with self._cv: + self._phase = phase + self._cv.notify_all() + + def begin(self, argv: list[Any]) -> tuple[bool, str | None]: + """Return (counted, error). Upgrade commands are not counted.""" + with self._cv: + upgrade = _is_upgrade_argv(argv) + if self._phase != _PHASE_OPEN and not upgrade: + return False, "serve_draining" + if upgrade: + return False, None + self._inflight += 1 + return True, None + + def end(self, counted: bool) -> None: + if not counted: + return + with self._cv: + self._inflight = max(0, self._inflight - 1) + if self._inflight == 0: + self._cv.notify_all() + + def enter_quiesce(self, retry_after_s: int) -> None: + with self._cv: + self._phase = _PHASE_QUIESCE + self.retry_after_s = retry_after_s + + def mark_ready(self) -> None: + with self._cv: + self._phase = _PHASE_READY + + def mark_open(self) -> None: + with self._cv: + self._phase = _PHASE_OPEN + self._cv.notify_all() + + def wait_idle(self, timeout_s: float) -> bool: + deadline = time.monotonic() + timeout_s + with self._cv: + while self._inflight > 0: + remaining = deadline - time.monotonic() + if remaining <= 0: + return False + self._cv.wait(remaining) + return True + + +drain_gate = DrainGate() + + +def reset_drain_gate() -> None: + drain_gate.reset() + + +def refuse_new_session() -> None: + """Refuse session_open while the serve is draining or ready to stop.""" + if drain_gate.phase == _PHASE_OPEN: + return + raise MemNetError( + "serve_draining", + f"retry_after_s={drain_gate.retry_after_s}", + exit_code=2, + ) + + +def draining_envelope(retry_after_s: int | None = None) -> dict[str, Any]: + seconds = drain_gate.retry_after_s if retry_after_s is None else retry_after_s + return { + "exit_code": 2, + "stdout": "", + "stderr": f"@ERR: serve_draining|retry_after_s={seconds}\n", + } + + +def _safe_name(session_id: str) -> str: + cleaned = "".join(ch if ch.isalnum() or ch in "._-" else "_" for ch in session_id) + return cleaned or "session" + + +def _edge_count(store: Any) -> int: + total = 0 + for hid in store.write_order: + rec = store._by_hid.get(hid) + if rec is not None and rec.tag == "EDG": + total += 1 + return total + + +def _house_tags(tag_map: Any) -> list[str]: + return sorted(tag_map.tags.keys()) + + +def _acl_to_dict(acl: SessionAcl | None) -> dict[str, Any]: + if acl is None: + return {"enabled": False, "callers": [], "bind": None} + callers = [] + for grant in acl.callers.values(): + scope = grant.write_scope.to_wire() if grant.write_scope else None + callers.append( + { + "caller": grant.caller, + "can_pin_map": grant.can_pin_map, + "can_mutate": grant.can_mutate, + "write_scope": scope, + } + ) + bind = None + if acl.bind is not None: + bind = {"mission_id": acl.bind.mission_id, "lease": acl.bind.lease} + return {"enabled": acl.enabled, "callers": callers, "bind": bind} + + +def _apply_acl(entry: Any, data: dict[str, Any] | None) -> None: + if not data: + return + acl = SessionAcl(enabled=bool(data.get("enabled"))) + for row in data.get("callers") or []: + acl.grant( + str(row.get("caller") or ""), + can_pin_map=bool(row.get("can_pin_map", True)), + can_mutate=bool(row.get("can_mutate", True)), + write_scope=parse_write_scope(row.get("write_scope")), + ) + bind = data.get("bind") + if isinstance(bind, dict) and bind.get("mission_id") and bind.get("lease"): + acl.set_bind(str(bind["mission_id"]), str(bind["lease"])) + if data.get("enabled"): + acl.enabled = True + entry.acl = acl + entry.meta.acl_enabled = acl.enabled + + +def _sha256(path: Path) -> str: + digest = hashlib.sha256() + digest.update(path.read_bytes()) + return digest.hexdigest() + + +def _atomic_write(path: Path, text: str) -> None: + path.parent.mkdir(parents=True, exist_ok=True) + tmp = path.with_suffix(path.suffix + ".tmp") + tmp.write_text(text, encoding="utf-8") + tmp.replace(path) + + +def _passport(entry: Any) -> dict[str, Any]: + meta = entry.meta + return { + "acl": _acl_to_dict(entry.acl), + "product": meta.product, + "created_at": meta.created_at, + "expires_at": meta.expires_at, + "ttl_minutes": meta.ttl_minutes, + "house_tags": _house_tags(entry.tag_map), + } + + +def _append_passport(path: Path, payload: dict[str, Any]) -> None: + line = PASSPORT_PREFIX + json.dumps(payload, separators=(",", ":"), sort_keys=True) + with path.open("a", encoding="utf-8") as handle: + handle.write(line + "\n") + + +def split_passport(text: str) -> tuple[str, dict[str, Any] | None]: + """Split an optional upgrade passport. v1 loaders skip the # line.""" + body: list[str] = [] + passport: dict[str, Any] | None = None + for line in text.splitlines(): + if line.startswith(PASSPORT_PREFIX): + passport = json.loads(line[len(PASSPORT_PREFIX) :]) + continue + body.append(line) + return "\n".join(body) + ("\n" if body else ""), passport + + +@dataclass +class PrepareResult: + ready_to_stop: bool + saved: int + unsaved: list[dict[str, str]] + manifest_path: Path | None + stat: str + exit_code: int + + def stdout_json(self) -> str: + payload = { + "ok": self.ready_to_stop, + "ready_to_stop": self.ready_to_stop, + "saved": self.saved, + "unsaved": self.unsaved, + "stat": self.stat, + "manifest": str(self.manifest_path) if self.manifest_path else None, + } + return json.dumps(payload, sort_keys=True) + "\n" + + +def _snapshot_one(entry: Any, snap_dir: Path) -> dict[str, Any]: + from memnet.session import SessionStore + + sid = entry.meta.session_id + path = snap_dir / f"{_safe_name(sid)}.snap" + with entry.lock: + rows = entry.store.row_count_non_law() + edges = _edge_count(entry.store) + house = _house_tags(entry.tag_map) + passport = _passport(entry) + write_snapshot(SessionStore(sid, entry.store.caps), path) + _append_passport(path, passport) + checksum = _sha256(path) + return { + "session_id": sid, + "file": f"{SNAPSHOT_DIRNAME}/{path.name}", + "rows": rows, + "edges": edges, + "checksum_sha256": checksum, + "expires_at": entry.meta.expires_at, + "created_at": entry.meta.created_at, + "ttl_minutes": entry.meta.ttl_minutes, + "house_tags": house, + "product": entry.meta.product, + "saved": True, + } + + +def prepare_upgrade( + presented_token: str | None, + *, + allow_unsaved: bool = False, + directory: Path | None = None, + drain_wait_s: float | None = None, +) -> PrepareResult: + """Quiesce, snapshot every loaded session, and write the manifest. + + On any unsaved session, ready_to_stop stays false unless allow_unsaved. + The serve returns to accepting work when the drain is not ready to stop. + """ + authenticate(presented_token) + root = directory or state_dir() + retry_after = _retry_after_s() + wait_s = DEFAULT_DRAIN_WAIT_S if drain_wait_s is None else drain_wait_s + drain_gate.enter_quiesce(retry_after) + if not drain_gate.wait_idle(wait_s): + drain_gate.mark_open() + raise MemNetError("upgrade_inflight", "in-flight commands still running", exit_code=2) + snap_dir = root / SNAPSHOT_DIRNAME + snap_dir.mkdir(parents=True, exist_ok=True) + saved_rows: list[dict[str, Any]] = [] + unsaved: list[dict[str, str]] = [] + for entry in list_entries(): + sid = entry.meta.session_id + try: + saved_rows.append(_snapshot_one(entry, snap_dir)) + except MemNetError as exc: + unsaved.append({"session_id": sid, "code": exc.code, "detail": exc.message}) + except OSError as exc: + unsaved.append( + {"session_id": sid, "code": "snapshot_io_error", "detail": type(exc).__name__} + ) + blocked = bool(unsaved) and not allow_unsaved + if blocked: + drain_gate.mark_open() + report = { + "manifest_format": MANIFEST_FORMAT, + "snapshot_format": snapshot_format_version(), + "serve_version": __version__, + "ready_to_stop": False, + "allow_unsaved": False, + "saved": saved_rows, + "unsaved": unsaved, + } + path = root / BLOCKED_NAME + _atomic_write(path, json.dumps(report, indent=2, sort_keys=True) + "\n") + stat = f"@STAT: upgrade_prepare|blocked|{len(saved_rows)}|unsaved|{len(unsaved)}" + return PrepareResult( + ready_to_stop=False, + saved=len(saved_rows), + unsaved=unsaved, + manifest_path=path, + stat=stat, + exit_code=2, + ) + manifest = { + "manifest_format": MANIFEST_FORMAT, + "snapshot_format": snapshot_format_version(), + "serve_version": __version__, + "ready_to_stop": True, + "retired": False, + "allow_unsaved": bool(allow_unsaved), + "sessions": saved_rows, + "unsaved": unsaved, + } + path = root / MANIFEST_NAME + _atomic_write(path, json.dumps(manifest, indent=2, sort_keys=True) + "\n") + drain_gate.mark_ready() + stat = f"@STAT: upgrade_prepare|ready|{len(saved_rows)}|unsaved|{len(unsaved)}" + return PrepareResult( + ready_to_stop=True, + saved=len(saved_rows), + unsaved=unsaved, + manifest_path=path, + stat=stat, + exit_code=0, + ) + + +@dataclass +class RestoreReport: + ok: int = 0 + failed: int = 0 + skipped: bool = False + errors: list[dict[str, str]] = field(default_factory=list) + session_ids: list[str] = field(default_factory=list) + details: list[dict[str, Any]] = field(default_factory=list) + + @property + def stat(self) -> str: + return f"@STAT: upgrade_restore|ok|{self.ok}|failed|{self.failed}" + + def as_dict(self) -> dict[str, Any]: + return { + "ok": self.ok, + "failed": self.failed, + "skipped": self.skipped, + "stat": self.stat, + "errors": self.errors, + "session_ids": self.session_ids, + "details": self.details, + } + + +def _load_json(path: Path) -> dict[str, Any]: + try: + data = json.loads(path.read_text(encoding="utf-8")) + except (OSError, json.JSONDecodeError) as exc: + raise MemNetError("upgrade_manifest", type(exc).__name__, exit_code=3) from exc + if not isinstance(data, dict): + raise MemNetError("upgrade_manifest", "manifest must be an object", exit_code=3) + return data + + +def _fail_session(report: RestoreReport, sid: str, code: str, detail: str) -> None: + report.failed += 1 + report.errors.append({"session_id": sid, "code": code, "detail": detail}) + + +def _restore_one(root: Path, row: dict[str, Any], report: RestoreReport) -> None: + sid = str(row.get("session_id") or "") + rel = str(row.get("file") or "") + path = root / rel + before = path.read_bytes() if path.is_file() else None + try: + if before is None: + _fail_session(report, sid, "snapshot_not_found", "missing snapshot file") + return + got = hashlib.sha256(before).hexdigest() + want = str(row.get("checksum_sha256") or "") + if got != want: + _fail_session(report, sid, "upgrade_checksum", "checksum mismatch") + return + text = before.decode("utf-8") + body, passport = split_passport(text) + if not body.startswith(SNAPSHOT_MAGIC): + _fail_session(report, sid, "bad_snapshot", "unsupported snapshot header") + return + from memnet.config import Caps + + ss = load_snapshot_text(body, caps=Caps(), keep_id=True, preserve_clocks=True) + entry = get_entry(ss.session_id) + if entry is None: + _fail_session(report, sid, "upgrade_restore", "session missing after load") + return + if passport: + if passport.get("expires_at") and passport.get("expires_at") != entry.meta.expires_at: + _fail_session(report, sid, "upgrade_clock", "ttl clock mismatch") + return + _apply_acl(entry, passport.get("acl")) + entry.meta.product = passport.get("product") + if ss.session_id != sid: + _fail_session(report, sid, "upgrade_session_id", "session id changed") + return + rows = entry.store.row_count_non_law() + edges = _edge_count(entry.store) + if rows != int(row.get("rows", -1)) or edges != int(row.get("edges", -1)): + _fail_session(report, sid, "upgrade_counts", "row or edge count mismatch") + return + house = _house_tags(entry.tag_map) + if house != list(row.get("house_tags") or []): + _fail_session(report, sid, "upgrade_house", "house tag map mismatch") + return + if path.read_bytes() != before: + _fail_session(report, sid, "upgrade_snapshot_mutated", "snapshot file changed") + return + report.ok += 1 + report.session_ids.append(ss.session_id) + bind = entry.acl.bind.mission_id if entry.acl and entry.acl.bind else None + report.details.append( + { + "session_id": ss.session_id, + "expires_at": entry.meta.expires_at, + "ttl_minutes": entry.meta.ttl_minutes, + "acl_enabled": bool(entry.acl and entry.acl.enabled), + "bind_mission": bind, + "rows": rows, + "edges": edges, + "product": entry.meta.product, + } + ) + except MemNetError as exc: + _fail_session(report, sid, exc.code, exc.message) + finally: + if before is not None and path.is_file() and path.read_bytes() != before: + path.write_bytes(before) + + +def restore_manifest(directory: Path | None = None) -> RestoreReport: + """Reload a ready manifest. Does not delete snapshot files.""" + root = directory or state_dir() + path = root / MANIFEST_NAME + if not path.is_file(): + return RestoreReport(skipped=True) + manifest = _load_json(path) + if manifest.get("retired"): + return RestoreReport(skipped=True) + if not manifest.get("ready_to_stop"): + raise MemNetError("upgrade_not_ready", "manifest is not ready_to_stop", exit_code=3) + fmt = int(manifest.get("snapshot_format") or 0) + if fmt not in SUPPORTED_SNAPSHOT_FORMATS: + raise MemNetError( + "upgrade_snapshot_format", + f"unsupported snapshot format {fmt}", + exit_code=3, + ) + if count() > 0: + raise MemNetError("upgrade_restore_registry_busy", "registry is not empty", exit_code=3) + report = RestoreReport() + sessions = manifest.get("sessions") or [] + if not isinstance(sessions, list): + raise MemNetError("upgrade_manifest", "sessions must be a list", exit_code=3) + for row in sessions: + if not isinstance(row, dict): + report.failed += 1 + report.errors.append( + {"session_id": "", "code": "upgrade_manifest", "detail": "bad row"} + ) + continue + _restore_one(root, row, report) + _atomic_write( + root / RESTORE_NAME, + json.dumps(report.as_dict(), indent=2, sort_keys=True) + "\n", + ) + return report + + +def startup_restore_or_exit(directory: Path | None = None) -> RestoreReport: + """Serve startup. Exit 3 on a loud restore failure. Leave files in place.""" + try: + report = restore_manifest(directory) + except MemNetError as exc: + line = f"@ERR: {exc.code}|{exc.message}\n" + import sys + + sys.stdout.write(line) + sys.stderr.write(line) + raise SystemExit(exc.exit_code) from exc + if report.skipped: + return report + import sys + + line = report.stat + "\n" + sys.stdout.write(line) + sys.stderr.write(line) + if report.failed: + raise SystemExit(3) + return report + + +def retire_manifest(directory: Path | None = None) -> None: + """Keep snapshot files, and stop replaying them on the next start. + + Refuses when the restore report is missing or not clean. + """ + root = directory or state_dir() + report_path = root / RESTORE_NAME + if not report_path.is_file(): + raise MemNetError("upgrade_retire_refused", "restore not verified", exit_code=2) + report = _load_json(report_path) + if int(report.get("failed") or 0) != 0 or int(report.get("ok") or 0) < 0: + raise MemNetError("upgrade_retire_refused", "restore not verified", exit_code=2) + if int(report.get("failed") or 0) != 0: + raise MemNetError("upgrade_retire_refused", "restore not verified", exit_code=2) + manifest_path = root / MANIFEST_NAME + if not manifest_path.is_file(): + raise MemNetError("upgrade_manifest_missing", "no manifest", exit_code=2) + manifest = _load_json(manifest_path) + if int(report.get("failed") or 0) != 0 or report.get("skipped"): + raise MemNetError("upgrade_retire_refused", "restore not verified", exit_code=2) + manifest["retired"] = True + _atomic_write(manifest_path, json.dumps(manifest, indent=2, sort_keys=True) + "\n") + + +def upgrade_prepare_envelope( + token: str | None, + *, + allow_unsaved: bool = False, +) -> dict[str, Any]: + try: + result = prepare_upgrade(token, allow_unsaved=allow_unsaved) + except MemNetError as exc: + return { + "exit_code": exc.exit_code, + "stdout": "", + "stderr": f"@ERR: {exc.code}|{exc.message}\n", + } + stderr = "" + if result.exit_code: + stderr = f"@ERR: upgrade_blocked|unsaved|{len(result.unsaved)}\n{result.stat}\n" + else: + stderr = result.stat + "\n" + return {"exit_code": result.exit_code, "stdout": result.stdout_json(), "stderr": stderr} + + +def upgrade_retire_envelope(token: str | None) -> dict[str, Any]: + try: + authenticate(token) + retire_manifest() + except MemNetError as exc: + return { + "exit_code": exc.exit_code, + "stdout": "", + "stderr": f"@ERR: {exc.code}|{exc.message}\n", + } + return { + "exit_code": 0, + "stdout": json.dumps({"ok": True, "retired": True}) + "\n", + "stderr": "@STAT: upgrade_retire|1|-\n", + } diff --git a/parts/common/memnet/memnet/upgrade_retry.py b/parts/common/memnet/memnet/upgrade_retry.py new file mode 100644 index 00000000..36dd5e49 --- /dev/null +++ b/parts/common/memnet/memnet/upgrade_retry.py @@ -0,0 +1,67 @@ +"""Bounded retry while a serve is draining or briefly refusing connections. + +Default window is 30 seconds (``MEMNET_UPGRADE_RETRY_S``). A window of 0 +tries once. Timeout is not retried: the serve may already have accepted +the command. ``serve_draining`` is retried because the serve refused it. +""" + +from __future__ import annotations + +import os +import re +import time +from collections.abc import Callable +from typing import Any + +ENV_RETRY_S = "MEMNET_UPGRADE_RETRY_S" +DEFAULT_RETRY_S = 30.0 +_RETRY_AFTER = re.compile(r"retry_after_s=(\d+(?:\.\d+)?)") +_RETRY_EXC = (ConnectionRefusedError, ConnectionResetError, ConnectionAbortedError, BrokenPipeError) + + +def retry_window_s() -> float: + raw = os.environ.get(ENV_RETRY_S) + if raw is None or raw.strip() == "": + return DEFAULT_RETRY_S + try: + return max(0.0, float(raw)) + except ValueError: + return DEFAULT_RETRY_S + + +def _retry_after(stderr: str) -> float: + match = _RETRY_AFTER.search(stderr or "") + if not match: + return 0.05 + return max(0.0, float(match.group(1))) + + +def call_with_upgrade_retry(fn: Callable[[], dict[str, Any]]) -> dict[str, Any]: + """Call ``fn`` until it returns or the retry window is spent. + + ``fn`` raises a connection-refusal error, or returns an envelope whose + stderr contains ``serve_draining``. + """ + window = retry_window_s() + deadline = time.monotonic() + window + delay = 0.05 + while True: + try: + raw = fn() + except _RETRY_EXC: + if window <= 0 or time.monotonic() >= deadline: + raise + remaining = deadline - time.monotonic() + time.sleep(min(delay, remaining)) + delay = min(delay * 2, 2.0) + continue + stderr = "" + if isinstance(raw, dict): + stderr = str(raw.get("stderr") or "") + if "serve_draining" in stderr and window > 0 and time.monotonic() < deadline: + remaining = deadline - time.monotonic() + pause = min(_retry_after(stderr), delay, remaining) + time.sleep(pause) + delay = min(delay * 2, 2.0) + continue + return raw diff --git a/parts/common/memnet/memnet/upgrade_run.py b/parts/common/memnet/memnet/upgrade_run.py new file mode 100644 index 00000000..4548edbb --- /dev/null +++ b/parts/common/memnet/memnet/upgrade_run.py @@ -0,0 +1,258 @@ +"""Operator helper for a side-by-side venv upgrade (MN-REQ-06.14). + +Order is fixed: the caller must confirm memnet-mcp and the gateway already +retry ``serve_draining``. This process then preflights the new interpreter, +drains, swaps ExecStart, updates the gateway pin, starts the new serve, and +rolls back to the unit backup if verification fails. + +No session ids are printed. Tokens are read from the environment. +""" + +from __future__ import annotations + +import argparse +import json +import os +import subprocess +import sys +import time +from collections.abc import Callable +from pathlib import Path +from typing import Any + +from memnet.exceptions import MemNetError +from memnet.upgrade import ( + MANIFEST_NAME, + RESTORE_NAME, + SUPPORTED_SNAPSHOT_FORMATS, + retire_manifest, +) + +_PROBE = ( + "import json,memnet\n" + "from memnet.upgrade import SUPPORTED_SNAPSHOT_FORMATS\n" + "print(json.dumps({'version': memnet.__version__, " + "'formats': list(SUPPORTED_SNAPSHOT_FORMATS)}))\n" +) + + +class UpgradeRunError(RuntimeError): + def __init__(self, code: str, message: str) -> None: + self.code = code + super().__init__(message) + + +def swap_execstart(text: str, new_exec: str) -> str: + """Replace every ExecStart= line. The previous text is the rollback copy.""" + lines = text.splitlines(keepends=True) + found = False + out: list[str] = [] + for line in lines: + if line.startswith("ExecStart="): + nl = "\n" if line.endswith("\n") else "" + out.append(f"ExecStart={new_exec}{nl}") + found = True + else: + out.append(line) + if not found: + raise UpgradeRunError("upgrade_unit", "unit has no ExecStart") + return "".join(out) + + +def set_pinned_version(config_text: str, product: str, version: str) -> str: + data = json.loads(config_text) + products = data.get("products") + if not isinstance(products, dict) or product not in products: + raise UpgradeRunError("upgrade_pin", "product is not in the gateway config") + row = products[product] + if not isinstance(row, dict): + raise UpgradeRunError("upgrade_pin", "product entry must be an object") + row["pinned_version"] = version + return json.dumps(data, indent=2) + "\n" + + +def preflight_new_python(python: str, directory: Path) -> dict[str, Any]: + """Import the new version and check snapshot format support.""" + proc = subprocess.run( + [python, "-c", _PROBE], + check=False, + capture_output=True, + text=True, + ) + if proc.returncode != 0: + raise UpgradeRunError("upgrade_preflight", "new interpreter failed to import memnet") + try: + info = json.loads(proc.stdout.strip().splitlines()[-1]) + except (json.JSONDecodeError, IndexError) as exc: + raise UpgradeRunError("upgrade_preflight", "new interpreter probe was not JSON") from exc + formats = set(info.get("formats") or []) + manifest_path = directory / MANIFEST_NAME + if manifest_path.is_file(): + manifest = json.loads(manifest_path.read_text(encoding="utf-8")) + fmt = int(manifest.get("snapshot_format") or 0) + if fmt not in formats: + raise UpgradeRunError( + "upgrade_snapshot_format", + f"new version cannot load snapshot format {fmt}", + ) + elif snapshot_local_format() not in formats: + raise UpgradeRunError( + "upgrade_snapshot_format", + "new version cannot load the current snapshot format", + ) + return info + + +def snapshot_local_format() -> int: + return int(SUPPORTED_SNAPSHOT_FORMATS[0]) + + +def _read(path: Path) -> str: + return path.read_text(encoding="utf-8") + + +def _write(path: Path, text: str) -> None: + path.write_text(text, encoding="utf-8") + + +def run_upgrade( + *, + new_python: str, + new_exec: str, + state: Path, + clients_ready: bool, + unit_path: Path | None = None, + gateway_config: Path | None = None, + product: str | None = None, + admin_token: str | None = None, + allow_unsaved: bool = False, + restarter: Callable[[], None] | None = None, + drainer: Callable[[], Any] | None = None, + verify_wait_s: float = 30.0, +) -> dict[str, Any]: + """Drain, swap, verify. Roll back the unit and pin if verify fails.""" + if not clients_ready: + raise UpgradeRunError( + "upgrade_clients", + "deploy memnet-mcp and gateway retry before the serve swap", + ) + info = preflight_new_python(new_python, state) + token = admin_token if admin_token is not None else os.environ.get("MEMNET_ADMIN_TOKEN") + if drainer is None: + _drain_via_serve(state, token, allow_unsaved=allow_unsaved) + else: + drainer() + unit_backup: str | None = None + pin_backup: str | None = None + if unit_path is not None: + unit_backup = _read(unit_path) + backup_path = unit_path.with_suffix(unit_path.suffix + ".bak") + _write(backup_path, unit_backup) + _write(unit_path, swap_execstart(unit_backup, new_exec)) + if gateway_config is not None: + if not product: + raise UpgradeRunError("upgrade_pin", "product is required to update the pin") + pin_backup = _read(gateway_config) + pin_path = gateway_config.with_suffix(gateway_config.suffix + ".bak") + _write(pin_path, pin_backup) + _write( + gateway_config, + set_pinned_version(pin_backup, product, str(info["version"])), + ) + if restarter is not None: + restarter() + report = _wait_restore(state, verify_wait_s) + if report is None or int(report.get("failed") or 0) != 0 or report.get("skipped"): + _rollback(unit_path, unit_backup, gateway_config, pin_backup, restarter) + raise UpgradeRunError("upgrade_verify", "restore did not verify; rolled back") + retire_manifest(state) + return {"ok": True, "version": info["version"], "restored": int(report.get("ok") or 0)} + + +def _drain_via_serve(state: Path, token: str | None, *, allow_unsaved: bool) -> None: + from memnet.serve import send_command + + args = ["admin", "upgrade-prepare", "--state-dir", str(state)] + if allow_unsaved: + args.append("--allow-unsaved") + raw = send_command(args, admin_token=token, timeout=120.0) + if int(raw.get("exit_code") or 1) != 0: + raise UpgradeRunError("upgrade_drain", "drain did not reach ready_to_stop") + + +def _wait_restore(state: Path, timeout_s: float) -> dict[str, Any] | None: + deadline = time.monotonic() + timeout_s + path = state / RESTORE_NAME + while time.monotonic() < deadline: + if path.is_file(): + try: + data = json.loads(path.read_text(encoding="utf-8")) + except json.JSONDecodeError: + time.sleep(0.05) + continue + if isinstance(data, dict): + return data + time.sleep(0.05) + return None + + +def _rollback( + unit_path: Path | None, + unit_backup: str | None, + gateway_config: Path | None, + pin_backup: str | None, + restarter: Callable[[], None] | None, +) -> None: + if unit_path is not None and unit_backup is not None: + _write(unit_path, unit_backup) + if gateway_config is not None and pin_backup is not None: + _write(gateway_config, pin_backup) + if restarter is not None: + restarter() + + +def main(argv: list[str] | None = None) -> int: + parser = argparse.ArgumentParser(prog="memnet-upgrade") + parser.add_argument("--new-python", required=True) + parser.add_argument("--new-exec", required=True, help="New unit ExecStart value") + parser.add_argument("--state-dir", required=True) + parser.add_argument("--unit", default="") + parser.add_argument("--gateway-config", default="") + parser.add_argument("--product", default="") + parser.add_argument("--allow-unsaved", action="store_true") + parser.add_argument( + "--clients-ready", + action="store_true", + help="Confirm memnet-mcp and the gateway already retry serve_draining", + ) + parser.add_argument("--restart-unit", default="", help="systemctl unit name to restart") + args = parser.parse_args(argv) + restarter = None + if args.restart_unit: + unit_name = args.restart_unit + + def restarter() -> None: + subprocess.run(["systemctl", "restart", unit_name], check=True) + + try: + result = run_upgrade( + new_python=args.new_python, + new_exec=args.new_exec, + state=Path(args.state_dir), + clients_ready=args.clients_ready, + unit_path=Path(args.unit) if args.unit else None, + gateway_config=Path(args.gateway_config) if args.gateway_config else None, + product=args.product or None, + allow_unsaved=args.allow_unsaved, + restarter=restarter, + ) + except (UpgradeRunError, MemNetError) as exc: + code = getattr(exc, "code", "upgrade") + sys.stderr.write(f"@ERR: {code}|{exc}\n") + return 2 + sys.stdout.write(json.dumps({"ok": True, "restored": result["restored"]}) + "\n") + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/parts/memnet-mcp/software/memnet_mcp/client.py b/parts/memnet-mcp/software/memnet_mcp/client.py index 32e0aa11..dacd1838 100644 --- a/parts/memnet-mcp/software/memnet_mcp/client.py +++ b/parts/memnet-mcp/software/memnet_mcp/client.py @@ -131,10 +131,18 @@ def run_memnet( mode = _transport() if mode == "tcp": - if probe(): - raw = send_command(full_argv, stdin=stdin) - return MemNetResponse.from_raw(raw, session_hint=session) - return MemNetResponse.serve_required(session_hint=session) + from memnet.upgrade_retry import call_with_upgrade_retry + + def _once() -> dict: + if not probe(): + raise ConnectionRefusedError("serve down") + return send_command(full_argv, stdin=stdin) + + try: + raw = call_with_upgrade_retry(_once) + except (ConnectionError, OSError, TimeoutError): + return MemNetResponse.serve_required(session_hint=session) + return MemNetResponse.from_raw(raw, session_hint=session) raw = run_argv(full_argv, stdin=stdin) return MemNetResponse.from_raw(raw, session_hint=session) diff --git a/parts/memnet-mcp/software/memnet_mcp/product_gateway.py b/parts/memnet-mcp/software/memnet_mcp/product_gateway.py index 6595d68c..369e793c 100644 --- a/parts/memnet-mcp/software/memnet_mcp/product_gateway.py +++ b/parts/memnet-mcp/software/memnet_mcp/product_gateway.py @@ -526,29 +526,24 @@ def _version(self, backend: Backend) -> str | None: ): return hit[1] version: str | None = None - if probe(host=backend.host, port=backend.port): - try: - raw = send_command( - ["version"], - host=backend.host, - port=backend.port, - timeout=self.timeout_s, - ) - except (OSError, json.JSONDecodeError): - raw = None - if isinstance(raw, dict): - for line in (raw.get("stdout") or "").splitlines(): - if line.startswith("@VER: memnet|"): - version = line.split("|", 1)[1].strip() - break + raw = self._command_with_retry(backend, ["version"], None) + if isinstance(raw, dict): + for line in (raw.get("stdout") or "").splitlines(): + if line.startswith("@VER: memnet|"): + version = line.split("|", 1)[1].strip() + break with self._lock: self._ver_cache[backend.id] = (now, version) return version - def _forward( + def _command_with_retry( self, backend: Backend, argv: list[str], stdin: str | None ) -> dict[str, Any] | None: - try: + from memnet.upgrade_retry import call_with_upgrade_retry + + def _once() -> dict[str, Any]: + if not probe(host=backend.host, port=backend.port): + raise ConnectionRefusedError("backend down") raw = send_command( list(argv), stdin=stdin, @@ -556,9 +551,20 @@ def _forward( port=backend.port, timeout=self.timeout_s, ) - except (OSError, json.JSONDecodeError): + if not isinstance(raw, dict): + raise ConnectionError("bad envelope") + return raw + + try: + return call_with_upgrade_retry(_once) + except (OSError, json.JSONDecodeError, ConnectionError): return None - if not isinstance(raw, dict): + + def _forward( + self, backend: Backend, argv: list[str], stdin: str | None + ) -> dict[str, Any] | None: + raw = self._command_with_retry(backend, argv, stdin) + if raw is None: return None err = (raw.get("stderr") or "").strip() if err == _CLIENT_TIMEOUT and not (raw.get("stdout") or "").strip(): diff --git a/parts/memnet-mcp/software/memnet_mcp/serve_bridge.py b/parts/memnet-mcp/software/memnet_mcp/serve_bridge.py index 8cd20252..5e74fe5d 100644 --- a/parts/memnet-mcp/software/memnet_mcp/serve_bridge.py +++ b/parts/memnet-mcp/software/memnet_mcp/serve_bridge.py @@ -12,4 +12,11 @@ def probe(self) -> bool: return probe() def send(self, argv: list[str], *, stdin: str | None = None) -> dict: - return send_command(argv, stdin=stdin) + from memnet.upgrade_retry import call_with_upgrade_retry + + def _once() -> dict: + if not probe(): + raise ConnectionRefusedError("serve down") + return send_command(argv, stdin=stdin) + + return call_with_upgrade_retry(_once) diff --git a/pyproject.toml b/pyproject.toml index 3a08449c..c1d49334 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -56,6 +56,7 @@ agensgraph = [ [project.scripts] memnet = "memnet.cli:main" memnet-mcp = "memnet_mcp.server:main" +memnet-upgrade = "memnet.upgrade_run:main" [tool.hatch.version] path = "parts/common/memnet/memnet/__init__.py" diff --git a/sysml-models/models/deploy.sysml b/sysml-models/models/deploy.sysml index 0d6d97b1..757daf6b 100644 --- a/sysml-models/models/deploy.sysml +++ b/sysml-models/models/deploy.sysml @@ -1308,6 +1308,59 @@ package MemNet { attribute onAgentMcp : Boolean = false; } + part def CmdAdminUpgradePrepare { + doc /* + CLI/serve: memnet admin upgrade-prepare. Optional envelope + {upgrade_prepare:true, admin_token}. NOT an agent MCP tool. + Allow-unsaved is an explicit operator override. + */ + attribute requiresStoreKeyId : Boolean = false; + attribute implemented : Boolean = true; + attribute productCommand : Boolean = false; + attribute onAgentMcp : Boolean = false; + attribute allowUnsavedOverride : Boolean = true; + } + + part def SafeServeUpgrade { + doc /* + Drain, manifest, and startup restore for a live memnet-serve + (MN-REQ-06.14). Lossless snapshot path. No silent session + loss. Older patch snapshot format v1 loads. Corrupt restore + fails loud and leaves snapshot files in place. Client retry + is on memnet-mcp and the product gateway, deployed first. + */ + attribute implemented : Boolean = true; + attribute noSemVerBump : Boolean = true; + attribute onAgentMcp : Boolean = false; + attribute adminCredEnv : String = "MEMNET_ADMIN_TOKEN"; + attribute refusesNewSessionOpen : Boolean = true; + attribute errDraining : String = "serve_draining"; + attribute finishesInFlight : Boolean = true; + attribute usesLosslessSnapshot : Boolean = true; + attribute manifestBeforeStop : Boolean = true; + attribute blocksOnUnsaveable : Boolean = true; + attribute allowUnsavedOverride : Boolean = true; + attribute silentSessionLoss : Boolean = false; + attribute restoreVerifiesChecksum : Boolean = true; + attribute restoreVerifiesCounts : Boolean = true; + attribute keepsFilesUntilVerified : Boolean = true; + attribute olderPatchSnapshotLoads : Boolean = true; + attribute corruptFailsLoud : Boolean = true; + attribute corruptDeletesSnapshots : Boolean = false; + attribute preservesSessionId : Boolean = true; + attribute preservesAcl : Boolean = true; + attribute preservesTtlClock : Boolean = true; + attribute preservesHouse : Boolean = true; + attribute retryDefaultS : Integer = 30; + attribute clientsBeforeServe : Boolean = true; + attribute helperRollsBack : Boolean = true; + attribute statUpgradeRestore : String = "upgrade_restore"; + + satisfy MN_REQ_00_MissionBridge::MN_REQ_06_ProcessBoundary::MN_REQ_06_14_SafeServeUpgrade; + satisfy MN_REQ_00_MissionBridge::MN_REQ_01_SessionLifecycle::MN_REQ_01_9_SnapshotLosslessRoundTrip; + satisfy MN_REQ_00_MissionBridge::MN_REQ_07_McpAgentBoundary::MN_REQ_07_3_NoDomainProductTools; + } + part def TcpServeBridge { doc /* Migration / multi-process fallback: length-prefixed JSON over TCP @@ -1316,6 +1369,7 @@ package MemNet { usageLookOut = serve_status (up/host/port + expire-save booleans; path redacted) for MemNetUsageDashboard. adminUsage = MN-REQ-06.11 admin JSON (opaque aliases; not 06.5). + safeUpgrade = MN-REQ-06.14 drain, manifest, and startup restore. Many mn_* sessions in one registry. Not a projectId map. Not a second graph. Admin report MUST NOT emit those mn_* ids. */ @@ -1325,6 +1379,7 @@ package MemNet { port adminCmdIn : AdminUsageCommandInPort; port adminReportOut : AdminUsageReportOutPort; part adminUsage : AdminUsageReport; + part safeUpgrade : SafeServeUpgrade; attribute host : String = "127.0.0.1"; attribute listenPort : Integer = 18765; attribute requestOutputIsolated : Boolean = true; @@ -1333,6 +1388,7 @@ package MemNet { satisfy MN_REQ_00_MissionBridge::MN_REQ_06_ProcessBoundary::MN_REQ_06_3_TcpServeMigration; satisfy MN_REQ_00_MissionBridge::MN_REQ_06_ProcessBoundary::MN_REQ_06_11_AdminServeUsageReport; satisfy MN_REQ_00_MissionBridge::MN_REQ_06_ProcessBoundary::MN_REQ_06_13_ServeRequestIsolation; + satisfy MN_REQ_00_MissionBridge::MN_REQ_06_ProcessBoundary::MN_REQ_06_14_SafeServeUpgrade; satisfy MN_REQ_00_MissionBridge::MN_REQ_12_MultitaskMode::MN_REQ_12_2_SharedStoreTransport; } @@ -1666,6 +1722,7 @@ package MemNet { part importSlice : CmdPathBImport; part catalogSnap : CmdCatalogSnap; part adminUsageReport : CmdAdminUsageReport; + part upgradePrepare : CmdAdminUpgradePrepare; attribute leftoverReadGetNested : Boolean = false; attribute leftoverNewNested : Boolean = false; attribute leftoverAssignedIdMapNested : Boolean = false; @@ -1682,9 +1739,11 @@ package MemNet { attribute emptyCueIsOutline : Boolean = true; attribute noRequiredIdTools : Boolean = true; attribute adminUsageOnAgentMcp : Boolean = false; + attribute upgradePrepareOnAgentMcp : Boolean = false; satisfy MN_REQ_00_MissionBridge::MN_REQ_01_SessionLifecycle::MN_REQ_01_2_SessionLifecycleOps; satisfy MN_REQ_00_MissionBridge::MN_REQ_06_ProcessBoundary::MN_REQ_06_11_AdminServeUsageReport; + satisfy MN_REQ_00_MissionBridge::MN_REQ_06_ProcessBoundary::MN_REQ_06_14_SafeServeUpgrade; satisfy MN_REQ_00_MissionBridge::MN_REQ_01_SessionLifecycle::MN_REQ_01_4_SaveSessionSnapshot; satisfy MN_REQ_00_MissionBridge::MN_REQ_01_SessionLifecycle::MN_REQ_01_5_LoadSessionSnapshot; satisfy MN_REQ_00_MissionBridge::MN_REQ_01_SessionLifecycle::MN_REQ_01_6_OptionalKeepSessionId; @@ -1777,6 +1836,10 @@ package MemNet { attribute noRequiredIdTools : Boolean = true; attribute adminUsageToolNested : Boolean = false; attribute productLabelOnSessionOpen : Boolean = false; + attribute upgradeRetry : Boolean = true; + attribute upgradeRetryDefaultS : Integer = 30; + attribute retryServeDraining : Boolean = true; + attribute retryConnectionRefusal : Boolean = true; satisfy MN_REQ_00_MissionBridge::MN_REQ_01_SessionLifecycle::MN_REQ_01_4_SaveSessionSnapshot; satisfy MN_REQ_00_MissionBridge::MN_REQ_01_SessionLifecycle::MN_REQ_01_5_LoadSessionSnapshot; @@ -1801,6 +1864,7 @@ package MemNet { satisfy MN_REQ_00_MissionBridge::MN_REQ_04_SliceEconomy::MN_REQ_04_9_SessionOutlineEmptyCue; satisfy MN_REQ_00_MissionBridge::MN_REQ_11_SnapshotInterop::MN_REQ_11_17_CatalogModelSnap; satisfy MN_REQ_00_MissionBridge::MN_REQ_11_SnapshotInterop::MN_REQ_11_17_CatalogModelSnap::MN_REQ_11_17_1_CrossCutSatisfyLocators; + satisfy MN_REQ_00_MissionBridge::MN_REQ_06_ProcessBoundary::MN_REQ_06_14_SafeServeUpgrade; } part def MemNetCoreLibrary { @@ -3577,8 +3641,13 @@ package MemNet { attribute adminCountsOmitSessionId : Boolean = true; attribute singleBackendUnchanged : Boolean = true; attribute hardcodedHost : Boolean = false; + attribute upgradeRetry : Boolean = true; + attribute upgradeRetryDefaultS : Integer = 30; + attribute retryServeDraining : Boolean = true; + attribute retryConnectionRefusal : Boolean = true; satisfy MN_REQ_00_MissionBridge::MN_REQ_06_ProcessBoundary::MN_REQ_06_12_ProductGateway; + satisfy MN_REQ_00_MissionBridge::MN_REQ_06_ProcessBoundary::MN_REQ_06_14_SafeServeUpgrade; satisfy MN_REQ_00_MissionBridge::MN_REQ_06_ProcessBoundary::MN_REQ_06_9_LanMcpFrontSeveralServes; } diff --git a/sysml-models/models/implementation.sysml b/sysml-models/models/implementation.sysml index b34ba787..b17c9ddc 100644 --- a/sysml-models/models/implementation.sysml +++ b/sysml-models/models/implementation.sysml @@ -122,6 +122,21 @@ package MemNetImplementation { attribute path : String = "parts/common/memnet/memnet/admin_usage.py"; attribute pyName : String = "memnet.admin_usage"; } + part def SafeUpgradeMod { + doc /* Serve drain, manifest, and startup restore (MN-REQ-06.14). */ + attribute path : String = "parts/common/memnet/memnet/upgrade.py"; + attribute pyName : String = "memnet.upgrade"; + } + part def UpgradeRetryMod { + doc /* Bounded retry for serve_draining and a brief connection refusal. */ + attribute path : String = "parts/common/memnet/memnet/upgrade_retry.py"; + attribute pyName : String = "memnet.upgrade_retry"; + } + part def UpgradeHelperMod { + doc /* Side-by-side venv helper: preflight, drain, swap, verify, rollback. */ + attribute path : String = "parts/common/memnet/memnet/upgrade_run.py"; + attribute pyName : String = "memnet.upgrade_run"; + } part def CliFacadeMod { attribute path : String = "parts/common/memnet/memnet/cli.py"; attribute pyName : String = "memnet.cli"; @@ -273,6 +288,9 @@ package MemNetImplementation { part tcpMod : TcpServeBridgeMod; part serveDaemonMod : ServeDaemonMod; part adminUsageMod : AdminUsageReportMod; + part safeUpgradeMod : SafeUpgradeMod; + part upgradeRetryMod : UpgradeRetryMod; + part upgradeHelperMod : UpgradeHelperMod; part cliMod : CliFacadeMod; part sessionsMod : SessionLifecycleMod; part storeMod : GraphStoreMod; @@ -354,6 +372,26 @@ package MemNetImplementation { end logical ::> memNetSystem.core.transport.tcp.adminUsage; end code ::> adminUsageMod; } + allocation safeUpgradeToMod : SoftwareAllocate { + doc /* Drain, manifest, and startup restore on memnet-serve. */ + end logical ::> memNetSystem.core.transport.tcp.safeUpgrade; + end code ::> safeUpgradeMod; + } + allocation upgradeRetryToMod : SoftwareAllocate { + doc /* Gateway retry across serve_draining and a brief refusal. */ + end logical ::> lanMcpFront.productGateway; + end code ::> upgradeRetryMod; + } + allocation mcpUpgradeRetryToMod : SoftwareAllocate { + doc /* memnet-mcp TCP client retry; same helper as the gateway. */ + end logical ::> memNetSystem.mcp.bridge; + end code ::> upgradeRetryMod; + } + allocation upgradeHelperToMod : SoftwareAllocate { + doc /* Operator helper memnet-upgrade. Not an agent MCP tool. */ + end logical ::> memNetSystem.core.cli; + end code ::> upgradeHelperMod; + } allocation cliToMod : SoftwareAllocate { end logical ::> memNetSystem.core.cli; end code ::> cliMod; diff --git a/sysml-models/models/requirements.sysml b/sysml-models/models/requirements.sysml index 251101a3..2532c4eb 100644 --- a/sysml-models/models/requirements.sysml +++ b/sysml-models/models/requirements.sysml @@ -36,6 +36,7 @@ MN-REQ-06.11 Admin-only serve usage report (counts/caps; opaque session alias; not agent MCP) MN-REQ-06.12 Product gateway on memnet-mcp: per-product route to owning serve (not mn_tip_) MN-REQ-06.13 Serve request isolation: each command's stdout, stderr, and exit_code + MN-REQ-06.14 Safe serve upgrade: drain, manifest, restore; no silent session loss MN-REQ-11.17 Catalog Snap / model Snap — session strata (0.15) MN-REQ-11.17.1 Cross-cut satisfy locators on the catalog; SysMLEdge pin_map is a different tool (not this catalog) @@ -1135,6 +1136,61 @@ package MemNetRequirements { */ attribute requirementId : String = "MN-REQ-06.13"; } + + requirement def MN_REQ_06_14_SafeServeUpgrade { + doc /* + SHALL upgrade a live memnet-serve (and the memnet-mcp / + product gateway in front of it) without silently dropping + a loaded session. Same usage method (cue then pin_map, + mutate). No SemVer bump. Patch-level c on 0.19. + Admin drain: credential MEMNET_ADMIN_TOKEN (same gate as + MN-REQ-06.11; unset admin_unconfigured; mismatch + admin_denied). Command `memnet admin upgrade-prepare` + and envelope {upgrade_prepare:true, admin_token}. Not an + agent MCP tool (onAgentMcp=false). + Drain refuses a new session_open with + @ERR: serve_draining|retry_after_s= (retryable). + It finishes commands already in flight, then snapshots + every loaded session through the lossless snapshot path + (MN-REQ-01.9). It writes a manifest in the state dir + (MEMNET_STATE_DIR): session ids, row and edge counts, + checksums, serve version, snapshot format version. + A session that cannot be saved (snapshot_unsaveable and + any other save failure) SHALL be named in the report. + ready_to_stop SHALL stay false unless the operator passes + the explicit allow-unsaved override. A session MUST NOT + be lost silently. + New serve startup reads that manifest, reloads every + saved session, and checks counts and checksums. It + reports @STAT: upgrade_restore|ok|n|failed|m. It keeps + the manifest and snapshot files until that restore is + verified. A snapshot written by an older patch (format + v1) SHALL load. If the format is unsupported or a + snapshot is corrupt, startup SHALL fail loudly and + SHALL NOT delete or rewrite the snapshot files, so a + rollback to the old venv can still reload them. + Restored sessions keep the same session id, CapsPolicy + ACL bindings, TTL expiry clock, and house (the session + tag map). Product label is kept with the passport. + memnet-mcp (TCP) and the product gateway SHALL treat + serve_draining and a brief connection refusal as + retryable with bounded backoff for MEMNET_UPGRADE_RETRY_S + (default 30 seconds). Order: deploy that client + tolerance first, then drain and swap the serve. The + gateway backend version pin is updated in the same + procedure, after the old serve has stopped and before + the new serve listens. + Helper `memnet-upgrade` encodes the side-by-side venv + procedure: preflight the new interpreter against the + snapshot format, drain, swap the systemd ExecStart from + a unit backup, restore, verify, and roll back to the + backup unit and the previous pin if verification fails. + Hosts: rpi5-syson (serve 18765, mcp 18766), the Endleaf + engine serve named by the gateway registry, and the + droplet gateway. No real session id in docs or tests. + */ + attribute requirementId : String = "MN-REQ-06.14"; + } } // ----- MCP agent boundary ----- @@ -2167,6 +2223,8 @@ package MemNetRequirements { : MN_REQ_00_MissionBridge::MN_REQ_06_ProcessBoundary::MN_REQ_06_12_ProductGateway; requirement serveRequestIsolationReq : MN_REQ_00_MissionBridge::MN_REQ_06_ProcessBoundary::MN_REQ_06_13_ServeRequestIsolation; + requirement safeServeUpgradeReq + : MN_REQ_00_MissionBridge::MN_REQ_06_ProcessBoundary::MN_REQ_06_14_SafeServeUpgrade; requirement mcpAgentBoundaryReq : MN_REQ_00_MissionBridge::MN_REQ_07_McpAgentBoundary; diff --git a/sysml-models/models/verify.sysml b/sysml-models/models/verify.sysml index 317d2aef..225cb2d8 100644 --- a/sysml-models/models/verify.sysml +++ b/sysml-models/models/verify.sysml @@ -50,6 +50,10 @@ memnet-mcp; per-product sha256 credential; one owning serve; no mn_tip_ reuse; pass-through; implemented. Parent #191 catalogue stays inventOnly. + Safe serve upgrade (MN-VER-06-S12): SafeServeUpgrade on + TcpServeBridge; admin drain; manifest; startup restore; + client retry on MCP and the product gateway; no silent + session loss; no SemVer bump. Optional id nickname TARGET (MN-VER-02-S01): GraphElement identity; CREATE () legal; HiddenStoreHandle off-wire; 0.9 leftover_by_id store. TARGET command list cue/pattern only (MN-VER-07-S01): no required-id @@ -1628,6 +1632,70 @@ package MemNetVerification { and gateway.adminCountsOmitSessionId == true and gateway.singleBackendUnchanged == true and gateway.hardcodedHost == false + and gateway.upgradeRetry == true + and gateway.retryServeDraining == true + and gateway.retryConnectionRefusal == true + } + } + } + + verification def MN_VER_06_S12_SafeServeUpgrade { + doc /* + MN-REQ-06.14 — drain, manifest, and startup restore. + Admin credential. New session_open refuses + serve_draining. Unsaveable sessions block ready_to_stop + unless allow-unsaved. Restore checks counts and + checksums. Corrupt snapshots stay on disk. Same session + id, ACL, TTL clock, and house. MCP and gateway retry + (default 30s) ship before the serve swap. Helper rolls + back. Not an agent MCP tool. No SemVer bump. + */ + attribute verificationId : String = "MN-VER-06-S12"; + + subject tcp : TcpServeBridge; + subject mcp : McpFacade; + subject cli : CliFacade; + subject gateway : MemNetProductGateway; + + objective safeServeUpgrade { + verify safeServeUpgradeReq; + require constraint { + tcp.safeUpgrade.implemented == true + and tcp.safeUpgrade.noSemVerBump == true + and tcp.safeUpgrade.onAgentMcp == false + and tcp.safeUpgrade.adminCredEnv == "MEMNET_ADMIN_TOKEN" + and tcp.safeUpgrade.refusesNewSessionOpen == true + and tcp.safeUpgrade.errDraining == "serve_draining" + and tcp.safeUpgrade.finishesInFlight == true + and tcp.safeUpgrade.usesLosslessSnapshot == true + and tcp.safeUpgrade.manifestBeforeStop == true + and tcp.safeUpgrade.blocksOnUnsaveable == true + and tcp.safeUpgrade.allowUnsavedOverride == true + and tcp.safeUpgrade.silentSessionLoss == false + and tcp.safeUpgrade.restoreVerifiesChecksum == true + and tcp.safeUpgrade.restoreVerifiesCounts == true + and tcp.safeUpgrade.keepsFilesUntilVerified == true + and tcp.safeUpgrade.olderPatchSnapshotLoads == true + and tcp.safeUpgrade.corruptFailsLoud == true + and tcp.safeUpgrade.corruptDeletesSnapshots == false + and tcp.safeUpgrade.preservesSessionId == true + and tcp.safeUpgrade.preservesAcl == true + and tcp.safeUpgrade.preservesTtlClock == true + and tcp.safeUpgrade.preservesHouse == true + and tcp.safeUpgrade.retryDefaultS == 30 + and tcp.safeUpgrade.clientsBeforeServe == true + and tcp.safeUpgrade.helperRollsBack == true + and mcp.upgradeRetry == true + and mcp.upgradeRetryDefaultS == 30 + and mcp.retryServeDraining == true + and mcp.retryConnectionRefusal == true + and mcp.adminUsageToolNested == false + and cli.upgradePrepareOnAgentMcp == false + and cli.upgradePrepare.onAgentMcp == false + and gateway.upgradeRetry == true + and gateway.upgradeRetryDefaultS == 30 + and gateway.retryServeDraining == true + and gateway.retryConnectionRefusal == true } } } @@ -2197,6 +2265,8 @@ package MemNetVerification { : MN_VER_06_S08_ClusterRouteVsSliceHandCarry; verification serveRequestIsolationVerify : MN_VER_06_S11_ServeRequestIsolation; + verification safeServeUpgradeVerify + : MN_VER_06_S12_SafeServeUpgrade; verification housekeepEndpointIdentityVerify : MN_VER_04_S06_HousekeepEndpointIdentity; } diff --git a/sysml-models/outputs/README.md b/sysml-models/outputs/README.md index 92c18bce..6fdfe53c 100644 --- a/sysml-models/outputs/README.md +++ b/sysml-models/outputs/README.md @@ -31,6 +31,7 @@ Keep product-canon and GQL application studies. Do not restore leftover `NEW` mi | [session-outline-case-study.md](session-outline-case-study.md) | Dark session empty q = Recall census of S (kinds + LIMIT exemplars) | MN-REQ-04.9; leftover skip leftover | | [usage-dashboard-case-study.md](usage-dashboard-case-study.md) | Human look at serve n/max + housekeep; manage = Memnetor+Devicor | MN-REQ-06.5; MN-VER-06-S02; HTTP parked | | [admin-usage-report-case-study.md](admin-usage-report-case-study.md) | Admin-only serve usage JSON; opaque alias; not agent MCP | MN-REQ-06.11; MN-VER-06-S09 | +| [safe-upgrade-case-study.md](safe-upgrade-case-study.md) | Drain, manifest, and startup restore so a serve swap does not drop sessions | MN-REQ-06.14; MN-VER-06-S12 | | [device-fleet-one-mcp-case-study.md](device-fleet-one-mcp-case-study.md) | Device MemNet services; one MemNet MCP (tip/ops) at the droplet; product face is sysmledge | MN-REQ-06.6; MN-VER-06-S03; tip≠face; not #47 | | [tip-memnet-access-portal-case-study.md](tip-memnet-access-portal-case-study.md) | Admin invite, Google login, Bearer for keyed tip MCP; unauthenticated WWW MemNet refused | MN-REQ-06.8; MN-VER-06-S05; tip≠face; portal sidecar | | [lan-mcp-front-case-study.md](lan-mcp-front-case-study.md) | One MCP catalogue over N LAN serves; ClusterRoute; registry owner; not shipped | MN-REQ-06.9; MN-VER-06-S07; tip≠face; inventOnly #191 | diff --git a/sysml-models/outputs/product-nest-one-page.md b/sysml-models/outputs/product-nest-one-page.md index 79727853..9a74dd05 100644 --- a/sysml-models/outputs/product-nest-one-page.md +++ b/sysml-models/outputs/product-nest-one-page.md @@ -26,6 +26,7 @@ ARCHIVE LOOK MemNetArchive (models/archive.sysml) — leftover_* / ProjectMemNet MUST NOT import OPS LOOK MemNetUsageDashboard — human look only; not agent wire SERVE ADMIN LOOK AdminUsageReport on TcpServeBridge — admin JSON; opaque alias; not agent MCP (MN-REQ-06.11) +SERVE UPGRADE SafeServeUpgrade on TcpServeBridge — drain, manifest, restore (MN-REQ-06.14) OPS FLEET MemNetOpsFleet — device MemNet services; one MemNet MCP at droplet (tip/ops; tip≠face; product face is sysmledge) OPS ACCESS TipMemNetAccessPortal — Szu-Wei invite, Google login, Bearer for keyed tip MCP at droplet WWW (tip≠face; not sysmledge; portal sidecar; look-only service/client status) OPS LAN FRONT MemNetLanMcpFront — ClusterRoute: one MCP catalogue over N LAN serves (inventOnly #191; tip≠face; not shipped) diff --git a/sysml-models/outputs/safe-upgrade-case-study.md b/sysml-models/outputs/safe-upgrade-case-study.md new file mode 100644 index 00000000..841ea92a --- /dev/null +++ b/sysml-models/outputs/safe-upgrade-case-study.md @@ -0,0 +1,20 @@ +# Case study: safe serve upgrade (MN-REQ-06.14) + +**Story:** A live `memnet-serve` has to move to a new patch without dropping sessions on the floor. The hand procedure (side-by-side venv, save, stop, start, reload) worked, and it also hid any session that failed to save. + +**Requirement:** MN-REQ-06.14 (`safeServeUpgradeReq`). **Verify:** MN-VER-06-S12. + +| | | +|--|--| +| Part | `SafeServeUpgrade` on `TcpServeBridge` | +| Command | `CmdAdminUpgradePrepare` on `CliFacade`; not on `McpFacade` | +| Retry | `McpFacade` and `MemNetProductGateway` (`upgradeRetry=true`, default 30s) | +| Helper | `memnet-upgrade` (`UpgradeHelperMod`) | +| Snapshot | lossless writer plus an optional `# upgrade-passport` line | +| Flag | `implemented=true`; `noSemVerBump=true`; `silentSessionLoss=false` | + +Drain refuses new `session_open` with `serve_draining`. It snapshots through the lossless path, writes a manifest (counts, checksums, format version), and stays not-ready when any session is `snapshot_unsaveable` unless `--allow-unsaved` is set. Startup restore checks counts and checksums and prints `@STAT: upgrade_restore|ok|n|failed|m`. A corrupt file or an unsupported format fails loud and leaves the files in place. + +Client tolerance ships first. Then the serve swaps. The gateway pin moves in that same procedure, after drain and before the new process listens. + +Hosts in the operator note: rpi5-syson serve `18765` and mcp `18766`; the Endleaf engine serve named by the gateway registry; the droplet gateway. Teach: [`docs/operations/safe-upgrade.md`](../../docs/operations/safe-upgrade.md). diff --git a/sysml-models/outputs/ssot-to-code-allocate-map.md b/sysml-models/outputs/ssot-to-code-allocate-map.md index 66e4c53d..714e9893 100644 --- a/sysml-models/outputs/ssot-to-code-allocate-map.md +++ b/sysml-models/outputs/ssot-to-code-allocate-map.md @@ -19,6 +19,10 @@ | tcpToMod | `…transport.tcp` | `tcpMod` | `parts/common/memnet/memnet/tcp_serve_bridge.py` | `memnet.tcp_serve_bridge` | true | | tcpDaemonToMod | `…transport.tcp` | `serveDaemonMod` | `parts/common/memnet/memnet/serve.py` | `memnet.serve` | true | | adminUsageToMod | `…transport.tcp.adminUsage` | `adminUsageMod` | `parts/common/memnet/memnet/admin_usage.py` | `memnet.admin_usage` | true | +| safeUpgradeToMod | `…transport.tcp.safeUpgrade` | `safeUpgradeMod` | `parts/common/memnet/memnet/upgrade.py` | `memnet.upgrade` | true | +| upgradeRetryToMod | `lanMcpFront.productGateway` | `upgradeRetryMod` | `parts/common/memnet/memnet/upgrade_retry.py` | `memnet.upgrade_retry` | true | +| mcpUpgradeRetryToMod | `memNetSystem.mcp.bridge` | `upgradeRetryMod` | `parts/common/memnet/memnet/upgrade_retry.py` | `memnet.upgrade_retry` | true | +| upgradeHelperToMod | `memNetSystem.core.cli` | `upgradeHelperMod` | `parts/common/memnet/memnet/upgrade_run.py` | `memnet.upgrade_run` | true | | cliToMod | `memNetSystem.core.cli` | `cliMod` | `parts/common/memnet/memnet/cli.py` | `memnet.cli` | true | | sessionsToMod | `…sessions` | `sessionsMod` | `parts/common/memnet/memnet/session_lifecycle.py` | `memnet.session_lifecycle` | true | | storeToMod | `…sessions.store` | `storeMod` | `parts/common/memnet/memnet/graph_store.py` | `memnet.graph_store` | true | diff --git a/sysml-models/outputs/system-design-notes.md b/sysml-models/outputs/system-design-notes.md index 7a9ebcbf..4576b028 100644 --- a/sysml-models/outputs/system-design-notes.md +++ b/sysml-models/outputs/system-design-notes.md @@ -56,6 +56,7 @@ Patterns on **SharedLlmMemory** — application shelf. Product-canon mechanism s | Multitask transport | [tcp-shared-multitask-case-study.md](tcp-shared-multitask-case-study.md) | | Lead imports member WM (path B) | [session-import-case-study.md](session-import-case-study.md) | | Snapshot passport | [snapshot-passport-case-study.md](snapshot-passport-case-study.md) | +| Safe serve upgrade | [safe-upgrade-case-study.md](safe-upgrade-case-study.md) | | Durable hydrate/flush | [durable-hydrate-flush-case-study.md](durable-hydrate-flush-case-study.md) | | Empty-cue session outline | [session-outline-case-study.md](session-outline-case-study.md) | | Human usage look (parked HTTP) | [usage-dashboard-case-study.md](usage-dashboard-case-study.md) | diff --git a/tests/conftest.py b/tests/conftest.py index 7be0aa10..a6bbc3b2 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -8,6 +8,12 @@ from memnet.session import purge_expired, reset_registry, set_now_override +@pytest.fixture(autouse=True) +def _upgrade_retry_off(monkeypatch: pytest.MonkeyPatch): + """Production default is a 30s upgrade retry. Tests opt in.""" + monkeypatch.setenv("MEMNET_UPGRADE_RETRY_S", "0") + + @pytest.fixture def memnet_temp(monkeypatch: pytest.MonkeyPatch): monkeypatch.setenv("MEMNET_TEST_INLINE", "1") @@ -18,11 +24,15 @@ def memnet_temp(monkeypatch: pytest.MonkeyPatch): from memnet.admin_usage import reset_pressure reset_pressure() + from memnet.upgrade import reset_drain_gate + + reset_drain_gate() yield set_now_override(None) reset_registry() purge_expired() reset_pressure() + reset_drain_gate() @pytest.fixture diff --git a/tests/test_safe_upgrade.py b/tests/test_safe_upgrade.py new file mode 100644 index 00000000..f23f1037 --- /dev/null +++ b/tests/test_safe_upgrade.py @@ -0,0 +1,556 @@ +"""Safe serve upgrade: drain, manifest, restore, retry, rollback (MN-REQ-06.14).""" + +from __future__ import annotations + +import json +import os +import socket +import subprocess +import sys +import time +from pathlib import Path + +import pytest + +from memnet.exceptions import MemNetError +from memnet.in_process_engine import run_argv +from memnet.registry import get_entry +from memnet.serve import probe, send_command +from memnet.session import open_session, reset_registry, utc_now +from memnet.snapshot import write_snapshot +from memnet.upgrade import ( + BLOCKED_NAME, + MANIFEST_NAME, + drain_gate, + prepare_upgrade, + reset_drain_gate, + restore_manifest, +) +from memnet.upgrade_retry import call_with_upgrade_retry +from memnet.upgrade_run import UpgradeRunError, run_upgrade, set_pinned_version, swap_execstart + +TOKEN = "test-admin-token" + + +def _mutate(sid: str, text: str) -> None: + raw = run_argv(["mutate", "--stdin", "--session", sid], stdin=text) + assert raw["exit_code"] == 0, raw["stderr"] + + +def _edge_poison(sid: str) -> None: + entry = get_entry(sid) + assert entry is not None + for hid in entry.store.write_order: + rec = entry.store._by_hid.get(hid) + if rec is not None and rec.tag == "EDG": + rec.fields["not_a_field"] = "x" + return + raise AssertionError("expected an edge to poison") + + +def test_round_trip_keeps_id_acl_ttl_and_near_cap(memnet_temp, monkeypatch, tmp_path, schema_file): + monkeypatch.setenv("MEMNET_ADMIN_TOKEN", TOKEN) + monkeypatch.setenv("MEMNET_MAX_ROWS", "80") + monkeypatch.setenv("MEMNET_STATE_DIR", str(tmp_path)) + plain = open_session(map_file=str(schema_file), product="gateapp") + _mutate(plain.session_id, "CREATE (:TSK {id: 'TSK_keep', goal: 'kept-fact', status: 'open'})\n") + acl = open_session(map_file=str(schema_file), ttl_minutes=2) + from datetime import timedelta + + _mutate(acl.session_id, "CREATE (:TSK {id: 'TSK_acl', goal: 'bound', status: 'open'})\n") + soon = (utc_now() + timedelta(seconds=90)).isoformat().replace("+00:00", "Z") + acl.meta.expires_at = soon + acl.enable_acl() + acl.grant_caller("caller-demo", can_pin_map=True, can_mutate=False, write_scope="labels=TSK") + acl.set_bind("mission-demo", "lease-demo") + wide = open_session(map_file=str(schema_file)) + lines = [f"CREATE (:TSK {{id: 'TSK_{i}', goal: 'g{i}', status: 'open'}})" for i in range(70)] + lines.append( + "MATCH (a {id: 'TSK_0'}), (b {id: 'TSK_1'}) CREATE (a)-[:owns {id: 'E_link'}]->(b)" + ) + _mutate(wide.session_id, "\n".join(lines) + "\n") + wide_rows = get_entry(wide.session_id).store.row_count_non_law() + assert wide_rows >= 71 + assert wide_rows >= int(0.85 * 80) + want = { + plain.session_id: {"expires": plain.meta.expires_at, "product": "gateapp"}, + acl.session_id: {"expires": soon}, + wide.session_id: {"rows": wide_rows}, + } + result = prepare_upgrade(TOKEN, directory=tmp_path) + assert result.ready_to_stop is True + assert result.unsaved == [] + assert (tmp_path / MANIFEST_NAME).is_file() + reset_registry() + reset_drain_gate() + report = restore_manifest(tmp_path) + assert report.failed == 0 + assert report.ok == 3 + assert set(report.session_ids) == set(want) + plain_back = get_entry(plain.session_id) + acl_back = get_entry(acl.session_id) + wide_back = get_entry(wide.session_id) + assert plain_back is not None and acl_back is not None and wide_back is not None + assert plain_back.meta.expires_at == want[plain.session_id]["expires"] + assert plain_back.meta.product == "gateapp" + assert acl_back.meta.expires_at == soon + assert acl_back.acl is not None and acl_back.acl.enabled is True + assert acl_back.acl.bind is not None + assert acl_back.acl.bind.mission_id == "mission-demo" + assert acl_back.acl.bind.lease == "lease-demo" + grant = acl_back.acl.callers["caller-demo"] + assert grant.can_pin_map is True + assert grant.can_mutate is False + assert grant.write_scope is not None + assert "TSK" in grant.write_scope.labels + assert wide_back.store.row_count_non_law() == wide_rows + snaps = list((tmp_path / "upgrade-snapshots").glob("*.snap")) + assert len(snaps) == 3 + + +def test_unsaveable_blocks_ready_unless_override(memnet_temp, monkeypatch, tmp_path, schema_file): + monkeypatch.setenv("MEMNET_ADMIN_TOKEN", TOKEN) + ss = open_session(map_file=str(schema_file)) + _mutate( + ss.session_id, + "CREATE (:TSK {id: 'TSK_a', goal: 'a', status: 'open'})\n" + "CREATE (:TSK {id: 'TSK_b', goal: 'b', status: 'open'})\n" + "MATCH (a {id: 'TSK_a'}), (b {id: 'TSK_b'}) " + "CREATE (a)-[:owns {id: 'E_bad'}]->(b)\n", + ) + _edge_poison(ss.session_id) + blocked = prepare_upgrade(TOKEN, directory=tmp_path) + assert blocked.ready_to_stop is False + assert blocked.exit_code == 2 + assert blocked.unsaved + assert blocked.unsaved[0]["code"] == "snapshot_unsaveable" + assert blocked.unsaved[0]["session_id"] == ss.session_id + assert not (tmp_path / MANIFEST_NAME).exists() + assert (tmp_path / BLOCKED_NAME).is_file() + assert drain_gate.phase == "open" + again = open_session(map_file=str(schema_file)) + assert again.session_id != ss.session_id + forced = prepare_upgrade(TOKEN, allow_unsaved=True, directory=tmp_path) + assert forced.ready_to_stop is True + manifest = json.loads((tmp_path / MANIFEST_NAME).read_text(encoding="utf-8")) + saved_ids = {row["session_id"] for row in manifest["sessions"]} + assert ss.session_id not in saved_ids + assert any(row["session_id"] == ss.session_id for row in manifest["unsaved"]) + + +def test_corrupt_snapshot_fails_loud_and_keeps_bytes( + memnet_temp, monkeypatch, tmp_path, schema_file +): + monkeypatch.setenv("MEMNET_ADMIN_TOKEN", TOKEN) + ss = open_session(map_file=str(schema_file)) + _mutate(ss.session_id, "CREATE (:TSK {id: 'TSK_c', goal: 'c', status: 'open'})\n") + assert prepare_upgrade(TOKEN, directory=tmp_path).ready_to_stop is True + snap = next((tmp_path / "upgrade-snapshots").glob("*.snap")) + original = snap.read_bytes() + snap.write_bytes(original[:40] + b"XXXX" + original[44:]) + corrupt = snap.read_bytes() + reset_registry() + reset_drain_gate() + report = restore_manifest(tmp_path) + assert report.failed == 1 + assert report.ok == 0 + assert snap.read_bytes() == corrupt + assert snap.is_file() + + +def test_older_patch_snapshot_loads(memnet_temp, monkeypatch, tmp_path, schema_file): + monkeypatch.setenv("MEMNET_ADMIN_TOKEN", TOKEN) + ss = open_session(map_file=str(schema_file)) + _mutate(ss.session_id, "CREATE (:TSK {id: 'TSK_old', goal: 'legacy', status: 'open'})\n") + entry = get_entry(ss.session_id) + assert entry is not None + snap_dir = tmp_path / "upgrade-snapshots" + snap_dir.mkdir() + path = snap_dir / "legacy.snap" + write_snapshot(ss, path) + assert b"upgrade-passport" not in path.read_bytes() + import hashlib + + rows = entry.store.row_count_non_law() + edges = sum(1 for hid in entry.store.write_order if entry.store._by_hid[hid].tag == "EDG") + manifest = { + "manifest_format": 1, + "snapshot_format": 1, + "serve_version": "0.19.19", + "ready_to_stop": True, + "retired": False, + "allow_unsaved": False, + "sessions": [ + { + "session_id": ss.session_id, + "file": "upgrade-snapshots/legacy.snap", + "rows": rows, + "edges": edges, + "checksum_sha256": hashlib.sha256(path.read_bytes()).hexdigest(), + "expires_at": entry.meta.expires_at, + "created_at": entry.meta.created_at, + "ttl_minutes": entry.meta.ttl_minutes, + "house_tags": sorted(entry.tag_map.tags), + "product": None, + "saved": True, + } + ], + "unsaved": [], + } + (tmp_path / MANIFEST_NAME).write_text(json.dumps(manifest), encoding="utf-8") + expires = entry.meta.expires_at + sid = ss.session_id + before = path.read_bytes() + reset_registry() + report = restore_manifest(tmp_path) + assert report.failed == 0, report.errors + assert report.ok == 1 + back = get_entry(sid) + assert back is not None + assert back.meta.session_id == sid + assert back.meta.expires_at == expires + assert back.store.row_count_non_law() == rows + assert path.read_bytes() == before + + +def test_unsupported_format_leaves_file(memnet_temp, tmp_path, schema_file): + ss = open_session(map_file=str(schema_file)) + snap_dir = tmp_path / "upgrade-snapshots" + snap_dir.mkdir() + path = snap_dir / "one.snap" + write_snapshot(ss, path) + before = path.read_bytes() + manifest = { + "snapshot_format": 99, + "ready_to_stop": True, + "retired": False, + "sessions": [], + } + (tmp_path / MANIFEST_NAME).write_text(json.dumps(manifest), encoding="utf-8") + reset_registry() + with pytest.raises(MemNetError) as exc: + restore_manifest(tmp_path) + assert exc.value.code == "upgrade_snapshot_format" + assert path.read_bytes() == before + + +def test_drain_refuses_session_open(memnet_temp, schema_file): + drain_gate.enter_quiesce(5) + with pytest.raises(MemNetError) as exc: + open_session(map_file=str(schema_file)) + assert exc.value.code == "serve_draining" + assert "retry_after_s=5" in exc.value.message + + +def test_inflight_blocks_ready_until_finished(): + reset_drain_gate() + counted, err = drain_gate.begin(["mutate", "--stdin"]) + assert counted is True and err is None + drain_gate.enter_quiesce(5) + _counted, refused = drain_gate.begin(["session", "open", "--map-file", "x"]) + assert refused == "serve_draining" + assert drain_gate.wait_idle(0.05) is False + drain_gate.end(counted) + assert drain_gate.wait_idle(0.2) is True + reset_drain_gate() + + +def test_client_retries_connection_refusal_and_draining(monkeypatch): + monkeypatch.setenv("MEMNET_UPGRADE_RETRY_S", "2") + state = {"n": 0} + + def flaky() -> dict: + state["n"] += 1 + if state["n"] < 3: + raise ConnectionRefusedError("down") + return {"exit_code": 0, "stdout": "", "stderr": ""} + + raw = call_with_upgrade_retry(flaky) + assert raw["exit_code"] == 0 + assert state["n"] == 3 + + state["n"] = 0 + + def draining() -> dict: + state["n"] += 1 + if state["n"] < 3: + return { + "exit_code": 2, + "stdout": "", + "stderr": "@ERR: serve_draining|retry_after_s=5\n", + } + return {"exit_code": 0, "stdout": "ok\n", "stderr": ""} + + raw = call_with_upgrade_retry(draining) + assert raw["stdout"] == "ok\n" + assert state["n"] == 3 + + +def test_client_no_retry_when_window_disabled(monkeypatch): + monkeypatch.setenv("MEMNET_UPGRADE_RETRY_S", "0") + + def down() -> dict: + raise ConnectionRefusedError("down") + + with pytest.raises(ConnectionRefusedError): + call_with_upgrade_retry(down) + + +def test_mcp_tcp_retries_until_serve_accepts(monkeypatch): + monkeypatch.setenv("MEMNET_UPGRADE_RETRY_S", "2") + monkeypatch.setenv("MEMNET_MCP_TRANSPORT", "tcp") + monkeypatch.delenv("MEMNET_TEST_INLINE", raising=False) + calls = {"n": 0} + + def fake_probe(*_a, **_k) -> bool: + calls["n"] += 1 + return calls["n"] >= 3 + + def fake_send(*_a, **_k) -> dict: + return {"exit_code": 0, "stdout": "@STAT: sessions|0|1024\n", "stderr": ""} + + monkeypatch.setattr("memnet_mcp.client.probe", fake_probe) + monkeypatch.setattr("memnet_mcp.client.send_command", fake_send) + from memnet_mcp.client import run_memnet + + resp = run_memnet(["session", "list"]) + assert resp.exit_code == 0 + assert calls["n"] >= 3 + + +def test_rollback_restores_unit_and_pin(memnet_temp, tmp_path): + unit = tmp_path / "memnet-serve.service" + unit.write_text("[Service]\nExecStart=/old/venv/bin/memnet serve\n", encoding="utf-8") + gateway = tmp_path / "gateway.json" + gateway.write_text( + json.dumps({"products": {"endleaf": {"pinned_version": "0.19.19", "backends": []}}}), + encoding="utf-8", + ) + state = tmp_path / "state" + state.mkdir() + calls = {"n": 0} + + def drainer() -> None: + manifest = { + "ready_to_stop": True, + "retired": False, + "snapshot_format": 1, + "sessions": [], + "unsaved": [], + } + (state / MANIFEST_NAME).write_text(json.dumps(manifest), encoding="utf-8") + + def restarter() -> None: + calls["n"] += 1 + if calls["n"] == 1: + (state / "upgrade-restore.json").write_text( + json.dumps({"ok": 0, "failed": 1, "skipped": False, "errors": []}), + encoding="utf-8", + ) + + with pytest.raises(UpgradeRunError) as exc: + run_upgrade( + new_python=sys.executable, + new_exec="/new/venv/bin/memnet serve", + state=state, + clients_ready=True, + unit_path=unit, + gateway_config=gateway, + product="endleaf", + drainer=drainer, + restarter=restarter, + verify_wait_s=1, + ) + assert exc.value.code == "upgrade_verify" + assert "ExecStart=/old/venv/bin/memnet serve" in unit.read_text(encoding="utf-8") + pin = json.loads(gateway.read_text(encoding="utf-8")) + assert pin["products"]["endleaf"]["pinned_version"] == "0.19.19" + assert (unit.with_suffix(".service.bak")).is_file() + assert calls["n"] == 2 + + +def test_clients_must_be_ready_before_swap(tmp_path): + unit = tmp_path / "memnet-serve.service" + unit.write_text("ExecStart=/old/venv/bin/memnet serve\n", encoding="utf-8") + with pytest.raises(UpgradeRunError) as exc: + run_upgrade( + new_python=sys.executable, + new_exec="/new/venv/bin/memnet serve", + state=tmp_path, + clients_ready=False, + unit_path=unit, + drainer=lambda: None, + restarter=lambda: None, + ) + assert exc.value.code == "upgrade_clients" + assert unit.read_text(encoding="utf-8").startswith("ExecStart=/old") + + +def test_swap_and_pin_helpers(): + swapped = swap_execstart("ExecStart=/old/bin/memnet serve\n", "/new/bin/memnet serve") + assert swapped == "ExecStart=/new/bin/memnet serve\n" + pinned = set_pinned_version( + json.dumps({"products": {"endleaf": {"pinned_version": "0.19.19"}}}), + "endleaf", + "0.19.20", + ) + assert json.loads(pinned)["products"]["endleaf"]["pinned_version"] == "0.19.20" + + +def _free_port() -> int: + with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as sock: + sock.bind(("127.0.0.1", 0)) + return int(sock.getsockname()[1]) + + +def _wait_probe(port: int, timeout_s: float = 8.0) -> bool: + deadline = time.monotonic() + timeout_s + while time.monotonic() < deadline: + if probe(host="127.0.0.1", port=port): + return True + time.sleep(0.05) + return False + + +def test_tcp_restart_restores_same_sessions(memnet_temp, monkeypatch, tmp_path, schema_file): + monkeypatch.setenv("MEMNET_ADMIN_TOKEN", TOKEN) + port = _free_port() + state = tmp_path / "state" + state.mkdir() + log = tmp_path / "serve.log" + env = os.environ.copy() + env.update( + { + "MEMNET_ADMIN_TOKEN": TOKEN, + "MEMNET_STATE_DIR": str(state), + "MEMNET_SERVE_HOST": "127.0.0.1", + "MEMNET_SERVE_PORT": str(port), + "MEMNET_MAX_ROWS": "80", + "MEMNET_UPGRADE_RETRY_S": "0", + } + ) + env.pop("MEMNET_TEST_INLINE", None) + env.pop("MEMNET_SESSION", None) + + def start() -> subprocess.Popen[bytes]: + handle = log.open("a", encoding="utf-8") + proc = subprocess.Popen( + [sys.executable, "-m", "memnet", "serve", "--host", "127.0.0.1", "--port", str(port)], + env=env, + stdout=handle, + stderr=subprocess.STDOUT, + ) + handle.close() + assert _wait_probe(port), log.read_text(encoding="utf-8") + return proc + + def stop(proc: subprocess.Popen[bytes]) -> None: + proc.terminate() + try: + proc.wait(timeout=5) + except subprocess.TimeoutExpired: + proc.kill() + proc.wait(timeout=5) + deadline = time.monotonic() + 5 + while time.monotonic() < deadline and probe(host="127.0.0.1", port=port): + time.sleep(0.05) + + proc = start() + try: + opened = send_command( + [ + "session", + "open", + "--map-file", + str(schema_file), + "--ttl", + "2", + "--product", + "gateapp", + ], + host="127.0.0.1", + port=port, + ) + assert opened["exit_code"] == 0, opened + sid = opened["stdout"].split("|", 1)[0].replace("@SESSION: ", "").strip() + mutated = send_command( + ["mutate", "--stdin", "--session", sid], + stdin="CREATE (:TSK {id: 'TSK_live', goal: 'kept-fact', status: 'open'})\n", + host="127.0.0.1", + port=port, + ) + assert mutated["exit_code"] == 0, mutated + granted = send_command( + [ + "session", + "acl-enable", + "--session", + sid, + ], + host="127.0.0.1", + port=port, + ) + assert granted["exit_code"] == 0, granted + bound = send_command( + [ + "session", + "acl-bind", + "--mission-id", + "mission-demo", + "--lease", + "lease-demo", + "--session", + sid, + ], + host="127.0.0.1", + port=port, + ) + assert bound["exit_code"] == 0, bound + prepared = send_command( + ["admin", "upgrade-prepare", "--state-dir", str(state)], + admin_token=TOKEN, + host="127.0.0.1", + port=port, + timeout=30, + ) + assert prepared["exit_code"] == 0, prepared + assert "upgrade_prepare|ready|1" in prepared["stderr"] + refused = send_command( + ["session", "open", "--map-file", str(schema_file)], + host="127.0.0.1", + port=port, + ) + assert "serve_draining" in refused["stderr"] + assert "retry_after_s" in refused["stderr"] + finally: + stop(proc) + + proc = start() + try: + report = json.loads((state / "upgrade-restore.json").read_text(encoding="utf-8")) + assert report["failed"] == 0 + assert report["ok"] == 1 + assert report["stat"] == "@STAT: upgrade_restore|ok|1|failed|0" + assert sid in report["session_ids"] + detail = report["details"][0] + assert detail["acl_enabled"] is True + assert detail["bind_mission"] == "mission-demo" + assert detail["product"] == "gateapp" + found = send_command( + ["query", "find", "--kind", "TSK", "--limit", "5", "--session", sid], + host="127.0.0.1", + port=port, + ) + assert found["exit_code"] == 0, found + assert "kept-fact" in found["stdout"] + assert (state / "upgrade-snapshots").is_dir() + assert list((state / "upgrade-snapshots").glob("*.snap")) + finally: + stop(proc) + + +def test_upgrade_doc_has_no_session_id(): + text = (Path(__file__).resolve().parents[1] / "docs/operations/safe-upgrade.md").read_text( + encoding="utf-8" + ) + assert "mn_" not in text diff --git a/tests/test_sysml_safe_upgrade.py b/tests/test_sysml_safe_upgrade.py new file mode 100644 index 00000000..b5bc7fae --- /dev/null +++ b/tests/test_sysml_safe_upgrade.py @@ -0,0 +1,69 @@ +"""Honesty-c: MN-REQ-06.14 safe serve upgrade is modelled before the code.""" + +from __future__ import annotations + +from pathlib import Path + +ROOT = Path(__file__).resolve().parents[1] +MODELS = ROOT / "sysml-models" / "models" +DEPLOY = MODELS / "deploy.sysml" +REQUIREMENTS = MODELS / "requirements.sysml" +VERIFY = MODELS / "verify.sysml" +IMPLEMENTATION = MODELS / "implementation.sysml" +STUDY = ROOT / "sysml-models" / "outputs" / "safe-upgrade-case-study.md" +NEST = ROOT / "sysml-models" / "outputs" / "product-nest-one-page.md" +TEACH = ROOT / "docs" / "operations" / "safe-upgrade.md" +MAP = ROOT / "sysml-models" / "outputs" / "ssot-to-code-allocate-map.md" + + +def test_safe_upgrade_parts(): + text = DEPLOY.read_text(encoding="utf-8") + assert "part def SafeServeUpgrade" in text + assert "part def CmdAdminUpgradePrepare" in text + assert "part safeUpgrade : SafeServeUpgrade" in text + assert "attribute silentSessionLoss : Boolean = false" in text + assert 'attribute errDraining : String = "serve_draining"' in text + assert "attribute retryDefaultS : Integer = 30" in text + assert "attribute clientsBeforeServe : Boolean = true" in text + mcp = text.split("part def McpFacade", 1)[1].split("part def ", 1)[0] + assert "part upgradePrepare" not in mcp + assert "attribute upgradeRetry : Boolean = true" in mcp + gateway = text.split("part def MemNetProductGateway", 1)[1].split("part def ", 1)[0] + assert "attribute retryConnectionRefusal : Boolean = true" in gateway + + +def test_requirement_verify_allocate(): + req = REQUIREMENTS.read_text(encoding="utf-8") + assert "requirement def MN_REQ_06_14_SafeServeUpgrade" in req + assert 'attribute requirementId : String = "MN-REQ-06.14"' in req + assert "safeServeUpgradeReq" in req + assert "serve_draining" in req + assert "allow-unsaved" in req + ver = VERIFY.read_text(encoding="utf-8") + assert "verification def MN_VER_06_S12_SafeServeUpgrade" in ver + assert "verify safeServeUpgradeReq" in ver + assert "tcp.safeUpgrade.corruptDeletesSnapshots == false" in ver + impl = IMPLEMENTATION.read_text(encoding="utf-8") + assert "part def SafeUpgradeMod" in impl + assert "allocation safeUpgradeToMod" in impl + assert "memnet.upgrade" in impl + assert "memnet.upgrade_retry" in impl + assert "memnet.upgrade_run" in impl + ledger = MAP.read_text(encoding="utf-8") + assert "| safeUpgradeToMod |" in ledger + assert "| upgradeHelperToMod |" in ledger + + +def test_outputs_and_docs(): + study = STUDY.read_text(encoding="utf-8") + assert "MN-REQ-06.14" in study + assert "MN-VER-06-S12" in study + nest = NEST.read_text(encoding="utf-8") + assert "SafeServeUpgrade" in nest + teach = TEACH.read_text(encoding="utf-8") + assert "memnet-upgrade" in teach + assert "MEMNET_ADMIN_TOKEN" in teach + assert "18765" in teach + assert "18766" in teach + assert "mn_" not in teach + assert "mn_" not in study