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 @@ -817,6 +817,29 @@ promotion retain their own acceptance. No new paid cohort or soak is authorized.
row. Todo/lease mutation admission, quota, automation, Goal Channel, and the
overall M3 activation hold remain separate.

### 2026-09-30: M3 Turn-journal exact commit candidate

- **Baseline:** `350f0f326`.
- **Proposed:** Keep `loopx_turn_journal_v0` and its existing path. Source-profile
plans carry matching exact GoalRef copies in the plan and transaction. Every
executor journal mutation uses one persistence callback that hands the
alias-scoped lifecycle guard from Python to TypeScript. TypeScript claims and
verifies that witness, reuses `decideFirstPartyHostRuntime(require_current)`,
then takes the existing journal mutation lock and commits before releasing
the source guard.
- **Evidence:** TypeScript owner tests reject missing, malformed, mismatched,
expired and stale source admission before journal mutation. A real
Python-to-TypeScript integration commits and replays Goal A, publishes
same-alias Goal B, proves a later A checkpoint leaves both files unchanged,
and commits B independently.
- **Compatibility:** Read-only inspection and recovery still read legacy
journals. Non-source writes retain the old RPC shape and persisted bytes.
Source admission facts and lock tokens are transport-only and never persist.
- **Remaining hold:** This qualifies only the `turn_journal` inventory row. It
does not complete `first_party_host_runtime`, downstream external-effect
drain, unsupported/warm binary coverage, or any other M3 row.
`execution_authority: false` and the overall activation hold remain.

## Appendix B: Decision log

| Date | Decision | Owner / approval | Alternatives | Normative sections changed |
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -740,6 +740,27 @@ service adoption、D1–D3 provider promotion 保留各自验收。不授权付
Todo/lease mutation admission、quota、automation、Goal Channel 与 M3 总
activation hold 仍分别验收。

### 2026-09-30:M3 Turn journal 精确提交候选

- **基线:** `350f0f326`。
- **候选实现:** 保持 `loopx_turn_journal_v0` 和既有路径。Source profile
plan 在 plan 与 transaction 中携带一致的精确 GoalRef。Executor 的每次
journal mutation 都通过同一个持久化回调,把 alias-scoped lifecycle guard
从 Python 交接给 TypeScript。TypeScript claim 并复核 witness,复用
`decideFirstPartyHostRuntime(require_current)`,再取得既有 journal mutation
lock 并提交,最后释放 source guard。
- **证据:** TypeScript owner 测试证明缺失、畸形、副本不一致、失效和 stale
source admission 均在 journal mutation 前拒绝。真实 Python-to-TypeScript
集成先提交并重放 Goal A,再发布同名 Goal B,证明迟到的 A checkpoint
不改变 A/B 文件,随后 B 可独立提交。
- **兼容性:** 只读 inspection 和 recovery 仍可读取 legacy journal。非 source
写入保留旧 RPC shape 与持久化字节。Source admission facts 和 lock token
只用于 transport,不落盘。
- **剩余 hold:** 本切片只资格化 `turn_journal` inventory 行,不代表
`first_party_host_runtime`、downstream external-effect drain、不支持的旧/常驻
binary 或其他 M3 行已完成。`execution_authority: false` 和总 activation hold
保持不变。

## 附录 B:决策日志

| 日期 | 决策 | Owner/批准 | 替代方案 | 变更的规范章节 |
Expand Down
34 changes: 32 additions & 2 deletions loopx/control_plane/goals/first_party_host_admission.py
Original file line number Diff line number Diff line change
@@ -1,11 +1,15 @@
from __future__ import annotations

from collections.abc import Callable, Mapping
from collections.abc import Callable, Iterator, Mapping
from contextlib import contextmanager
from dataclasses import dataclass
from pathlib import Path
from typing import Any, TypeVar

from ...file_lock import exclusive_cross_runtime_file_lock
from ...file_lock import (
cross_runtime_lock_witness,
exclusive_cross_runtime_file_lock,
)
from ..effect_runtime import effect_runtime_result
from ..projects.registry_codec import (
SOURCE_SESSION_PROFILE_ID,
Expand Down Expand Up @@ -197,6 +201,32 @@ def require_current(self) -> None:
):
self._decision("require_current")

@contextmanager
def source_journal_admission(
self,
) -> Iterator[dict[str, Any] | None]:
"""Hand one journal mutation to the TS owner under the source guard."""

if not self.source_profile:
yield None
return
target = guard_path(self.registry_path, self.goal_id)
with exclusive_cross_runtime_file_lock(
target,
operation="first_party_host_journal_commit",
):
yield {
"schema_version": "loopx_turn_journal_source_admission_v0",
"profile_id": SOURCE_SESSION_PROFILE_ID,
"registry_path": str(self.registry_path),
"planned_goal_ref": self.planned_goal_ref,
"authority": _source_authority(
self.registry_path,
self.goal_id,
),
"lock": cross_runtime_lock_witness(target),
}

