Skip to content
Open
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 @@ -840,6 +840,33 @@ promotion retain their own acceptance. No new paid cohort or soak is authorized.
drain, unsupported/warm binary coverage, or any other M3 row.
`execution_authority: false` and the overall activation hold remain.

### 2026-09-30: M3 source Turn effect admission and drain candidate

- **Baseline:** `3ec049e138917a8cce4f84197ba196d26445b2b0`.
- **Delivered:** Source-profile Turn settlement now records a per-effect
admission before durable writeback, quota spend, or terminal closeout. The
existing alias guard covers admission plus the prepared journal checkpoint,
and later covers the committed or aborted checkpoint plus admission release.
The final provider admission remains held through scheduler apply and the
optional post-settlement observer. A durable `source_effect_hold` marker keeps
crashes and `scheduler_action_required` resumable, and the admission releases
only with the final journal checkpoint. Provider calls, readbacks, and tail
callbacks remain outside the alias guard.
- **Retirement:** Recreation first closes the exact Goal A gate. It returns
`drain_required` while an admitted effect still needs provider readback and
publishes Goal B only after the admission set is empty. A committed readback
checkpoints once, an absent readback aborts without provider re-execution,
and an unknown readback keeps Goal A current.
- **Evidence:** Deterministic thread and crash tests cover provider commit
before checkpoint, close versus next-step admission, recreation during
scheduler and post-settlement callbacks, and committed, absent, and unknown
readbacks. Existing source-session recreation and non-source Turn paths
retain their schemas and behavior.
- **Remaining hold:** This qualifies the built-in source Turn settlement slice
of `first_party_host_runtime`. Other Host effects, unsupported or warm
binaries, and every other partial inventory owner remain blocked.
`execution_authority: false` and the overall M3 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 @@ -761,6 +761,30 @@ service adoption、D1–D3 provider promotion 保留各自验收。不授权付
binary 或其他 M3 行已完成。`execution_authority: false` 和总 activation hold
保持不变。

### 2026-09-30:M3 source Turn effect 准入与 drain 候选

- **基线:** `3ec049e138917a8cce4f84197ba196d26445b2b0`。
- **已交付:** Source profile Turn settlement 在 durable writeback、quota spend
或 terminal closeout 前持久化逐 effect admission。既有 alias guard 覆盖
admission 与 prepared journal checkpoint,之后再覆盖 committed/aborted
checkpoint 与 admission release。最后一个 provider admission 会保持到
scheduler apply 和可选 post-settlement observer 完成。持久化
`source_effect_hold` 标记让 crash 和 `scheduler_action_required` 可恢复;该
admission 只在最终 journal checkpoint 时释放。Provider 调用、readback 和
tail callback 都不持 alias guard。
- **Retirement:** Recreation 先关闭精确 Goal A 的 gate。仍有 effect 需要
provider readback 时返回 `drain_required`;只有 admission 集合为空后才发布
Goal B。Committed readback 只 checkpoint 一次;absent readback 直接 abort,
不重新调用 provider;unknown readback 保持 Goal A 为当前实例。
- **证据:** 确定性的线程与 crash 测试覆盖 provider commit 先于 checkpoint、
close 与下一 settlement step admission 的竞争、scheduler 与 post-settlement
callback 期间的 recreation,以及 committed、absent、unknown 三类 readback。
既有 source-session recreation 与非 source Turn 的 schema 和行为保持不变。
- **剩余 hold:** 本切片只资格化 `first_party_host_runtime` 中内置 source Turn
settlement 的部分。其他 Host effect、不支持或常驻 binary,以及其余 partial
inventory owner 仍处于 hold。`execution_authority: false` 和 M3 总 activation
hold 保持不变。

## 附录 B:决策日志

| 日期 | 决策 | Owner/批准 | 替代方案 | 变更的规范章节 |
Expand Down
20 changes: 19 additions & 1 deletion loopx/cli_commands/project.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@
from ..control_plane.coordination.shadow_management import ShadowManagementError

import argparse
from collections.abc import Callable
from collections.abc import Callable, Mapping
from pathlib import Path

