From bff2dda4e8ef870f5cbcb14a3054d7ba6d3ffb85 Mon Sep 17 00:00:00 2001 From: root Date: Sun, 30 Aug 2026 17:10:48 +0000 Subject: [PATCH] Repair surrendered Azure worker reconciliation --- bin/fm-azure-worker-provider.py | 62 ++++++- bin/fm-worker-lifecycle.py | 213 ++++++++++++++++++++--- tests/fm-worker-lifecycle.test.sh | 278 +++++++++++++++++++++++++++++- 3 files changed, 520 insertions(+), 33 deletions(-) diff --git a/bin/fm-azure-worker-provider.py b/bin/fm-azure-worker-provider.py index 26f912e76f8..189520a1a5d 100755 --- a/bin/fm-azure-worker-provider.py +++ b/bin/fm-azure-worker-provider.py @@ -55,6 +55,7 @@ REQUEST_SCHEMA = "fm.worker-provider-request/v1" RESPONSE_SCHEMA = "fm.worker-provider-response/v1" INVENTORY_SCHEMA = "fm.worker-provider-inventory/v1" +RELEASED_ABSENCE_SCHEMA = "fm.worker-released-absence/v1" EXECUTION_TERMINAL_SCHEMA = "fm.worker-execution-terminal/v1" EXECUTION_RESULT_SCHEMA = "fm.worker-execution-result/v1" EXECUTE_DISPOSITION_SUBMIT = "submit" @@ -273,6 +274,10 @@ "monitor-extension", "bootstrap-command", "task-command", "ttl-schedule", }) READY_CHILD_KINDS = frozenset({"monitor-extension", "bootstrap-command"}) +RELEASED_ABSENT_RESOURCE_KINDS = ( + "vm", "nic", "os-disk", "task-disk", "account-disk", + "monitor-extension", "bootstrap-command", "task-command", "ttl-schedule", +) class ProviderError(RuntimeError): @@ -3046,11 +3051,37 @@ def mutate_delete_compute(controller, action): return final +def released_absent_data_reset(action): + """Validate the controller's narrow, idempotency-bound absence adoption.""" + evidence = action.get("released_absent_data") + if evidence is None: + return False + expected = { + "schema": RELEASED_ABSENCE_SCHEMA, + "assignment_generation": (action.get("bindings") or {}).get( + "assignment_generation" + ), + "release_proof_digest": action.get("release_proof_digest"), + "absent_resources": list(RELEASED_ABSENT_RESOURCE_KINDS), + } + if evidence != expected: + raise ProviderError("released-absence reset evidence is not exact") + return True + + +def reset_worker_inventory(controller, action, released_absent_data): + """Read one reset snapshot, preserving ambiguity refusal on every replay step.""" + snapshot = inventory_slot(controller, action["slot"]) + if released_absent_data and snapshot.get("conflicts"): + raise ProviderError("released-absence reset inventory is ambiguous") + return worker_by_slot(snapshot, action["slot"]) + + def mutate_reset(controller, action): if not re.match(r"^[0-9a-f]{64}$", str(action.get("release_proof_digest", ""))): raise ProviderError("reset requires the exact ordinary release-proof digest") - snapshot = inventory_slot(controller, action["slot"]) - worker = worker_by_slot(snapshot, action["slot"]) + released_absent_data = released_absent_data_reset(action) + worker = reset_worker_inventory(controller, action, released_absent_data) if worker is None: return None resources = worker.get("resources") or {} @@ -3065,20 +3096,37 @@ def mutate_reset(controller, action): "vm", "nic", "os-disk", "monitor-extension", "bootstrap-command", "task-command", "ttl-schedule", ) + allowed_missing = disposable + if released_absent_data: + allowed_missing += RELEASED_ABSENT_RESOURCE_KINDS + ( + "staging-request", "staging-result", + ) resources = cleanup_recorded_exact( - action, worker, allow_missing=disposable, require_ready_children=False + action, worker, allow_missing=allowed_missing, require_ready_children=False ) + if released_absent_data and any( + kind in resources for kind in RELEASED_ABSENT_RESOURCE_KINDS + ): + raise ProviderError( + "released-absence reset refuses while data or compute still exists" + ) if any(kind in resources for kind in disposable): raise ProviderError("reset refuses while disposable compute still exists") mark_cleanup_container( controller, action, "reset-action", action["idempotency_key"] ) - worker = worker_by_slot(inventory_slot(controller, action["slot"]), action["slot"]) + worker = reset_worker_inventory(controller, action, released_absent_data) allow_missing = tuple(kind for kind in REQUIRED_RESOURCE_KINDS if kind != "state-container") resources = cleanup_recorded_exact( action, worker, allow_missing=allow_missing, skip_immutable=("state-container",), require_ready_children=False, ) + if released_absent_data and any( + kind in resources for kind in RELEASED_ABSENT_RESOURCE_KINDS + ): + raise ProviderError( + "released-absence reset refuses while data or compute still exists" + ) if any(kind in resources for kind in ( "vm", "nic", "os-disk", "monitor-extension", "bootstrap-command", "task-command", "ttl-schedule", @@ -3138,7 +3186,7 @@ def delete_archive(blob_name=blob_name): "staging archive deletion failed for {}: {}".format(blob_name, stderr)) independent.append((blob_name, delete_archive)) run_independent_cleanup(independent) - refreshed = worker_by_slot(inventory_slot(controller, action["slot"]), action["slot"]) + refreshed = reset_worker_inventory(controller, action, released_absent_data) if refreshed is None: raise ProviderError("cleanup marker container disappeared before exact reset completed") state_container = (refreshed.get("resources") or {}).get("state-container") @@ -3164,8 +3212,8 @@ def delete_archive(blob_name=blob_name): ], check=False) if rc != 0: raise ProviderError("exact worker state-container deletion failed: {}".format(stderr)) - final = inventory_slot(controller, action["slot"]) - if worker_by_slot(final, action["slot"]) is not None: + final = reset_worker_inventory(controller, action, released_absent_data) + if final is not None: raise ProviderError("released worker capacity remains after exact reset") return None diff --git a/bin/fm-worker-lifecycle.py b/bin/fm-worker-lifecycle.py index b05c2b78b61..f55b10aadf1 100755 --- a/bin/fm-worker-lifecycle.py +++ b/bin/fm-worker-lifecycle.py @@ -61,6 +61,7 @@ EXECUTION_TERMINAL_SCHEMA = "fm.worker-execution-terminal/v1" EXECUTE_ABANDON_MARKER = "execute-abandon-action" RELEASE_SCHEMA = "fm.worker-release/v2" +RELEASED_ABSENCE_SCHEMA = "fm.worker-released-absence/v1" AUTHORITY_SCHEMA = "fm.worker-authority/v1" CAPACITY_RESERVATION_SCHEMA = "fm.capacity-reservation/v1" CAPACITY_FENCE_RETIREMENT_SCHEMA = "fm.capacity-fence-retirement/v1" @@ -141,6 +142,17 @@ MUTABLE_PROVISIONING_CHILD_KINDS = frozenset({ "monitor-extension", "bootstrap-command", "task-command", "ttl-schedule", }) +MISSING_COMPUTE_RESOURCE_KINDS = frozenset({ + "vm", "nic", "os-disk", "monitor-extension", "bootstrap-command", + "task-command", "ttl-schedule", "staging-request", "staging-result", +}) +# This is deliberately a complete observed-absence shape, not another generic +# missing-resource allowance. Only an exact surrendered release can use it, +# and the provider re-observes the same absence before touching residuals. +RELEASED_ABSENT_RESOURCE_KINDS = ( + "vm", "nic", "os-disk", "task-disk", "account-disk", + "monitor-extension", "bootstrap-command", "task-command", "ttl-schedule", +) REVIEWED_SKU_FAMILY = { "Standard_D4as_v6": "standardDav6Family", "Standard_D4as_v7": "StandardDasv7Family", @@ -1502,7 +1514,7 @@ def resource_identity(resource): } -def resources_exact(worker, cloud, allow_missing_compute=False): +def _resources_exact(worker, cloud, allowed_missing): resources = cloud.get("resources") or {} recorded = worker.get("resources") or {} missing = [] @@ -1510,10 +1522,7 @@ def resources_exact(worker, cloud, allow_missing_compute=False): current = resources.get(kind) prior = recorded.get(kind) if current is None: - if allow_missing_compute and kind in ( - "vm", "nic", "os-disk", "monitor-extension", "bootstrap-command", - "task-command", "ttl-schedule", "staging-request", "staging-result", - ): + if kind in allowed_missing: continue missing.append(kind) continue @@ -1553,6 +1562,33 @@ def resources_exact(worker, cloud, allow_missing_compute=False): return True, "" +def resources_exact(worker, cloud, allow_missing_compute=False): + allowed_missing = MISSING_COMPUTE_RESOURCE_KINDS if allow_missing_compute else () + return _resources_exact(worker, cloud, allowed_missing) + + +def released_absent_data_resources_exact(worker, cloud): + """Prove the one reset-adoption inventory shape without widening absence.""" + if not isinstance(cloud, dict) or cloud.get("slot") != worker.get("slot"): + return False, "released absence is not bound to the exact worker slot" + resources = cloud.get("resources") or {} + unknown = sorted(set(resources) - set(REQUIRED_RESOURCE_KINDS)) + if unknown: + return False, "released absence inventory has unknown resources: {}".format( + ", ".join(unknown) + ) + present = [kind for kind in RELEASED_ABSENT_RESOURCE_KINDS if resources.get(kind) is not None] + if present: + return False, "released absence still has data or compute resources: {}".format( + ", ".join(present) + ) + return _resources_exact( + worker, + cloud, + MISSING_COMPUTE_RESOURCE_KINDS.union({"task-disk", "account-disk"}), + ) + + def classify_worker(worker, cloud, now=None): now = now or now_utc() if cloud is None: @@ -1587,6 +1623,93 @@ def classify_worker(worker, cloud, now=None): return "assigned", "one exact active task generation" +def recorded_unlanded_executions(state, worker): + """Return exact-assignment executions whose repository outcome is not clean.""" + produced = [] + bindings = worker.get("bindings") or {} + for request_digest, execution in sorted((state.get("executions") or {}).items()): + if not isinstance(execution, dict): + continue + if ( + execution.get("task") != bindings.get("task") + or execution.get("task_generation") != bindings.get("task_generation") + or execution.get("assignment_generation") != worker.get("assignment_generation") + ): + continue + commits = execution.get("outcome_commits") + if ( + execution.get("outcome_present") not in (None, False) + or execution.get("outcome_uncommitted_changes") not in (None, False) + or commits not in (None, 0) + ): + produced.append(request_digest) + return produced + + +def released_absent_data_adoption(state, worker, cloud): + """Recognize only an exact surrender whose data and compute are already gone. + + The release proof remains bound to the durable pre-absence identities. A + fresh provider inventory must still prove every residual identity and tag; + absence supplies no authority for an unreleased or outcome-bearing worker. + """ + exact, reason = released_absent_data_resources_exact(worker, cloud) + if not exact: + return False, reason + if worker.get("phase") != "release-proved": + return False, "released absence worker is not in release-proved phase" + proof = worker.get("release_proof") + try: + verify_surrender_release_against_worker(proof, worker) + except LifecycleError as exc: + return False, "released absence proof is not exact: {}".format(exc) + item = state.get("queue", {}).get(worker.get("queue_key")) + bindings = worker.get("bindings") or {} + if ( + not isinstance(item, dict) + or worker.get("queue_key") != request_key( + bindings.get("task", ""), bindings.get("task_generation", "") + ) + or item.get("status") != "releasing" + or item.get("slot") != worker.get("slot") + or item.get("assignment_generation") != worker.get("assignment_generation") + ): + return False, "released absence has no exact releasing queue owner" + recorded = recorded_unlanded_executions(state, worker) + discarded = (proof.get("surrender") or {}).get("discarded_unlanded_executions") + if recorded or discarded: + return False, "released absence has recorded unlanded outcome: {}".format( + ", ".join(recorded or discarded) + ) + return True, "exact surrendered worker data and compute are authoritatively absent" + + +def reconciled_worker_classification(state, worker, cloud, now=None): + classification, note = classify_worker(worker, cloud, now=now) + if classification != "retained-for-investigation" or cloud is None: + return classification, note + adopted, adoption_note = released_absent_data_adoption(state, worker, cloud) + if adopted: + return "orphaned-safe-to-delete", adoption_note + # Name a failed candidate precisely, but preserve the ordinary classifier's + # reason for every shape that still has compute or data. + if not any( + (cloud.get("resources") or {}).get(kind) is not None + for kind in RELEASED_ABSENT_RESOURCE_KINDS + ): + return classification, adoption_note + return classification, note + + +def released_absence_evidence(worker): + return { + "schema": RELEASED_ABSENCE_SCHEMA, + "assignment_generation": worker["assignment_generation"], + "release_proof_digest": worker["release_proof"]["proof_digest"], + "absent_resources": list(RELEASED_ABSENT_RESOURCE_KINDS), + } + + def adopt_cloud_resources(worker, cloud): resources = cloud.get("resources") or {} if set(resources) != set(REQUIRED_RESOURCE_KINDS): @@ -2603,7 +2726,9 @@ def refresh_classifications(state, inventory, now=None): # An unapplied mutation owns this slot; a display value derived # from a record whose durable phase is not yet true would lie. continue - classification, note = classify_worker(worker, cloud.get(worker["slot"]), now=now) + classification, note = reconciled_worker_classification( + state, worker, cloud.get(worker["slot"]), now=now + ) worker["last_classification"] = classification worker["classification_note"] = note @@ -2690,7 +2815,9 @@ def next_reconcile_action(env, state, inventory, now=None): ) continue current = cloud.get(worker["slot"]) - classification, note = classify_worker(worker, current, now=now) + classification, note = reconciled_worker_classification( + state, worker, current, now=now + ) worker["last_classification"] = classification worker["classification_note"] = note if worker.get("role") in RETIRED_WORKER_ROLES and not worker.get("release_proof"): @@ -2731,10 +2858,13 @@ def next_reconcile_action(env, state, inventory, now=None): ) continue if classification == "orphaned-safe-to-delete": - return make_action( - env, "reset", worker=worker, - release_proof_digest=worker["release_proof"]["proof_digest"], - ) + fields = { + "release_proof_digest": worker["release_proof"]["proof_digest"], + } + adopted, _ = released_absent_data_adoption(state, worker, current) + if adopted: + fields["released_absent_data"] = released_absence_evidence(worker) + return make_action(env, "reset", worker=worker, **fields) if not waiting: return None @@ -2891,10 +3021,56 @@ def verify_release_against_worker(proof, worker): "resources": worker["resources"], } for field, expected in checks.items(): - if proof.get(field) != expected: + if not isinstance(proof, dict) or proof.get(field) != expected: raise LifecycleError("worker release proof {} binding is not exact".format(field)) +def verify_surrender_release_against_worker(proof, worker): + """Validate the stored surrender bytes before absence can authorize reset.""" + if not isinstance(proof, dict) or proof.get("schema") != RELEASE_SCHEMA: + raise LifecycleError("surrender release proof schema is not exact") + surrender = proof.get("surrender") + if not isinstance(surrender, dict): + raise LifecycleError("release proof is not a surrender proof") + supplied = proof.get("proof_digest") + unsigned = dict(proof) + unsigned.pop("proof_digest", None) + if supplied != digest_value(unsigned): + raise LifecycleError("surrender release proof digest is not exact") + verify_release_against_worker(proof, worker) + for field in ("home_binding", "account_binding", "worktree_binding", "repository_binding"): + require_binding(field, proof.get(field)) + for field in ("task", "task_generation", "assignment_generation", "repository_generation"): + require_id(field, proof.get(field)) + resources = proof.get("resources") + if not isinstance(resources, dict) or set(resources) != set(REQUIRED_RESOURCE_KINDS): + raise LifecycleError("surrender release proof does not bind every worker resource") + for kind, identity in resources.items(): + if not isinstance(identity, dict) or not identity.get("id") or not identity.get("immutable_id"): + raise LifecycleError("surrender release proof {} identity is incomplete".format(kind)) + authorities = proof.get("authorities") + expected_authorities = {"endpoint", "report", "landing", "account", "worktree"} + if not isinstance(authorities, dict) or set(authorities) != expected_authorities: + raise LifecycleError("surrender release proof authority set is not exact") + for name, receipt in authorities.items(): + if not isinstance(receipt, dict) or receipt.get("schema") != AUTHORITY_SCHEMA: + raise LifecycleError("{} surrender authority schema is not exact".format(name)) + receipt_unsigned = dict(receipt) + receipt_digest = receipt_unsigned.pop("receipt_digest", None) + if receipt_digest != digest_value(receipt_unsigned): + raise LifecycleError("{} surrender authority digest is not exact".format(name)) + if ( + receipt.get("authority") != name + or receipt.get("task") != proof.get("task") + or receipt.get("task_generation") != proof.get("task_generation") + or receipt.get("assignment_generation") != proof.get("assignment_generation") + or receipt.get("verdict") != "surrendered" + or receipt.get("evidence_digest") + != digest_value({"authority": name, "surrender": surrender}) + ): + raise LifecycleError("{} surrender authority binding is not exact".format(name)) + + def proof_template(state, task, generation): key = request_key(task, generation) item = state["queue"].get(key) @@ -4961,18 +5137,7 @@ def command_surrender(env, args): raise LifecycleError("worker already has an ordinary release proof; reconcile releases it") if item.get("status") != "assigned": raise LifecycleError("surrender requires one exact assigned task generation") - produced = [] - for request_digest, execution in sorted((state.get("executions") or {}).items()): - if not isinstance(execution, dict): - continue - if (execution.get("task") != args.task - or execution.get("task_generation") != args.task_generation - or execution.get("assignment_generation") != worker["assignment_generation"]): - continue - if (execution.get("outcome_present") is True - or execution.get("outcome_uncommitted_changes") is True - or (execution.get("outcome_commits") or 0) > 0): - produced.append(request_digest) + produced = recorded_unlanded_executions(state, worker) if produced and not args.confirm_discard_unlanded: # The controller's own durable record says this worker produced # repository work whose landing is unproven; surrendering it leads diff --git a/tests/fm-worker-lifecycle.test.sh b/tests/fm-worker-lifecycle.test.sh index 25269c7a866..72690e9a38a 100755 --- a/tests/fm-worker-lifecycle.test.sh +++ b/tests/fm-worker-lifecycle.test.sh @@ -960,6 +960,102 @@ except module.ProviderIdentityRefusal as exc: else: raise AssertionError("deallocate cleanup accepted a foreign monitor-extension ID") module.inventory = original_inventory + +# A reset may adopt only the controller's exact surrendered-absence shape. +# The real reset reaches its durable marker with exact identity/storage +# residuals and both data disks absent. Missing/malformed evidence, a disk that +# still exists, residual identity drift, and a provider conflict all refuse +# before that marker can authorize deletion. +absence_action = copy.deepcopy(action) +absence_action.update({ + "type": "reset", + "release_proof_digest": "6" * 64, +}) +absence_action["released_absent_data"] = { + "schema": module.RELEASED_ABSENCE_SCHEMA, + "assignment_generation": absence_action["bindings"]["assignment_generation"], + "release_proof_digest": absence_action["release_proof_digest"], + "absent_resources": list(module.RELEASED_ABSENT_RESOURCE_KINDS), +} +absence_action["idempotency_key"] = hashlib.sha256(module.canonical_bytes( + {key: value for key, value in absence_action.items() if key != "idempotency_key"} +)).hexdigest() +residual = copy.deepcopy(worker) +for kind in module.RELEASED_ABSENT_RESOURCE_KINDS: + residual["resources"].pop(kind, None) + +class MarkerReached(Exception): + pass + +original_inventory_slot = module.inventory_slot +original_mark_cleanup_container = module.mark_cleanup_container +module.inventory_slot = lambda *_a, **_k: { + "workers": [copy.deepcopy(residual)], "conflicts": [], "metrics": {}, +} +module.mark_cleanup_container = lambda *_a, **_k: (_ for _ in ()).throw(MarkerReached()) +try: + module.mutate_reset(controller, absence_action) +except MarkerReached: + pass +else: + raise AssertionError("released-absence reset did not reach its fenced residual marker") + +without_evidence = copy.deepcopy(absence_action) +without_evidence.pop("released_absent_data") +try: + module.mutate_reset(controller, without_evidence) +except module.ProviderError as exc: + assert "task-disk resource is absent" in str(exc), exc +else: + raise AssertionError("data-disk absence authorized reset without release evidence") + +malformed_evidence = copy.deepcopy(absence_action) +malformed_evidence["released_absent_data"]["assignment_generation"] = "asg-foreign" +try: + module.mutate_reset(controller, malformed_evidence) +except module.ProviderError as exc: + assert "evidence is not exact" in str(exc), exc +else: + raise AssertionError("malformed released-absence evidence authorized reset") + +attached_residual = copy.deepcopy(residual) +attached_residual["resources"]["task-disk"] = copy.deepcopy(worker["resources"]["task-disk"]) +module.inventory_slot = lambda *_a, **_k: { + "workers": [copy.deepcopy(attached_residual)], "conflicts": [], "metrics": {}, +} +try: + module.mutate_reset(controller, absence_action) +except module.ProviderError as exc: + assert "data or compute still exists" in str(exc), exc +else: + raise AssertionError("released-absence reset accepted an attached task disk") + +foreign_residual = copy.deepcopy(residual) +foreign_residual["resources"]["identity"]["immutable_id"] = "foreign-principal" +module.inventory_slot = lambda *_a, **_k: { + "workers": [copy.deepcopy(foreign_residual)], "conflicts": [], "metrics": {}, +} +try: + module.mutate_reset(controller, absence_action) +except module.ProviderIdentityRefusal as exc: + assert "identity immutable identity differs" in str(exc), exc +else: + raise AssertionError("released-absence reset accepted residual identity drift") + +module.inventory_slot = lambda *_a, **_k: { + "workers": [copy.deepcopy(residual)], + "conflicts": [{"slot": 1, "kind": "unknown", "reason": "ambiguous"}], + "metrics": {}, +} +try: + module.mutate_reset(controller, absence_action) +except module.ProviderError as exc: + assert "inventory is ambiguous" in str(exc), exc +else: + raise AssertionError("released-absence reset accepted ambiguous provider inventory") +module.inventory_slot = original_inventory_slot +module.mark_cleanup_container = original_mark_cleanup_container + # Ordinary provider failures stay on exit 2. If every ProviderError were # mislabeled as a permanent identity refusal, abandon-claim could clear a # transiently failed claim. @@ -4247,14 +4343,52 @@ assert controller_state()["workers"][slot]["release_proof"] == proof rewritten = json.loads((Path(env["FM_HOME"]) / "surrender-1.json").read_text()) assert rewritten == proof, "the idempotent rerun did not re-issue the exact stored proof" -# Reconcile now converges the surrendered slot to nothing through the ordinary -# fenced machinery: delete-compute, then reset. +# The deliberately discarded worker keeps its existing ordinary convergence +# path: deallocate/delete-compute/reset. for _ in range(4): run("reconcile", "--apply", "--confirm-subscription", env["FM_AZURE_SUBSCRIPTION_ID"]) state = controller_state() assert state["workers"] == {}, state["workers"] assert state["queue"]["task-1@gen-1"]["status"] == "complete" assert slot not in json.loads(Path(fixture_path).read_text())["workers"] + +# A second, clean surrender models the live recovery seam: Azure has already +# removed the exact VM, NIC, OS disk, both data disks, and every VM child while +# exact identity/storage residuals remain. Reconcile adopts that +# authority-bound absence and resets directly, without inventing or replaying +# a compute deletion. +run( + "request", "--task", "task-2", "--task-generation", "gen-2", + "--home-binding", binding(1002), "--account-binding", binding(2002), + "--worktree-binding", binding(3002), "--repository-binding", binding(4002), + "--repository-generation", "repo-2", "--owner-kind", "primary", "--eligible", +) +run("reconcile", "--apply", "--confirm-subscription", env["FM_AZURE_SUBSCRIPTION_ID"]) +state = controller_state() +slot2 = str(state["queue"]["task-2@gen-2"]["slot"]) +fixture = json.loads(Path(fixture_path).read_text()) +fixture["workers"][slot2]["resources"]["vm"]["power_state"] = "VM deallocated" +Path(fixture_path).write_text(json.dumps(fixture, sort_keys=True, separators=(",", ":")) + "\n") +run( + "surrender", "--task", "task-2", "--task-generation", "gen-2", + "--reason", "local metadata was consumed after the replacement landed", + "--output", str(Path(env["FM_HOME"]) / "surrender-2.json"), + "--confirm-surrender", *confirm, +) +fixture = json.loads(Path(fixture_path).read_text()) +for kind in ( + "vm", "nic", "os-disk", "task-disk", "account-disk", + "monitor-extension", "bootstrap-command", "task-command", "ttl-schedule", +): + fixture["workers"][slot2]["resources"].pop(kind, None) +calls_before = len(fixture["calls"]) +Path(fixture_path).write_text(json.dumps(fixture, sort_keys=True, separators=(",", ":")) + "\n") +run("reconcile", "--apply", "--confirm-subscription", env["FM_AZURE_SUBSCRIPTION_ID"]) +state = controller_state() +fixture = json.loads(Path(fixture_path).read_text()) +assert state["queue"]["task-2@gen-2"]["status"] == "complete" +assert slot2 not in state["workers"] and slot2 not in fixture["workers"] +assert [call["type"] for call in fixture["calls"][calls_before:]] == ["reset"], fixture["calls"] PY pass "surrender releases an authority-less worker through refusal-first gates and ordinary reset" } @@ -8134,6 +8268,146 @@ def build_assigned(slot): } return build_state, worker, item, cloud, inventory +def surrender_proof(worker, task=None): + bindings = worker["bindings"] + task = task or bindings["task"] + surrender = { + "reason": "ordinary authority is gone after a landed replacement", + "ordinary_refusal": "WORKER AUTHORITY REFUSED: fixture metadata is absent", + "surrendered_at": module.iso_utc(T0 - dt.timedelta(hours=1)), + "power_state": "vm deallocated", + "last_execution_digest": worker.get("last_execution_digest"), + "discarded_unlanded_executions": [], + } + proof = { + "schema": module.RELEASE_SCHEMA, + "home_binding": bindings["home_binding"], + "task": task, + "task_generation": bindings["task_generation"], + "assignment_generation": worker["assignment_generation"], + "account_binding": bindings["account_binding"], + "worktree_binding": bindings["worktree_binding"], + "repository_binding": bindings["repository_binding"], + "repository_generation": bindings["repository_generation"], + "cloud_instance_id": worker["cloud_instance_id"], + "resources": copy.deepcopy(worker["resources"]), + "surrender": surrender, + "authorities": {}, + } + for name in ("endpoint", "report", "landing", "account", "worktree"): + receipt = { + "schema": module.AUTHORITY_SCHEMA, + "authority": name, + "task": task, + "task_generation": bindings["task_generation"], + "assignment_generation": worker["assignment_generation"], + "verdict": "surrendered", + "evidence_digest": module.digest_value({"authority": name, "surrender": surrender}), + } + receipt["receipt_digest"] = module.digest_value(receipt) + proof["authorities"][name] = receipt + proof["proof_digest"] = module.digest_value(proof) + return proof + +def mark_surrendered(state, worker, item, task=None): + worker["release_proof"] = surrender_proof(worker, task=task) + worker["released_at"] = module.iso_utc(T0 - dt.timedelta(hours=1)) + worker["phase"] = "release-proved" + item["status"] = "releasing" + return state, worker, item + +def released_absent_fixture(slot=3): + state, worker, item, cloud, inventory = build_assigned(slot) + mark_surrendered(state, worker, item) + for kind in module.RELEASED_ABSENT_RESOURCE_KINDS: + cloud["resources"].pop(kind, None) + state["executions"] = {} + inventory = dict(inventory, workers=[cloud]) + return state, worker, item, cloud, inventory + +# A release-proved surrender whose exact VM, NIC, OS disk, both data disks, +# and compute children are already absent adopts only the exact residuals and +# plans reset directly. The action carries a proof/assignment-bound shape that +# the provider re-observes before deleting anything. +rstate, rworker, ritem, rcloud, rinventory = released_absent_fixture() +classification, note = module.reconciled_worker_classification( + rstate, rworker, rcloud, now=T0 +) +assert classification == "orphaned-safe-to-delete", (classification, note) +rreset = module.next_reconcile_action(penv, rstate, rinventory, now=T0) +assert rreset["type"] == "reset", rreset +assert rreset["released_absent_data"] == { + "schema": module.RELEASED_ABSENCE_SCHEMA, + "assignment_generation": rworker["assignment_generation"], + "release_proof_digest": rworker["release_proof"]["proof_digest"], + "absent_resources": list(module.RELEASED_ABSENT_RESOURCE_KINDS), +}, rreset + +# Absence is never authority on its own. Unreleased state, a missing proof, +# any durable unlanded outcome, residual identity drift, live compute, an +# attached data disk, a proof bound to another task, and provider ambiguity all +# stay outside the released-absence reset path. +refusal_state, refusal_worker, refusal_item, refusal_cloud, refusal_inventory = ( + released_absent_fixture(4) +) +refusal_worker["release_proof"] = None +refusal_worker["phase"] = "assigned" +refusal_item["status"] = "assigned" +assert module.next_reconcile_action( + penv, refusal_state, refusal_inventory, now=T0 +) is None + +missing_state, missing_worker, missing_item, _missing_cloud, missing_inventory = ( + released_absent_fixture(5) +) +missing_worker["release_proof"] = None +assert module.next_reconcile_action(penv, missing_state, missing_inventory, now=T0) is None + +outcome_state, outcome_worker, _outcome_item, _outcome_cloud, outcome_inventory = ( + released_absent_fixture(6) +) +outcome_state["executions"]["e" * 64] = { + "task": outcome_worker["bindings"]["task"], + "task_generation": outcome_worker["bindings"]["task_generation"], + "assignment_generation": outcome_worker["assignment_generation"], + "outcome_present": False, + "outcome_commits": 0, + "outcome_uncommitted_changes": True, +} +assert module.next_reconcile_action(penv, outcome_state, outcome_inventory, now=T0) is None + +identity_state, _identity_worker, _identity_item, identity_cloud, identity_inventory = ( + released_absent_fixture(7) +) +identity_cloud["resources"]["identity"]["immutable_id"] = "foreign-principal" +assert module.next_reconcile_action(penv, identity_state, identity_inventory, now=T0) is None + +live_state, live_worker, live_item, _live_cloud, live_inventory = build_assigned(8) +mark_surrendered(live_state, live_worker, live_item) +live_action = module.next_reconcile_action(penv, live_state, live_inventory, now=T0) +assert live_action["type"] == "deallocate", live_action +assert "released_absent_data" not in live_action + +attached_state, attached_worker, attached_item, attached_cloud, attached_inventory = build_assigned(9) +mark_surrendered(attached_state, attached_worker, attached_item) +for kind in module.RELEASED_ABSENT_RESOURCE_KINDS: + if kind != "task-disk": + attached_cloud["resources"].pop(kind, None) +assert attached_cloud["resources"]["task-disk"]["attached_to"] +assert module.next_reconcile_action(penv, attached_state, attached_inventory, now=T0) is None + +proof_state, proof_worker, proof_item, _proof_cloud, proof_inventory = released_absent_fixture(10) +mark_surrendered(proof_state, proof_worker, proof_item, task="another-task") +assert module.next_reconcile_action(penv, proof_state, proof_inventory, now=T0) is None + +conflict_state, _conflict_worker, _conflict_item, _conflict_cloud, conflict_inventory = ( + released_absent_fixture(11) +) +conflict_inventory["conflicts"] = [ + {"slot": 11, "kind": "unknown-child", "reason": "ambiguous residual"} +] +assert module.next_reconcile_action(penv, conflict_state, conflict_inventory, now=T0) is None + # --- one exact conflicted compartment cannot freeze unrelated cleanup or be # selected for fresh admission. The conflicting slot itself stays parked. cstate, cworker, citem, ccloud, cinventory = build_assigned(2)