Skip to content

Commit db7aa05

Browse files
committed
fix(azure): recover no-mistakes worker completion
1 parent 381eb91 commit db7aa05

6 files changed

Lines changed: 246 additions & 14 deletions

bin/fm-azure-worker-provider.py

Lines changed: 21 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -3109,10 +3109,29 @@ def execute_terminal_disposition(controller, action, resources):
31093109
)
31103110
names = expected_names(controller, action["slot"])
31113111
view = view or run_command_instance_view(controller, names["vm"], names["task-command"])
3112-
if view.get("executionState") != "Succeeded":
3112+
execution_state = view.get("executionState")
3113+
if str(execution_state).lower() in ("failed", "canceled"):
3114+
exit_code = view.get("exitCode")
3115+
if isinstance(exit_code, bool) or not isinstance(exit_code, int):
3116+
raise ProviderError(
3117+
"exact terminal worker execution has no integer exit code: state={}".format(
3118+
execution_state
3119+
)
3120+
)
3121+
return EXECUTE_DISPOSITION_TERMINAL, {
3122+
"schema": EXECUTION_TERMINAL_SCHEMA,
3123+
"request_digest": request_digest,
3124+
"idempotency_key": action.get("idempotency_key"),
3125+
"disposition": "provider-terminal",
3126+
"provisioning_state": provisioning_state,
3127+
"execution_state": execution_state,
3128+
"exit_code": exit_code,
3129+
"task_command_id": task_command_resource["id"],
3130+
}
3131+
if execution_state != "Succeeded":
31133132
raise ProviderError(
31143133
"exact worker execution has no recoverable terminal result: state={}".format(
3115-
view.get("executionState")
3134+
execution_state
31163135
)
31173136
)
31183137
execution = exact_execution_marker(view)

bin/fm-worker-lifecycle.py