def accept_result(self, commit_result: Callable[[], T]) -> T:
if not self.source_profile:
return commit_result()
Expand Down
91 changes: 50 additions & 41 deletions loopx/control_plane/turn_driver/executor.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@
from ...authority import validate_public_safe_text
from ...file_lock import LockAcquireTimeoutError, exclusive_file_lock
from ...runtime import validate_goal_id_path_segment
from ..effect_runtime import EffectRuntimeConflict
from ..effect_program import (
SettlementStepKind,
interpret_turn_result_packet,
Expand Down Expand Up @@ -127,6 +128,7 @@
Spend = Callable[..., dict[str, Any]]
Scheduler = Callable[[dict[str, Any]], dict[str, Any]]
HostRunner = Callable[[Mapping[str, Any]], dict[str, Any]]
JournalPersist = Callable[[Mapping[str, Any]], None]


def build_loopx_turn_host_request(plan: Mapping[str, Any]) -> dict[str, Any]:
Expand Down Expand Up @@ -758,12 +760,11 @@ def _host_result_stage(
project: Path,
timeout_seconds: float,
journal: dict[str, Any],
journal_path: Path,
persist_journal: JournalPersist,
effects: dict[str, bool],
confirm_start: Callable[[], None] | None = None,
usage_runtime_root: Path | None = None,
usage_goal_id: str = "",
goal_admission: FirstPartyHostGoalAdmission | None = None,
) -> tuple[dict[str, Any] | None, list[str], dict[str, Any] | None]:
completed_phases = list(journal.get("completed_phases") or [])
result = (
Expand All @@ -773,7 +774,7 @@ def _host_result_stage(
)
if "typed_result" not in completed_phases:
journal["host_attempt_count"] = int(journal.get("host_attempt_count") or 0) + 1
_write_journal(journal_path, journal)
persist_journal(journal)
# The attempt is durable now, so a later restart must not resume this
# reservation. Confirmation failure stops before the host starts.
if confirm_start is not None:
Expand Down Expand Up @@ -815,12 +816,7 @@ def _host_result_stage(
journal["host_recovery"] = build_host_recovery_record(recovery_kind)
else:
journal.pop("host_recovery", None)
if goal_admission is None:
_write_journal(journal_path, journal)
else:
goal_admission.accept_result(
lambda: _write_journal(journal_path, journal)
)
persist_journal(journal)
return (
None,
[],
Expand Down Expand Up @@ -858,12 +854,7 @@ def _host_result_stage(
result_kind=LoopXTurnResultKind.VALIDATION_FAILED.value,
validation_stage="host_result_contract",
)
if goal_admission is None:
_write_journal(journal_path, journal)
else:
goal_admission.accept_result(
lambda: _write_journal(journal_path, journal)
)
persist_journal(journal)
return (
None,
list(TRANSACTION_PHASES[:2]),
Expand All @@ -884,12 +875,7 @@ def _host_result_stage(
result_kind=normalized.get("result_kind"),
completed_phases=completed_phases,
)
if goal_admission is None:
_write_journal(journal_path, journal)
else:
goal_admission.accept_result(
lambda: _write_journal(journal_path, journal)
)
persist_journal(journal)
return normalized, completed_phases, None


Expand All @@ -900,7 +886,7 @@ def _task_validation_stage(
task_validator: TaskValidator | None,
completed_phases: list[str],
journal: dict[str, Any],
journal_path: Path,
persist_journal: JournalPersist,
effects: dict[str, bool],
) -> tuple[list[str], dict[str, Any] | None]:
turn = interpret_turn_result_packet(result)
Expand All @@ -918,7 +904,7 @@ def _task_validation_stage(
receipt=_receipt(plan, result, completed_phases=completed_phases),
scheduler={"disposition": "not_applicable"},
)
_write_journal(journal_path, journal)
persist_journal(journal)
return completed_phases, execution_payload(
plan,
journal,
Expand Down Expand Up @@ -963,7 +949,7 @@ def _task_validation_stage(
result_kind=LoopXTurnResultKind.VALIDATION_FAILED.value,
validation_stage="task_postcondition",
)
_write_journal(journal_path, journal)
persist_journal(journal)
return list(TRANSACTION_PHASES[:2]), execution_payload(
plan,
journal,
Expand All @@ -979,7 +965,7 @@ def _task_validation_stage(
completed_phases=completed_phases,
validation_stage="task_postcondition",
)
_write_journal(journal_path, journal)
persist_journal(journal)
return completed_phases, None


Expand Down Expand Up @@ -1024,7 +1010,7 @@ def _typed_settlement_stage(
*,
completed_phases: list[str],
journal: dict[str, Any],
journal_path: Path,
persist_journal: JournalPersist,
effects: dict[str, bool],
writeback: Writeback,
completion_writeback: CompletionWriteback | None,
Expand Down Expand Up @@ -1104,7 +1090,7 @@ def writeback_effect(effect_ref: str) -> Mapping[str, Any]:
journal_adapter = TurnSettlementJournalAdapter(
journal,
effects,
lambda: _write_journal(journal_path, journal),
lambda: persist_journal(journal),
_compact_callback,
)

Expand Down Expand Up @@ -1168,7 +1154,7 @@ def writeback_effect(effect_ref: str) -> Mapping[str, Any]:
reason=failure["reason"],
receipt=failure["receipt"],
)
_write_journal(journal_path, journal)
persist_journal(journal)
return execution_payload(
plan,
journal,
Expand All @@ -1188,7 +1174,7 @@ def writeback_effect(effect_ref: str) -> Mapping[str, Any]:
result = {**result, "result_kind": outcome["result_kind"]}
completed_phases = [str(phase) for phase in outcome["completed_phases"]]
spend_payload = dict(settlement_state.quota_spend)
_write_journal(journal_path, journal)
persist_journal(journal)

scheduler_payload = scheduler(spend_payload)
journal["scheduler"] = scheduler_payload
Expand All @@ -1197,14 +1183,14 @@ def writeback_effect(effect_ref: str) -> Mapping[str, Any]:
result=result,
post_settlement=post_settlement,
journal=journal,
journal_path=journal_path,
persist_journal=persist_journal,
)
if scheduler_payload.get("completed") is not True:
journal.update(
status="scheduler_action_required",
receipt=_receipt(plan, result, completed_phases=completed_phases),
)
_write_journal(journal_path, journal)
persist_journal(journal)
return execution_payload(
plan,
journal,
Expand All @@ -1220,7 +1206,7 @@ def writeback_effect(effect_ref: str) -> Mapping[str, Any]:
completed_phases=completed_phases,
receipt=_receipt(plan, result, completed_phases=completed_phases),
)
_write_journal(journal_path, journal)
persist_journal(journal)
return execution_payload(
plan,
journal,
Expand Down Expand Up @@ -1312,6 +1298,30 @@ def run_loopx_turn_once(

turn_key = str(request["turn_key"])
journal_path = turn_journal_path(runtime_root, goal_id=goal_id, turn_key=turn_key)

def persist_journal(snapshot: Mapping[str, Any]) -> None:
if goal_admission is None or not goal_admission.enabled:
_write_journal(journal_path, snapshot)
return
try:
with goal_admission.source_journal_admission() as source_admission:
if source_admission is None:
raise RuntimeError("source journal admission was not produced")
_write_journal(
journal_path,
snapshot,
source_admission=source_admission,
)
except EffectRuntimeConflict as exc:
if exc.diagnostic_code in {
"goal_not_registered",
"goal_authority_unavailable",
"goal_instance_id_missing",
"stale_goal_instance",
}:
raise FirstPartyHostRuntimeRejected(exc.diagnostic_code) from exc
raise

with exclusive_file_lock(journal_path):
journal = _load_journal(journal_path)
recovery_decision: dict[str, Any] | None = None
Expand Down Expand Up @@ -1390,7 +1400,7 @@ def run_loopx_turn_once(
journal["recovery_audit"] = build_turn_recovery_audit(
recovery_decision, journal, status="started", host_invoked=None,
)
_write_journal(journal_path, journal)
persist_journal(journal)

if journal and journal.get("status") == "failed":
receipt = (
Expand All @@ -1414,7 +1424,7 @@ def run_loopx_turn_once(
journal.pop("host_recovery", None)
journal.pop("host_failure", None)
journal["status"] = "in_progress"
_write_journal(journal_path, journal)
persist_journal(journal)
if journal is None:
journal = {
"schema_version": LOOPX_TURN_JOURNAL_SCHEMA_VERSION,
Expand All @@ -1425,10 +1435,10 @@ def run_loopx_turn_once(
"completed_phases": [],
"plan": dict(plan),
}
_write_journal(journal_path, journal)
persist_journal(journal)
if admission is not None and admission.get("reserved") is True:
journal["admission"] = admission
_write_journal(journal_path, journal)
persist_journal(journal)

effects = dict(empty_effects)

Expand All @@ -1441,7 +1451,7 @@ def finish_recovery(payload: dict[str, Any]) -> dict[str, Any]:
status="finished",
host_invoked=effects.get("host_invoked") is True,
)
_write_journal(journal_path, journal)
persist_journal(journal)
payload["recovery"] = dict(journal["recovery_audit"])
return payload

Expand All @@ -1463,14 +1473,13 @@ def finish_recovery(payload: dict[str, Any]) -> dict[str, Any]:
project=project,
timeout_seconds=timeout_seconds,
journal=journal,
journal_path=journal_path,
persist_journal=persist_journal,
effects=effects,
confirm_start=(
confirm_start
if admission is not None and admission.get("reserved") is True
else None
),
goal_admission=goal_admission,
)
if terminal is not None:
return finish_recovery(terminal)
Expand All @@ -1482,7 +1491,7 @@ def finish_recovery(payload: dict[str, Any]) -> dict[str, Any]:
task_validator=task_validator,
completed_phases=completed_phases,
journal=journal,
journal_path=journal_path,
persist_journal=persist_journal,
effects=effects,
)
if terminal is not None:
Expand All @@ -1493,7 +1502,7 @@ def finish_recovery(payload: dict[str, Any]) -> dict[str, Any]:
result,
completed_phases=completed_phases,
journal=journal,
journal_path=journal_path,
persist_journal=persist_journal,
effects=effects,
writeback=writeback,
completion_writeback=completion_writeback,
Expand Down
Loading
Loading