Skip to content

Commit 066b5bf

Browse files
authored
fix(delegation): complete under the original lease and preserve replan handoff (#5466)
1 parent 760f039 commit 066b5bf

6 files changed

Lines changed: 277 additions & 9 deletions

File tree

‎docs/architecture/rfcs/ledger/shared-goal-authority-state-provider-v0/2026-09-28-retirement-cadence.md‎

Lines changed: 12 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -575,8 +575,13 @@ Renewal uses the latest proved version, one unchanged-intent retry for a lost
575575
transport reply, and the last proved expiry even when renewal hangs. A rejected
576576
proof cancels the CLI; its TERM adapter unwinds nested managed Hosts before
577577
returning. Ordinary non-hard delegation keeps the existing subprocess route.
578-
Completion reads the current claim, persists its terminal CAS intent before the
579-
effect, and replays that exact completion after an ambiguous reply. Canonical
578+
After the supervised CLI returns, completion renews the original execution
579+
using its acquisition TTL before starting independent Todo acceptance. It
580+
journals the renewal version before the effect, then proves current authority
581+
and freezes the terminal CAS version. Lost renewal and completion replies each
582+
replay their original intent; neither replay can revive an expired or replaced
583+
execution. The validation phase makes no lease writes that would invalidate
584+
its provider-revision witness. Canonical
580585
completion releases the execution lease; subsequent original-Turn accounting
581586
uses its terminal receipt rather than reacquiring an open-work lease.
582587

@@ -592,7 +597,11 @@ their own runtime instead of inheriting operator state.
592597
The acceptance slice uses disposable File/SQLite providers, real CLI/Turn
593598
execution and a synthetic model process: crossing the initial expiry, canonical
594599
release, a new execution epoch, lost completion/renewal replies and rejected or
595-
hung renewal, including control-pipe loss with a TERM-resistant process. It
600+
hung renewal, including control-pipe loss with a TERM-resistant process. Final
601+
Todo acceptance also crosses a short remaining lease deadline, with lost-reply,
602+
expiry and replacement controls. The Python delegation adapter only sequences
603+
the existing TS-owned claim, renew and terminal contracts; no provider rule,
604+
public setting or frontend permission changes. It
596605
does not qualify a paid model, remote job cancellation or
597606
Windows process-tree cleanup. Stop acknowledgements and interrupted-Turn
598607
no-progress settlement remain with the existing delegation-stop work (#5308);

‎docs/architecture/rfcs/ledger/shared-goal-authority-state-provider-v0/2026-09-28-retirement-cadence.zh-CN.md‎

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -404,3 +404,15 @@ Archive restore/audit 已使用 provider 通用的 1–64 个操作批量回
404404
验收。File 每恢复一笔仍重写保留的文件;此前一次完整历史恢复超出调用方的
405405
300 秒超时,随后才发布精确匹配的确认。批量回执优化没有闭合这项恢复成本。
406406
#4224 的 soak 已启动;其最终证据和对当前候选的适用性仍待核对。
407+
408+
### 委派完成阶段的原租约衔接
409+
410+
#5436 的 TS Host 监督覆盖模型执行和 Turn 验证,最终 Todo 独立验收位于其后。
411+
因此完成前先按原获取 TTL 续租同一执行,分别持久化续租与完成的 CAS 版本,再
412+
进入验收。两种回复丢失都重放原请求;历史回执不能恢复已过期或被替换的执行。
413+
验收期间不写租约,避免自己改变验收所依赖的 provider revision。
414+
415+
隔离 File/SQLite 回归覆盖剩余短租约跨越最终验收,以及续租/完成回复丢失、
416+
过期和新 epoch。Python 仅衔接既有 TS claim、renew、terminal 权威合同,没有
417+
新增 provider 规则、公开配置或前端权限。Stop ACK 与中断后的无进展结算仍归
418+
#5308;持续 D2、默认准入和旧 writer 删除仍需各自证据。

‎loopx/collaboration_mcp.py‎

Lines changed: 22 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1174,6 +1174,28 @@ def _complete_delegated_todo(self, row: dict, binding: dict) -> None:
11741174
# key/epoch; it cannot reacquire an expired execution. Renewal has
11751175
# changed its version, so the historical acquisition is not CAS.
11761176
if "completion_lease_version" not in row:
1177+
# The Host supervisor has stopped. Renew the original execution
1178+
# before validation captures its provider revision; renewing
1179+
# during validation would invalidate that source witness. The
1180+
# canonical TS lease owner decides admission and replay. This
1181+
# adapter journals one intent, not a new lease or a longer TTL.
1182+
if "completion_lease_renewal_version" not in row:
1183+
proof = self._cli(binding, *self._delegation_claim_arguments(row, binding))
1184+
if proof.get("ok") is not True:
1185+
raise ValueError("delegation current execution proof lost before completion")
1186+
row["completion_lease_renewal_version"] = proof["lease"]["version"]
1187+
_write(self.path(row["identity"]["operation_id"]), row)
1188+
renewed = self._cli(
1189+
binding, "task-lease", "renew", "--goal-id", self.goal_id,
1190+
"--todo-id", binding["todo_id"], "--owner", binding["agent_id"],
1191+
"--idempotency-key", lease["lease"]["idempotency_key"],
1192+
"--expected-version", str(row["completion_lease_renewal_version"]),
1193+
"--ttl-seconds", str(lease["lease"]["acquire_ttl_seconds"]),
1194+
)
1195+
if renewed.get("ok") is not True:
1196+
raise ValueError("delegation original lease renewal rejected before completion")
1197+
# A renewal receipt can be historical after a lost reply. Read
1198+
# current authority before freezing the terminal intent below.
11771199
proof = self._cli(binding, *self._delegation_claim_arguments(row, binding))
11781200
current = proof.get("lease", {})
11791201
if (proof.get("ok") is not True or current.get("owner") != binding["agent_id"]

‎loopx/control_plane/turn_driver/driver.py‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -460,6 +460,11 @@ def build_loopx_turn_plan(
460460
else LoopXTurnRoute.CONTRACT_ERROR
461461
)
462462
selected_todo = selected_turn_todo(envelope)
463+
# Replanning an empty frontier is valid, but this Host executor settles
464+
# against a Todo. Preserve the replan packet for the controller to resolve;
465+
# do not turn the absence of a successor into a malformed Host session.
466+
if route is LoopXTurnRoute.REPLAN_REQUIRED and not selected_todo:
467+
route = LoopXTurnRoute.BLOCKED
463468
lineage = _turn_lineage(envelope, selected_todo=selected_todo)
464469
session, session_error = _session_plan(
465470
route=route,

‎tests/test_delegation_lease_lifetime.py‎

Lines changed: 178 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -13,7 +13,7 @@
1313

1414
import pytest
1515

16-
from test_local_delegation import HOST, brief, service # noqa: F401
16+
from test_local_delegation import HOST, brief, demo, service # noqa: F401
1717
from loopx.control_plane.collaboration.inbox import _read
1818
from loopx.control_plane.coordination.local_authority import read_canonical_todos_if_promoted
1919
from tests.control_plane.host_process_fixture import COUNTER_PROCESS_SOURCE
@@ -42,6 +42,8 @@ def prepare_lease(root, runner, monkeypatch, *, ttl=20):
4242
binding = runner.binding("analysis")
4343
runner._acquire_delegation_lease(runner.path("lease-lifetime"), row, binding)
4444
lease = row["task_lease"]["lease"]
45+
if ttl is None:
46+
return lease
4547
renewed = runner._cli(binding, "task-lease", "renew", "--goal-id", runner.goal_id,
4648
"--todo-id", binding["todo_id"], "--owner", binding["agent_id"],
4749
"--idempotency-key", lease["idempotency_key"], "--expected-version", str(lease["version"]),
@@ -54,6 +56,164 @@ def inspect(runner):
5456
"--todo-id", "todo_analyst-initial")
5557

5658

59+
@pytest.fixture(params=["file", "sqlite"])
60+
def completion_service(tmp_path, request, monkeypatch):
61+
"""Pin a real slow final acceptance command before preparing authority."""
62+
validator = tmp_path / "completion-validator.py"
63+
validator.write_text('''import sys, time, runpy
64+
from datetime import datetime
65+
from pathlib import Path
66+
root = Path(sys.argv[1])
67+
counter = root / 'validator-calls'
68+
calls = int(counter.read_text()) + 1 if counter.exists() else 1
69+
counter.write_text(str(calls))
70+
if calls == 2:
71+
deadline = datetime.fromisoformat((root / 'completion-prior-expiry').read_text())
72+
time.sleep(max(0, deadline.timestamp() - time.time()) + 0.2)
73+
(root / 'completion-validation-ended').touch()
74+
actual = root / 'project/validation/acceptance.py'
75+
sys.path.insert(0, str(actual.parent))
76+
sys.argv[0] = str(actual)
77+
runpy.run_path(str(actual), run_name='__main__')
78+
''')
79+
write = demo.write
80+
81+
def configure(path, value):
82+
if path.name == "bootstrap.json":
83+
for criterion in value["document"]["criteria"]:
84+
criterion["validation_timeout_seconds"] = 1
85+
if criterion["id"] == "analyst-initial":
86+
# 21 + four one-second criteria stays within the public
87+
# 25-second completion budget, including a 20s deadline.
88+
criterion["validation_timeout_seconds"] = 21
89+
criterion["validation_argv"][1] = str(validator)
90+
return write(path, value)
91+
92+
monkeypatch.setattr(demo, "write", configure)
93+
return service.__wrapped__(tmp_path, request, monkeypatch)
94+
95+
96+
@pytest.mark.parametrize("lost_reply", [None, "renewal", "completion"])
97+
def test_completion_renews_before_validation_and_replays_each_intent(completion_service, monkeypatch, lost_reply):
98+
root, runner = completion_service
99+
# This case shortens the lease at the completion boundary below. An
100+
# unrelated short Host deadline can cancel execution before that boundary
101+
# under load; the running-Host cases separately exercise that deadline.
102+
original = prepare_lease(root, runner, monkeypatch, ttl=None)
103+
complete, cli = runner._complete_delegated_todo, runner._cli
104+
shortened = None
105+
renewal_calls, completion_calls = [], []
106+
dropped = False
107+
108+
def enter_completion(row, binding):
109+
nonlocal shortened
110+
if shortened is None:
111+
current = inspect(runner)["lease"]
112+
# Fix the phase boundary, independent of whether the Host happened
113+
# to cross its earlier renewal timer. Use the real canonical API.
114+
shortened = cli(binding, "task-lease", "renew", "--goal-id", runner.goal_id,
115+
"--todo-id", binding["todo_id"], "--owner", binding["agent_id"],
116+
"--idempotency-key", current["idempotency_key"],
117+
"--expected-version", str(current["version"]), "--ttl-seconds", "20")["lease"]
118+
# Cross the actual pre-renewal deadline, not an assumed amount of
119+
# CLI startup time. Allow cold claim/renew commands to reach the
120+
# boundary; the independent validator still outlives that lease.
121+
(root / "completion-prior-expiry").write_text(shortened["expires_at"])
122+
return complete(row, binding)
123+
124+
def observe_reply(binding, *args, **kwargs):
125+
nonlocal dropped
126+
result = cli(binding, *args, **kwargs)
127+
phase = None
128+
if args[:2] == ("task-lease", "renew"):
129+
renewal_calls.append((args, result))
130+
phase = "renewal"
131+
elif args[:2] == ("todo", "complete"):
132+
completion_calls.append((args, result))
133+
phase = "completion"
134+
if phase == lost_reply and phase is not None and not dropped:
135+
dropped = True
136+
raise ValueError("fixture dropped the committed " + phase + " reply")
137+
return result
138+
139+
monkeypatch.setattr(runner, "_complete_delegated_todo", enter_completion)
140+
monkeypatch.setattr(runner, "_cli", observe_reply)
141+
runner.execute("lease-lifetime")
142+
if lost_reply:
143+
uncertain = _read(runner.path("lease-lifetime"))
144+
assert uncertain["status"] == "turn_returned", uncertain
145+
assert "fixture dropped" in uncertain["error"]
146+
assert uncertain["completion_lease_renewal_version"] == shortened["version"]
147+
assert ("completion_lease_version" in uncertain) == (lost_reply == "completion")
148+
runner.execute("lease-lifetime")
149+
row = _read(runner.path("lease-lifetime"))
150+
assert row["status"] == "accepted", row
151+
assert row["turn_result"]["result_kind"] == "validated_progress"
152+
assert (root / "completion-validation-ended").exists()
153+
assert time.time() > datetime.fromisoformat(shortened["expires_at"].replace("Z", "+00:00")).timestamp()
154+
assert (root / "analyst/initial/host-invocations").read_text() == "1"
155+
final = inspect(runner)["lease"]
156+
assert final["status"] == "released"
157+
assert final["lease_epoch"] == original["lease_epoch"]
158+
assert final["idempotency_key"] == original["idempotency_key"]
159+
assert final["version"] == shortened["version"] + 1
160+
assert len(renewal_calls) == (2 if lost_reply == "renewal" else 1)
161+
assert len(completion_calls) == (2 if lost_reply == "completion" else 1)
162+
for calls in (renewal_calls, completion_calls):
163+
assert all(args == calls[0][0] for args, _ in calls)
164+
assert all(result["provider_revision"] == calls[0][1]["provider_revision"] for _, result in calls)
165+
166+
167+
@pytest.mark.parametrize("authority_loss", ["expiry", "replacement"])
168+
def test_completion_renewal_receipt_cannot_revive_lost_execution(service, monkeypatch, authority_loss):
169+
root, runner = service
170+
prepare_lease(root, runner, monkeypatch)
171+
cli = runner._cli
172+
dropped = False
173+
completions = []
174+
175+
def lose_renewal_reply(binding, *args, **kwargs):
176+
nonlocal dropped
177+
if args[:2] == ("todo", "complete"):
178+
completions.append(args)
179+
result = cli(binding, *args, **kwargs)
180+
if args[:2] == ("task-lease", "renew") and not dropped:
181+
dropped = True
182+
raise ValueError("fixture lost renewal response before terminal intent")
183+
return result
184+
185+
monkeypatch.setattr(runner, "_cli", lose_renewal_reply)
186+
runner.execute("lease-lifetime")
187+
row = _read(runner.path("lease-lifetime"))
188+
assert row["status"] == "turn_returned", row
189+
assert "completion_lease_renewal_version" in row
190+
assert "completion_lease_version" not in row
191+
current = inspect(runner)["lease"]
192+
binding = runner.binding("analysis")
193+
args = ("--goal-id", runner.goal_id, "--todo-id", binding["todo_id"],
194+
"--owner", binding["agent_id"], "--idempotency-key", current["idempotency_key"],
195+
"--expected-version", str(current["version"]))
196+
if authority_loss == "expiry":
197+
expired = cli(binding, "task-lease", "renew", *args, "--ttl-seconds", "1")["lease"]
198+
time.sleep(max(0, datetime.fromisoformat(expired["expires_at"].replace("Z", "+00:00")).timestamp()
199+
- time.time()) + 0.1)
200+
else:
201+
assert cli(binding, "task-lease", "release", *args)["ok"] is True
202+
replacement = cli(binding, "task-lease", "acquire", "--goal-id", runner.goal_id,
203+
"--todo-id", binding["todo_id"], "--owner", binding["agent_id"],
204+
"--idempotency-key", "replacement", "--expected-version", str(current["version"]))
205+
assert replacement["lease"]["lease_epoch"] > current["lease_epoch"]
206+
runner.execute("lease-lifetime")
207+
rejected = _read(runner.path("lease-lifetime"))
208+
assert rejected["status"] != "accepted", rejected
209+
assert "completion_lease_version" not in rejected
210+
assert not completions
211+
assert (root / "analyst/initial/host-invocations").read_text() == "1"
212+
snapshot = read_canonical_todos_if_promoted(runtime_root=runner.root, goal_id=runner.goal_id)
213+
todo = next(item for item in snapshot["todos"] if item["todo_id"] == binding["todo_id"])
214+
assert todo["done"] is False
215+
216+
57217
def await_started(root, future, runner):
58218
deadline = time.monotonic() + 45
59219
while time.monotonic() < deadline and not (root / "host-started").exists():
@@ -125,11 +285,24 @@ def test_real_revocation_or_new_execution_stops_nested_host_without_acceptance(s
125285
with ThreadPoolExecutor(max_workers=1) as pool:
126286
future = pool.submit(runner.execute, "lease-lifetime")
127287
await_started(root, future, runner)
128-
current = inspect(runner)["lease"]
129288
binding = runner.binding("analysis")
130-
runner._cli(binding, "task-lease", "release", "--goal-id", runner.goal_id,
131-
"--todo-id", "todo_analyst-initial", "--owner", "analyst",
132-
"--idempotency-key", original["idempotency_key"], "--expected-version", str(current["version"]))
289+
# The live supervisor may renew between inspection and revocation.
290+
# Retry only that CAS race, never a replacement execution or rejection.
291+
for _ in range(5):
292+
current = inspect(runner)["lease"]
293+
assert current["lease_epoch"] == original["lease_epoch"]
294+
assert current["idempotency_key"] == original["idempotency_key"]
295+
try:
296+
runner._cli(binding, "task-lease", "release", "--goal-id", runner.goal_id,
297+
"--todo-id", "todo_analyst-initial", "--owner", "analyst",
298+
"--idempotency-key", original["idempotency_key"],
299+
"--expected-version", str(current["version"]))
300+
break
301+
except ValueError as exc:
302+
if str(exc) != "canonical task lease release rejected: version_mismatch":
303+
raise
304+
else:
305+
pytest.fail("could not revoke the original execution during concurrent renewal")
133306
if reclaim:
134307
acquired = runner._cli(binding, "task-lease", "acquire", "--goal-id", runner.goal_id,
135308
"--todo-id", "todo_analyst-initial", "--owner", "analyst", "--idempotency-key", "new-execution",

0 commit comments

Comments
 (0)