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
1 change: 1 addition & 0 deletions loopx/control_plane/quota/cli_projection.py
Original file line number Diff line number Diff line change
Expand Up @@ -134,6 +134,7 @@
"spend_after_validation",
"spend_policy",
"delivery_workspace_causality",
"settlement_resume_ref",
),
}
_RETAINED_MONITOR_POLL_RESPONSE_PLAN_FIELDS = (
Expand Down
12 changes: 12 additions & 0 deletions loopx/control_plane/quota/monitor_poll.py
Original file line number Diff line number Diff line change
Expand Up @@ -698,6 +698,7 @@ def record_quota_monitor_poll_for_decision(
task_lease_idempotency_key: str | None = None,
task_lease_expected_version: int | None = None,
use_current_task_lease: bool = False,
auxiliary_settlement_todo: Mapping[str, Any] | None = None,
turn_instance_id: str | None = None,
_index_lock_held: bool = False,
status_reloader: Callable[[], dict[str, Any]] | None = None,
Expand Down Expand Up @@ -750,6 +751,17 @@ def record_quota_monitor_poll_for_decision(
registry_path=registry_path,
runtime_root=runtime_root,
)
if auxiliary_settlement_todo is not None:
decision["auxiliary_settlement_todo"] = {
key: auxiliary_settlement_todo.get(key)
for key in (
"todo_id",
"task_class",
"status",
"claimed_by",
"excluded_agents",
)
}
observation = _observation_packet(
before=before,
agent_id=agent_id,
Expand Down
64 changes: 59 additions & 5 deletions loopx/control_plane/quota/monitor_poll_commit.ts
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@ import { basename, dirname, join, resolve } from "node:path";
import { monitorSuccessorIntent, monitorSuccessorRoute } from "../scheduler/monitor_successor.ts";
import { normalizeTodoCapabilities } from "../todos/work_requirements.ts";
import { parseProjectionDelivery } from "../todos/projection_delivery.ts";
import { readQuotaSettlement, QUOTA_SETTLEMENT_READBACK_REQUEST_SCHEMA } from "./settlement_readback.ts";

import type { JsonObject } from "../effect_program.ts";
import { EffectRuntimeRequestError } from "../effect_runtime_errors.ts";
Expand Down Expand Up @@ -89,6 +90,7 @@ interface MonitorDecision extends JsonObject {
vision_wait_state: JsonObject;
due_monitor_candidates: JsonObject[];
registry_due_monitor: JsonObject;
auxiliary_settlement_todo: JsonObject | null;
}

interface MonitorObservation extends JsonObject {
Expand Down Expand Up @@ -311,6 +313,8 @@ function decisionObject(value: unknown): MonitorDecision {
decision.registry_due_monitor,
"decision.registry_due_monitor",
),
auxiliary_settlement_todo: decision.auxiliary_settlement_todo == null ? null :
requiredObject(decision.auxiliary_settlement_todo, "decision.auxiliary_settlement_todo"),
};
}

Expand Down Expand Up @@ -537,10 +541,60 @@ interface Admission {
external: boolean;
}

function admission(request: MonitorRequest): Admission {
async function auxiliaryMonitorAllowed(
request: MonitorRequest, historicalAdmission: boolean,
): Promise<boolean | null> {
const { decision, observation } = request;
const settlementTodo = observation.settlement_todo_id;
if (!settlementTodo || settlementTodo === observation.todo_id) return null;
const conflict = () => new EffectRuntimeRequestError(
"turn-scoped monitor-poll conflicts with the committed advancement settlement identity: " +
`settlement Todo ${settlementTodo}, observation Todo ${observation.todo_id}`,
"heartbeat_receipt_identity_conflict",
);
const todo = decision.auxiliary_settlement_todo;
// A verified pending v1 receipt already admitted this exact observation.
// Earlier receipts predate the lifecycle-fact field; preserve their recovery
// basis, never reuse it for a new effect or a changed settlement binding.
if (historicalAdmission && todo === null) return null;
if (!request.runtime_root || !request.turn_instance_id || !decision.agent_id ||
observation.actor_agent_id !== decision.agent_id ||
todo === null || todo.todo_id !== settlementTodo || todo.task_class !== "advancement_task" ||
!["open", "done"].includes(String(todo.status)) ||
(todo.claimed_by != null && todo.claimed_by !== decision.agent_id) ||
(todo.excluded_agents != null && (!Array.isArray(todo.excluded_agents) ||
todo.excluded_agents.some(value => typeof value !== "string") ||
todo.excluded_agents.includes(decision.agent_id)))) throw conflict();
if (!historicalAdmission) {
const settlement = await readQuotaSettlement({
schema_version: QUOTA_SETTLEMENT_READBACK_REQUEST_SCHEMA,
runtime_root: request.runtime_root, goal_id: request.goal_id,
agent_id: decision.agent_id, todo_id: settlementTodo,
turn_instance_id: request.turn_instance_id, replan_obligation_id: null,
infer_turn_instance_id: false, allow_unbound_binding: false,
});
const progress = jsonObject(settlement.progress);
if (settlement.found !== true || jsonObject(jsonObject(settlement.identity)?.result)?.failure !== null ||
progress?.schema_version !== "quota_settlement_progress_v0" ||
progress.state === "identity_required" || progress.state === "settled") throw conflict();
}
const monitor = decision.registry_due_monitor;
// Retain ordinary quota/due-work admission and capability/gate projections;
// lifecycle lookup must not become a second should-run bypass.
return !decision.requires_user_action && dueMonitorAllowed(decision, observation) &&
monitor.due === true && candidateMatches(monitor, observation) &&
(monitor.claimed_by == null || monitor.claimed_by === decision.agent_id);
}

async function admission(request: MonitorRequest, historicalAdmission = false): Promise<Admission> {
const blocked = blockedSuccessorAllowed(request.decision);
const external = externalMonitorAllowed(request.decision);
const due = dueMonitorAllowed(request.decision, request.observation);
const auxiliary = await auxiliaryMonitorAllowed(request, historicalAdmission);
const due = auxiliary ?? dueMonitorAllowed(request.decision, request.observation);
if (auxiliary === false) {
throw new EffectRuntimeRequestError("auxiliary monitor-poll requires its own due Monitor target",
"monitor_poll_admission_rejected");
}
if (
request.decision.effective_action !== EffectiveAction.MONITOR_QUIET_SKIP &&
!external && !due && !blocked
Expand Down Expand Up @@ -1925,7 +1979,7 @@ export async function evaluateQuotaMonitorPollCommit(
throw new EffectRuntimeRequestError("provider rejection recovery requires execute");
}
if (request.phase === "event") {
const record = buildRecord(request, admission(request));
const record = buildRecord(request, await admission(request));
return result(
request,
fingerprint,
Expand All @@ -1946,7 +2000,7 @@ export async function evaluateQuotaMonitorPollCommit(
if (!request.execute) {
// A provider may mutate the Todo registry, so previews must pass admission
// before returning a provider plan.
const allowed = admission(request);
const allowed = await admission(request);
if (request.phase === "preflight") {
const providerPlan = providerPlanFor(request);
return result(
Expand Down Expand Up @@ -2076,7 +2130,7 @@ export async function evaluateQuotaMonitorPollCommit(
}
let allowed: Admission;
try {
allowed = admission(admittedRequest);
allowed = await admission(admittedRequest, existing?.schema_version === MONITOR_PENDING_ADMISSION_SCHEMA);
} catch (error) {
if (existing?.schema_version === QUOTA_MONITOR_POLL_COMMIT_RECEIPT_SCHEMA &&
error instanceof EffectRuntimeRequestError && error.code === "monitor_poll_admission_rejected") {
Expand Down
172 changes: 97 additions & 75 deletions loopx/quota.py
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@
from collections.abc import Callable, Mapping
from datetime import datetime, timedelta, timezone
from pathlib import Path
import shlex
from typing import Any

from .control_plane import compact_control_plane_policy
Expand All @@ -23,7 +24,6 @@
ReceiptBoundReplayPhase,
ReceiptBoundTerminalPhase,
)
from .control_plane.quota.error_codes import HeartbeatReceiptIdentityConflictError
from .control_plane.quota.goal_boundary import (
registry_goal_by_id as _registry_goal_by_id,
)
Expand All @@ -34,7 +34,6 @@
from .control_plane.quota.monitor_poll import (
QUOTA_MONITOR_POLL_CLASSIFICATION as QUOTA_MONITOR_POLL_CLASSIFICATION,
build_quota_monitor_poll_event as build_quota_monitor_poll_event,
find_quota_monitor_poll_turn,
record_quota_monitor_poll_for_decision,
resolve_due_monitor_candidate,
)
Expand Down Expand Up @@ -98,8 +97,6 @@
CODEX_APP_SURFACE,
)
from .control_plane.todos.contract import (
TODO_TASK_CLASS_ADVANCEMENT,
TODO_TASK_CLASS_MONITOR,
normalize_todo_claimed_by,
normalize_todo_id,
)
Expand Down Expand Up @@ -1138,87 +1135,31 @@ def should_run(current_status: dict[str, Any]) -> dict[str, Any]:
)

before = should_run(status_payload)
auxiliary_settlement_todo = None
if (
normalized_receipt_todo_id
and normalized_observation_todo_id
and normalized_observation_todo_id != normalized_receipt_todo_id
):
selected = (
before.get("selected_todo")
if isinstance(before.get("selected_todo"), Mapping)
else {}
)
lane = (
before.get("work_lane_contract")
if isinstance(before.get("work_lane_contract"), Mapping)
else {}
)
summary = (
before.get("agent_todo_summary")
if isinstance(before.get("agent_todo_summary"), Mapping)
else {}
)
candidate_values = [
*(lane.get("monitor_due_items") or []),
*(summary.get("monitor_due_items") or []),
]
normalized_agent_id = normalize_todo_claimed_by(agent_id)
auxiliary_due_monitor = any(
isinstance(candidate, Mapping)
and normalize_todo_id(candidate.get("todo_id"))
== normalized_observation_todo_id
and candidate.get("task_class") == TODO_TASK_CLASS_MONITOR
and normalize_todo_claimed_by(candidate.get("claimed_by"))
in {None, normalized_agent_id}
for candidate in candidate_values
)
auxiliary_registry_due = bool(
resolved_monitor
and normalize_todo_claimed_by(resolved_monitor.get("claimed_by"))
in {None, normalized_agent_id}
)
existing_observation = (
find_quota_monitor_poll_turn(
Path(str(raw_runtime_root)).expanduser(),
# Discovery/selection is an open-work projection, not the committed
# Turn's lifecycle authority. Read the exact bound record, including
# completed history; TS monitor admission owns its interpretation.
from .todos import list_goal_todos

if registry_path is not None:
bound_records = list_goal_todos(
registry_path=registry_path,
runtime_root_arg=str(runtime_root) if runtime_root else None,
goal_id=safe_goal_id,
agent_id=normalized_agent_id or "",
turn_instance_id=str(turn_instance_id or ""),
todo_id=normalized_observation_todo_id,
)
if raw_runtime_root and agent_id and turn_instance_id
else None
)
auxiliary_replay = bool(
isinstance(existing_observation, Mapping)
and normalize_todo_id(existing_observation.get("todo_id"))
== normalized_observation_todo_id
and normalize_todo_id(
existing_observation.get("settlement_todo_id")
)
== normalized_receipt_todo_id
)
auxiliary_observation_allowed = bool(
normalize_todo_id(selected.get("todo_id"))
== normalized_receipt_todo_id
and selected.get("task_class") == TODO_TASK_CLASS_ADVANCEMENT
and selected.get("selection_binding") == "heartbeat_receipt"
and (
auxiliary_due_monitor
or auxiliary_registry_due
or auxiliary_replay
)
)
if not auxiliary_observation_allowed:
raise HeartbeatReceiptIdentityConflictError(
"turn-scoped monitor-poll Todo conflicts with the committed "
"heartbeat receipt: expected settlement Todo "
f"{normalized_receipt_todo_id}, requested observation Todo "
f"{normalized_observation_todo_id}"
todo_id=normalized_receipt_todo_id,
role="agent",
)
items = bound_records.get("todos") or []
auxiliary_settlement_todo = items[0] if len(items) == 1 else None
effective_todo_id = normalized_observation_todo_id or (
normalized_receipt_todo_id if not target_key else None
)
return record_quota_monitor_poll_for_decision(
result = record_quota_monitor_poll_for_decision(
before,
status_payload,
goal_id=safe_goal_id,
Expand All @@ -1230,6 +1171,7 @@ def should_run(current_status: dict[str, Any]) -> dict[str, Any]:
reason_summary=reason_summary,
agent_id=agent_id,
settlement_todo_id=normalized_receipt_todo_id,
auxiliary_settlement_todo=auxiliary_settlement_todo,
todo_id=effective_todo_id,
target_key=target_key,
result_hash=result_hash,
Expand All @@ -1251,6 +1193,86 @@ def should_run(current_status: dict[str, Any]) -> dict[str, Any]:
turn_instance_id=turn_instance_id,
status_reloader=status_reloader,
)
continuation = result.get("turn_continuation") or {}
if (
result.get("ok") is True
and continuation.get("settlement_binding_matches_observation") is False
):
# The open-work readback may now select another Todo. Render the
# original receipt's closeout separately; never borrow that selection
# to construct this Turn's refresh/spend commands.
from .control_plane.agents.capability_gate import (
runtime_capabilities_for_cli_projection,
)
from .control_plane.quota.settlement import (
attach_settlement_progress,
build_turn_scoped_cli_settlement_plan,
)

readback = read_heartbeat_settlement(
runtime_root,
goal_id=safe_goal_id,
agent_id=agent_id,
todo_id=normalized_receipt_todo_id,
turn_instance_id=turn_instance_id,
)
if readback is None or readback.identity.value is None:
raise RuntimeError(
"auxiliary monitor receipt lost its original settlement readback"
)
attach_settlement_progress(
result, readback, registry_path=registry_path, runtime_root=runtime_root,
)
identity = readback.identity.value
prefix = "loopx"
if registry_path is not None:
prefix += f" --registry {shlex.quote(str(registry_path))}"
prefix += f" --runtime-root {shlex.quote(str(runtime_root))}"
scoped_args = "".join(
f" --available-capability {shlex.quote(capability)}"
for capability in runtime_capabilities_for_cli_projection(
available_capabilities
)
)
plan = build_turn_scoped_cli_settlement_plan(
goal_id=identity.goal_id,
agent_id=identity.agent_id,
todo_id=identity.todo_id,
turn_instance_id=identity.turn_instance_id,
replan_obligation_id=identity.replan_obligation_id,
command_prefix=prefix,
scoped_cli_args=scoped_args,
lifecycle_actor_args="",
quota_spend_source=readback.progress["quota_spend_source"],
)
result["settlement_resume"] = {
"schema_version": "auxiliary_monitor_settlement_resume_v0",
"identity": identity.as_dict(),
"progress_ref": "$.settlement_progress",
"next_step": next(
(
step for step in plan.as_dict()["ordered_steps"]
if step["kind"] == readback.progress["next_step"]
),
None,
),
"grants_new_delivery": False,
}
# A later discovery projection is not a second settlement plan. The
# top-level resume/readback above owns this original Turn's closeout.
after = result.get("after") or {}
cli = (after.get("interaction_contract") or {}).get("cli_channel")
if (
isinstance(cli, dict)
and (after.get("selected_todo") or {}).get("todo_id") != identity.todo_id
):
cli.pop("settlement_plan", None)
cli.pop("next_cli_actions", None)
cli.update(
spend_allowed_now=False, spend_after_validation=False,
settlement_resume_ref="$.settlement_resume",
)
return result


def build_quota_slot_void_preview(
Expand Down
Loading
Loading