Lines changed: 57 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -2384,12 +2384,20 @@ def apply_action_result(env, state, action, result):
23842384
expected_task_command_id = (
23852385
((action.get("resources") or {}).get("task-command") or {}).get("id")
23862386
)
2387+
provisioning_state = str(execution.get("provisioning_state", "")).lower()
2388+
execution_state = str(execution.get("execution_state", "")).lower()
2389+
provisioning_terminal = provisioning_state in ("failed", "canceled")
2390+
execution_terminal = (
2391+
provisioning_state == "succeeded"
2392+
and execution_state in ("failed", "canceled")
2393+
and isinstance(execution.get("exit_code"), int)
2394+
and not isinstance(execution.get("exit_code"), bool)
2395+
)
23872396
if (
23882397
execution.get("request_digest") != action.get("request_digest")
23892398
or execution.get("idempotency_key") != action.get("idempotency_key")
23902399
or execution.get("disposition") != "provider-terminal"
2391-
or str(execution.get("provisioning_state", "")).lower()
2392-
not in ("failed", "canceled")
2400+
or not (provisioning_terminal or execution_terminal)
23932401
or not isinstance(execution.get("task_command_id"), str)
23942402
or not execution["task_command_id"]
23952403
):
@@ -2398,9 +2406,10 @@ def apply_action_result(env, state, action, result):
23982406
raise ProviderResultIdentityRefused(
23992407
"provider terminal execution task-command identity differs from the claimed action"
24002408
)
2409+
terminal_state = execution_state if execution_terminal else provisioning_state
24012410
raise LifecycleError(
24022411
"provider-terminal {}: exact execution is {} and cannot be applied".format(
2403-
execution["request_digest"], execution["provisioning_state"]
2412+
execution["request_digest"], terminal_state
24042413
)
24052414
)
24062415
if not isinstance(execution, dict) or execution.get("schema") != EXECUTION_RESULT_SCHEMA:
@@ -4523,16 +4532,53 @@ def command_service_complete(env, args):
45234532
require_id("task generation", args.task_generation),
45244533
)
45254534
item = state["queue"].get(key)
4526-
worker = state["workers"].get(str((item or {}).get("slot")))
4527-
if item is None or item.get("status") != "assigned" or worker is None:
4528-
raise LifecycleError("service completion requires one exact assigned worker")
4529-
if item.get("role") != "no-mistakes" or worker.get("role") != "no-mistakes":
4535+
if item is None or item.get("role") != "no-mistakes":
45304536
raise LifecycleError("service completion is owned by no-mistakes workers only")
4531-
if worker.get("assignment_generation") != args.assignment_generation:
4532-
raise LifecycleError("service completion assignment generation is not exact")
45334537
execution = state["executions"].get(request_digest)
45344538
if not isinstance(execution, dict):
45354539
raise LifecycleError("service completion has no exact recorded execution")
4540+
receipt = item.get("service_completion_receipt")
4541+
if receipt is not None:
4542+
if item.get("status") not in ("releasing", "complete"):
4543+
raise LifecycleError("service completion receipt exists in an invalid queue state")
4544+
if not isinstance(receipt, dict):
4545+
raise LifecycleError("service completion receipt is malformed")
4546+
unsigned_receipt = dict(receipt)
4547+
receipt_digest = unsigned_receipt.pop("proof_digest", None)
4548+
if (
4549+
receipt.get("schema") != "fm.worker-service-release/v1"
4550+
or receipt.get("task") != args.task
4551+
or receipt.get("task_generation") != args.task_generation
4552+
or receipt.get("assignment_generation") != args.assignment_generation
4553+
or receipt.get("request_digest") != request_digest
4554+
or receipt.get("result_digest") != execution.get("result_digest")
4555+
or execution.get("request_digest") != request_digest
4556+
or execution.get("assignment_generation") != args.assignment_generation
4557+
or receipt.get("verdict") != "proved"
4558+
or receipt_digest != digest_value(unsigned_receipt)
4559+
):
4560+
raise LifecycleError("service completion receipt identity differs")
4561+
worker = state["workers"].get(str(item.get("slot")))
4562+
worker_owns_item = (
4563+
worker is not None and worker.get("queue_key") == key
4564+
)
4565+
if item.get("status") == "releasing" and not worker_owns_item:
4566+
raise LifecycleError("releasing service completion lost its exact worker")
4567+
if worker_owns_item and (
4568+
worker.get("role") != "no-mistakes"
4569+
or worker.get("assignment_generation") != args.assignment_generation
4570+
or worker.get("release_proof") != receipt
4571+
):
4572+
raise LifecycleError("service completion worker receipt identity differs")
4573+
print("service release proof already recorded with exact identity")
4574+
return
4575+
worker = state["workers"].get(str(item.get("slot")))
4576+
if item.get("status") != "assigned" or worker is None:
4577+
raise LifecycleError("service completion requires one exact assigned worker")
4578+
if worker.get("role") != "no-mistakes":
4579+
raise LifecycleError("service completion is owned by no-mistakes workers only")
4580+
if worker.get("assignment_generation") != args.assignment_generation:
4581+
raise LifecycleError("service completion assignment generation is not exact")
45364582
if (
45374583
execution.get("request_digest") != request_digest
45384584
or execution.get("result_digest") != worker.get("last_execution_digest")
@@ -4556,9 +4602,10 @@ def command_service_complete(env, args):
45564602
raise LifecycleError("worker already has a different service release proof")
45574603
print("service release proof already recorded with exact identity")
45584604
return
4559-
worker["release_proof"] = proof
4605+
worker["release_proof"] = copy.deepcopy(proof)
45604606
worker["released_at"] = iso_utc()
45614607
worker["phase"] = "release-proved"
4608+
item["service_completion_receipt"] = copy.deepcopy(proof)
45624609
item["status"] = "releasing"
45634610
save_state(env, state)
45644611
print("service execution proved; exact idle capacity is eligible for cleanup")

bin/fm-worker-lifecycle.sh

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -35,6 +35,7 @@
3535
# fm-worker-lifecycle.sh capacity-release <exact fence and cleanup receipt>
3636
# fm-worker-lifecycle.sh execute <exact assignment flags> [--existing-task-disk --return-kind <ship|scout> --outcome-dir <dir>] -- <argv...>
3737
# fm-worker-lifecycle.sh authority-receipt <exact assignment flags> --output <json>
38+
# fm-worker-lifecycle.sh service-complete <exact no-mistakes execution binding>
3839
# fm-worker-lifecycle.sh proof-template --task <id> --task-generation <id>
3940
# fm-worker-lifecycle.sh release --task <id> --task-generation <id> --proof-file <json>
4041
# fm-worker-lifecycle.sh withdraw --task <id> --task-generation <id> --confirm-withdraw --confirm-subscription <uuid>
@@ -138,7 +139,7 @@ fm_worker_receipt_credential_remains() { # <task home file> <task id> <controll
138139
}
139140

140141
case "${1:-}" in
141-
request|release|resume|steer|execute|authority-receipt|capacity-reserve|capacity-reserve-shape|capacity-release|capacity-retire-fence|abandon-claim|message-put|message-collect|compartment-chain-tip)
142+
request|release|resume|steer|execute|authority-receipt|service-complete|capacity-reserve|capacity-reserve-shape|capacity-release|capacity-retire-fence|abandon-claim|message-put|message-collect|compartment-chain-tip)
142143
fm_refuse_if_gate_agent
143144
exec python3 "$SCRIPT_DIR/fm-worker-lifecycle.py" "$@"
144145
;;

