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
Original file line number Diff line number Diff line change
Expand Up @@ -475,8 +475,13 @@ Renewal uses the latest proved version, one unchanged-intent retry for a lost
transport reply, and the last proved expiry even when renewal hangs. A rejected
proof cancels the CLI; its TERM adapter unwinds nested managed Hosts before
returning. Ordinary non-hard delegation keeps the existing subprocess route.
Completion reads the current claim, persists its terminal CAS intent before the
effect, and replays that exact completion after an ambiguous reply. Canonical
After the supervised CLI returns, completion renews the original execution
using its acquisition TTL before starting independent Todo acceptance. It
journals the renewal version before the effect, then proves current authority
and freezes the terminal CAS version. Lost renewal and completion replies each
replay their original intent; neither replay can revive an expired or replaced
execution. The validation phase makes no lease writes that would invalidate
its provider-revision witness. Canonical
completion releases the execution lease; subsequent original-Turn accounting
uses its terminal receipt rather than reacquiring an open-work lease.

Expand All @@ -492,7 +497,11 @@ their own runtime instead of inheriting operator state.
The acceptance slice uses disposable File/SQLite providers, real CLI/Turn
execution and a synthetic model process: crossing the initial expiry, canonical
release, a new execution epoch, lost completion/renewal replies and rejected or
hung renewal, including control-pipe loss with a TERM-resistant process. It
hung renewal, including control-pipe loss with a TERM-resistant process. Final
Todo acceptance also crosses a short remaining lease deadline, with lost-reply,
expiry and replacement controls. The Python delegation adapter only sequences
the existing TS-owned claim, renew and terminal contracts; no provider rule,
public setting or frontend permission changes. It
does not qualify a paid model, remote job cancellation or
Windows process-tree cleanup. Stop acknowledgements and interrupted-Turn
no-progress settlement remain with the existing delegation-stop work (#5308);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -321,3 +321,15 @@ Archive restore/audit 已使用 provider 通用的 1–64 个操作批量回
验收。File 每恢复一笔仍重写保留的文件;此前一次完整历史恢复超出调用方的
300 秒超时,随后才发布精确匹配的确认。批量回执优化没有闭合这项恢复成本。
#4224 的 soak 已启动;其最终证据和对当前候选的适用性仍待核对。

### 委派完成阶段的原租约衔接

#5436 的 TS Host 监督覆盖模型执行和 Turn 验证,最终 Todo 独立验收位于其后。
因此完成前先按原获取 TTL 续租同一执行,分别持久化续租与完成的 CAS 版本,再
进入验收。两种回复丢失都重放原请求;历史回执不能恢复已过期或被替换的执行。
验收期间不写租约,避免自己改变验收所依赖的 provider revision。

隔离 File/SQLite 回归覆盖剩余短租约跨越最终验收,以及续租/完成回复丢失、
过期和新 epoch。Python 仅衔接既有 TS claim、renew、terminal 权威合同,没有
新增 provider 规则、公开配置或前端权限。Stop ACK 与中断后的无进展结算仍归
#5308;持续 D2、默认准入和旧 writer 删除仍需各自证据。
22 changes: 22 additions & 0 deletions loopx/collaboration_mcp.py
Original file line number Diff line number Diff line change
Expand Up @@ -1167,6 +1167,28 @@ def _complete_delegated_todo(self, row: dict, binding: dict) -> None:
# key/epoch; it cannot reacquire an expired execution. Renewal has
# changed its version, so the historical acquisition is not CAS.
if "completion_lease_version" not in row:
# The Host supervisor has stopped. Renew the original execution
# before validation captures its provider revision; renewing
# during validation would invalidate that source witness. The
# canonical TS lease owner decides admission and replay. This
# adapter journals one intent, not a new lease or a longer TTL.
if "completion_lease_renewal_version" not in row:
proof = self._cli(binding, *self._delegation_claim_arguments(row, binding))
if proof.get("ok") is not True:
raise ValueError("delegation current execution proof lost before completion")
row["completion_lease_renewal_version"] = proof["lease"]["version"]
_write(self.path(row["identity"]["operation_id"]), row)
renewed = self._cli(
binding, "task-lease", "renew", "--goal-id", self.goal_id,
"--todo-id", binding["todo_id"], "--owner", binding["agent_id"],
"--idempotency-key", lease["lease"]["idempotency_key"],
"--expected-version", str(row["completion_lease_renewal_version"]),
"--ttl-seconds", str(lease["lease"]["acquire_ttl_seconds"]),
)
if renewed.get("ok") is not True:
raise ValueError("delegation original lease renewal rejected before completion")
# A renewal receipt can be historical after a lost reply. Read
# current authority before freezing the terminal intent below.
proof = self._cli(binding, *self._delegation_claim_arguments(row, binding))
current = proof.get("lease", {})
if (proof.get("ok") is not True or current.get("owner") != binding["agent_id"]
Expand Down
5 changes: 5 additions & 0 deletions loopx/control_plane/turn_driver/driver.py
Original file line number Diff line number Diff line change
Expand Up @@ -460,6 +460,11 @@ def build_loopx_turn_plan(
else LoopXTurnRoute.CONTRACT_ERROR
)
selected_todo = selected_turn_todo(envelope)
# Replanning an empty frontier is valid, but this Host executor settles
# against a Todo. Preserve the replan packet for the controller to resolve;
# do not turn the absence of a successor into a malformed Host session.
if route is LoopXTurnRoute.REPLAN_REQUIRED and not selected_todo:
route = LoopXTurnRoute.BLOCKED
lineage = _turn_lineage(envelope, selected_todo=selected_todo)
session, session_error = _session_plan(
route=route,
Expand Down
183 changes: 178 additions & 5 deletions tests/test_delegation_lease_lifetime.py
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,7 @@

import pytest

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


@pytest.fixture(params=["file", "sqlite"])
def completion_service(tmp_path, request, monkeypatch):
"""Pin a real slow final acceptance command before preparing authority."""
validator = tmp_path / "completion-validator.py"
validator.write_text('''import sys, time, runpy
from datetime import datetime
from pathlib import Path
root = Path(sys.argv[1])
counter = root / 'validator-calls'
calls = int(counter.read_text()) + 1 if counter.exists() else 1
counter.write_text(str(calls))
if calls == 2:
deadline = datetime.fromisoformat((root / 'completion-prior-expiry').read_text())
time.sleep(max(0, deadline.timestamp() - time.time()) + 0.2)
(root / 'completion-validation-ended').touch()
actual = root / 'project/validation/acceptance.py'
sys.path.insert(0, str(actual.parent))
sys.argv[0] = str(actual)
runpy.run_path(str(actual), run_name='__main__')
''')
write = demo.write

def configure(path, value):
if path.name == "bootstrap.json":
for criterion in value["document"]["criteria"]:
criterion["validation_timeout_seconds"] = 1
if criterion["id"] == "analyst-initial":
# 21 + four one-second criteria stays within the public
# 25-second completion budget, including a 20s deadline.
criterion["validation_timeout_seconds"] = 21
criterion["validation_argv"][1] = str(validator)
return write(path, value)

monkeypatch.setattr(demo, "write", configure)
return service.__wrapped__(tmp_path, request, monkeypatch)


@pytest.mark.parametrize("lost_reply", [None, "renewal", "completion"])
def test_completion_renews_before_validation_and_replays_each_intent(completion_service, monkeypatch, lost_reply):
root, runner = completion_service
# This case shortens the lease at the completion boundary below. An
# unrelated short Host deadline can cancel execution before that boundary
# under load; the running-Host cases separately exercise that deadline.
original = prepare_lease(root, runner, monkeypatch, ttl=None)
complete, cli = runner._complete_delegated_todo, runner._cli
shortened = None
renewal_calls, completion_calls = [], []
dropped = False

def enter_completion(row, binding):
nonlocal shortened
if shortened is None:
current = inspect(runner)["lease"]
# Fix the phase boundary, independent of whether the Host happened
# to cross its earlier renewal timer. Use the real canonical API.
shortened = cli(binding, "task-lease", "renew", "--goal-id", runner.goal_id,
"--todo-id", binding["todo_id"], "--owner", binding["agent_id"],
"--idempotency-key", current["idempotency_key"],
"--expected-version", str(current["version"]), "--ttl-seconds", "20")["lease"]
# Cross the actual pre-renewal deadline, not an assumed amount of
# CLI startup time. Allow cold claim/renew commands to reach the
# boundary; the independent validator still outlives that lease.
(root / "completion-prior-expiry").write_text(shortened["expires_at"])
return complete(row, binding)

def observe_reply(binding, *args, **kwargs):
nonlocal dropped
result = cli(binding, *args, **kwargs)
phase = None
if args[:2] == ("task-lease", "renew"):
renewal_calls.append((args, result))
phase = "renewal"
elif args[:2] == ("todo", "complete"):
completion_calls.append((args, result))
phase = "completion"
if phase == lost_reply and phase is not None and not dropped:
dropped = True
raise ValueError("fixture dropped the committed " + phase + " reply")
return result

monkeypatch.setattr(runner, "_complete_delegated_todo", enter_completion)
monkeypatch.setattr(runner, "_cli", observe_reply)
runner.execute("lease-lifetime")
if lost_reply:
uncertain = _read(runner.path("lease-lifetime"))
assert uncertain["status"] == "turn_returned", uncertain
assert "fixture dropped" in uncertain["error"]
assert uncertain["completion_lease_renewal_version"] == shortened["version"]
assert ("completion_lease_version" in uncertain) == (lost_reply == "completion")
runner.execute("lease-lifetime")
row = _read(runner.path("lease-lifetime"))
assert row["status"] == "accepted", row
assert row["turn_result"]["result_kind"] == "validated_progress"
assert (root / "completion-validation-ended").exists()
assert time.time() > datetime.fromisoformat(shortened["expires_at"].replace("Z", "+00:00")).timestamp()
assert (root / "analyst/initial/host-invocations").read_text() == "1"
final = inspect(runner)["lease"]
assert final["status"] == "released"
assert final["lease_epoch"] == original["lease_epoch"]
assert final["idempotency_key"] == original["idempotency_key"]
assert final["version"] == shortened["version"] + 1
assert len(renewal_calls) == (2 if lost_reply == "renewal" else 1)
assert len(completion_calls) == (2 if lost_reply == "completion" else 1)
for calls in (renewal_calls, completion_calls):
assert all(args == calls[0][0] for args, _ in calls)
assert all(result["provider_revision"] == calls[0][1]["provider_revision"] for _, result in calls)


@pytest.mark.parametrize("authority_loss", ["expiry", "replacement"])
def test_completion_renewal_receipt_cannot_revive_lost_execution(service, monkeypatch, authority_loss):
root, runner = service
prepare_lease(root, runner, monkeypatch)
cli = runner._cli
dropped = False
completions = []

def lose_renewal_reply(binding, *args, **kwargs):
nonlocal dropped
if args[:2] == ("todo", "complete"):
completions.append(args)
result = cli(binding, *args, **kwargs)
if args[:2] == ("task-lease", "renew") and not dropped:
dropped = True
raise ValueError("fixture lost renewal response before terminal intent")
return result

monkeypatch.setattr(runner, "_cli", lose_renewal_reply)
runner.execute("lease-lifetime")
row = _read(runner.path("lease-lifetime"))
assert row["status"] == "turn_returned", row
assert "completion_lease_renewal_version" in row
assert "completion_lease_version" not in row
current = inspect(runner)["lease"]
binding = runner.binding("analysis")
args = ("--goal-id", runner.goal_id, "--todo-id", binding["todo_id"],
"--owner", binding["agent_id"], "--idempotency-key", current["idempotency_key"],
"--expected-version", str(current["version"]))
if authority_loss == "expiry":
expired = cli(binding, "task-lease", "renew", *args, "--ttl-seconds", "1")["lease"]
time.sleep(max(0, datetime.fromisoformat(expired["expires_at"].replace("Z", "+00:00")).timestamp()
- time.time()) + 0.1)
else:
assert cli(binding, "task-lease", "release", *args)["ok"] is True
replacement = cli(binding, "task-lease", "acquire", "--goal-id", runner.goal_id,
"--todo-id", binding["todo_id"], "--owner", binding["agent_id"],
"--idempotency-key", "replacement", "--expected-version", str(current["version"]))
assert replacement["lease"]["lease_epoch"] > current["lease_epoch"]
runner.execute("lease-lifetime")
rejected = _read(runner.path("lease-lifetime"))
assert rejected["status"] != "accepted", rejected
assert "completion_lease_version" not in rejected
assert not completions
assert (root / "analyst/initial/host-invocations").read_text() == "1"
snapshot = read_canonical_todos_if_promoted(runtime_root=runner.root, goal_id=runner.goal_id)
todo = next(item for item in snapshot["todos"] if item["todo_id"] == binding["todo_id"])
assert todo["done"] is False


def await_started(root, future, runner):
deadline = time.monotonic() + 45
while time.monotonic() < deadline and not (root / "host-started").exists():
Expand Down Expand Up @@ -125,11 +285,24 @@ def test_real_revocation_or_new_execution_stops_nested_host_without_acceptance(s
with ThreadPoolExecutor(max_workers=1) as pool:
future = pool.submit(runner.execute, "lease-lifetime")
await_started(root, future, runner)
current = inspect(runner)["lease"]
binding = runner.binding("analysis")
runner._cli(binding, "task-lease", "release", "--goal-id", runner.goal_id,
"--todo-id", "todo_analyst-initial", "--owner", "analyst",
"--idempotency-key", original["idempotency_key"], "--expected-version", str(current["version"]))
# The live supervisor may renew between inspection and revocation.
# Retry only that CAS race, never a replacement execution or rejection.
for _ in range(5):
current = inspect(runner)["lease"]
assert current["lease_epoch"] == original["lease_epoch"]
assert current["idempotency_key"] == original["idempotency_key"]
try:
runner._cli(binding, "task-lease", "release", "--goal-id", runner.goal_id,
"--todo-id", "todo_analyst-initial", "--owner", "analyst",
"--idempotency-key", original["idempotency_key"],
"--expected-version", str(current["version"]))
break
except ValueError as exc:
if str(exc) != "canonical task lease release rejected: version_mismatch":
raise
else:
pytest.fail("could not revoke the original execution during concurrent renewal")
if reclaim:
acquired = runner._cli(binding, "task-lease", "acquire", "--goal-id", runner.goal_id,
"--todo-id", "todo_analyst-initial", "--owner", "analyst", "--idempotency-key", "new-execution",
Expand Down
Loading
Loading