Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
52 changes: 45 additions & 7 deletions bin/fm-no-mistakes-worker
Original file line number Diff line number Diff line change
Expand Up @@ -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():
Expand Down Expand Up @@ -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")
Expand Down
35 changes: 35 additions & 0 deletions tests/fm-no-mistakes-worker.test.sh
Original file line number Diff line number Diff line change
Expand Up @@ -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"))
Expand Down Expand Up @@ -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"
Expand Down
Loading