bin/fm-worker-supervisor.py

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -29,6 +29,7 @@
2929
MAX_RESULT_BYTES = 8 * 1024 * 1024
3030
MAX_WALL_SECONDS = 6 * 60 * 60
3131
MAX_OUTPUT_BYTES = 4 * 1024 * 1024
32+
MAX_NO_MISTAKES_RUNTIME_FILES = 20000
3233

3334

3435
class SupervisorError(RuntimeError):
@@ -386,7 +387,7 @@ def stage_no_mistakes_runtime(source, target, enforce_linux=True):
386387
try:
387388
with tarfile.open(source, mode="r:gz") as archive:
388389
members = archive.getmembers()
389-
if not members or len(members) > 4096:
390+
if not members or len(members) > MAX_NO_MISTAKES_RUNTIME_FILES:
390391
raise SupervisorError("no-mistakes runtime member inventory is unbounded")
391392
for member in members:
392393
parts = Path(member.name).parts

tests/fm-no-mistakes-runtime.test.sh

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -74,6 +74,8 @@ def load(name, path):
7474
7575
builder = load("runtime_builder", root / "bin/fm-no-mistakes-runtime.py")
7676
supervisor = load("worker_supervisor", root / "bin/fm-worker-supervisor.py")
77+
assert supervisor.MAX_NO_MISTAKES_RUNTIME_FILES == builder.MAX_FILES, (
78+
supervisor.MAX_NO_MISTAKES_RUNTIME_FILES, builder.MAX_FILES)
7779
bundle = temporary / "runtime.tar.gz"
7880
args = argparse.Namespace(
7981
no_mistakes=str(temporary / "no-mistakes"), node=str(temporary / "node"),

tests/fm-worker-lifecycle.test.sh

Lines changed: 162 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,149 @@ AUTHORITY="$ROOT/bin/fm-worker-authority.py"
1616
DOC="$ROOT/docs/azure-workers.md"
1717
SUB=11111111-1111-4111-8111-111111111111
1818

19+
service_complete_front_door() {
20+
local tmp help
21+
fm_test_tmproot_into tmp fm-worker-service-complete-front-door
22+
mkdir -p "$tmp/home"
23+
help=$(FM_HOME="$tmp/home" "$WRAPPER" service-complete --help) \
24+
|| fail "supported lifecycle wrapper rejected service-complete"
25+
case "$help" in
26+
*--request-digest*--confirm-subscription*) ;;
27+
*) fail "service-complete help lost its exact execution binding" ;;
28+
esac
29+
pass "supported lifecycle wrapper exposes service-complete"
30+
}
31+
32+
service_complete_replay_contract() {
33+
python3 - "$CONTROLLER" <<'PY' || fail "service completion replay contract failed"
34+
import contextlib
35+
import importlib.util
36+
from types import SimpleNamespace
37+
import sys
38+
39+
spec = importlib.util.spec_from_file_location("lifecycle", sys.argv[1])
40+
module = importlib.util.module_from_spec(spec)
41+
spec.loader.exec_module(module)
42+
43+
bindings = {
44+
"home_binding": "1" * 64,
45+
"task": "service-task",
46+
"task_generation": "service-generation",
47+
"assignment_generation": "asg-00000001",
48+
"account_binding": "2" * 64,
49+
"worktree_binding": "3" * 64,
50+
"repository_binding": "4" * 64,
51+
"repository_generation": "repository-generation",
52+
}
53+
request_digest = "5" * 64
54+
result_digest = "6" * 64
55+
item = {
56+
**bindings,
57+
"role": "no-mistakes",
58+
"status": "assigned",
59+
"slot": 1,
60+
}
61+
worker = {
62+
"role": "no-mistakes",
63+
"queue_key": "service-task@service-generation",
64+
"assignment_generation": bindings["assignment_generation"],
65+
"bindings": bindings,
66+
"cloud_instance_id": "worker-instance",
67+
"resources": {"vm": {"id": "/exact/vm"}},
68+
"last_execution_digest": result_digest,
69+
"release_proof": None,
70+
}
71+
execution = {
72+
"request_digest": request_digest,
73+
"result_digest": result_digest,
74+
"assignment_generation": bindings["assignment_generation"],
75+
}
76+
state = {
77+
"queue": {"service-task@service-generation": item},
78+
"workers": {"1": worker},
79+
"executions": {request_digest: execution},
80+
}
81+
module.controller_lock = lambda _env: contextlib.nullcontext()
82+
module.load_state = lambda _env: state
83+
module.save_state = lambda _env, _state: None
84+
args = SimpleNamespace(
85+
task="service-task",
86+
task_generation="service-generation",
87+
assignment_generation=bindings["assignment_generation"],
88+
request_digest=request_digest,
89+
confirm_subscription="subscription",
90+
)
91+
env = {"subscription": "subscription"}
92+
93+
# Crash window one: the first call durably moved the item to releasing, but
94+
# the caller died before observing success. The exact retry is idempotent.
95+
module.command_service_complete(env, args)
96+
assert item["status"] == "releasing", item
97+
assert item["service_completion_receipt"] == worker["release_proof"], item
98+
module.command_service_complete(env, args)
99+
100+
# Crash window two: reconcile completed the reset and removed the worker, but
101+
# the caller died before publishing the cached result. The queue-owned exact
102+
# receipt survives reset and admits only the same bound completion request.
103+
item["status"] = "complete"
104+
# The released slot may already belong to a later task. Its presence cannot
105+
# invalidate the old queue item's exact, self-digested completion receipt.
106+
state["workers"] = {"1": {
107+
"queue_key": "later-task@later-generation",
108+
"role": "author",
109+
"assignment_generation": "asg-00000002",
110+
"release_proof": None,
111+
}}
112+
module.command_service_complete(env, args)
113+
wrong = SimpleNamespace(**vars(args))
114+
wrong.assignment_generation = "asg-99999999"
115+
try:
116+
module.command_service_complete(env, wrong)
117+
except module.LifecycleError as exc:
118+
assert "identity differs" in str(exc), exc
119+
else:
120+
raise AssertionError("completed service receipt admitted a foreign assignment")
121+
122+
# The provider-side regression below emits this exact terminal shape. Prove
123+
# the lifecycle consumer recognizes it as terminal (rather than malformed),
124+
# fails closed, and leaves the durable execute claim available for explicit
125+
# abandonment/recovery.
126+
terminal_action = {
127+
"type": "execute",
128+
"slot": 1,
129+
"request_digest": "7" * 64,
130+
"idempotency_key": "8" * 64,
131+
"resources": {"task-command": {"id": "/exact/task-command"}},
132+
}
133+
terminal_state = {
134+
"queue": {"terminal-task@terminal-generation": {"status": "assigned"}},
135+
"workers": {"1": {"queue_key": "terminal-task@terminal-generation"}},
136+
"executions": {},
137+
"pending_actions": {"1": terminal_action},
138+
"completed_worker_seconds": 0.0,
139+
}
140+
state = terminal_state
141+
terminal_result = {"execution": {
142+
"schema": "fm.worker-execution-terminal/v1",
143+
"request_digest": terminal_action["request_digest"],
144+
"idempotency_key": terminal_action["idempotency_key"],
145+
"disposition": "provider-terminal",
146+
"provisioning_state": "Succeeded",
147+
"execution_state": "Failed",
148+
"exit_code": 2,
149+
"task_command_id": "/exact/task-command",
150+
}}
151+
try:
152+
module.apply_pending(env, terminal_action, terminal_result)
153+
except module.LifecycleError as exc:
154+
assert "provider-terminal" in str(exc) and "failed" in str(exc), exc
155+
else:
156+
raise AssertionError("failed guest execution was applied as a successful result")
157+
assert terminal_state["pending_actions"]["1"] == terminal_action, terminal_state
158+
PY
159+
pass "service completion replays across releasing and completed crash windows"
160+
}
161+
19162
static_contract() {
20163
python3 - "$CONTROLLER" "$AZURE" "$SUPERVISOR" "$AUTHORITY" "$DOC" <<'PY' || fail "elastic worker static contract failed"
21164
from pathlib import Path
@@ -859,6 +1002,23 @@ recovered_kind, recovered_execution = module.execute_terminal_disposition(
8591002
)
8601003
assert recovered_kind == module.EXECUTE_DISPOSITION_RECOVERED, recovered_kind
8611004
assert recovered_execution == execution, recovered_execution
1005+
module.run_command_instance_view = lambda *_args, **_kwargs: {
1006+
"executionState": "Failed", "exitCode": 2, "output": "", "error": "guest failed",
1007+
}
1008+
failed_kind, failed_execution = module.execute_terminal_disposition(
1009+
controller, execute_action, worker["resources"]
1010+
)
1011+
assert failed_kind == module.EXECUTE_DISPOSITION_TERMINAL, failed_kind
1012+
assert failed_execution == {
1013+
"schema": "fm.worker-execution-terminal/v1",
1014+
"request_digest": execute_action["request_digest"],
1015+
"idempotency_key": execute_action["idempotency_key"],
1016+
"disposition": "provider-terminal",
1017+
"provisioning_state": "Succeeded",
1018+
"execution_state": "Failed",
1019+
"exit_code": 2,
1020+
"task_command_id": task_command["id"],
1021+
}, failed_execution
8621022
8631023
# Once the exact request owns the Run Command, a replay may only recover its
8641024
# terminal disposition or exact result. Updating/Running and a Succeeded
@@ -8692,6 +8852,8 @@ PY
86928852
}
86938853

86948854
static_contract
8855+
service_complete_front_door
8856+
service_complete_replay_contract
86958857
compartment_payload_contract
86968858
classification_and_admission_matrix
86978859
azure_provider_refusal_matrix

0 commit comments

Comments
 (0)