from ..control_plane.projects.registry import (
Expand Down Expand Up @@ -102,14 +102,32 @@ def render_project_command_markdown(payload: dict[str, object]) -> str:
f"- registry: `{payload.get('registry')}`",
]
for field in (
"status",
"changed",
"replayed",
"gate_state",
"resolution",
"source",
"project_id",
"foreground_goal_id",
):
if field in payload:
lines.append(f"- {field}: `{payload.get(field)}`")
if payload.get("recovery_action"):
lines.append(f"- recovery_action: {payload.get('recovery_action')}")
pending_effects = payload.get("pending_effects")
if isinstance(pending_effects, list) and pending_effects:
lines.extend(["", "## Pending Turn effects", ""])
for pending in pending_effects:
if not isinstance(pending, Mapping):
continue
lines.append(
"- "
f"turn_key=`{pending.get('turn_key')}`; "
f"step_kind=`{pending.get('step_kind')}`; "
f"reason=`{pending.get('reason')}`; "
f"recovery_action={pending.get('recovery_action')}"
)
if payload.get("error"):
lines.append(f"- error: {payload.get('error')}")
return "\n".join(lines)
Expand Down
26 changes: 25 additions & 1 deletion loopx/control_plane/effect_program.py
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,7 @@
from enum import StrEnum
from functools import lru_cache
from pathlib import Path
from typing import Any, Generic, TypeVar
from typing import Any, Generic, Literal, TypeAlias, TypeVar

from .effect_runtime import EffectRuntimeRejected, effect_runtime_result

Expand Down Expand Up @@ -176,6 +176,30 @@ class SettlementStepKind(StrEnum):
TERMINAL_CLOSEOUT = "terminal_closeout"


TurnProviderStepKind: TypeAlias = Literal[
SettlementStepKind.DURABLE_WRITEBACK,
SettlementStepKind.QUOTA_SPEND,
SettlementStepKind.TERMINAL_CLOSEOUT,
]
TURN_PROVIDER_STEP_KINDS: tuple[TurnProviderStepKind, ...] = (
SettlementStepKind.DURABLE_WRITEBACK,
SettlementStepKind.QUOTA_SPEND,
SettlementStepKind.TERMINAL_CLOSEOUT,
)


def require_turn_provider_step_kind(
step_kind: SettlementStepKind,
) -> TurnProviderStepKind:
if step_kind is SettlementStepKind.DURABLE_WRITEBACK:
return step_kind
if step_kind is SettlementStepKind.QUOTA_SPEND:
return step_kind
if step_kind is SettlementStepKind.TERMINAL_CLOSEOUT:
return step_kind
raise ValueError("validation is not a provider effect step")


class SettlementBindingKind(StrEnum):
TODO = "todo"
AUTONOMOUS_REPLAN = "autonomous_replan"
Expand Down
4 changes: 4 additions & 0 deletions loopx/control_plane/effect_program.ts
Original file line number Diff line number Diff line change
Expand Up @@ -97,6 +97,10 @@ export const SETTLEMENT_STEP_KINDS = [
"terminal_closeout",
] as const;
export type SettlementStepKind = (typeof SETTLEMENT_STEP_KINDS)[number];
export type TurnProviderStepKind = Exclude<SettlementStepKind, "validation">;
export const TURN_PROVIDER_STEP_KINDS = SETTLEMENT_STEP_KINDS.filter(
(kind): kind is TurnProviderStepKind => kind !== "validation",
);

