From 11fb003f4b2a8832ead92db7f5a51eededa895cf Mon Sep 17 00:00:00 2001 From: Igor Apresov Date: Fri, 25 Sep 2026 17:39:26 +0300 Subject: [PATCH 1/2] Support bounded stop-start MCP rollout with predecessor recovery --- .github/workflows/ci.yml | 3 +- delivery/ci/publish_mcp_artifacts.py | 11 ++-- delivery/vps/v8std_mcp_release.py | 96 ++++++++++++++++++++++------ tests/mcp_release_fixture.py | 16 ++++- tests/test_mcp_publication.py | 28 ++++++++ tests/test_v8std_mcp_release.py | 91 ++++++++++++++++++++++++++ 6 files changed, 218 insertions(+), 27 deletions(-) diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 9f16384..9aa785d 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -267,7 +267,7 @@ jobs: needs.validate.result == 'success' && needs.publish.result == 'success' && vars.MCP_RUNTIME_DEPLOY_ENABLED == 'true' runs-on: ubuntu-24.04 - timeout-minutes: 10 + timeout-minutes: 20 environment: mcp-production permissions: contents: read @@ -276,6 +276,7 @@ jobs: env: GH_TOKEN: ${{ github.token }} MCP_RUNTIME_DEPLOY_ENABLED: ${{ vars.MCP_RUNTIME_DEPLOY_ENABLED }} + MCP_STOP_START_DEPLOY_ENABLED: ${{ vars.MCP_STOP_START_DEPLOY_ENABLED }} MCP_GATES: >- {"build":"${{ needs.validate.outputs.build }}","tests":"${{ needs.validate.outputs.tests }}","benchmark":"${{ needs.validate.outputs.benchmark }}"} steps: diff --git a/delivery/ci/publish_mcp_artifacts.py b/delivery/ci/publish_mcp_artifacts.py index a41ca50..48279a7 100644 --- a/delivery/ci/publish_mcp_artifacts.py +++ b/delivery/ci/publish_mcp_artifacts.py @@ -928,9 +928,10 @@ def verify_running_runtime(source_sha): return runtime_smoke(url, record, deadline) -def deploy_runtime(context, adapter, accepted, *, enabled, configuration_digest, platform): +def deploy_runtime(context, adapter, accepted, *, enabled, configuration_digest, platform, stop_start=False): authorized(context) require(type(enabled) is bool, "activation_type") + require(type(stop_start) is bool, "activation_type") if not enabled: return {"state": "DISABLED"} runtime = accepted.get("runtime") @@ -956,15 +957,16 @@ def deploy_runtime(context, adapter, accepted, *, enabled, configuration_digest, and item["platform"].get("os") == os_name and item["platform"].get("architecture") == architecture and item["platform"].get("variant", "") in ({"", "v8"} if architecture == "arm64" else {""})] require(len(descriptors) == 1, "platform_descriptor") + duration = 900 if stop_start else 300 envelope = {"schema_version": 1, "release_id": f"ci-{context['run_id']}-{context['attempt']}", "sequence": context["run_number"] * 1000 + context["attempt"], "trigger_sha": context["sha"], "runtime_source_sha": runtime["source_sha"], "image": IMAGE, "image_digest": runtime["image_digest"], "platform_digest": descriptors[0]["digest"], "configuration_digest": configuration_digest, "corpus_id": manifest["corpus_id"], "archive_sha256": manifest["archive"]["sha256"], - "deadline": int(adapter.time()) + 300} + "deadline": int(adapter.time()) + duration} validate_envelope(canonical_json(envelope), now=adapter.time()) wire = canonical_json(envelope) + b"\n" - deadline = adapter.monotonic() + 300 + deadline = adapter.monotonic() + duration require(adapter.command("validate-envelope", wire) == envelope, "validated_envelope_identity") initial = adapter.command("deploy", wire) require(isinstance(initial, dict) and initial.get("release_id") == envelope["release_id"], "release_queue_identity") @@ -1063,7 +1065,8 @@ def main(argv=None): require(accepted["trigger_sha"] == context["sha"], "accepted_trigger") result = deploy_runtime(context, adapter, accepted, enabled=activated(os.environ, "MCP_RUNTIME_DEPLOY_ENABLED"), - configuration_digest=os.environ.get("MCP_CONFIGURATION_DIGEST"), platform=os.environ.get("MCP_PLATFORM")) + configuration_digest=os.environ.get("MCP_CONFIGURATION_DIGEST"), platform=os.environ.get("MCP_PLATFORM"), + stop_start=activated(os.environ, "MCP_STOP_START_DEPLOY_ENABLED")) print("runtime " + result["state"]) except (ValueError, OSError, subprocess.SubprocessError) as error: code = str(error) if isinstance(error, PublicationError) else getattr(error, "code", type(error).__name__) diff --git a/delivery/vps/v8std_mcp_release.py b/delivery/vps/v8std_mcp_release.py index a101b9a..de3b9f8 100644 --- a/delivery/vps/v8std_mcp_release.py +++ b/delivery/vps/v8std_mcp_release.py @@ -64,12 +64,15 @@ "OnBootSec=5s\nOnUnitInactiveSec=15s\nUnit=v8std-bootstrap-recover.service\n\n" "[Install]\nWantedBy=timers.target\n") TRANSACTION = 300 +STOP_START_TRANSACTION = 900 READINESS = 90 SMOKE = DRAIN = 30 STOP = 45 # Rollback gets a live budget even if preparation uses its entire work allowance. # 90 ready + 30 smoke + 45 stop + 15 nginx/control overhead. RECOVERY_RESERVE = 180 +STOP_START_RECOVERY_RESERVE = STOP + READINESS + 2 * SMOKE + 15 +STOP_START_ACCEPTANCE_BUDGET = STOP + 2 * READINESS + 3 * SMOKE + DRAIN + 20 INITIAL_RECOVERY_RESERVE = 60 TERMINAL = {"FAILED", "ROLLED_BACK", "COMMITTED", "RECOVERY_REQUIRED"} ID = re.compile(r"[a-z0-9][a-z0-9-]{0,63}\Z") @@ -133,7 +136,7 @@ def validate_envelope(raw, *, now=None, expired=False): require(type(value["deadline"]) is int, "deadline") if not expired: now = time.time() if now is None else now - require(now < value["deadline"] <= now + TRANSACTION, "deadline") + require(now < value["deadline"] <= now + STOP_START_TRANSACTION, "deadline") return value @@ -446,8 +449,9 @@ def backup_inventory(root, window, deadline): def validate_policy(policy): - require(set(policy) == {"schema_version", "enabled", "runtime_enabled", "platform", "public_url", "configs", - "capacity", "ports", "nginx_include", "static_root"}, "policy_fields") + fields = {"schema_version", "enabled", "runtime_enabled", "platform", "public_url", "configs", + "capacity", "ports", "nginx_include", "static_root"} + require(fields <= set(policy) <= fields | {"stop_start"}, "policy_fields") require(type(policy["schema_version"]) is int and policy["schema_version"] == 1 and policy["enabled"] is True, "not_activated") require(type(policy["runtime_enabled"]) is bool, "runtime_activation") @@ -475,9 +479,32 @@ def validate_policy(policy): "capacity_reserve") else: require(capacity["network_evidence"] is None or matches(HEX, capacity["network_evidence"]), "capacity_evidence") + if "stop_start" in policy: + target = policy["stop_start"] + keys = {"trigger_sha", "runtime_source_sha", "image_digest", "platform_digest", + "configuration_digest", "corpus_id", "archive_sha256", "expires_at"} + require(isinstance(target, dict) and set(target) == keys, "stop_start_fields") + for key in ("trigger_sha", "runtime_source_sha"): + require(matches(SHA, target[key]), "stop_start_target") + for key in ("image_digest", "platform_digest"): + require(matches(DIGEST, target[key]), "stop_start_target") + for key in ("configuration_digest", "corpus_id", "archive_sha256"): + require(matches(HEX, target[key]), "stop_start_target") + require(type(target["expires_at"]) is int and target["expires_at"] > 0, "stop_start_expiry") return policy +def rollout_mode(policy, envelope, *, now=None): + target = policy.get("stop_start") + if target is None: + require(envelope["deadline"] - (time.time() if now is None else now) <= TRANSACTION, "deadline") + return "overlap" + require(all(envelope[key] == value for key, value in target.items() if key != "expires_at"), + "stop_start_target") + require((time.time() if now is None else now) < target["expires_at"], "stop_start_expired") + return "stop-start" + + def verify_descriptors(index_raw, child_raw, envelope, platform): def decoded(raw, expected): # buildx adds a display newline on some versions; only discard it if @@ -666,9 +693,13 @@ def verify(self, envelope, deadline): child = run(["docker", "buildx", "imagetools", "inspect", "--raw", IMAGE + "@" + envelope["platform_digest"]], deadline) return verify_descriptors(index, child, envelope, self.policy["platform"]) - def capacity(self, deadline, *, reclaim_bytes=0): + def capacity(self, deadline, *, reclaim_bytes=0, pre_stop=False): limits = self.policy["capacity"] - self._basic_capacity(deadline, limits["available_memory_bytes"], reclaim_bytes=reclaim_bytes) + require(not pre_stop or reclaim_bytes == 0, "capacity_mode") + # The old container still owns its memory during stop-start preflight. + # Check disk, descriptors and real network evidence before disruption; + # memory must be checked again after the old container has stopped. + self._basic_capacity(deadline, 0 if pre_stop else limits["available_memory_bytes"], reclaim_bytes=reclaim_bytes) evidence = read_file(self.root / "capacity" / (limits["network_evidence"] + ".json")) require(digest(evidence) == limits["network_evidence"], "network_capacity") remaining(deadline) @@ -876,7 +907,8 @@ def switch(self, record, deadline): run(["nginx", "-s", "reload"], deadline) def stop(self, record, deadline): - if self.inspect(record, deadline) is not None: + info = self.inspect(record, deadline) + if info is not None and info["State"]["Running"]: # No automatic restart policy; deterministic named objects survive # controller death and are reconciled, not replaced by unrelated IDs. run(["docker", "stop", "--time", str(max(0, min(DRAIN, int(remaining(deadline)) - 5))), record["name"]], deadline) @@ -1115,20 +1147,23 @@ def deploy(self, raw): if existing: return self.result(existing) require(self.adapter.policy.get("runtime_enabled") is True, "runtime_not_activated") + mode = rollout_mode(self.adapter.policy, envelope) self.adapter.config(envelope) validate_envelope(raw) - require(envelope["deadline"] - time.time() > RECOVERY_RESERVE, "insufficient_transaction_budget") + reserve = STOP_START_RECOVERY_RESERVE if mode == "stop-start" else RECOVERY_RESERVE + require(envelope["deadline"] - time.time() > reserve, "insufficient_transaction_budget") require((self.root / "active.json").is_file(), "predecessor_required") previous = parse(read_file(self.root / "active.json"), 65536) self.adapter.config(previous) - deadline = time.monotonic() + min(TRANSACTION, envelope["deadline"] - time.time()) - work = deadline - RECOVERY_RESERVE + transaction = STOP_START_TRANSACTION if mode == "stop-start" else TRANSACTION + deadline = time.monotonic() + min(transaction, envelope["deadline"] - time.time()) + work = deadline - reserve journal = {"envelope": envelope, "state": "RECEIVED", "intent": "verify", "predecessor": previous, - "candidate": None, "cleanup_complete": False} + "candidate": None, "cleanup_complete": False, "rollout_mode": mode} self.save(journal) try: descriptors = self.adapter.verify(envelope, work) - self.adapter.capacity(work) + self.adapter.capacity(work, pre_stop=mode == "stop-start") manifest = self.adapter.manifest(envelope) self.save(journal, "VERIFIED", "hold_predecessor") token = digest(canonical_json(envelope))[:32] @@ -1144,6 +1179,13 @@ def deploy(self, raw): journal["candidate"] = candidate self.save(journal, intent="pull_candidate") self.adapter.pull(candidate, work) + if mode == "stop-start": + require(remaining(work) >= STOP_START_ACCEPTANCE_BUDGET, + "insufficient_stop_start_budget") + self.save(journal, intent="stop_predecessor") + self.adapter.stop(previous, min(work, time.monotonic() + STOP)) + self.save(journal, intent="capacity_after_stop") + self.adapter.capacity(work) self.save(journal, intent="start_candidate") candidate = self.adapter.hold(candidate, token, min(work, time.monotonic() + READINESS), manifest) journal["candidate"] = candidate @@ -1154,8 +1196,11 @@ def deploy(self, raw): self.adapter.switch(candidate, work) self.save(journal, "SWITCHED", "public_smoke") self.adapter.check(candidate, min(work, time.monotonic() + SMOKE), public=True) - self.save(journal, "COMMITTED", "accept_pointer") - self.cleanup(journal, deadline) + # In stop-start the predecessor is already down. Keep rollback + # possible through the second public check and candidate resume. + if mode != "stop-start": + self.save(journal, "COMMITTED", "accept_pointer") + self.cleanup(journal, work if mode == "stop-start" else deadline) except Exception as error: journal["error_code"] = error.code if isinstance(error, ReleaseError) else "host_failure" self.save(journal) @@ -1197,12 +1242,18 @@ def cleanup(self, journal, deadline): self.adapter.resume(candidate, min(deadline, time.monotonic() + SMOKE)) journal["cleanup_complete"] = True journal.pop("error_code", None) - self.save(journal, intent="complete") + self.save(journal, "COMMITTED" if journal.get("rollout_mode") == "stop-start" else None, "complete") def rollback(self, journal, deadline): switched = journal.get("switch_attempted", False) try: previous = journal["predecessor"] + if journal.get("rollout_mode") == "stop-start" and journal["candidate"]: + # A stopped predecessor cannot be restarted alongside the + # candidate on a small host. Persist before stopping it, also + # when recovery is replaying a crash after candidate start. + self.save(journal, intent="stop_candidate") + self.adapter.stop(journal["candidate"], min(deadline, time.monotonic() + STOP)) self.save(journal, intent="restore_predecessor") # A crash during capture might leave no acknowledged identity. Select # the persisted previously accepted generation, never a newer cache pointer. @@ -1214,7 +1265,7 @@ def rollback(self, journal, deadline): self.adapter.switch(previous, deadline - STOP - SMOKE) self.adapter.check(previous, min(deadline - STOP, time.monotonic() + SMOKE), public=True) write_json(self.root / "active.json", previous) - if journal["candidate"]: + if journal["candidate"] and journal.get("rollout_mode") != "stop-start": self.save(journal, intent="stop_candidate") self.adapter.stop(journal["candidate"], min(deadline, time.monotonic() + STOP)) self.adapter.resume(previous, deadline) @@ -2037,14 +2088,16 @@ def gc(self, *, now=None): return removed -def schedule(kind): +def schedule(kind, *, runtime_max=TRANSACTION): require(kind in {"deploy", "index", "recover", "bootstrap", "initial-install"}, "job_kind") + require(runtime_max in {TRANSACTION, STOP_START_TRANSACTION} + and (runtime_max == TRANSACTION or kind == "deploy"), "job_budget") if kind == "initial-install": require(os.geteuid() == 0, "host_privilege") # Shared unit name and controller lock serialize all host effects. No --pipe, # --wait or inherited SSH stdin; timer recovers a crash before enqueue. return run(["systemd-run", "--unit=v8std-release-job", "--collect", "--no-block", - "--property=Type=exec", "--property=RuntimeMaxSec=300s", "--property=TimeoutStopSec=5s", + "--property=Type=exec", f"--property=RuntimeMaxSec={runtime_max}s", "--property=TimeoutStopSec=5s", "--property=KillMode=control-group", "/usr/bin/python3", "-I", INSTALL, "_" + kind], time.monotonic() + 10) @@ -2057,17 +2110,19 @@ def submit(root, adapter, raw): if existing: return controller.result(existing) require(adapter.policy.get("runtime_enabled") is True, "runtime_not_activated") + mode = rollout_mode(adapter.policy, envelope) adapter.config(envelope) require((Path(root) / "active.json").is_file(), "predecessor_required") validate_envelope(raw) - require(envelope["deadline"] - time.time() > RECOVERY_RESERVE, "insufficient_transaction_budget") + reserve = STOP_START_RECOVERY_RESERVE if mode == "stop-start" else RECOVERY_RESERVE + require(envelope["deadline"] - time.time() > reserve, "insufficient_transaction_budget") pending = Path(root) / "pending-deploy.json" if pending.exists(): previous = parse(read_file(pending)) require(previous == envelope or any(item["envelope"] == previous and item.get("cleanup_complete") for item in controller.journals()), "busy") write_json(pending, envelope) - schedule("deploy") + schedule("deploy", runtime_max=STOP_START_TRANSACTION if mode == "stop-start" else TRANSACTION) return {"state": "QUEUED", "release_id": envelope["release_id"]} @@ -2135,7 +2190,8 @@ def main(): pending.unlink() sync_dir(ROOT) except ReleaseError as error: - if error.code not in {"deadline", "insufficient_transaction_budget", "runtime_not_activated", "predecessor_required", "stale_sequence"}: + if error.code not in {"deadline", "insufficient_transaction_budget", "runtime_not_activated", "predecessor_required", + "stale_sequence", "stop_start_target", "stop_start_expired"}: raise with locked(ROOT): write_json(ROOT / "rejected" / (envelope["release_id"] + ".json"), diff --git a/tests/mcp_release_fixture.py b/tests/mcp_release_fixture.py index ee1f613..05e853d 100644 --- a/tests/mcp_release_fixture.py +++ b/tests/mcp_release_fixture.py @@ -82,8 +82,8 @@ def verify(self, envelope, deadline): return {envelope["image_digest"]: next(iter(release.INDEX_TYPES)), envelope["platform_digest"]: next(iter(release.MANIFEST_TYPES))} - def capacity(self, deadline): - self.record("capacity") + def capacity(self, deadline, *, reclaim_bytes=0, pre_stop=False): + self.record("capacity_pre_stop" if pre_stop else "capacity") def pull(self, record, deadline): self.record("pull", record) @@ -110,6 +110,10 @@ def inspect(self, record, deadline): def start(self, record, deadline): self.record("start", record) + if self.policy.get("stop_start"): + live = [path.stem for path in (self.root / "processes").glob("*.json") + if self.inspect(release.read_record(path)["record"], deadline)["State"]["Running"]] + release.require(not live or live == [record["release_id"]], "fixture_overlap") info = self.inspect(record, deadline) if info and info["State"]["Running"]: return @@ -145,6 +149,10 @@ def check(self, record, deadline, *, public=False): # fail a real HTTP request before rollback can be claimed. if "public_dead" in self.fault.split(",") and public and record["release_id"] != "predecessor": self.stop(record, deadline) + if self.fault == "cleanup_public_dead" and operation == "public": + calls = [json.loads(line) for line in (self.root / "calls.jsonl").read_text().splitlines()] + if sum(call["operation"] == "public" for call in calls) == 2: + self.stop(record, deadline) result = super().check(record, min(deadline, time.monotonic() + 5), public=public) if self.fault == "kill_active_after_smoke" and public: os.kill(os.getpid(), signal.SIGKILL) @@ -153,6 +161,8 @@ def check(self, record, deadline, *, public=False): def resume(self, record, deadline): if self.fault == "kill_active_before_resume": os.kill(os.getpid(), signal.SIGKILL) + if self.fault == "cleanup_resume" and record["release_id"] != "predecessor": + raise release.ReleaseError("injected_resume") return super().resume(record, deadline) def stop(self, record, deadline): @@ -167,6 +177,8 @@ def stop(self, record, deadline): os.waitpid(info["pid"], os.WNOHANG) except (ChildProcessError, TypeError): pass + if self.fault == "crash_after_predecessor_stop" and record["release_id"] == "predecessor": + os._exit(93) class BootstrapAdapter(ProcessAdapter): diff --git a/tests/test_mcp_publication.py b/tests/test_mcp_publication.py index 157514d..e830c4b 100644 --- a/tests/test_mcp_publication.py +++ b/tests/test_mcp_publication.py @@ -1159,6 +1159,34 @@ def smoke(source): self.assertEqual(effects, [] if error else ["live-smoke"]) self.assertEqual(len(limits), 150 if elapsed == 300 else 4 if not error else 1) + def test_stop_start_runtime_uses_extended_bounded_deadline(self): + context = dict(event="push", repository="zeegin/v8std", ref="refs/heads/main", + sha="c" * 40, main_sha="c" * 40, run_id=1, run_number=2, attempt=1, + gates={name: "success" for name in self.p.GATES}) + archive, manifest = fixture.snapshot_fixture() + manifest["archive"]["path"] = "https://ai.v8std.ru/indexes/v1/" + manifest["archive"]["path"] + accepted = {"runtime": {"source_sha": "b" * 40, "image_digest": "sha256:" + "a" * 64}, "manifest": manifest} + adapter = Transport(archive, manifest) + submitted = {} + def command(name, payload, **kwargs): + if name == "status" and not payload: + return {"state": "COMMITTED", "cleanup_complete": True, "image_digest": "sha256:" + "f" * 64} + if name == "validate-envelope": + return json.loads(payload) + if name == "deploy": + submitted.update(json.loads(payload)) + return {"state": "QUEUED", "release_id": submitted["release_id"]} + return {**submitted, "state": "COMMITTED", "cleanup_complete": True, "error_code": None} + adapter.command = command + index = {"manifests": [{"digest": "sha256:" + "e" * 64, + "platform": {"os": "linux", "architecture": "amd64"}}]} + with patch.object(self.p, "registry_manifest", return_value=fixture.json_bytes(index)), \ + patch.object(self.p, "verify_running_runtime"): + result = self.p.deploy_runtime(context, adapter, accepted, enabled=True, + configuration_digest="d" * 64, platform="linux/amd64", stop_start=True) + self.assertEqual(submitted["deadline"], 1_800_000_900) + self.assertEqual(result["state"], "COMMITTED") + def test_immutable_tag_never_overwrites_conflict_and_rechecks_main_before_write(self): digest = "sha256:" + "a" * 64 image = self.p.IMAGE + "@" + digest diff --git a/tests/test_v8std_mcp_release.py b/tests/test_v8std_mcp_release.py index f4a5865..b12def9 100644 --- a/tests/test_v8std_mcp_release.py +++ b/tests/test_v8std_mcp_release.py @@ -1060,6 +1060,15 @@ def health(self): except release.ReleaseError: return None + def enable_stop_start(self): + self.env["deadline"] = int(time.time()) + release.STOP_START_TRANSACTION + self.policy["stop_start"] = {key: self.env[key] for key in ( + "trigger_sha", "runtime_source_sha", "image_digest", "platform_digest", + "configuration_digest", "corpus_id", "archive_sha256")} + self.policy["stop_start"]["expires_at"] = int(time.time()) + 3600 + release.write_json(self.root / "policy.json", self.policy) + release.write_json(self.root / "envelope.json", self.env) + def invoke(self, fault="", mode="deploy"): result = subprocess.run([sys.executable, "-m", "tests.mcp_release_fixture", mode, str(self.root), fault], cwd=ROOT, capture_output=True, timeout=25) @@ -1079,6 +1088,88 @@ def test_success_real_mcp_health_and_independent_runtime_corpus_shas(self): self.assertEqual(self.health()["corpus_source_sha"], self.source.manifest["source_sha"]) self.assertFalse(self.adapter.inspect(self.previous, time.monotonic() + 1)["State"]["Running"]) + def test_stop_start_updates_without_overlapping_runtimes(self): + self.enable_stop_start() + before = len((self.root / "calls.jsonl").read_text().splitlines()) + result = self.invoke() + self.assertEqual(result["state"], "COMMITTED", result) + self.assertTrue(result["cleanup_complete"]) + self.assertEqual(self.health()["runtime_sha"], self.env["runtime_source_sha"]) + calls = [json.loads(line) for line in (self.root / "calls.jsonl").read_text().splitlines()[before:]] + operations = [call["operation"] for call in calls] + self.assertLess(operations.index("capacity_pre_stop"), operations.index("pull")) + self.assertLess(next(i for i, call in enumerate(calls) if call == + {"operation": "stop", "release_id": "predecessor"}), operations.index("capacity")) + self.assertLess(operations.index("capacity"), next(i for i, call in enumerate(calls) if call == + {"operation": "start", "release_id": "release-1"})) + + def test_stop_start_post_stop_capacity_failure_restores_predecessor(self): + self.enable_stop_start() + result = self.invoke("capacity") + self.assertEqual(result["state"], "FAILED", result) + self.assertEqual(result["error_code"], "injected_capacity") + self.assertEqual(self.health()["runtime_sha"], self.previous["runtime_source_sha"]) + + def test_stop_start_public_failure_stops_candidate_before_restoring_predecessor(self): + self.enable_stop_start() + result = self.invoke("public_dead") + self.assertEqual(result["state"], "ROLLED_BACK", result) + self.assertEqual(self.health()["runtime_sha"], self.previous["runtime_source_sha"]) + + def test_stop_start_cleanup_public_failure_still_restores_predecessor(self): + self.enable_stop_start() + result = self.invoke("cleanup_public_dead") + self.assertEqual(result["state"], "ROLLED_BACK", result) + self.assertEqual(self.health()["runtime_sha"], self.previous["runtime_source_sha"]) + + def test_stop_start_resume_failure_after_pointer_restores_predecessor(self): + self.enable_stop_start() + result = self.invoke("cleanup_resume") + self.assertEqual(result["state"], "ROLLED_BACK", result) + self.assertEqual(result["error_code"], "injected_resume") + self.assertEqual(release.read_record(self.root / "active.json")["release_id"], "predecessor") + self.assertEqual(self.health()["runtime_sha"], self.previous["runtime_source_sha"]) + + def test_stop_start_crash_after_old_stop_recovers_predecessor(self): + self.enable_stop_start() + self.invoke("crash_after_predecessor_stop") + self.assertIsNone(self.health()) + result = self.invoke(mode="recover") + self.assertEqual(result["state"], "FAILED", result) + self.assertEqual(self.health()["runtime_sha"], self.previous["runtime_source_sha"]) + + def test_stop_start_crash_after_candidate_start_stops_it_before_recovery(self): + self.enable_stop_start() + self.invoke("crash_after_start") + result = self.invoke(mode="recover") + self.assertEqual(result["state"], "FAILED", result) + self.assertEqual(self.health()["runtime_sha"], self.previous["runtime_source_sha"]) + + def test_stop_start_requires_exact_unexpired_target_before_scheduling(self): + self.enable_stop_start() + with patch.object(release, "schedule") as schedule: + self.policy["stop_start"]["image_digest"] = "sha256:" + "f" * 64 + with self.assertRaisesRegex(release.ReleaseError, "stop_start_target"): + release.submit(self.root, self.adapter, release.canonical_json(self.env)) + self.policy["stop_start"]["image_digest"] = self.env["image_digest"] + self.policy["stop_start"]["expires_at"] = int(time.time()) - 1 + with self.assertRaisesRegex(release.ReleaseError, "stop_start_expired"): + release.submit(self.root, self.adapter, release.canonical_json(self.env)) + schedule.assert_not_called() + + def test_stop_start_target_expiring_after_submit_is_durably_rejected(self): + self.enable_stop_start() + with patch.object(release, "schedule"): + self.assertEqual(release.submit(self.root, self.adapter, release.canonical_json(self.env))["state"], "QUEUED") + self.policy["stop_start"]["expires_at"] = int(time.time()) + 60 + release.write_json(self.root / "policy.json", self.policy) + self.invoke(mode="queued_expired") + query = {"schema_version": 1, "kind": "release", "id": self.env["release_id"]} + result = release.query_status(self.root, self.adapter, query) + self.assertEqual(result["state"], "REJECTED", result) + self.assertEqual(result["error_code"], "stop_start_expired") + self.assertFalse((self.root / "pending-deploy.json").exists()) + def test_smoke_accepts_real_held_tools_only_runtime(self): health = release.smoke(self.adapter.url(self.previous), self.previous, time.monotonic() + 30) self.assertEqual(health["hold_token"], self.previous["hold_token"]) From 363f040da7429df62a71b664d53e27b8d7d76e64 Mon Sep 17 00:00:00 2001 From: Igor Apresov Date: Fri, 25 Sep 2026 17:52:22 +0300 Subject: [PATCH 2/2] Preserve ordinary release adapters and dynamic job budgets --- delivery/vps/v8std_mcp_release.py | 8 ++++++-- 1 file changed, 6 insertions(+), 2 deletions(-) diff --git a/delivery/vps/v8std_mcp_release.py b/delivery/vps/v8std_mcp_release.py index de3b9f8..d512d6c 100644 --- a/delivery/vps/v8std_mcp_release.py +++ b/delivery/vps/v8std_mcp_release.py @@ -1163,7 +1163,10 @@ def deploy(self, raw): self.save(journal) try: descriptors = self.adapter.verify(envelope, work) - self.adapter.capacity(work, pre_stop=mode == "stop-start") + if mode == "stop-start": + self.adapter.capacity(work, pre_stop=True) + else: + self.adapter.capacity(work) manifest = self.adapter.manifest(envelope) self.save(journal, "VERIFIED", "hold_predecessor") token = digest(canonical_json(envelope))[:32] @@ -2088,8 +2091,9 @@ def gc(self, *, now=None): return removed -def schedule(kind, *, runtime_max=TRANSACTION): +def schedule(kind, *, runtime_max=None): require(kind in {"deploy", "index", "recover", "bootstrap", "initial-install"}, "job_kind") + runtime_max = TRANSACTION if runtime_max is None else runtime_max require(runtime_max in {TRANSACTION, STOP_START_TRANSACTION} and (runtime_max == TRANSACTION or kind == "deploy"), "job_budget") if kind == "initial-install":