diff --git a/docs/architecture/rfcs/goal-instance-identity-and-orphan-recovery-v0.md b/docs/architecture/rfcs/goal-instance-identity-and-orphan-recovery-v0.md index e63344dda..23b84bf36 100644 --- a/docs/architecture/rfcs/goal-instance-identity-and-orphan-recovery-v0.md +++ b/docs/architecture/rfcs/goal-instance-identity-and-orphan-recovery-v0.md @@ -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 | diff --git a/docs/architecture/rfcs/goal-instance-identity-and-orphan-recovery-v0.zh-CN.md b/docs/architecture/rfcs/goal-instance-identity-and-orphan-recovery-v0.zh-CN.md index 52b049de3..8ccb2d1c9 100644 --- a/docs/architecture/rfcs/goal-instance-identity-and-orphan-recovery-v0.zh-CN.md +++ b/docs/architecture/rfcs/goal-instance-identity-and-orphan-recovery-v0.zh-CN.md @@ -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/批准 | 替代方案 | 变更的规范章节 | diff --git a/loopx/control_plane/goals/first_party_host_admission.py b/loopx/control_plane/goals/first_party_host_admission.py index 0d2b45e03..ae315c30b 100644 --- a/loopx/control_plane/goals/first_party_host_admission.py +++ b/loopx/control_plane/goals/first_party_host_admission.py @@ -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, @@ -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() diff --git a/loopx/control_plane/turn_driver/executor.py b/loopx/control_plane/turn_driver/executor.py index 57cc53a97..99e941cfc 100644 --- a/loopx/control_plane/turn_driver/executor.py +++ b/loopx/control_plane/turn_driver/executor.py @@ -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, @@ -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]: @@ -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 = ( @@ -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: @@ -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, [], @@ -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]), @@ -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 @@ -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) @@ -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, @@ -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, @@ -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 @@ -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, @@ -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, ) @@ -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, @@ -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 @@ -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, @@ -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, @@ -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 @@ -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 = ( @@ -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, @@ -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) @@ -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 @@ -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) @@ -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: @@ -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, diff --git a/loopx/control_plane/turn_driver/journal_store.py b/loopx/control_plane/turn_driver/journal_store.py index ccc072509..142928965 100644 --- a/loopx/control_plane/turn_driver/journal_store.py +++ b/loopx/control_plane/turn_driver/journal_store.py @@ -58,11 +58,17 @@ def journal_committed_effect_id(journal: Mapping[str, Any]) -> str | None: return effect_id or None -def write_turn_journal_checkpoint(path: Path, journal: Mapping[str, Any]) -> None: +def write_turn_journal_checkpoint( + path: Path, + journal: Mapping[str, Any], + *, + source_admission: Mapping[str, Any] | None = None, +) -> None: write_turn_journal( str(path), journal, expected_effect_id=journal_committed_effect_id(journal), + source_admission=source_admission, ) diff --git a/loopx/control_plane/turn_driver/post_settlement.py b/loopx/control_plane/turn_driver/post_settlement.py index b1630eebc..34458715c 100644 --- a/loopx/control_plane/turn_driver/post_settlement.py +++ b/loopx/control_plane/turn_driver/post_settlement.py @@ -3,11 +3,8 @@ from __future__ import annotations from collections.abc import Callable, Mapping -from pathlib import Path from typing import Any -from .journal_store import write_turn_journal_checkpoint - PostSettlement = Callable[ [Mapping[str, Any], Mapping[str, Any], Mapping[str, Any]], @@ -21,7 +18,7 @@ def run_post_settlement_callback( result: Mapping[str, Any], post_settlement: PostSettlement | None, journal: dict[str, Any], - journal_path: Path, + persist_journal: Callable[[Mapping[str, Any]], None], ) -> None: if post_settlement is None: return @@ -56,4 +53,4 @@ def run_post_settlement_callback( "fail_open": True, "external_writes_performed": False, } - write_turn_journal_checkpoint(journal_path, journal) + persist_journal(journal) diff --git a/loopx/control_plane/turn_driver/turn_journal.ts b/loopx/control_plane/turn_driver/turn_journal.ts index 3b10f9e0b..1e31a6ce0 100644 --- a/loopx/control_plane/turn_driver/turn_journal.ts +++ b/loopx/control_plane/turn_driver/turn_journal.ts @@ -11,6 +11,7 @@ import { type EffectTurn, } from "../effect_program.ts"; import { EffectRuntimeRequestError } from "../effect_runtime_errors.ts"; +import { parseExactGoalRef } from "../goals/goal_instance_identity.ts"; import { preparedAttemptViolation } from "./turn_journal_attempt_contract.ts"; import { recordedTurnEffects, type RecordedTurnEffects } from "./turn_journal_effect_readback.ts"; @@ -19,6 +20,22 @@ export const TURN_JOURNAL_INSPECTION_SCHEMA_VERSION = type JsonObject = Record; +type WireGoalRef = Readonly<{ + goal_id: string; + goal_instance_id: string; +}>; + +export type TurnJournalGoalBinding = + | Readonly<{ kind: "legacy" }> + | Readonly<{ kind: "exact"; goal_ref: WireGoalRef }> + | Readonly<{ + kind: "invalid"; + violation: + | "goal_ref_binding_incomplete" + | "goal_ref_binding_invalid" + | "goal_ref_binding_mismatch"; + }>; + export interface TurnJournalInspectionRequest { schema_version: "loopx_turn_journal_interpretation_request_v0"; journal: JsonObject; @@ -145,6 +162,36 @@ function asObject(value: unknown): JsonObject { : {}; } +export function parseTurnJournalGoalBinding( + journal: JsonObject, +): TurnJournalGoalBinding { + const plan = asObject(journal.plan); + const transaction = asObject(plan.transaction); + const hasPlanRef = Object.hasOwn(plan, "goal_ref"); + const hasTransactionRef = Object.hasOwn(transaction, "goal_ref"); + if (!hasPlanRef && !hasTransactionRef) return { kind: "legacy" }; + if (!hasPlanRef || !hasTransactionRef) { + return { kind: "invalid", violation: "goal_ref_binding_incomplete" }; + } + const planned = parseExactGoalRef(plan.goal_ref); + const transactionPlanned = parseExactGoalRef(transaction.goal_ref); + if (planned.kind === "invalid" || transactionPlanned.kind === "invalid") { + return { kind: "invalid", violation: "goal_ref_binding_invalid" }; + } + const goalRef = { + goal_id: planned.value.goalId.value, + goal_instance_id: planned.value.goalInstanceId.value, + }; + if ( + goalRef.goal_id !== transactionPlanned.value.goalId.value + || goalRef.goal_instance_id !== + transactionPlanned.value.goalInstanceId.value + ) { + return { kind: "invalid", violation: "goal_ref_binding_mismatch" }; + } + return { kind: "exact", goal_ref: goalRef }; +} + function isValidIdentity(value: unknown): value is string { return typeof value === "string" && value.trim().length > 0; } @@ -559,6 +606,7 @@ export function interpretTurnJournalEffect( const identity = asObject(settlement.identity); const hostResult = asObject(journal.host_result); const receipt = asObject(journal.receipt); + const goalBinding = parseTurnJournalGoalBinding(journal); const [goalComplete, goalMatches] = identityState( [journal.goal_id, envelope.goal_id, identity.goal_id], @@ -587,6 +635,14 @@ export function interpretTurnJournalEffect( const violations: string[] = []; if (!goalComplete) violations.push("goal_identity_missing"); else if (!goalMatches) violations.push("goal_mismatch"); + if (goalBinding.kind === "invalid") { + violations.push(goalBinding.violation); + } else if ( + goalBinding.kind === "exact" + && goalBinding.goal_ref.goal_id !== request.goal_id + ) { + violations.push("goal_ref_binding_mismatch"); + } if (!ownerComplete) violations.push("owner_identity_missing"); else if (!ownerMatches) violations.push("owner_mismatch"); if (!settlementIdentityValid) violations.push("settlement_identity_invalid"); @@ -625,6 +681,11 @@ export function interpretTurnJournalEffect( const lineageConsistent = goalMatches && + goalBinding.kind !== "invalid" && + ( + goalBinding.kind !== "exact" + || goalBinding.goal_ref.goal_id === request.goal_id + ) && ownerMatches && settlementIdentityValid && settlementTurnInstanceMatches && diff --git a/loopx/control_plane/turn_driver/turn_journal_effects.ts b/loopx/control_plane/turn_driver/turn_journal_effects.ts index c86da287d..09c203d27 100644 --- a/loopx/control_plane/turn_driver/turn_journal_effects.ts +++ b/loopx/control_plane/turn_driver/turn_journal_effects.ts @@ -1,6 +1,6 @@ import { createHash } from "node:crypto"; import { readFile } from "node:fs/promises"; -import { isAbsolute } from "node:path"; +import { dirname, isAbsolute, join, resolve } from "node:path"; import type { JsonObject } from "../effect_program.ts"; import { settlementIdentityFromPlan } from "../effect_program.ts"; @@ -10,16 +10,27 @@ import { } from "../effect_runtime_errors.ts"; import { atomicWriteJson, + claimFileMutationLock, + mutationLockOwner, + releaseFileMutationLock, + releaseFileMutationLockClaim, withFileMutationLock, } from "../effect_runtime_io.ts"; +import { decideFirstPartyHostRuntime } from "../goals/first_party_host_runtime.ts"; +import { parseExactGoalRef } from "../goals/goal_instance_identity.ts"; import { requireNonEmptyString as requiredString } from "../runtime_decode.ts"; import { preparedAttemptViolation } from "./turn_journal_attempt_contract.ts"; import { interpretTurnJournalEffect, + parseTurnJournalGoalBinding, supportedJournalStatuses, transactionPhases, + type TurnJournalGoalBinding, } from "./turn_journal.ts"; +const SOURCE_ADMISSION_SCHEMA_VERSION = + "loopx_turn_journal_source_admission_v0"; +const SOURCE_SESSION_PROFILE_ID = "source_session_v1"; const terminalStatuses = new Set(["committed", "stopped"]); const statusTransitions: Readonly>> = { in_progress: new Set([ @@ -78,6 +89,20 @@ interface JournalState { failedPhase: string | null; } +interface SourceAdmission { + registryPath: string; + plannedGoalRef: { + goal_id: string; + goal_instance_id: string; + }; + authority: unknown; + lock: { + target: string; + pid: number; + token: string; + }; +} + function conflict(message: string, code = "journal_transition_conflict"): never { throw new EffectRuntimeConflictError(message, code); } @@ -161,6 +186,101 @@ function requireJournalState(journal: JsonObject): JournalState { return state; } +function sameGoalRef( + left: SourceAdmission["plannedGoalRef"], + right: SourceAdmission["plannedGoalRef"], +): boolean { + return left.goal_id === right.goal_id + && left.goal_instance_id === right.goal_instance_id; +} + +function sourceGuardTarget(registryPath: string, goalId: string): string { + return join( + dirname(registryPath), + ".loopx", + "lifecycle", + "goal-instance", + "guards", + `${sha256(goalId)}.guard`, + ); +} + +function requireSourceAdmission( + value: unknown, + binding: Extract, +): SourceAdmission { + const admission = asObject(value); + if ( + admission.schema_version !== SOURCE_ADMISSION_SCHEMA_VERSION + || admission.profile_id !== SOURCE_SESSION_PROFILE_ID + ) { + throw new EffectRuntimeRequestError( + "Turn journal source admission is malformed", + "journal_source_admission_invalid", + ); + } + const registryPath = requiredString( + admission.registry_path, + "source admission registry_path", + ); + if (!isAbsolute(registryPath) || resolve(registryPath) !== registryPath) { + throw new EffectRuntimeRequestError( + "Turn journal source admission registry path must be absolute and normalized", + "journal_source_admission_invalid", + ); + } + const planned = parseExactGoalRef(admission.planned_goal_ref); + if (planned.kind === "invalid") { + throw new EffectRuntimeRequestError( + "Turn journal source admission has an invalid planned GoalRef", + "journal_source_admission_invalid", + ); + } + const plannedGoalRef = { + goal_id: planned.value.goalId.value, + goal_instance_id: planned.value.goalInstanceId.value, + }; + if (!sameGoalRef(plannedGoalRef, binding.goal_ref)) { + throw new EffectRuntimeRequestError( + "Turn journal source admission does not match the journal GoalRef", + "journal_source_admission_invalid", + ); + } + const lock = asObject(admission.lock); + const target = requiredString(lock.target, "source admission lock target"); + const expectedTarget = sourceGuardTarget(registryPath, plannedGoalRef.goal_id); + if ( + !isAbsolute(target) + || resolve(target) !== target + || target !== expectedTarget + ) { + throw new EffectRuntimeRequestError( + "Turn journal source admission lock target mismatch", + "journal_source_admission_invalid", + ); + } + if ( + typeof lock.pid !== "number" + || !Number.isSafeInteger(lock.pid) + || lock.pid <= 0 + ) { + throw new EffectRuntimeRequestError( + "Turn journal source admission lock owner is invalid", + "journal_source_admission_invalid", + ); + } + return { + registryPath, + plannedGoalRef, + authority: admission.authority, + lock: { + target, + pid: lock.pid, + token: requiredString(lock.token, "source admission lock token"), + }, + }; +} + function requireJournalTransition( existing: JsonObject, incoming: JsonObject, @@ -217,6 +337,13 @@ export async function commitTurnJournal( throw new EffectRuntimeRequestError("Turn journal path must be absolute"); } const journal = asObject(params.journal); + const goalBinding = parseTurnJournalGoalBinding(journal); + if (goalBinding.kind === "invalid") { + throw new EffectRuntimeRequestError( + `Turn journal GoalRef binding is invalid: ${goalBinding.violation}`, + "journal_snapshot_invalid", + ); + } const incomingState = requireJournalState(journal); const expectedEffectId = typeof params.expected_effect_id === "string" ? params.expected_effect_id.trim() @@ -228,48 +355,122 @@ export async function commitTurnJournal( ); } const incomingOperationId = operationId(journal); - return await withFileMutationLock(path, async () => { - let existing: JsonObject | null = null; - try { - const encoded = await readFile(path, "utf8"); - existing = asObject(JSON.parse(encoded)); - const existingState = requireJournalState(existing); - const existingEffectId = existingState.effectId; - if ( - existingEffectId && - existingEffectId !== incomingEffectId - ) { - throw new EffectRuntimeConflictError( - "Turn journal belongs to another settlement effect", - "journal_effect_conflict", - ); + const commit = async (): Promise => + await withFileMutationLock(path, async () => { + let existing: JsonObject | null = null; + try { + const encoded = await readFile(path, "utf8"); + existing = asObject(JSON.parse(encoded)); + const existingState = requireJournalState(existing); + const existingEffectId = existingState.effectId; + if ( + existingEffectId && + existingEffectId !== incomingEffectId + ) { + throw new EffectRuntimeConflictError( + "Turn journal belongs to another settlement effect", + "journal_effect_conflict", + ); + } + if (operationId(existing) === incomingOperationId) { + return { + ok: true, + appended: false, + replayed: true, + effect_id: incomingEffectId, + operation_id: incomingOperationId, + }; + } + requireJournalTransition(existing, journal, existingState, incomingState); + } catch (error) { + if ((error as NodeJS.ErrnoException).code !== "ENOENT") throw error; } - if (operationId(existing) === incomingOperationId) { - return { - ok: true, - appended: false, - replayed: true, - effect_id: incomingEffectId, - operation_id: incomingOperationId, - }; + if (existing === null && ( + incomingState.status !== "in_progress" || + incomingState.completedPhases.length !== 0 + )) { + conflict("A new Turn journal must begin in progress with no completed phases"); } - requireJournalTransition(existing, journal, existingState, incomingState); - } catch (error) { - if ((error as NodeJS.ErrnoException).code !== "ENOENT") throw error; + await atomicWriteJson(path, journal); + return { + ok: true, + appended: true, + replayed: false, + effect_id: incomingEffectId, + operation_id: incomingOperationId, + }; + }); + if (goalBinding.kind === "legacy") { + if (params.source_admission !== undefined) { + throw new EffectRuntimeRequestError( + "Legacy Turn journals cannot carry source admission", + "journal_source_admission_invalid", + ); } - if (existing === null && ( - incomingState.status !== "in_progress" || - incomingState.completedPhases.length !== 0 - )) { - conflict("A new Turn journal must begin in progress with no completed phases"); + return await commit(); + } + if (params.source_admission === undefined) { + throw new EffectRuntimeRequestError( + "Exact GoalRef Turn journals require source admission", + "journal_source_admission_required", + ); + } + const admission = requireSourceAdmission( + params.source_admission, + goalBinding, + ); + const claim = await claimFileMutationLock( + admission.lock.target, + admission.lock.token, + ); + if (!claim) { + conflict( + "Turn journal source admission lock handoff expired", + "journal_source_admission_expired", + ); + } + let adopted = false; + try { + const owner = await mutationLockOwner(admission.lock.target); + if ( + owner?.pid !== admission.lock.pid + || owner.token !== admission.lock.token + ) { + conflict( + "Turn journal source admission lock owner changed", + "journal_source_admission_expired", + ); } - await atomicWriteJson(path, journal); - return { - ok: true, - appended: true, - replayed: false, - effect_id: incomingEffectId, - operation_id: incomingOperationId, - }; - }); + adopted = true; + const decision = decideFirstPartyHostRuntime({ + profile_id: SOURCE_SESSION_PROFILE_ID, + operation: "require_current", + planned_goal_ref: admission.plannedGoalRef, + authority: admission.authority, + }); + if (decision.kind === "reject") { + conflict( + `Turn journal source admission rejected: ${decision.code}`, + decision.code, + ); + } + if (decision.kind !== "resume") { + throw new EffectRuntimeRequestError( + "Turn journal source admission did not resume the exact GoalRef", + "journal_source_admission_invalid", + ); + } + return await commit(); + } finally { + if (adopted) { + await releaseFileMutationLock( + admission.lock.target, + admission.lock.token, + claim, + true, + ); + } else { + await releaseFileMutationLockClaim(claim); + } + } } diff --git a/loopx/control_plane/turn_driver/turn_journal_runtime.py b/loopx/control_plane/turn_driver/turn_journal_runtime.py index 0065faaf3..226723d91 100644 --- a/loopx/control_plane/turn_driver/turn_journal_runtime.py +++ b/loopx/control_plane/turn_driver/turn_journal_runtime.py @@ -166,16 +166,20 @@ def write_turn_journal( journal: Mapping[str, Any], *, expected_effect_id: str | None = None, + source_admission: Mapping[str, Any] | None = None, ) -> dict[str, object]: """Commit a Turn-journal transition through the TS semantic owner.""" + request: dict[str, Any] = { + "path": path, + "journal": dict(journal), + "expected_effect_id": expected_effect_id, + } + if source_admission is not None: + request["source_admission"] = dict(source_admission) payload = effect_runtime_result( "turn_journal.write", - { - "path": path, - "journal": dict(journal), - "expected_effect_id": expected_effect_id, - }, + request, retry_safe=True, ) if ( diff --git a/loopx/semantics/goal_instance_binding_inventory_v1.json b/loopx/semantics/goal_instance_binding_inventory_v1.json index f5dfe3651..9a3a42ad9 100644 --- a/loopx/semantics/goal_instance_binding_inventory_v1.json +++ b/loopx/semantics/goal_instance_binding_inventory_v1.json @@ -217,20 +217,22 @@ { "owner_id": "turn_journal", "locator": "/goals//turns/.json", - "revision_signal": "turn_key and turn_instance_id", + "revision_signal": "turn_key, turn_instance_id, and exact GoalRef", "content_digest_signal": "turn key SHA-256", - "observed_reference": "goal_id", + "observed_reference": "goal_id + goal_instance_id for source_session_v1; goal_id otherwise", "producer_sites": [ + "loopx/control_plane/goals/first_party_host_admission.py::FirstPartyHostGoalAdmission.source_journal_admission", "loopx/control_plane/turn_driver/turn_journal_runtime.py::write_turn_journal" ], "consumer_sites": [ - "loopx/control_plane/turn_driver/journal_store.py::load_loopx_turn_plan_from_journal" + "loopx/control_plane/turn_driver/journal_store.py::load_loopx_turn_plan_from_journal", + "loopx/control_plane/turn_driver/turn_journal.ts::parseTurnJournalGoalBinding" ], - "effect_boundary": "loopx/control_plane/turn_driver/turn_journal.ts::interpretTurnJournalEffect", + "effect_boundary": "loopx/control_plane/turn_driver/turn_journal_effects.ts::commitTurnJournal", "authority_role": "turn_execution_journal", - "current_identity_strength": "goal_alias_only", + "current_identity_strength": "exact_goal_ref_enforced", "cleanup_support": "journal_recovery_and_terminal_closeout", - "m1_disposition": "alias_only_inventory", + "m1_disposition": "m3_qualified", "target_milestone": "M3" } ] diff --git a/loopx/semantics/project_registry_io_manifest_v1.json b/loopx/semantics/project_registry_io_manifest_v1.json index 78b65b58d..79e84adb4 100644 --- a/loopx/semantics/project_registry_io_manifest_v1.json +++ b/loopx/semantics/project_registry_io_manifest_v1.json @@ -1279,7 +1279,7 @@ }, { "site": "loopx/control_plane/goals/first_party_host_admission.py::.FirstPartyHostGoalAdmission.for_plan::codec_read:load_project_registry#1", - "line": 120, + "line": 124, "column": 17, "kind": "codec_read", "api": "load_project_registry", @@ -1287,7 +1287,7 @@ }, { "site": "loopx/control_plane/goals/first_party_host_admission.py::._source_authority::codec_read:load_project_registry#1", - "line": 36, + "line": 40, "column": 20, "kind": "codec_read", "api": "load_project_registry", @@ -1295,7 +1295,7 @@ }, { "site": "loopx/control_plane/goals/first_party_host_admission.py::.capture_first_party_host_goal_ref::codec_read:load_project_registry#1", - "line": 71, + "line": 75, "column": 16, "kind": "codec_read", "api": "load_project_registry", diff --git a/tests/architecture/test_goal_instance_binding_inventory.py b/tests/architecture/test_goal_instance_binding_inventory.py index d4d6a8f95..466461766 100644 --- a/tests/architecture/test_goal_instance_binding_inventory.py +++ b/tests/architecture/test_goal_instance_binding_inventory.py @@ -53,7 +53,11 @@ PARTIALLY_ENFORCED_OWNER_IDS = { "first_party_host_runtime", } -QUALIFIED_OWNER_IDS = {"attached_host_chat_session", "handoff_inbox_outbox"} +QUALIFIED_OWNER_IDS = { + "attached_host_chat_session", + "handoff_inbox_outbox", + "turn_journal", +} TYPESCRIPT_DECLARATION = re.compile( r"^(?:export\s+)?(?:async\s+)?(?:function|class)\s+([A-Za-z_$][A-Za-z0-9_$]*)", re.MULTILINE, diff --git a/tests/control_plane/test_turn_journal_goal_instance.py b/tests/control_plane/test_turn_journal_goal_instance.py new file mode 100644 index 000000000..f2cb96e55 --- /dev/null +++ b/tests/control_plane/test_turn_journal_goal_instance.py @@ -0,0 +1,194 @@ +from __future__ import annotations + +import copy +import json +from pathlib import Path +from typing import Any + +import pytest + +from loopx.control_plane.effect_runtime import EffectRuntimeConflict +from loopx.control_plane.goals.first_party_host_admission import ( + FirstPartyHostGoalAdmission, +) +from loopx.control_plane.goals.source_session_registry_state import guard_path +from loopx.control_plane.projects.registry_codec import ( + source_session_registry_transaction, +) +from loopx.control_plane.turn_driver.journal_store import ( + write_turn_journal_checkpoint, +) +from loopx.file_lock import exclusive_cross_runtime_file_lock + + +INSTANCE_A = "ginst_aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa" +INSTANCE_B = "ginst_bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb" +TURN_A = f"sha256:{'a' * 64}" +TURN_B = f"sha256:{'b' * 64}" + + +def _registry_payload(path: Path, instance_id: str) -> dict[str, Any]: + return { + "schema_version": "0.2", + "registry_role": "project-local", + "profile_id": "source_session_v1", + "common_runtime_root": str(path.parent), + "projects": [], + "goals": [ + { + "id": "fixture-goal", + "goal_instance_id": instance_id, + "status": "active", + "execution_authority": False, + } + ], + "session_bindings": [], + "session_receipts": [], + "lifetime_receipts": [], + "retired_goal_instances": [], + } + + +def _write_registry(path: Path, instance_id: str) -> None: + expected = _registry_payload(path, instance_id) + create = None if path.exists() else lambda: expected + with source_session_registry_transaction( + path, + operation="turn_journal_goal_instance_test", + create=create, + ) as transaction: + payload = transaction.payload_copy() + payload["goals"] = expected["goals"] + transaction.commit(payload) + + +def _replace_goal(path: Path, instance_id: str) -> None: + with exclusive_cross_runtime_file_lock( + guard_path(path, "fixture-goal"), + operation="turn_journal_goal_instance_test_recreate", + ): + _write_registry(path, instance_id) + + +def _journal( + instance_id: str | None, + turn_key: str, + *, + completed_phases: list[str] | None = None, +) -> dict[str, Any]: + effect_id = ( + f"fixture-goal:fixture-agent:todo_fixture0001:{turn_key}" + ) + transaction: dict[str, Any] = { + "turn_key": turn_key, + "settlement_plan": { + "schema_version": "quota_settlement_plan_v1", + "identity": { + "schema_version": "quota_settlement_identity_v0", + "effect_id": effect_id, + "goal_id": "fixture-goal", + "agent_id": "fixture-agent", + "todo_id": "todo_fixture0001", + "turn_instance_id": turn_key, + }, + }, + } + plan: dict[str, Any] = { + "turn_envelope": { + "goal_id": "fixture-goal", + "agent_id": "fixture-agent", + "action": {"selected_todo": {"todo_id": "todo_fixture0001"}}, + }, + "transaction": transaction, + } + if instance_id is not None: + goal_ref = { + "goal_id": "fixture-goal", + "goal_instance_id": instance_id, + } + plan["goal_ref"] = goal_ref + transaction["goal_ref"] = goal_ref + return { + "schema_version": "loopx_turn_journal_v0", + "goal_id": "fixture-goal", + "turn_key": turn_key, + "status": "in_progress", + "completed_phases": completed_phases or [], + "plan": plan, + } + + +def _commit_source( + path: Path, + journal: dict[str, Any], + admission: FirstPartyHostGoalAdmission, +) -> None: + with admission.source_journal_admission() as source_admission: + assert source_admission is not None + write_turn_journal_checkpoint( + path, + journal, + source_admission=source_admission, + ) + + +def test_source_journal_commit_is_fenced_by_current_exact_goal_ref( + tmp_path: Path, +) -> None: + registry = tmp_path / "project" / ".loopx" / "registry.json" + runtime = tmp_path / "runtime" / "goals" / "fixture-goal" / "turns" + path_a = runtime / f"{'a' * 64}.json" + path_b = runtime / f"{'b' * 64}.json" + _write_registry(registry, INSTANCE_A) + admission_a = FirstPartyHostGoalAdmission.for_plan( + registry_path=registry, + goal_id="fixture-goal", + planned_goal_ref={ + "goal_id": "fixture-goal", + "goal_instance_id": INSTANCE_A, + }, + ) + journal_a = _journal(INSTANCE_A, TURN_A) + + _commit_source(path_a, journal_a, admission_a) + _commit_source(path_a, journal_a, admission_a) + before_a = path_a.read_bytes() + assert "source_admission" not in json.loads(before_a) + + _replace_goal(registry, INSTANCE_B) + advanced_a = copy.deepcopy(journal_a) + advanced_a["completed_phases"] = ["host_execute", "typed_result"] + with pytest.raises(EffectRuntimeConflict) as exc_info: + _commit_source(path_a, advanced_a, admission_a) + assert exc_info.value.diagnostic_code == "stale_goal_instance" + assert path_a.read_bytes() == before_a + assert not path_b.exists() + + journal_b = _journal(INSTANCE_B, TURN_B) + admission_b = FirstPartyHostGoalAdmission.for_plan( + registry_path=registry, + goal_id="fixture-goal", + planned_goal_ref={ + "goal_id": "fixture-goal", + "goal_instance_id": INSTANCE_B, + }, + ) + _commit_source(path_b, journal_b, admission_b) + + assert json.loads(path_a.read_text(encoding="utf-8")) == journal_a + assert json.loads(path_b.read_text(encoding="utf-8")) == journal_b + assert not Path( + f"{guard_path(registry, 'fixture-goal')}.ts-effect.lock" + ).exists() + + +def test_legacy_journal_commit_keeps_the_existing_wire_shape( + tmp_path: Path, +) -> None: + path = tmp_path / "runtime" / "legacy.json" + journal = _journal(None, TURN_A) + + write_turn_journal_checkpoint(path, journal) + + assert json.loads(path.read_text(encoding="utf-8")) == journal + assert "source_admission" not in path.read_text(encoding="utf-8") diff --git a/tests/control_plane_ts/turn_journal_effects.test.ts b/tests/control_plane_ts/turn_journal_effects.test.ts index 11d6bc6b0..4636e8c46 100644 --- a/tests/control_plane_ts/turn_journal_effects.test.ts +++ b/tests/control_plane_ts/turn_journal_effects.test.ts @@ -1,14 +1,22 @@ import assert from "node:assert/strict"; +import { createHash } from "node:crypto"; import { mkdtemp, readFile, rm } from "node:fs/promises"; import { tmpdir } from "node:os"; -import { join } from "node:path"; +import { dirname, join } from "node:path"; import test from "node:test"; +import { EffectRuntimeConflictError } from "../../loopx/control_plane/effect_runtime_errors.ts"; +import { + acquireFileMutationLock, + releaseFileMutationLock, +} from "../../loopx/control_plane/effect_runtime_io.ts"; import { commitTurnJournal } from "../../loopx/control_plane/turn_driver/turn_journal_effects.ts"; import { interpretTurnJournal } from "../../loopx/control_plane/turn_driver/turn_journal.ts"; const turnKey = `sha256:${"a".repeat(64)}`; const todoId = "todo_fixture0001"; +const instanceA = "ginst_aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa"; +const instanceB = "ginst_bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb"; const phases = [ "host_execute", "typed_result", @@ -19,19 +27,20 @@ const phases = [ "scheduler_ack", ] as const; -function effectId(agentId = "fixture-agent"): string { - return `fixture-goal:${agentId}:${todoId}:${turnKey}`; +function effectId(agentId = "fixture-agent", key = turnKey): string { + return `fixture-goal:${agentId}:${todoId}:${key}`; } function journal( status = "in_progress", completedPhases: readonly string[] = [], agentId = "fixture-agent", + key = turnKey, ): Record { return { schema_version: "loopx_turn_journal_v0", goal_id: "fixture-goal", - turn_key: turnKey, + turn_key: key, status, completed_phases: [...completedPhases], plan: { @@ -41,16 +50,16 @@ function journal( action: { selected_todo: { todo_id: todoId } }, }, transaction: { - turn_key: turnKey, + turn_key: key, settlement_plan: { schema_version: "quota_settlement_plan_v1", identity: { schema_version: "quota_settlement_identity_v0", - effect_id: effectId(agentId), + effect_id: effectId(agentId, key), goal_id: "fixture-goal", agent_id: agentId, todo_id: todoId, - turn_instance_id: turnKey, + turn_instance_id: key, }, }, }, @@ -58,6 +67,78 @@ function journal( }; } +function sourceJournal( + goalInstanceId: string, + key: string, + completedPhases: readonly string[] = [], +): Record { + const snapshot = journal("in_progress", completedPhases, "fixture-agent", key); + const plan = snapshot.plan as Record; + const goalRef = { + goal_id: "fixture-goal", + goal_instance_id: goalInstanceId, + }; + plan.goal_ref = goalRef; + (plan.transaction as Record).goal_ref = goalRef; + return snapshot; +} + +function sourceGuardPath(registryPath: string, goalId: string): string { + const alias = createHash("sha256").update(goalId, "utf8").digest("hex"); + return join( + dirname(registryPath), + ".loopx", + "lifecycle", + "goal-instance", + "guards", + `${alias}.guard`, + ); +} + +async function commitSource( + path: string, + snapshot: Record, + plannedInstanceId: string, + currentInstanceId = plannedInstanceId, +) { + const registryPath = join(dirname(path), "project", ".loopx", "registry.json"); + const target = sourceGuardPath(registryPath, "fixture-goal"); + const lock = await acquireFileMutationLock(target); + try { + return await commitTurnJournal({ + path, + journal: snapshot, + expected_effect_id: effectId( + "fixture-agent", + String(snapshot.turn_key), + ), + source_admission: { + schema_version: "loopx_turn_journal_source_admission_v0", + profile_id: "source_session_v1", + registry_path: registryPath, + planned_goal_ref: { + goal_id: "fixture-goal", + goal_instance_id: plannedInstanceId, + }, + authority: { + kind: "present", + goal_ref: { + goal_id: "fixture-goal", + goal_instance_id: currentInstanceId, + }, + }, + lock: { + target, + pid: process.pid, + token: lock.token, + }, + }, + }); + } finally { + await releaseFileMutationLock(target, lock.token, null, true); + } +} + async function withJournalPath( run: (path: string) => Promise, ): Promise { @@ -91,6 +172,176 @@ test("journal checkpoint retry is idempotent and operation-scoped", async () => }); }); +test("exact GoalRef journals require source admission", async () => { + await withJournalPath(async (path) => { + const snapshot = sourceJournal(instanceA, turnKey); + await assert.rejects( + commitTurnJournal({ + path, + journal: snapshot, + expected_effect_id: effectId(), + }), + /source admission/, + ); + }); +}); + +test("source admission must be well formed and hold a live guard", async () => { + await withJournalPath(async (path) => { + const snapshot = sourceJournal(instanceA, turnKey); + await assert.rejects( + commitTurnJournal({ + path, + journal: snapshot, + expected_effect_id: effectId(), + source_admission: {}, + }), + /source admission is malformed/, + ); + + const registryPath = join( + dirname(path), + "project", + ".loopx", + "registry.json", + ); + const target = sourceGuardPath(registryPath, "fixture-goal"); + const lock = await acquireFileMutationLock(target); + await releaseFileMutationLock(target, lock.token); + await assert.rejects( + commitTurnJournal({ + path, + journal: snapshot, + expected_effect_id: effectId(), + source_admission: { + schema_version: "loopx_turn_journal_source_admission_v0", + profile_id: "source_session_v1", + registry_path: registryPath, + planned_goal_ref: { + goal_id: "fixture-goal", + goal_instance_id: instanceA, + }, + authority: { + kind: "present", + goal_ref: { + goal_id: "fixture-goal", + goal_instance_id: instanceA, + }, + }, + lock: { + target, + pid: process.pid, + token: lock.token, + }, + }, + }), + (error: unknown) => { + assert.ok(error instanceof EffectRuntimeConflictError); + assert.equal(error.code, "journal_source_admission_expired"); + return true; + }, + ); + await assert.rejects(readFile(path), { code: "ENOENT" }); + }); +}); + +test("source admission rejects stale Goal A without mutating A or B", async () => { + const directory = await mkdtemp(join(tmpdir(), "loopx-ts-source-journal-")); + const aKey = `sha256:${"a".repeat(64)}`; + const bKey = `sha256:${"b".repeat(64)}`; + const aPath = join( + directory, + "runtime", + "goals", + "fixture-goal", + "turns", + `${aKey.slice("sha256:".length)}.json`, + ); + const bPath = join( + directory, + "runtime", + "goals", + "fixture-goal", + "turns", + `${bKey.slice("sha256:".length)}.json`, + ); + try { + const initialA = sourceJournal(instanceA, aKey); + const first = await commitSource(aPath, initialA, instanceA); + const replay = await commitSource(aPath, initialA, instanceA); + assert.equal(first.appended, true); + assert.equal(replay.replayed, true); + + const beforeA = await readFile(aPath, "utf8"); + await assert.rejects( + commitSource( + aPath, + sourceJournal(instanceA, aKey, phases.slice(0, 2)), + instanceA, + instanceB, + ), + (error: unknown) => { + assert.ok(error instanceof EffectRuntimeConflictError); + assert.equal(error.code, "stale_goal_instance"); + return true; + }, + ); + assert.equal(await readFile(aPath, "utf8"), beforeA); + + const currentB = sourceJournal(instanceB, bKey); + const committedB = await commitSource(bPath, currentB, instanceB); + assert.equal(committedB.appended, true); + assert.deepEqual(JSON.parse(await readFile(bPath, "utf8")), currentB); + assert.equal(await readFile(aPath, "utf8"), beforeA); + } finally { + await rm(directory, { recursive: true, force: true }); + } +}); + +test("source journal plan GoalRef copies must agree", async () => { + await withJournalPath(async (path) => { + const snapshot = sourceJournal(instanceA, turnKey); + const plan = snapshot.plan as Record; + (plan.transaction as Record).goal_ref = { + goal_id: "fixture-goal", + goal_instance_id: instanceB, + }; + const inspection = interpretTurnJournal({ + schema_version: "loopx_turn_journal_interpretation_request_v0", + journal: snapshot, + goal_id: "fixture-goal", + agent_id: "fixture-agent", + turn_key: turnKey, + }); + assert.ok( + inspection.violations.includes("goal_ref_binding_mismatch"), + ); + assert.equal(inspection.journal_consistent, false); + await assert.rejects( + commitSource(path, snapshot, instanceA), + /GoalRef/, + ); + }); +}); + +test("legacy journals reject source admission without changing their wire", async () => { + await withJournalPath(async (path) => { + const snapshot = journal(); + await assert.rejects( + commitTurnJournal({ + path, + journal: snapshot, + expected_effect_id: effectId(), + source_admission: {}, + }), + /Legacy Turn journals cannot carry source admission/, + ); + const committed = await commit(path, snapshot); + assert.equal(committed.appended, true); + assert.deepEqual(JSON.parse(await readFile(path, "utf8")), snapshot); + }); +}); + test("TS journal owner accepts the complete monotonic transaction", async () => { await withJournalPath(async (path) => { await commit(path, journal()); diff --git a/tests/test_loopx_turn_executor.py b/tests/test_loopx_turn_executor.py index f4671360e..d533eb111 100644 --- a/tests/test_loopx_turn_executor.py +++ b/tests/test_loopx_turn_executor.py @@ -400,7 +400,10 @@ def test_task_validation_stage_reads_result_kind_through_effect_turn( task_validator=None, completed_phases=list(TRANSACTION_PHASES[:2]), journal=journal, - journal_path=journal_path, + persist_journal=lambda snapshot: turn_executor._write_journal( + journal_path, + snapshot, + ), effects={}, ) @@ -645,18 +648,28 @@ def test_cached_host_result_cannot_resume_after_goal_recreation( "completed_phases": [], "plan": plan, } - turn_executor._write_journal(path, journal) - journal.update( - completed_phases=list(TRANSACTION_PHASES[:2]), - host_result=_host_result(plan), - result_kind="validated_progress", - ) - turn_executor._write_journal(path, journal) admission = FirstPartyHostGoalAdmission.for_plan( registry_path=registry, goal_id="fixture-goal", planned_goal_ref=plan["goal_ref"], ) + + def write_source_journal() -> None: + with admission.source_journal_admission() as source_admission: + assert source_admission is not None + turn_executor._write_journal( + path, + journal, + source_admission=source_admission, + ) + + write_source_journal() + journal.update( + completed_phases=list(TRANSACTION_PHASES[:2]), + host_result=_host_result(plan), + result_kind="validated_progress", + ) + write_source_journal() _replace_source_goal(registry, INSTANCE_B) calls = {"writeback": 0, "spend": 0, "scheduler": 0} writeback, spend, scheduler = _callbacks(calls)