export const SETTLEMENT_BINDING_KINDS = [
"todo",
Expand Down
11 changes: 11 additions & 0 deletions loopx/control_plane/effect_runtime_handlers.ts
Original file line number Diff line number Diff line change
Expand Up @@ -136,6 +136,10 @@ import {
decideGoalRecreation,
decideProjectSessionBind,
decideProjectSessionUnbind,
decideSourceTurnEffectAbsentResolution,
decideSourceTurnEffectAdmission,
decideSourceTurnEffectGate,
decideSourceTurnEffectRelease,
} from "./goals/source_session_lifetime.ts";
import { decideFirstPartyHostRuntime } from "./goals/first_party_host_runtime.ts";
import { decideChatSessionLifecycle } from "./goals/chat_session_lifecycle.ts";
Expand Down Expand Up @@ -555,6 +559,13 @@ export function createEffectRuntimeHandlers(
["goal.source_session.bind.decide", decideProjectSessionBind],
["goal.source_session.unbind.decide", decideProjectSessionUnbind],
["goal.source_session.recreate.decide", decideGoalRecreation],
["goal.source_session.turn_effect.admit", decideSourceTurnEffectAdmission],
[
"goal.source_session.turn_effect.resolve_absent",
decideSourceTurnEffectAbsentResolution,
],
["goal.source_session.turn_effect.release", decideSourceTurnEffectRelease],
["goal.source_session.turn_effect.gate", decideSourceTurnEffectGate],
["goal.first_party_host_runtime.decide", decideFirstPartyHostRuntime],
["goal.chat_session.lifecycle.decide", decideChatSessionLifecycle],
["goal.acceptance.inspect", inspectLocalGoalAcceptance],
Expand Down
130 changes: 119 additions & 11 deletions loopx/control_plane/goals/first_party_host_admission.py
Original file line number Diff line number Diff line change
Expand Up @@ -15,11 +15,20 @@
SOURCE_SESSION_PROFILE_ID,
load_project_registry,
)
from ..effect_program import TurnProviderStepKind
from .source_session_registry_state import (
exact_goal_ref,
guard_path,
require_goal_id,
)
from .source_session_turn_effects import (
JournalPersist,
SourceTurnEffect,
SourceTurnEffectRejected,
prepare_source_turn_effect,
release_source_turn_effect,
source_turn_effect_allows_absent_reexecute,
)


T = TypeVar("T")
Expand All @@ -33,6 +42,89 @@ def __init__(self, code: str) -> None:
self.code = code


@dataclass(frozen=True, slots=True)
class FirstPartyHostTurnEffectAdmission:
goal_admission: FirstPartyHostGoalAdmission
turn_key: str
journal_path: Path

def _effect(
self,
step_kind: TurnProviderStepKind,
effect_ref: str,
) -> SourceTurnEffect:
goal_ref = self.goal_admission.planned_goal_ref
if not isinstance(goal_ref, Mapping):
raise FirstPartyHostRuntimeRejected("goal_instance_id_missing")
goal_id = goal_ref.get("goal_id")
goal_instance_id = goal_ref.get("goal_instance_id")
if not isinstance(goal_id, str) or not isinstance(goal_instance_id, str):
raise FirstPartyHostRuntimeRejected("goal_instance_id_missing")
return SourceTurnEffect(
goal_ref=exact_goal_ref(goal_id, goal_instance_id),
turn_key=self.turn_key,
step_kind=step_kind,
effect_ref=effect_ref,
journal_path=self.journal_path,
)

def prepare(
self,
step_kind: TurnProviderStepKind,
effect_ref: str,
persist_journal: JournalPersist,
) -> None:
try:
prepare_source_turn_effect(
registry_path=self.goal_admission.registry_path,
goal_id=self.goal_admission.goal_id,
effect=self._effect(step_kind, effect_ref),
source_admission=self.goal_admission.source_journal_admission_locked,
persist_journal=persist_journal,
)
except SourceTurnEffectRejected as exc:
raise FirstPartyHostRuntimeRejected(exc.code) from exc

def hold(
self,
step_kind: TurnProviderStepKind,
effect_ref: str,
persist_journal: JournalPersist,
) -> None:
self.prepare(step_kind, effect_ref, persist_journal)

def release(
self,
step_kind: TurnProviderStepKind,
effect_ref: str,
persist_journal: JournalPersist,
) -> None:
try:
release_source_turn_effect(
registry_path=self.goal_admission.registry_path,
goal_id=self.goal_admission.goal_id,
effect=self._effect(step_kind, effect_ref),
source_admission=self.goal_admission.source_journal_admission_locked,
persist_journal=persist_journal,
)
except SourceTurnEffectRejected as exc:
raise FirstPartyHostRuntimeRejected(exc.code) from exc

def allows_absent_reexecute(
self,
step_kind: TurnProviderStepKind,
effect_ref: str,
) -> bool:
try:
return source_turn_effect_allows_absent_reexecute(
registry_path=self.goal_admission.registry_path,
goal_id=self.goal_admission.goal_id,
effect=self._effect(step_kind, effect_ref),
)
except SourceTurnEffectRejected as exc:
raise FirstPartyHostRuntimeRejected(exc.code) from exc


def _source_authority(registry_path: Path, goal_id: str) -> dict[str, Any]:
if not registry_path.is_file():
return {"kind": "unavailable", "reason": "registry_missing"}
Expand Down Expand Up @@ -215,17 +307,33 @@ def source_journal_admission(
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),
}
yield self.source_journal_admission_locked()

def source_journal_admission_locked(self) -> dict[str, Any]:
"""Build a TS handoff while the caller holds this Goal's source guard."""

target = guard_path(self.registry_path, self.goal_id)
return {
"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 turn_effect_admission(
self,
*,
turn_key: str,
journal_path: Path,
) -> FirstPartyHostTurnEffectAdmission | None:
if not self.source_profile:
return None
return FirstPartyHostTurnEffectAdmission(self, turn_key, journal_path)

def accept_result(self, commit_result: Callable[[], T]) -> T:
if not self.source_profile:
Expand Down
Loading
Loading