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 @@ -445,8 +445,11 @@ shipped Stage 2B cutovers are in place:
reduction. A real caller-approved validation command remains an explicit
Python provider between two reductions. Todo and policy-source snapshots are
compared after the mutation lock so a receipt for one declaration or agent
registry cannot authorize changed facts. Materialized and event-projected
writes consume the same typed result.
registry cannot authorize changed facts. Policy admission failures are
returned as typed data by that same reduction and consumed only after Python
actor/lease admission, preserving legacy error priority without a leaf
runtime call inside the writer critical section. Materialized and
event-projected writes consume the same typed result.
- Scheduler heartbeat/state: TypeScript owns receipt freshness, ACK and
host-failure validation, identity-aware progression, failure-cache
retention/counting, replay and CAS fencing, preview reduction, the locked
Expand Down
2 changes: 0 additions & 2 deletions loopx/control_plane/effect_runtime_handlers.ts
Original file line number Diff line number Diff line change
Expand Up @@ -64,7 +64,6 @@ import {
selectTodoCompletionContinuation,
} from "./todos/completion_state.ts";
import { reduceTodoCompletionTransaction } from "./todos/completion_transaction.ts";
import { resolveTodoCompletionPolicy } from "./todos/completion_policy.ts";
import { transitionTodoNextAction } from "./todos/next_action.ts";
import {
evaluateTodoResumeConditions,
Expand Down Expand Up @@ -375,7 +374,6 @@ export function createEffectRuntimeHandlers(
),
],
["todo.completion.reduce", reduceTodoCompletionTransaction],
["todo.completion_policy.resolve", resolveTodoCompletionPolicy],
["todo.next_action.transition", transitionTodoNextAction],
["todo.resume_condition.normalize", normalizeTodoResumeWhen],
["todo.resume_condition.evaluate", evaluateTodoResumeConditions],
Expand Down
48 changes: 18 additions & 30 deletions loopx/control_plane/todos/completion_policy.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,6 @@
load_goal_from_registry,
registered_agent_ids_for_goal,
)
from ..effect_runtime import EffectRuntimeRejected, effect_runtime_result
from .active_state_editing import find_todo_block
from .contract import (
normalize_todo_claimed_by,
Expand All @@ -19,6 +18,7 @@

TODO_COMPLETION_POLICY_REQUEST_SCHEMA = "loopx_todo_completion_policy_request_v0"
TODO_COMPLETION_POLICY_RESULT_SCHEMA = "loopx_todo_completion_policy_result_v0"
TODO_COMPLETION_POLICY_FAILURE_SCHEMA = "loopx_todo_completion_policy_failure_v0"


@dataclass(frozen=True)
Expand Down Expand Up @@ -148,6 +148,23 @@ def completion_policy_from_transaction(
# These fields are dead on the event-projected replay path; the TS
# completion fence has already prohibited every write.
return CompletionPolicy(None, [], None, [], False)
if transaction.get("decision") == "policy_reject":
failure = transaction.get("completion_policy_failure")
if not (
isinstance(failure, Mapping)
and failure.get("schema_version")
== TODO_COMPLETION_POLICY_FAILURE_SCHEMA
and failure.get("kind") == "completion_policy_rejected"
and isinstance(failure.get("diagnostic_code"), str)
and bool(failure.get("diagnostic_code"))
and isinstance(failure.get("summary"), str)
and bool(failure.get("summary"))
and transaction.get("completion_policy") is None
):
raise RuntimeError(
"TypeScript Todo completion policy failure shape mismatch"
)
raise ValueError(str(failure["summary"]))
policy = transaction.get("completion_policy")
if not isinstance(policy, Mapping) or (
policy.get("schema_version") != TODO_COMPLETION_POLICY_RESULT_SCHEMA
Expand Down Expand Up @@ -179,32 +196,3 @@ def completion_policy_from_transaction(
self_merged=bool(policy["self_merged"]),
linked_successor_id=linked_successor_id,
)


def bind_completion_policy_to_transaction(
transaction: Mapping[str, Any],
completion_policy_request: Mapping[str, Any],
) -> dict[str, Any]:
"""Attach the TS-owned policy after actor and lease admission.

External validation and completion-state reduction stay single-shot. Only
the pure policy reducer runs under the locked authority boundary, which
preserves the legacy actor -> lease -> policy error priority.
"""

bound = dict(transaction)
if bound.get("decision") != "commit":
return bound
try:
result = effect_runtime_result(
"todo.completion_policy.resolve",
dict(completion_policy_request),
)
except EffectRuntimeRejected as exc:
raise ValueError(str(exc)) from None
if not isinstance(result, Mapping):
raise RuntimeError("TypeScript Todo completion policy result must be an object")
bound["completion_policy"] = dict(result)
# Reuse the public adapter as the exact result-shape guard.
completion_policy_from_transaction(bound)
return bound
62 changes: 52 additions & 10 deletions loopx/control_plane/todos/completion_transaction.py
Original file line number Diff line number Diff line change
Expand Up @@ -14,7 +14,7 @@
"loopx_todo_completion_transaction_result_v0"
)

_DECISIONS = {"execute_validation", "commit", "replay", "reject"}
_DECISIONS = {"execute_validation", "commit", "policy_reject", "replay", "reject"}
_IDENTITY_SOURCES = {
"turn_settlement",
"unscoped_completion",
Expand All @@ -37,6 +37,7 @@
"validation_timeout_seconds",
)
_COMPLETION_POLICY_RESULT_SCHEMA = "loopx_todo_completion_policy_result_v0"
_COMPLETION_POLICY_FAILURE_SCHEMA = "loopx_todo_completion_policy_failure_v0"


def _json_sequence(value: Any) -> Any:
Expand Down Expand Up @@ -335,6 +336,18 @@ def _valid_completion_policy(value: Any) -> bool:
)


def _valid_completion_policy_failure(value: Any) -> bool:
return (
isinstance(value, Mapping)
and value.get("schema_version") == _COMPLETION_POLICY_FAILURE_SCHEMA
and value.get("kind") == "completion_policy_rejected"
and isinstance(value.get("diagnostic_code"), str)
and bool(value.get("diagnostic_code"))
and isinstance(value.get("summary"), str)
and bool(value.get("summary"))
)


def _valid_execute_validation_result(result: Mapping[str, Any]) -> bool:
effect = result.get("validation_effect")
return (
Expand Down Expand Up @@ -374,15 +387,10 @@ def _valid_execute_validation_result(result: Mapping[str, Any]) -> bool:
)


def _valid_commit_result(
result: Mapping[str, Any],
*,
completion_policy_required: bool,
) -> bool:
def _valid_completion_settlement(result: Mapping[str, Any]) -> bool:
state = result.get("completion_state")
updates = result.get("metadata_updates")
receipt = result.get("validation_receipt")
policy = result.get("completion_policy")
return (
isinstance(state, Mapping)
and state.get("continuation") in _CONTINUATIONS
Expand All @@ -396,10 +404,38 @@ def _valid_commit_result(
and updates.get("completion_continuation") == state.get("continuation")
and updates.get("completion_recovery") == state.get("recovery")
and (receipt is None or _valid_receipt(receipt))
)


def _valid_commit_result(
result: Mapping[str, Any],
*,
completion_policy_required: bool,
) -> bool:
policy = result.get("completion_policy")
return (
_valid_completion_settlement(result)
and result.get("completion_policy_failure") is None
and (
_valid_completion_policy(policy)
if completion_policy_required or policy is not None
else True
if completion_policy_required
else policy is None
)
)


def _valid_policy_reject_result(
result: Mapping[str, Any],
*,
completion_policy_required: bool,
) -> bool:
if not completion_policy_required:
return False
return (
_valid_completion_settlement(result)
and result.get("completion_policy") is None
and _valid_completion_policy_failure(
result.get("completion_policy_failure")
)
)

Expand Down Expand Up @@ -431,6 +467,11 @@ def _valid_result(
result,
completion_policy_required=completion_policy_required,
)
if decision == "policy_reject":
return _valid_policy_reject_result(
result,
completion_policy_required=completion_policy_required,
)
if decision == "replay":
return result["fence"].get("outcome") == "replay"
return _valid_reject_result(result)
Expand Down Expand Up @@ -482,7 +523,8 @@ def reduce_todo_completion_transaction(
if not _valid_result(
result,
completion_policy_required=(
completion_policy_request is not None and result.get("decision") == "commit"
completion_policy_request is not None
and result.get("decision") in {"commit", "policy_reject"}
),
):
raise RuntimeError(
Expand Down
78 changes: 68 additions & 10 deletions loopx/control_plane/todos/completion_transaction.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,8 @@
import type { JsonObject } from "../effect_program.ts";
import { EffectRuntimeRequestError } from "../effect_runtime_errors.ts";
import {
effectRuntimeErrorPayload,
EffectRuntimeRequestError,
} from "../effect_runtime_errors.ts";
import {
optionalNonEmptyString,
requireBoolean,
Expand Down Expand Up @@ -38,6 +41,8 @@ export const TODO_COMPLETION_TRANSACTION_REQUEST_SCHEMA =
"loopx_todo_completion_transaction_v0";
export const TODO_COMPLETION_TRANSACTION_RESULT_SCHEMA =
"loopx_todo_completion_transaction_result_v0";
export const TODO_COMPLETION_POLICY_FAILURE_SCHEMA =
"loopx_todo_completion_policy_failure_v0";
const CALLER_VALIDATION_RECEIPT_SCHEMA = "issue_fix_validation_command_v0";

const PROJECTION_SOURCES = ["materialized", "event_log"] as const;
Expand Down Expand Up @@ -85,14 +90,29 @@ export interface TodoCompletionExecuteValidation
validation_effect: TodoCompletionValidationEffect;
}

export interface TodoCompletionCommit extends CompletionTransactionBase {
decision: "commit";
interface TodoCompletionSettlement extends CompletionTransactionBase {
completion_state: CompletionStateProjection;
metadata_updates: JsonObject;
validation_receipt: CallerValidationReceipt | null;
}

export interface TodoCompletionCommit extends TodoCompletionSettlement {
decision: "commit";
completion_policy?: TodoCompletionPolicyResult;
}

export interface TodoCompletionPolicyFailure extends JsonObject {
schema_version: typeof TODO_COMPLETION_POLICY_FAILURE_SCHEMA;
kind: "completion_policy_rejected";
diagnostic_code: string;
summary: string;
}

export interface TodoCompletionPolicyReject extends TodoCompletionSettlement {
decision: "policy_reject";
completion_policy_failure: TodoCompletionPolicyFailure;
}

export interface TodoCompletionReplay extends CompletionTransactionBase {
decision: "replay";
}
Expand All @@ -109,6 +129,7 @@ export interface TodoCompletionReject extends CompletionTransactionBase {
export type TodoCompletionTransactionResult =
| TodoCompletionExecuteValidation
| TodoCompletionCommit
| TodoCompletionPolicyReject
| TodoCompletionReplay
| TodoCompletionReject;

Expand Down Expand Up @@ -294,6 +315,33 @@ function baseResult(
};
}

function evaluateCompletionPolicy(
request: JsonObject | null,
):
| { outcome: "not_requested" }
| { outcome: "accepted"; policy: TodoCompletionPolicyResult }
| { outcome: "rejected"; failure: TodoCompletionPolicyFailure } {
if (request === null) return { outcome: "not_requested" };
try {
return {
outcome: "accepted",
policy: resolveTodoCompletionPolicy(request),
};
} catch (error) {
if (!(error instanceof EffectRuntimeRequestError)) throw error;
const failure = effectRuntimeErrorPayload(error);
return {
outcome: "rejected",
failure: {
schema_version: TODO_COMPLETION_POLICY_FAILURE_SCHEMA,
kind: "completion_policy_rejected",
diagnostic_code: failure.code,
summary: failure.message,
},
};
}
}

/**
* Reduce one Todo completion admission/settlement transaction.
*
Expand Down Expand Up @@ -414,20 +462,30 @@ export function reduceTodoCompletionTransaction(
metadataResult.updates,
"completion metadata updates",
);
const completionPolicy = request.completion_policy_request === null
? null
: resolveTodoCompletionPolicy(request.completion_policy_request);
return {
const completionPolicy = evaluateCompletionPolicy(
request.completion_policy_request,
);
const settlement = {
...base,
decision: "commit",
completion_state: {
continuation: completionStateResult.continuation,
recovery: completionStateResult.recovery,
},
metadata_updates: updates,
validation_receipt: request.validation_receipt,
...(completionPolicy === null
};
if (completionPolicy.outcome === "rejected") {
return {
...settlement,
decision: "policy_reject",
completion_policy_failure: completionPolicy.failure,
};
}
return {
...settlement,
decision: "commit",
...(completionPolicy.outcome === "not_requested"
? {}
: { completion_policy: completionPolicy }),
: { completion_policy: completionPolicy.policy }),
};
}
10 changes: 5 additions & 5 deletions loopx/control_plane/todos/completion_validation.py
Original file line number Diff line number Diff line change
Expand Up @@ -256,10 +256,10 @@ def run_completion_validation_gate_with_source(
todo_id=todo_id,
requested_has_successor=requested_has_successor,
validation_receipt=None,
# Preserve the legacy error priority: policy admission is evaluated
# only after actor authority and the task-lease fence are established
# under the write lock. The source is still captured here for CAS.
completion_policy_request=None,
# The coarse reducer returns policy success or typed failure as data.
# The public writer consumes that projection only after actor and lease
# admission, preserving legacy error priority without a second IPC.
completion_policy_request=completion_policy_source,
)
completion_validation = None
if transaction["decision"] == "execute_validation":
Expand Down Expand Up @@ -302,7 +302,7 @@ def run_completion_validation_gate_with_source(
todo_id=todo_id,
requested_has_successor=requested_has_successor,
validation_receipt=completion_validation,
completion_policy_request=None,
completion_policy_request=completion_policy_source,
)
if transaction["decision"] != "reject":
return {
Expand Down
4 changes: 0 additions & 4 deletions loopx/todos.py
Original file line number Diff line number Diff line change
Expand Up @@ -62,7 +62,6 @@
from .control_plane.todos.addition import matching_todo_block, require_replan_successor_rebinding, require_replan_successor_scope
from .control_plane.todos.completed_archive import archive_completed_todo_lines
from .control_plane.todos.completion_policy import (
bind_completion_policy_to_transaction,
completion_policy_from_transaction,
)
from .control_plane.todos.completion_transaction import (
Expand Down Expand Up @@ -1768,9 +1767,6 @@ def complete_goal_todo(
runtime_root=shadow_runtime_root,
)
)
completion_transaction = bind_completion_policy_to_transaction(
completion_transaction, locked_completion_policy_source
)
completion_fence = completion_transaction["fence"]
completion_state = completion_transaction.get("completion_state")
completion_policy = completion_policy_from_transaction(completion_transaction)
Expand Down
Loading