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
6 changes: 6 additions & 0 deletions docs/architecture/rfcs/loopx-overall-roadmap-v0.md
Original file line number Diff line number Diff line change
Expand Up @@ -512,6 +512,12 @@ preflight through the actual Turn dry-run and selected executor/profile. Task
admission, current pinned acceptance and runtime availability remain distinct;
unprobed generic/cloud availability stays unknown. The same inspection is available
to CLI and enabled MCP/new Chat tools, without a new state store or launch effect.
Turn exception readback now distinguishes uncertain invocation effects from the
original journal's checkpointed effects and prepared-effect recovery. This is a
bounded R2 recovery slice, not runtime provisioning or G1 completion; caller
admission, fresh task derivation and frontend/Lark qualification remain open.
中文:异常读回区分本次调用未知副作用与原 Turn 的持久观察/待核对效果,只完成
R2 的一个恢复切片;调用方准入、新任务生成与前端/Lark 验收仍未闭环。
The same owner-local panel now opens current validated artifact text and its
version/source identifiers, accepts feedback through the original coordinator
inbox and exposes coordinator pause with its actual scope. Stale reads clear prior
Expand Down
36 changes: 36 additions & 0 deletions docs/reference/protocols/loopx-turn-v0.md
Original file line number Diff line number Diff line change
Expand Up @@ -383,6 +383,42 @@ quota construction, scheduler context, planning, host invocation, settlement,
spend, or state writeback. Its `effects` field is therefore always an empty
list.

`recorded_effects` separately reports lower-bound observations of the **original
Turn**, not effects performed by this inspection or a recovery invocation.
Boolean values mean checkpointed execution or no recorded attempt; `null` means
unknown. A saved Host attempt precedes launch confirmation; a `prepared` writeback
or spend is not a commit receipt. Inconsistent/foreign identity or phase lineage
supplies only unknown effect facts. A scheduler phase does not imply host acknowledgement.

The inspector and journal writer share the original TS prepared-intent contract:
one supported step, object shape, `prepared` status, exact settlement effect ref,
the next phase and a nonterminal status. Unknown or malformed intents block the
existing recovery decision, including Host reinvocation. Valid lineage's already
completed effects remain proved; other effects are unknown, not `false`. The
writer's empty-map rejection is retained: no pending intent uses an absent field;
committed history is carried by completed checkpoints, not retained intents.

If executing `run-once` raises unexpectedly, its error response retains the
original error and resume key, marks uncertain current-invocation `effects` as
`null`, and adds `journal_observation` from this same read-only TS owner. Its
`scope=original_turn` prevents a saved Host result being mistaken for another
launch. Unavailable inspection remains explicit; there is no Python fallback.
Use the original `recovery_decision` and provider readback to recover, never a
fresh task or repeated model call inferred from a failed CLI reply. This readback
does not bypass controller completion validation, lease conflicts or quota gates.
Pre-execution failures retain known Turn-start hook writes without claiming Host
execution. Normal successful/replayed `effects` remain invocation-scoped.

中文:`recorded_effects` 是原 Turn 的持久观察,不是本次检查或恢复又发生了副作用。
`null` 表示未知:已登记 Host attempt 不等于模型已启动,`prepared` 不等于写回或扣额
已提交。异常返回保留原错误/恢复身份,以同一 TS owner 的只读 `journal_observation`
区分本次调用与原 Turn;读回失败不猜测“没有执行”。按原恢复判定和 provider 回执
继续,不能因 CLI 报错新建任务重跑模型,也不放松验收、租约或扣额门禁。
检查与写入共用原 TS prepared-intent 合同,校验步骤、形状、状态、effect 身份和
阶段绑定。未知或畸形 intent 阻断原恢复判定,不能建议重调 Host;保留合法身份与
阶段已经证明的执行事实,其余返回未知而非 `false`。不放松原 writer 的空 map
拒绝语义:无待决 intent 应省略该字段,已提交历史由完成 checkpoint 表达。

