Skip to content

Commit 2403e45

Browse files
committed
refactor(quota): make TypeScript own prior-Turn closeout recovery
Move the prior-host-Turn closeout transaction behind one TypeScript request/response boundary so the Python coordinator keeps only transport and public-payload projection. The typed owner now decides which prior must-attempt Turn still needs a closeout, whether its settlement already validates, which closeout is accepted, and what the recovery obligation is. Python reads exactly the two provider facts the preflight names - the bound Todo and the committed monitor-poll receipt - and projects the typed verdict. Two requests, not one per receipt: a fail-closed preflight that reads the persisted guards and the settlement readback itself, then one final reduction over the checkpointed bound facts. The preflight reads the rollout log from runtime_root because a real goal persists 3.25 MiB of receipt facts for its busiest Agent, and the runtime bridge rejects any request over its 2 MiB bound. Semantics tightened while sharing the rule with the settlement readback: - the heartbeat receipt identity rule has one owner instead of one TypeScript and one Python copy; - a receipt that declares a Todo binding the settlement authority cannot address now fails the read instead of minting a recovery obligation for an identity no other reader can reproduce; - a receipt that names a settlement effect without a Todo or autonomous replan binding keeps failing closed, and the declared effect id is no longer silently replaced by the derived one. Python deletion: the closeout policy, the receipt-selection rule, the identity rule, and the settlement-readback transport for this path. Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com>
1 parent c3cee40 commit 2403e45

14 files changed

Lines changed: 1799 additions & 311 deletions

