diff --git a/bin/fm-no-mistakes-worker b/bin/fm-no-mistakes-worker index e0f333f55cc..b997c0dbb7f 100755 --- a/bin/fm-no-mistakes-worker +++ b/bin/fm-no-mistakes-worker @@ -403,6 +403,48 @@ def request_assignment(config, task, generation): raise WrapperError("Azure worker assignment timed out{}".format(detail), "assignment_timeout", True) +def execute_service(config, task, generation, assignment, account_home, staged, outcome_dir, argv): + subscription = config["lifecycle_env"]["FM_AZURE_SUBSCRIPTION_ID"] + deadline = time.monotonic() + config["wall_seconds"] + 1800 + last_error = None + while time.monotonic() < deadline: + remaining = deadline - time.monotonic() + try: + return lifecycle( + config, "execute", "--task", task, "--task-generation", generation, + "--assignment-generation", assignment, + "--wall-seconds", str(config["wall_seconds"]), + "--payload-dir", str(staged), "--account-dir", account_home, + "--outcome-dir", str(outcome_dir), "--confirm-execute", + "--confirm-subscription", subscription, "--", *argv, + timeout=max(1, math.ceil(remaining)), + ) + except WrapperError as error: + if not error.retryable: + raise + last_error = error + remaining = deadline - time.monotonic() + if remaining <= 0: + break + try: + lifecycle( + config, "reconcile", "--apply", "--confirm-subscription", subscription, + timeout=max(1, math.ceil(remaining)), + ) + except WrapperError as error: + if not error.retryable: + raise + last_error = error + remaining = deadline - time.monotonic() + if remaining > 0: + time.sleep(min(config["poll_seconds"], remaining)) + detail = ": {}".format(last_error) if last_error is not None else "" + raise WrapperError( + "Azure worker execution exhausted its exact-identity retry window{}".format(detail), + "execution_timeout", False, + ) + + def extract_service_return(root, execution, request): bundle = root / "outcome" / "outcome.bundle" if execution.get("service_return_present") is not True or not bundle.is_file(): @@ -659,13 +701,9 @@ def execute(args): assignment, account_home = request_assignment(config, task, generation) outcome_dir = root / "outcome" outcome_dir.mkdir(exist_ok=True, mode=0o700) - output = lifecycle( - config, "execute", "--task", task, "--task-generation", generation, - "--assignment-generation", assignment, "--wall-seconds", str(config["wall_seconds"]), - "--payload-dir", str(staged), "--account-dir", account_home, - "--outcome-dir", str(outcome_dir), "--confirm-execute", - "--confirm-subscription", config["lifecycle_env"]["FM_AZURE_SUBSCRIPTION_ID"], - "--", *request["guest_argv"], timeout=config["wall_seconds"] + 1800, + output = execute_service( + config, task, generation, assignment, account_home, staged, outcome_dir, + request["guest_argv"], ) execution = last_json(output, "worker execution") request_digest = execution.get("request_digest") diff --git a/tests/fm-no-mistakes-worker.test.sh b/tests/fm-no-mistakes-worker.test.sh index abcb9fd76c8..7ecb31807fa 100755 --- a/tests/fm-no-mistakes-worker.test.sh +++ b/tests/fm-no-mistakes-worker.test.sh @@ -91,6 +91,13 @@ elif command == "status": }] print(json.dumps({"account_placements": placements}, separators=(",", ":"))) elif command == "execute": + failures = root / "execute-failures" + if failures.exists(): + remaining = int(failures.read_text().strip()) + if remaining > 0: + failures.write_text(str(remaining - 1)) + print("fixture transient execute refusal", file=sys.stderr) + raise SystemExit(75) request_digest = "d" * 64 payload = Path(value("--payload-dir")) outcome_dir = Path(value("--outcome-dir")) @@ -443,6 +450,34 @@ assert len(reconciles) >= 3, reconciles PY pass "wrapper resumes transient reconciliation under one durable Azure task identity" +printf 'ok\n' > "$TMP_ROOT/mode" +printf '2\n' > "$TMP_ROOT/execute-failures" +: > "$TMP_ROOT/calls.log" +write_request "$TMP_ROOT/request-transient-execute.json" job-transient-execute +"$WRAPPER" --config "$TMP_ROOT/config.json" execute \ + --request "$TMP_ROOT/request-transient-execute.json" --payload "$PAYLOAD" \ + --result "$TMP_ROOT/result-transient-execute.json" \ + --outcome "$TMP_ROOT/outcome-transient-execute.bundle" \ + --step-outcome "$TMP_ROOT/step-outcome-transient-execute.json" \ + || fail "transient execute refusals escaped the durable worker invocation" +python3 - "$TMP_ROOT/result-transient-execute.json" "$TMP_ROOT/calls.log" <<'PY' \ + || fail "transient execute retries duplicated or lost the assigned Azure job" +import json, pathlib, re, sys +result = json.loads(pathlib.Path(sys.argv[1]).read_text()) +calls = pathlib.Path(sys.argv[2]).read_text().splitlines() +requests = [line for line in calls if line.startswith("request ")] +executes = [line for line in calls if line.startswith("execute ")] +identities = { + re.search(r"--task ([^ ]+) --task-generation ([^ ]+)", line).groups() + for line in requests + executes +} +assert result["outcome"] == "succeeded", result +assert len(requests) == 1, requests +assert len(executes) == 3, executes +assert len(identities) == 1, identities +PY +pass "wrapper resumes transient execution under one assigned Azure task identity" + printf 'ok\n' > "$TMP_ROOT/mode" printf '2\n' > "$TMP_ROOT/cleanup-reconcile-failures" : > "$TMP_ROOT/calls.log"