Exit zero means that inspection completed, including when `decision` is
`replay_blocked`. A non-zero exit means the command could not inspect the
requested journal because a selector, file, JSON document, or schema was
Expand Down
18 changes: 17 additions & 1 deletion loopx/cli_commands/turn.py
Original file line number Diff line number Diff line change
Expand Up @@ -50,6 +50,7 @@
load_loopx_turn_plan_from_journal,
run_codex_cli_host,
run_loopx_turn_once,
inspect_loopx_turn_journal,
selected_turn_todo,
)
from ..control_plane.turn_driver.host_binding import managed_executor_binding
Expand Down Expand Up @@ -114,6 +115,7 @@ def handle_turn_command(
output_format=output_format, print_payload=print_payload,
)
payload: dict[str, Any] = {}
execution_started = False
try:
if getattr(args, "todo_id", None) is not None and (
getattr(args, "resume_turn_key", None)
Expand Down Expand Up @@ -1056,6 +1058,7 @@ def on_managed_start_admitted() -> None:
on_admitted=on_managed_start_admitted,
)

execution_started = bool(args.execute)
payload = run_loopx_turn_once(
payload,
host_argv=raw_argv,
Expand Down Expand Up @@ -1090,7 +1093,20 @@ def on_managed_start_admitted() -> None:
else:
raise ValueError("turn requires the `plan` or `run-once` subcommand")
except Exception as exc: # noqa: BLE001 - CLI boundary renders typed JSON failure
payload = build_turn_error_payload(payload, exc, turn_command=args.turn_command)
journal_readback = None
if execution_started:
transaction = payload.get("transaction") or {}
try:
journal_readback = inspect_loopx_turn_journal(
runtime_root, goal_id=args.goal_id, agent_id=args.agent_id,
turn_key=str(transaction.get("turn_key") or ""),
)
except Exception: # noqa: BLE001 - retain original error and unknown effects
pass
payload = build_turn_error_payload(
payload, exc, turn_command=args.turn_command,
execution_started=execution_started, journal_readback=journal_readback,
)
renderer = (
_render_loopx_turn_execution_markdown
if args.turn_command == "run-once"
Expand Down
38 changes: 32 additions & 6 deletions loopx/cli_commands/turn_rendering.py
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,8 @@

def build_turn_error_payload(
planned: dict[str, Any], exc: Exception, *, turn_command: str,
execution_started: bool = False,
journal_readback: Mapping[str, Any] | None = None,
) -> dict[str, Any]:
"""Preserve the failed Turn's identity and effect readback at the CLI edge."""

Expand All @@ -24,6 +26,15 @@ def build_turn_error_payload(
run_once = turn_command == "run-once"
error_code = getattr(exc, "code", None)
error_payload = getattr(exc, "payload", None)
planned_effects = planned.get("effects")
planned_effects = planned_effects if isinstance(planned_effects, Mapping) else {}
# Once the executor was entered an exception can follow a protected effect.
# A missing reply is not proof of non-execution. The original journal is
# projected separately so a replay never looks like a second host launch.
effects = {
key: True if planned_effects.get(key) is True else None if execution_started else False
for key in ("host_invoked", "state_written", "scheduler_acknowledged", "quota_spent")
}
return {
**({"error_code": error_code, **(error_payload if isinstance(error_payload, Mapping) else {})}
if isinstance(error_code, str) else {}),
Expand All @@ -33,12 +44,16 @@ def build_turn_error_payload(
),
"mode": "run_once" if run_once else "plan",
"error": str(exc),
"effects": {
"host_invoked": False,
"state_written": False,
"scheduler_acknowledged": False,
"quota_spent": False,
},
"effects": effects,
"effects_scope": "current_invocation",
**({"journal_observation": {
"scope": "original_turn",
"status": "observed" if journal_readback is not None else "unavailable",
**({key: journal_readback[key] for key in (
"journal_consistent", "journal_status", "completed_phases",
"recorded_effects", "recovery_decision",
)} if journal_readback is not None else {}),
}} if execution_started else {}),
**({
"resume_turn_key": turn_key,
"journal_ref": f"turn:{turn_key.removeprefix('sha256:')[:16]}",
Expand Down Expand Up @@ -114,6 +129,10 @@ def render_loopx_turn_execution_markdown(payload: dict[str, object]) -> str:
if isinstance(payload.get("host_failure"), dict)
else {}
)
journal_observation = (
payload.get("journal_observation")
if isinstance(payload.get("journal_observation"), dict) else {}
)
return "\n".join(
[
"# LoopX Turn Run Once",
Expand Down Expand Up @@ -145,6 +164,12 @@ def render_loopx_turn_execution_markdown(payload: dict[str, object]) -> str:
f"- host_invoked: {effects.get('host_invoked')}",
f"- state_written: {effects.get('state_written')}",
f"- quota_spent: {effects.get('quota_spent')}",
*([f"- error: {payload['error']}"] if payload.get("error") else []),
*([f"- journal_readback: {journal_observation.get('status')}",
f"- recorded_effects: {journal_observation.get('recorded_effects')}",
f"- recovery_from: {(journal_observation.get('recovery_decision') or {}).get('resume_from')}",
f"- recovery_reinvoke_host: {(journal_observation.get('recovery_decision') or {}).get('reinvoke_host')}"]
if journal_observation else []),
*(
[
f"- recovery_plan: {planned.get('action')}",
Expand Down Expand Up @@ -209,6 +234,7 @@ def render_loopx_turn_journal_inspection_markdown(
f"- recovery_reason: {recovery.get('reason')}",
f"- recovery_checks: {rendered_checks or 'none'}",
f"- journal_consistent: {payload.get('journal_consistent')}",
f"- original_turn_recorded_effects: {payload.get('recorded_effects')}",
*(
[
f"- last_recovery_plan: {last_planned.get('action')}",
Expand Down
19 changes: 17 additions & 2 deletions loopx/control_plane/turn_driver/turn_journal.ts
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,8 @@ import {
type EffectTurn,
} from "../effect_program.ts";
import { EffectRuntimeRequestError } from "../effect_runtime_errors.ts";
import { preparedAttemptViolation } from "./turn_journal_attempt_contract.ts";
import { recordedTurnEffects, type RecordedTurnEffects } from "./turn_journal_effect_readback.ts";

export const TURN_JOURNAL_INSPECTION_SCHEMA_VERSION =
"loopx_turn_journal_inspection_v1";
Expand Down Expand Up @@ -76,6 +78,7 @@ export interface TurnJournalInspection {
journal_consistent: boolean;
recovery_decision: TurnRecoveryDecision;
last_recovery: TurnRecoveryAudit | null;
recorded_effects: RecordedTurnEffects;
effects: [];
}

Expand All @@ -92,6 +95,7 @@ export interface TurnJournalEffectContext {
journal_consistent: boolean;
recovery_decision: TurnRecoveryDecision;
last_recovery: TurnRecoveryAudit | null;
recorded_effects: RecordedTurnEffects;
}

// Replay has its own verdict. It is not a quota decision and must not manufacture
Expand Down Expand Up @@ -619,8 +623,7 @@ export function interpretTurnJournalEffect(
violations.push("journal_status_unsupported");
}

const replayLegal = violations.length === 0;
const journalConsistent =
const lineageConsistent =
goalMatches &&
ownerMatches &&
settlementIdentityValid &&
Expand All @@ -629,6 +632,14 @@ export function interpretTurnJournalEffect(
turnKeyMatches &&
phasesFormOrderedPrefix &&
supportedJournalStatuses.has(journalStatus);
const attemptViolation = preparedAttemptViolation(journal, {
status: journalStatus,
completedPhases,
effectId: settlementIdentityFromPlan(transaction).value?.effect_id ?? "",
});
if (attemptViolation) violations.push(attemptViolation.code);
const replayLegal = violations.length === 0;
const journalConsistent = lineageConsistent && attemptViolation === null;
const decision = replayLegal ? "replay_legal" : "replay_blocked";
const turnRecoveryDecision = recoveryDecision(
request,
Expand Down Expand Up @@ -657,6 +668,9 @@ export function interpretTurnJournalEffect(
journal_consistent: journalConsistent,
recovery_decision: turnRecoveryDecision,
last_recovery: projectRecoveryAudit(journal.recovery_audit),
recorded_effects: recordedTurnEffects(
journal, completedPhases, lineageConsistent, attemptViolation === null,
),
},
},
interpretation: {
Expand Down Expand Up @@ -709,6 +723,7 @@ export function projectTurnJournalInspection(
journal_consistent: context.journal_consistent,
recovery_decision: context.recovery_decision,
last_recovery: context.last_recovery,
recorded_effects: context.recorded_effects,
effects: [],
};
}
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,57 @@
/** The journal owner's prepared-intent contract, shared by writes and reads. */
import transactionContract from "../turn_transaction_contract.json" with { type: "json" };
import { SETTLEMENT_STEP_KINDS, type JsonObject } from "../effect_program.ts";

const preparedStepKinds: ReadonlySet<string> = new Set(
SETTLEMENT_STEP_KINDS.filter((kind) => kind !== "validation"),
);

interface AttemptViolation {
code: string;
message: string;
}

function isObject(value: unknown): value is JsonObject {
return typeof value === "object" && value !== null && !Array.isArray(value);
}

export function preparedAttemptViolation(
journal: JsonObject,
state: { status: string; completedPhases: readonly string[]; effectId: string },
): AttemptViolation | null {
if (journal.effect_attempts === undefined) return null;
if (!isObject(journal.effect_attempts)) {
return { code: "prepared_effect_attempts_invalid",
message: "Turn journal prepared effects must be an object" };
}
const entries = Object.entries(journal.effect_attempts);
// Preserve the writer contract: no pending intent is represented by an
// absent field, not by an empty map or a set of historical committed intents.
if (entries.length !== 1) {
return { code: "prepared_effect_count_invalid",
message: "Turn journal must carry at most one prepared effect" };
}
const [stepKind, attempt] = entries[0];
if (!preparedStepKinds.has(stepKind)) {
return { code: "prepared_effect_step_unsupported",
message: "Turn journal carries an unsupported prepared effect step" };
}
if (!isObject(attempt) || attempt.status !== "prepared"
|| attempt.effect_ref !== `${state.effectId}#${stepKind}`) {
return { code: "prepared_effect_identity_invalid",
message: "Turn journal prepared effect does not match settlement identity" };
}
const phases: readonly string[] = transactionContract.phases;
const phaseIndex = phases.indexOf(stepKind === "terminal_closeout"
? "scheduler_apply" : stepKind);
if (phaseIndex < 0 || state.completedPhases.length !== phaseIndex
|| !state.completedPhases.every((phase, index) => phase === phases[index])) {
return { code: "prepared_effect_phase_invalid",
message: "Turn journal prepared effect is not the next settlement step" };
}
if (!["in_progress", "failed"].includes(state.status)) {
return { code: "prepared_effect_status_invalid",
message: "Turn journal terminal state cannot retain a prepared effect" };
}
return null;
}
55 changes: 55 additions & 0 deletions loopx/control_plane/turn_driver/turn_journal_effect_readback.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,55 @@
/** Lower-bound observations of the original Turn, never new execution effects. */
export interface RecordedTurnEffects {
host_invoked: boolean | null;
state_written: boolean | null;
quota_spent: boolean | null;
scheduler_acknowledged: boolean | null;
}

type JsonObject = Record<string, unknown>;

function object(value: unknown): JsonObject {
return typeof value === "object" && value !== null && !Array.isArray(value)
? value as JsonObject : {};
}

export function recordedTurnEffects(
journal: JsonObject,
completedPhases: readonly string[],
lineageConsistent: boolean,
attemptsValid: boolean,
): RecordedTurnEffects {
const unknown: RecordedTurnEffects = {
host_invoked: null, state_written: null,
quota_spent: null, scheduler_acknowledged: null,
};
// Foreign, corrupt or contradictory lineage cannot supply effect facts.
if (!lineageConsistent) return unknown;
const completed = new Set(completedPhases);
const scheduler = object(journal.scheduler);
// Unknown or contradictory intents cannot prove non-execution. Keep only
// facts proved by the valid lineage's already completed checkpoints.
if (!attemptsValid) return {
host_invoked: completed.has("host_execute") ? true : null,
state_written: completed.has("durable_writeback") ? true : null,
quota_spent: completed.has("quota_spend") ? true : null,
scheduler_acknowledged: scheduler.acknowledged === true ? true : null,
};
const attempts = object(journal.effect_attempts);
const pending = (step: string) => Object.hasOwn(attempts, step);
// An attempt is persisted BEFORE confirmation/host launch. It proves neither
// launch nor non-launch until the host checkpoint is durable.
const hostCount = journal.host_attempt_count;
const hostUncertain = hostCount !== undefined
&& (!Number.isInteger(hostCount) || Number(hostCount) !== 0);
return {
host_invoked: completed.has("host_execute") ? true : hostUncertain ? null : false,
state_written: completed.has("durable_writeback") ? true
: pending("durable_writeback") || pending("terminal_closeout") ? null : false,
quota_spent: completed.has("quota_spend") ? true
: pending("quota_spend") ? null : false,
// A completed outer-controller phase need not acknowledge any host cadence.
scheduler_acknowledged: typeof scheduler.acknowledged === "boolean"
? scheduler.acknowledged : completed.has("quota_spend") ? null : false,
};
}
Loading
Loading