‎loopx/control_plane/effect_runtime_handlers.ts‎

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -52,6 +52,10 @@ import { evaluateDeliveryWorkspaceCausality } from "./quota/settlement_workspace
5252
import { evaluateQuotaSpendCommit } from "./quota/spend_commit.ts";
5353
import { evaluateQuotaVoidCommit } from "./quota/void_commit.ts";
5454
import { readQuotaSettlement } from "./quota/settlement_readback.ts";
55+
import {
56+
preflightPriorHostTurnCloseout,
57+
reduceUnsettledHostTurnRecovery,
58+
} from "./quota/unsettled_host_turn_recovery.ts";
5559
import { evaluateTurnEnvelope } from "./quota/turn_envelope.ts";
5660
import { evaluateQuotaMonitorPollCommit } from "./quota/monitor_poll_commit.ts";
5761
import { planMonitorSuccessor, selectMonitorTodoRequest } from "./scheduler/monitor_successor.ts";
@@ -463,6 +467,14 @@ export function createEffectRuntimeHandlers(
463467
["quota.spend.commit", evaluateQuotaSpendCommit],
464468
["quota.void.commit", evaluateQuotaVoidCommit],
465469
["quota.settlement.read", readQuotaSettlement],
470+
[
471+
"quota.prior_host_turn_closeout.preflight",
472+
preflightPriorHostTurnCloseout,
473+
],
474+
[
475+
"quota.unsettled_host_turn_recovery.reduce",
476+
reduceUnsettledHostTurnRecovery,
477+
],
466478
["quota.turn_envelope.evaluate", evaluateTurnEnvelope],
467479
["task_lease.owner_eligibility", evaluateTaskLeaseOwnerEligibility],
468480
["task_lease.acquire.decide", evaluateTaskLeaseAcquireDecision],

‎loopx/control_plane/quota/heartbeat_receipt.py‎

Lines changed: 0 additions & 45 deletions
Original file line numberDiff line numberDiff line change
@@ -153,51 +153,6 @@ def find_heartbeat_receipt(
153153
)
154154

155155

156-
def prior_closeout_required_heartbeat_receipts(
157-
runtime_root: Path,
158-
*,
159-
goal_id: str,
160-
agent_id: str,
161-
exclude_turn_instance_id: str | None = None,
162-
) -> list[dict[str, object]]:
163-
"""Return prior heartbeat guards that explicitly require host closeout.
164-
165-
The expectation is opt-in on the persisted guard so older receipts cannot
166-
become false-positive recovery obligations after an upgrade. Results are
167-
newest first and contain at most one effective receipt per Turn.
168-
"""
169-
170-
events = load_rollout_events(rollout_event_log_path(runtime_root, goal_id))
171-
matching_by_turn: dict[str, list[dict[str, object]]] = {}
172-
newest_turns: list[str] = []
173-
excluded = str(exclude_turn_instance_id or "").strip()
174-
for event in events:
175-
if (
176-
event.get("event_kind") != "quota_should_run"
177-
or str(event.get("goal_id") or "") != goal_id
178-
or str(event.get("agent_id") or "") != agent_id
179-
):
180-
continue
181-
turn_id = str(event.get("run_id") or "").strip()
182-
if not turn_id or turn_id == excluded:
183-
continue
184-
matching_by_turn.setdefault(turn_id, []).append(event)
185-
if turn_id in newest_turns:
186-
newest_turns.remove(turn_id)
187-
newest_turns.append(turn_id)
188-
189-
required: list[dict[str, object]] = []
190-
for turn_id in reversed(newest_turns):
191-
effective = _effective_heartbeat_receipt(matching_by_turn[turn_id])
192-
if effective is None:
193-
continue
194-
details_value = effective.get("details")
195-
details = details_value if isinstance(details_value, Mapping) else {}
196-
if details.get("closeout_required") is True:
197-
required.append(effective)
198-
return required
199-
200-
201156
def ensure_turn_heartbeat_settlement_receipt(
202157
runtime_root: Path,
203158
identity: SettlementIdentity,
Lines changed: 182 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,182 @@
1+
/**
2+
* One owner for the persisted heartbeat-receipt settlement identity rule.
3+
*
4+
* A Turn can persist several `quota_should_run` events. The effective receipt
5+
* is the one that binds a settlement identity; a Turn whose events disagree on
6+
* that binding is an identity conflict and must fail closed instead of letting
7+
* a caller infer, upgrade, or silently prefer one binding.
8+
*
9+
* Both the settlement readback and the prior-host-Turn closeout selection read
10+
* this rule, so it lives here rather than in either caller.
11+
*/
12+
import type { JsonObject } from "../effect_program.ts";
13+
import {settlementIdentity, type SettlementIdentity} from "../effect_program.ts";
14+
import { EffectRuntimeRequestError } from "../effect_runtime_errors.ts";
15+
import { jsonObject } from "../runtime_decode.ts";
16+
17+
const TODO_ID_PATTERN = /^todo_[a-z0-9_-]{3,64}$/;
18+
const REPLAN_OBLIGATION_ID_PATTERN = /^replan-[a-f0-9]{16}$/;
19+
/** Public diagnostic preserved from the retired Python identity rule. */
20+
const IDENTITY_CONFLICT_CODE = "heartbeat_receipt_identity_conflict";
21+
22+
function receiptIdentityConflict(message: string): EffectRuntimeRequestError {
23+
return new EffectRuntimeRequestError(message, IDENTITY_CONFLICT_CODE);
24+
}
25+
26+
export function heartbeatReceiptDetails(event: JsonObject | null): JsonObject {
27+
return jsonObject(event?.details) ?? {};
28+
}
29+
30+
export function optionalHeartbeatString(value: unknown): string | null {
31+
if (value === null || value === undefined) return null;
32+
return String(value).trim() || null;
33+
}
34+
35+
export function normalizeHeartbeatTodoId(value: unknown): string | null {
36+
const candidate = String(value ?? "").trim().toLowerCase();
37+
return candidate && TODO_ID_PATTERN.test(candidate) ? candidate : null;
38+
}
39+
40+
export function normalizeHeartbeatReplanObligationId(
41+
value: unknown,
42+
): string | null {
43+
const candidate = String(value ?? "").trim();
44+
return candidate && REPLAN_OBLIGATION_ID_PATTERN.test(candidate)
45+
? candidate
46+
: null;
47+
}
48+
49+
export interface HeartbeatReceiptBinding extends JsonObject {
50+
binding_kind: "todo" | "autonomous_replan";
51+
/** The Todo id or autonomous replan obligation id, per `binding_kind`. */
52+
binding_id: string;
53+
/**
54+
* The effect id the receipt declares, kept verbatim for projection.
55+
*
56+
* It is intentionally not repaired from the identity's derived effect id: a
57+
* receipt that names an effect the settlement authority could not derive
58+
* stays visible in the projected payload instead of being silently rewritten.
59+
*/
60+
settlement_effect_id: string | null;
61+
/** The conflict-detection key: one Turn may declare only one of these. */
62+
identity_key: string;
63+
}
64+
65+
/**
66+
* The only receipt fields the settlement binding rule can read.
67+
*
68+
* The persisted rollout event and the settlement readback's own projection both
69+
* decode into this shape, so one rule serves both readers. The goal and Agent
70+
* a receipt belongs to are not receipt facts: they are the scope the reader
71+
* already selected.
72+
*/
73+
export interface HeartbeatReceiptFact extends JsonObject {
74+
event_id: string | null;
75+
run_id: string | null;
76+
todo_id: string | null;
77+
replan_obligation_id: string | null;
78+
settlement_effect_id: string | null;
79+
closeout_required: boolean;
80+
}
81+
82+
export function heartbeatReceiptFactFromEvent(
83+
event: JsonObject,
84+
): HeartbeatReceiptFact {
85+
const eventDetails = heartbeatReceiptDetails(event);
86+
return {
87+
event_id: optionalHeartbeatString(event.event_id),
88+
run_id: optionalHeartbeatString(event.run_id),
89+
todo_id: optionalHeartbeatString(eventDetails.todo_id),
90+
replan_obligation_id: optionalHeartbeatString(
91+
eventDetails.replan_obligation_id,
92+
),
93+
settlement_effect_id: optionalHeartbeatString(
94+
eventDetails.settlement_effect_id,
95+
),
96+
closeout_required: eventDetails.closeout_required === true,
97+
};
98+
}
99+
100+
/**
101+
* Resolve the settlement binding a persisted receipt declares.
102+
*
103+
* A receipt without a binding is not a settlement identity; it is reported as
104+
* `null` so the caller can decide whether an unbound receipt is admissible.
105+
*/
106+
export function heartbeatReceiptBinding(
107+
goalId: string,
108+
agentId: string,
109+
fact: HeartbeatReceiptFact,
110+
): HeartbeatReceiptBinding | null {
111+
const declaredTodoId = optionalHeartbeatString(fact.todo_id);
112+
const replanObligationId = normalizeHeartbeatReplanObligationId(
113+
fact.replan_obligation_id,
114+
);
115+
const declaredEffectId = optionalHeartbeatString(fact.settlement_effect_id);
116+
if (declaredTodoId && replanObligationId) {
117+
throw receiptIdentityConflict(
118+
"heartbeat receipt has conflicting Todo and autonomous replan bindings",
119+
);
120+
}
121+
if (declaredEffectId && !declaredTodoId && !replanObligationId) {
122+
throw receiptIdentityConflict(
123+
"heartbeat receipt has an effect identity without a Todo or autonomous replan binding; refuse to infer or upgrade it",
124+
);
125+
}
126+
if (!declaredTodoId && !replanObligationId) return null;
127+
// A Todo-binding the settlement authority cannot address is not a binding we
128+
// may act on: acting on the raw string would create a recovery obligation
129+
// whose identity no other reader can reproduce, and dropping it would let a
130+
// required closeout disappear. Refuse the read instead.
131+
const todoId = declaredTodoId;
132+
if (todoId !== null && normalizeHeartbeatTodoId(todoId) !== todoId) {
133+
throw receiptIdentityConflict(
134+
"heartbeat receipt declares a Todo binding that is not a legal Todo id",
135+
);
136+
}
137+
const identity = settlementIdentity({
138+
goal_id: goalId,
139+
agent_id: agentId,
140+
todo_id: todoId,
141+
turn_instance_id: fact.run_id ?? "",
142+
replan_obligation_id: replanObligationId,
143+
});
144+
return {
145+
binding_kind: identity.binding_kind === "todo"
146+
? "todo"
147+
: "autonomous_replan",
148+
binding_id: identity.binding_id,
149+
settlement_effect_id: declaredEffectId,
150+
identity_key: `${identity.binding_kind}\u0000${identity.binding_id}\u0000${
151+
declaredEffectId ?? identity.effect_id
152+
}`,
153+
};
154+
}
155+
156+
/**
157+
* Reduce one Turn's receipts to its effective receipt.
158+
*
159+
* Receipts without a settlement binding are admissible only while the Turn
160+
* declares no binding at all; they cannot outrank a bound receipt, and they
161+
* cannot silently turn a conflicting Turn into a valid one.
162+
*/
163+
export function selectEffectiveHeartbeatReceipt<Value>(
164+
goalId: string,
165+
agentId: string,
166+
entries: readonly { fact: HeartbeatReceiptFact; value: Value }[],
167+
): Value | null {
168+
if (entries.length === 0) return null;
169+
const identities = new Map<string, Value>();
170+
for (const entry of entries) {
171+
const binding = heartbeatReceiptBinding(goalId, agentId, entry.fact);
172+
if (binding) identities.set(binding.identity_key, entry.value);
173+
}
174+
if (identities.size > 1) {
175+
throw receiptIdentityConflict(
176+
"heartbeat receipt has conflicting settlement identities for the same goal, agent, and turn",
177+
);
178+
}
179+
return identities.size === 1
180+
? [...identities.values()][0]!
181+
: entries.at(-1)!.value;
182+
}

0 commit comments

Comments
 (0)