diff --git a/bin/fm-azure-worker-provider.py b/bin/fm-azure-worker-provider.py index 0d4b05adba8..64912628ad6 100755 --- a/bin/fm-azure-worker-provider.py +++ b/bin/fm-azure-worker-provider.py @@ -1085,7 +1085,12 @@ def metrics(controller, vms, capacity_reservations, specialized_active_by_family } -def inventory(controller, include_metrics=True): +def inventory(controller, include_metrics=True, target_slot=None): + if target_slot is not None and ( + not isinstance(target_slot, int) or isinstance(target_slot, bool) + or not 1 <= target_slot <= 16 + ): + raise ProviderError("target worker slot is outside the reviewed fleet") account, rc, stderr = az(controller, ["account", "show"], check=False) if rc != 0 or not isinstance(account, dict): raise ProviderError("Azure scope is unreadable: {}".format(stderr)) @@ -1142,15 +1147,15 @@ def add(kind, value, slot, power=None, tags_override=None): for vm in vms: slot = slot_from_name(vm.get("name"), r"^vm-{}-wkr-".format(prefix)) - if slot is not None: + if slot is not None and (target_slot is None or slot == target_slot): add("vm", vm, slot, vm.get("powerState") or vm.get("power_state") or "unknown") for nic in nics: slot = slot_from_name(nic.get("name"), r"^nic-{}-wkr-".format(prefix)) - if slot is not None: + if slot is not None and (target_slot is None or slot == target_slot): add("nic", nic, slot) for disk in disks: slot = slot_from_name(disk.get("name"), r"^disk-{}-wkr-".format(prefix)) - if slot is None: + if slot is None or (target_slot is not None and slot != target_slot): continue name = str(disk.get("name")) if name.endswith("-os"): @@ -1166,7 +1171,7 @@ def add(kind, value, slot, power=None, tags_override=None): identity_principals = {} for identity in identities: slot = slot_from_name(identity.get("name"), r"^id-{}-wkr-".format(prefix)) - if slot is not None: + if slot is not None and (target_slot is None or slot == target_slot): add("identity", identity, slot) principal = immutable_id("identity", identity) if principal: @@ -1180,7 +1185,7 @@ def add(kind, value, slot, power=None, tags_override=None): if not match: continue slot = int(match.group(1)) - if not 1 <= slot <= 16: + if not 1 <= slot <= 16 or (target_slot is not None and slot != target_slot): continue metadata = partial_container_metadata( metadata_to_tags(container.get("metadata") or {}), slot, workers, @@ -1218,7 +1223,10 @@ def add(kind, value, slot, power=None, tags_override=None): for extension in extensions: slot = slot_from_name(extension.get("name"), r"^vm-{}-wkr-".format(prefix)) - if slot is None or not str(extension.get("name", "")).endswith("/AzureMonitorLinuxAgent"): + if ( + slot is None or (target_slot is not None and slot != target_slot) + or not str(extension.get("name", "")).endswith("/AzureMonitorLinuxAgent") + ): continue vm_id = exact_id( controller, "Microsoft.Compute", "virtualMachines", @@ -1235,7 +1243,7 @@ def add(kind, value, slot, power=None, tags_override=None): add("monitor-extension", value, slot) for command in run_commands: slot = slot_from_name(command.get("name"), r"^vm-{}-wkr-".format(prefix)) - if slot is None: + if slot is None or (target_slot is not None and slot != target_slot): continue child = str(command.get("name", "")).rsplit("/", 1)[-1] kind = {"bootstrap": "bootstrap-command", "execute": "task-command"}.get(child) @@ -1254,7 +1262,7 @@ def add(kind, value, slot, power=None, tags_override=None): add(kind, value, slot) for schedule in schedules: slot = slot_from_name(schedule.get("name"), r"^shutdown-computevm-vm-{}-wkr-".format(prefix)) - if slot is not None: + if slot is not None and (target_slot is None or slot == target_slot): value = show_full( controller, schedule["id"], api_version="2018-09-15", inventory_missing_ok=True, @@ -1274,7 +1282,7 @@ def add(kind, value, slot, power=None, tags_override=None): if not match: continue slot = int(match.group(1)) - if not 1 <= slot <= 16: + if not 1 <= slot <= 16 or (target_slot is not None and slot != target_slot): continue principal = role.get("principalId") if identity_principals.get(slot) != principal: @@ -1313,9 +1321,12 @@ def add(kind, value, slot, power=None, tags_override=None): if nic and nic.get("attached_to") and vm and nic["attached_to"].lower() != vm["id"].lower(): conflicts.append({"kind": "nic", "slot": slot, "reason": "NIC is attached to another VM"}) - capacity_reservations, specialized_active_by_family = specialized_capacity_inventory( - controller, vms, identities - ) + if include_metrics: + capacity_reservations, specialized_active_by_family = specialized_capacity_inventory( + controller, vms, identities + ) + else: + capacity_reservations, specialized_active_by_family = [], {} result = { "schema": INVENTORY_SCHEMA, "observed_at": iso_utc(), @@ -1340,6 +1351,11 @@ def add(kind, value, slot, power=None, tags_override=None): return result +def inventory_slot(controller, slot): + """Read one exact worker compartment without expanding peer children.""" + return inventory(controller, include_metrics=False, target_slot=slot) + + def worker_by_slot(snapshot, slot): matches = [worker for worker in snapshot["workers"] if worker["slot"] == slot] if len(matches) > 1: @@ -2439,7 +2455,7 @@ def run_pilot_create(controller, action): def converge_create_tags(controller, action): - snapshot = inventory(controller, include_metrics=False) + snapshot = inventory_slot(controller, action["slot"]) if snapshot["conflicts"]: raise ProviderError("new worker inventory contains foreign or unsafe resources") worker = worker_by_slot(snapshot, action["slot"]) @@ -2454,7 +2470,7 @@ def converge_create_tags(controller, action): continue tag_resource(controller, resource["id"], tags) tag_container(controller, expected_names(controller, action["slot"])["state-container"], tags) - snapshot = inventory(controller, include_metrics=False) + snapshot = inventory_slot(controller, action["slot"]) if snapshot["conflicts"]: raise ProviderError("tagged worker inventory contains foreign or unsafe resources") worker = worker_by_slot(snapshot, action["slot"]) @@ -2465,7 +2481,7 @@ def converge_create_tags(controller, action): def create_or_resume(controller, action): - snapshot = inventory(controller, include_metrics=False) + snapshot = inventory_slot(controller, action["slot"]) if snapshot["conflicts"]: raise ProviderError("same-name foreign worker resources refuse create/adopt") existing = worker_by_slot(snapshot, action["slot"]) @@ -2486,7 +2502,7 @@ def create_or_resume(controller, action): # for, and it is the one that skips the create path entirely. if ensure_worker_running(controller, expected_names(controller, action["slot"])["vm"]): existing = worker_by_slot( - inventory(controller, include_metrics=False), action["slot"] + inventory_slot(controller, action["slot"]), action["slot"] ) return existing except ProviderError: @@ -2594,7 +2610,7 @@ def run_independent_cleanup(operations): def mutate_deallocate(controller, action): - snapshot = inventory(controller, include_metrics=False) + snapshot = inventory_slot(controller, action["slot"]) resources = cleanup_recorded_exact( action, worker_by_slot(snapshot, action["slot"]), require_ready_children=False ) @@ -2606,7 +2622,7 @@ def mutate_deallocate(controller, action): ], check=False) if rc != 0: raise ProviderError("exact worker deallocation failed: {}".format(stderr)) - final = worker_by_slot(inventory(controller, include_metrics=False), action["slot"]) + final = worker_by_slot(inventory_slot(controller, action["slot"]), action["slot"]) final_resources = cleanup_recorded_exact(action, final, require_ready_children=False) if "deallocated" not in str(final_resources["vm"].get("power_state", "")).lower(): raise ProviderError("worker compute did not reach Azure deallocated state") @@ -2632,7 +2648,7 @@ def service_cancel_allows_missing_task_command(action): def mutate_delete_compute(controller, action): - snapshot = inventory(controller, include_metrics=False) + snapshot = inventory_slot(controller, action["slot"]) worker = worker_by_slot(snapshot, action["slot"]) if worker is None: raise ProviderError("compute cleanup lost exact retained task/account ownership") @@ -2669,7 +2685,7 @@ def mutate_delete_compute(controller, action): mark_cleanup_container( controller, action, "compute-action", action["idempotency_key"] ) - worker = worker_by_slot(inventory(controller, include_metrics=False), action["slot"]) + worker = worker_by_slot(inventory_slot(controller, action["slot"]), action["slot"]) # ttl-schedule may be absent only on re-entry after the VM already # cascaded away (Azure deletes shutdown-computevm schedules with their # target VM); the fresh-entry path above still required it alongside the @@ -2683,7 +2699,7 @@ def mutate_delete_compute(controller, action): ) if resources.get("ttl-schedule") is None and resources.get("vm") is not None: raise ProviderError("TTL disappeared while the worker VM still exists") - worker = worker_by_slot(inventory(controller, include_metrics=False), action["slot"]) + worker = worker_by_slot(inventory_slot(controller, action["slot"]), action["slot"]) if worker is None: raise ProviderError("VM deletion also lost exact retained task/account capacity") remaining = worker.get("resources") or {} @@ -2707,7 +2723,7 @@ def mutate_delete_compute(controller, action): wait_absent(controller, resource["id"]) # NIC/disk attach relations only clear once the VM is gone, so the # detach proof below must read a fresh snapshot. - refreshed = worker_by_slot(inventory(controller, include_metrics=False), action["slot"]) + refreshed = worker_by_slot(inventory_slot(controller, action["slot"]), action["slot"]) if refreshed is None: raise ProviderError("VM deletion also lost exact retained task/account capacity") remaining = refreshed.get("resources") or {} @@ -2729,7 +2745,7 @@ def mutate_delete_compute(controller, action): continue conditional_delete(controller, kind, resource) run_independent_cleanup(detached) - final = worker_by_slot(inventory(controller, include_metrics=False), action["slot"]) + final = worker_by_slot(inventory_slot(controller, action["slot"]), action["slot"]) if final is None: raise ProviderError("compute cleanup lost retained task/account ownership") final_resources = final.get("resources") or {} @@ -2741,7 +2757,7 @@ def mutate_delete_compute(controller, action): ttl = final_resources.get("ttl-schedule") if ttl is not None: conditional_delete(controller, "ttl-schedule", ttl) - final = worker_by_slot(inventory(controller, include_metrics=False), action["slot"]) + final = worker_by_slot(inventory_slot(controller, action["slot"]), action["slot"]) if final is None: raise ProviderError("TTL cleanup lost retained task/account ownership") # An absent TTL here is the Azure cascade outcome: shutdown-computevm @@ -2758,7 +2774,7 @@ def mutate_delete_compute(controller, action): 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(controller, include_metrics=False) + snapshot = inventory_slot(controller, action["slot"]) worker = worker_by_slot(snapshot, action["slot"]) if worker is None: return None @@ -2782,7 +2798,7 @@ def mutate_reset(controller, action): mark_cleanup_container( controller, action, "reset-action", action["idempotency_key"] ) - worker = worker_by_slot(inventory(controller, include_metrics=False), action["slot"]) + worker = worker_by_slot(inventory_slot(controller, action["slot"]), action["slot"]) 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",), @@ -2847,7 +2863,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(controller, include_metrics=False), action["slot"]) + refreshed = worker_by_slot(inventory_slot(controller, action["slot"]), action["slot"]) if refreshed is None: raise ProviderError("cleanup marker container disappeared before exact reset completed") state_container = (refreshed.get("resources") or {}).get("state-container") @@ -2873,7 +2889,7 @@ def delete_archive(blob_name=blob_name): ], check=False) if rc != 0: raise ProviderError("exact worker state-container deletion failed: {}".format(stderr)) - final = inventory(controller, include_metrics=False) + final = inventory_slot(controller, action["slot"]) if worker_by_slot(final, action["slot"]) is not None: raise ProviderError("released worker capacity remains after exact reset") return None @@ -3125,7 +3141,7 @@ def abandon_execute(controller, action): if not isinstance(expected_task, str) or not expected_task: raise ProviderError("execute abandonment action carries no exact task-command identity") - worker = worker_by_slot(inventory(controller, include_metrics=False), action["slot"]) + worker = worker_by_slot(inventory_slot(controller, action["slot"]), action["slot"]) if worker is None: raise ProviderError("execute abandonment lost exact retained worker ownership") current = worker.get("resources") or {} @@ -3160,7 +3176,7 @@ def abandon_execute(controller, action): mark_cleanup_container( controller, action, EXECUTE_ABANDON_MARKER, action["idempotency_key"] ) - bracketed = worker_by_slot(inventory(controller, include_metrics=False), action["slot"]) + bracketed = worker_by_slot(inventory_slot(controller, action["slot"]), action["slot"]) bracketed_resources = recorded_exact( action, bracketed, skip_immutable=("state-container",), require_ready_children=False, ) @@ -3340,7 +3356,7 @@ def persist_execute_result(controller, action, names, tags, execution): def mutate_execute(controller, action): - snapshot = inventory(controller, include_metrics=False) + snapshot = inventory_slot(controller, action["slot"]) worker = worker_by_slot(snapshot, action["slot"]) resources = recorded_exact(action, worker) durable_retired_key = action.get("retired_execute_key") @@ -3378,7 +3394,7 @@ def mutate_execute(controller, action): if disposition == EXECUTE_DISPOSITION_RECOVERED: persist_execute_result(controller, action, names, tags, recovered) worker = worker_by_slot( - inventory(controller, include_metrics=False), action["slot"] + inventory_slot(controller, action["slot"]), action["slot"] ) if worker is None: raise ProviderError("execution result persistence lost its exact worker") @@ -3482,14 +3498,14 @@ def mutate_execute(controller, action): if execution is None: raise ProviderError("private worker execution returned no exact result") persist_execute_result(controller, action, names, tags, execution) - worker = worker_by_slot(inventory(controller, include_metrics=False), action["slot"]) + worker = worker_by_slot(inventory_slot(controller, action["slot"]), action["slot"]) if worker is None: raise ProviderError("execution result persistence lost its exact worker") return worker, execution def mutate_steer(controller, action): - snapshot = inventory(controller, include_metrics=False) + snapshot = inventory_slot(controller, action["slot"]) resources = recorded_exact(action, worker_by_slot(snapshot, action["slot"])) if "deallocated" in str(resources["vm"].get("power_state", "")).lower(): raise ProviderError("steer refuses deallocated worker compute") @@ -3526,7 +3542,7 @@ def mutate_steer(controller, action): or ack.get("assignment_generation") != bindings["assignment_generation"] ): raise ProviderError("guest supervisor did not acknowledge the exact steer request") - return worker_by_slot(inventory(controller, include_metrics=False), action["slot"]) + return worker_by_slot(inventory_slot(controller, action["slot"]), action["slot"]) def validate_mutation_action(controller, action): @@ -3620,6 +3636,19 @@ def main(): operation = request.get("operation") if operation == "inventory": value = response(controller, operation, inventory=inventory(controller)) + elif operation == "inventory-slot": + selector = request.get("action") + if ( + not isinstance(selector, dict) or set(selector) != {"slot"} + or not isinstance(selector.get("slot"), int) + or isinstance(selector.get("slot"), bool) + or not 1 <= selector["slot"] <= 16 + ): + raise ProviderError("slot inventory requires one exact reviewed slot") + value = response( + controller, operation, slot=selector["slot"], + inventory=inventory_slot(controller, selector["slot"]), + ) elif operation == "mutate": require_landed_code() value = response(controller, operation, result=mutate(controller, request.get("action"))) diff --git a/bin/fm-worker-lifecycle.py b/bin/fm-worker-lifecycle.py index ef9ca72ca37..b90df7767b1 100755 --- a/bin/fm-worker-lifecycle.py +++ b/bin/fm-worker-lifecycle.py @@ -1282,6 +1282,17 @@ def provider_call(env, operation, action=None): return _provider_call_raw(env, operation, action) +def provider_slot_inventory(env, slot): + response = provider_call(env, "inventory-slot", {"slot": slot}) + if response.get("slot") != slot: + raise LifecycleError("provider slot inventory binding is not exact") + value = response.get("inventory") + verify_inventory(value) + if any(worker.get("slot") != slot for worker in value["workers"]): + raise LifecycleError("provider slot inventory escaped its exact compartment") + return value + + def _provider_call_raw(env, operation, action=None): request = { "schema": PROVIDER_REQUEST_SCHEMA, @@ -1332,7 +1343,7 @@ def verify_provider_response(env, operation, response): expected = env[field] if controller.get(field) != expected: raise LifecycleError("provider response {} binding is not exact".format(field)) - if operation == "inventory": + if operation in ("inventory", "inventory-slot"): verify_inventory(response.get("inventory")) @@ -2679,7 +2690,6 @@ def next_service_reconcile_action(env, state, inventory, task, generation, now=N raise LifecycleError( "provider found same-fleet worker-name conflicts; unrelated resources were not adopted" ) - roll_daily_baseline(state, inventory["metrics"].get("actual_usd"), now) key = request_key(task, generation) item = state["queue"].get(key) if item is None or item.get("role") != "no-mistakes": @@ -2688,6 +2698,7 @@ def next_service_reconcile_action(env, state, inventory, task, generation, now=N if status == "complete": return None if status == "queued": + roll_daily_baseline(state, inventory["metrics"].get("actual_usd"), now) if not item.get("eligible"): raise LifecycleError("service reconcile refuses an ineligible request") if active_count(state, inventory) >= env["max_workers"]: @@ -3898,12 +3909,18 @@ def command_service_reconcile(env, args): print("pending: slot={} task={}@{}".format(slot_key, task, generation)) return - inventory = provider_call(env, "inventory")["inventory"] + if item.get("status") == "queued": + inventory = provider_call(env, "inventory")["inventory"] + else: + if worker is None or slot_key is None: + raise LifecycleError("service reconcile task has no exact durable worker owner") + inventory = provider_slot_inventory(env, worker["slot"]) action = None with contextlib.ExitStack() as stack: with controller_lock(env): state = load_state(env) - state["last_metrics"] = metrics_from_inventory(inventory) + if item.get("status") == "queued": + state["last_metrics"] = metrics_from_inventory(inventory) action = next_service_reconcile_action( env, state, inventory, task, generation ) diff --git a/docs/azure-workers.md b/docs/azure-workers.md index 39aa29dbaa0..9e38c12a241 100644 --- a/docs/azure-workers.md +++ b/docs/azure-workers.md @@ -404,7 +404,7 @@ The guest verifies the runtime's exact file inventory, runs the role command wit The wrapper verifies the semantic bytes and head binding before writing the controller-facing result; a process exit, missing outcome, malformed outcome, or changed read-only head is a failed result, never `CLEAR` by inference. Repair results return one digest-bound single-ref bundle whose head must descend from the requested head, while review and test return no code bundle and must keep the exact requested head. The wrapper records a retryable local candidate before cleanup, releases through `service-complete` only after the lifecycle owns the exact execution result, and replays the candidate after a lost response instead of executing the step again. -Admission, execute recovery, and cleanup use `service-reconcile`, which advances only the caller's exact task generation or replays that task's own pending slot claim rather than converging unrelated fleet work. +Admission, execute recovery, and cleanup use `service-reconcile`, which advances only the caller's exact task generation or replays that task's own pending slot claim rather than converging unrelated fleet work. Once a service task owns a slot, its execution and cleanup inventory expands only that exact slot's Azure children; queued admission retains the whole-fleet quota, spend, and conflict census. The guest supervisor marks a no-mistakes Azure execution as the already-isolated test boundary, so the repository test command runs the focused service suite directly instead of recursively provisioning the general validation fleet or a Herdr lab. The root-owned supervisor stages the job, then runs the no-mistakes process as the dedicated non-root `fmworker` user with no supplementary groups; the sealed runtime remains root-owned and read-only while the exact repository and projected account are writable only by that service identity. `bin/fm-azure-service-test-scope.py` owns that focused inventory and the narrow source set eligible for focused pull-request CI; an empty, mixed, or unknown diff and every push to `main` retain the complete behavior suite. diff --git a/tests/fm-azure-pilot.test.sh b/tests/fm-azure-pilot.test.sh index 58fe2df45f5..2a011cf4cca 100755 --- a/tests/fm-azure-pilot.test.sh +++ b/tests/fm-azure-pilot.test.sh @@ -984,7 +984,7 @@ execution["result_digest"] = provider.hashlib.sha256( ).hexdigest() worker = {"slot": 1, "resources": {"vm": {"power_state": "VM running"}}} -provider.inventory = lambda controller, include_metrics=True: {"workers": [worker]} +provider.inventory = lambda controller, include_metrics=True, target_slot=None: {"workers": [worker]} provider.worker_by_slot = lambda snapshot, slot: worker provider.recorded_exact = lambda action, worker, **kwargs: worker["resources"] provider.action_tags = lambda controller, action: {} @@ -1299,7 +1299,7 @@ assert [c for c in calls if c[:2] == ["vm", "deallocate"]], ( # the case the change exists for unguarded. existing_worker = {"slot": 1, "resources": {"vm": {"tags": {}}}} provider.az, calls = make_az(["PowerState/deallocated", "PowerState/running"]) -provider.inventory = lambda controller, include_metrics=True: { +provider.inventory = lambda controller, include_metrics=True, target_slot=None: { "conflicts": [], "workers": [existing_worker], } provider.worker_by_slot = lambda snapshot, slot: existing_worker diff --git a/tests/fm-worker-lifecycle.test.sh b/tests/fm-worker-lifecycle.test.sh index bab9168375f..01247bb5423 100755 --- a/tests/fm-worker-lifecycle.test.sh +++ b/tests/fm-worker-lifecycle.test.sh @@ -1085,13 +1085,13 @@ else: deallocated_transition = copy.deepcopy(transitioned) deallocated_transition["resources"]["vm"]["power_state"] = "VM deallocated" original_inventory = module.inventory -module.inventory = lambda controller_arg, include_metrics=False: { +module.inventory = lambda controller_arg, include_metrics=False, target_slot=None: { "workers": [deallocated_transition], "conflicts": [], "metrics": {}, } assert module.mutate_deallocate(controller, action) == deallocated_transition foreign_transition = copy.deepcopy(deallocated_transition) foreign_transition["resources"]["monitor-extension"]["id"] = "/foreign/monitor-extension" -module.inventory = lambda controller_arg, include_metrics=False: { +module.inventory = lambda controller_arg, include_metrics=False, target_slot=None: { "workers": [foreign_transition], "conflicts": [], "metrics": {}, } try: @@ -1639,7 +1639,7 @@ partial_module.create_lifecycle_children=lambda controller, action: partial_call partial_module.converge_create_tags=lambda controller, action: partial_calls.append("converge") or {"slot": 1, "resources": {}} partial_action={"slot":1,"bindings":{"home_binding":"h"*64,"task":"task-a","task_generation":"g","assignment_generation":"asg-1","account_binding":"a"*64,"worktree_binding":"w"*64,"repository_binding":"r"*64,"repository_generation":"rg"},"sku":"Standard_D4as_v6","sku_family":"standardDav6Family","shared_admission_digest":"x","type":"create"} partial_worker={"slot":1,"resources":{"nic":{"id":"/nic","immutable_id":"n","tags":{"home-binding":"h"*64,"task-binding":"task-a","invocation-binding":"asg-1"}},"task-disk":{"id":"/d","immutable_id":"d","tags":{}}}} -partial_module.inventory=lambda controller, include_metrics=True: {"workers":[partial_worker],"conflicts":[],"capacity_reservations":[],"metrics":{}} +partial_module.inventory=lambda controller, include_metrics=True, target_slot=None: {"workers":[partial_worker],"conflicts":[],"capacity_reservations":[],"metrics":{}} partial_module.worker_by_slot=lambda snapshot, slot: partial_worker partial_module.create_or_resume({"prefix":"fmtest"}, partial_action) assert partial_calls==["create","children","converge"], partial_calls @@ -2241,13 +2241,25 @@ else: "workers": [state["workers"][key] for key in sorted(state["workers"], key=int)], "capacity_reservations": state.get("capacity_reservations", []), "conflicts": [], "metrics": metrics, } + if request["operation"] == "inventory-slot": + slot = request.get("action", {}).get("slot") + inventory["workers"] = [ + worker for worker in inventory["workers"] if worker.get("slot") == slot + ] + inventory["capacity_reservations"] = [] + response_slot = slot result = inventory response = { "schema": "fm.worker-provider-response/v1", "operation": request["operation"], "controller": controller, } -response["inventory" if request["operation"] == "inventory" else "result"] = result +if request["operation"] in ("inventory", "inventory-slot"): + response["inventory"] = result + if request["operation"] == "inventory-slot": + response["slot"] = response_slot +else: + response["result"] = result print(json.dumps(response, sort_keys=True, separators=(",", ":"))) PY chmod +x "$1" @@ -9191,8 +9203,94 @@ PY pass "independent exact Azure cleanup mutations overlap and fail deterministically" } +targeted_slot_inventory_contract() { + python3 - "$AZURE" "$CONTROLLER" <<'PY' || fail "targeted Azure slot inventory contract failed" +import importlib.util +import sys + +provider_spec = importlib.util.spec_from_file_location("azure_provider_slot", sys.argv[1]) +provider = importlib.util.module_from_spec(provider_spec) +provider_spec.loader.exec_module(provider) +lifecycle_spec = importlib.util.spec_from_file_location("lifecycle_slot", sys.argv[2]) +lifecycle = importlib.util.module_from_spec(lifecycle_spec) +lifecycle_spec.loader.exec_module(lifecycle) + +controller = { + "subscription": "11111111-1111-4111-8111-111111111111", + "resource_group": "rg", "prefix": "fixture", + "deployment_generation": "dep", "owner": "owner", "home_binding": "h" * 64, +} +tags = { + "workload": "firstmate", "deployment-generation": "dep", "cleanup-owner": "owner", +} +vms = [ + {"name": "vm-fixture-wkr-01", "id": "/vm/1", "vmId": "vm-id-1", + "powerState": "VM running", "tags": tags}, + {"name": "vm-fixture-wkr-02", "id": "/vm/2", "vmId": "vm-id-2", + "powerState": "VM running", "tags": tags}, +] +extensions = [ + {"name": "vm-fixture-wkr-01/AzureMonitorLinuxAgent", "id": "/vm/1/ext"}, + {"name": "vm-fixture-wkr-02/AzureMonitorLinuxAgent", "id": "/vm/2/ext"}, +] +expanded = [] + +provider.az = lambda _controller, args, check=False, timeout=provider.AZ_TIMEOUT_SECONDS: ( + ({"id": controller["subscription"], "state": "Enabled"}, 0, "") + if args[:2] == ["account", "show"] else (None, 1, "unexpected") +) +def listing(_controller, args, transient_not_found_attempts=1): + if args[:2] == ["vm", "list"]: + return vms + if args[:2] == ["resource", "list"] and args[-1] == "Microsoft.Compute/virtualMachines/extensions": + return extensions + return [] +provider.list_json = listing +def show(_controller, resource_id, api_version=None, inventory_missing_ok=False): + expanded.append(resource_id) + return { + "id": resource_id, "tags": tags, + "properties": {"provisioningState": "Succeeded"}, + } +provider.show_full = show + +snapshot = provider.inventory_slot(controller, 2) +assert [worker["slot"] for worker in snapshot["workers"]] == [2], snapshot +assert expanded == ["/vm/2/ext"], expanded +assert snapshot["capacity_reservations"] == [] +assert snapshot["metrics"]["actual_usd"] is None +try: + provider.inventory_slot(controller, 17) +except provider.ProviderError: + pass +else: + raise AssertionError("targeted inventory accepted a slot outside the fleet") + +inventory = dict(snapshot) +response = { + "schema": lifecycle.PROVIDER_RESPONSE_SCHEMA, "operation": "inventory-slot", + "controller": {key: controller[key] for key in ( + "home_binding", "subscription", "deployment_generation", "owner", "prefix", "resource_group" + )}, + "slot": 2, "inventory": inventory, +} +env = dict(response["controller"]) +lifecycle.provider_call = lambda _env, operation, action=None: response +assert lifecycle.provider_slot_inventory(env, 2) == inventory +response["slot"] = 1 +try: + lifecycle.provider_slot_inventory(env, 2) +except lifecycle.LifecycleError: + pass +else: + raise AssertionError("controller accepted a differently bound slot inventory") +PY + pass "service lifecycle expands and validates only its exact Azure slot" +} + static_contract independent_cleanup_contract +targeted_slot_inventory_contract service_complete_front_door service_reconcile_scope_contract service_complete_replay_contract diff --git a/tests/fm-worker-outcome-transport.test.sh b/tests/fm-worker-outcome-transport.test.sh index 9e3e37f7227..95556d53a58 100755 --- a/tests/fm-worker-outcome-transport.test.sh +++ b/tests/fm-worker-outcome-transport.test.sh @@ -331,7 +331,7 @@ execution["result_digest"] = hashlib.sha256( # Only the collaborators that read live Azure state are substituted; every # line of mutate_execute own body runs for real. worker = {"slot": 1, "resources": {"vm": {"power_state": "VM running"}}} -provider.inventory = lambda controller, include_metrics=True: {"workers": [worker]} +provider.inventory = lambda controller, include_metrics=True, target_slot=None: {"workers": [worker]} provider.worker_by_slot = lambda snapshot, slot: worker provider.recorded_exact = lambda action, worker, **kwargs: worker["resources"] provider.action_tags = lambda controller, action: {}