distributed-systems-iteration-budget-design-handoff
The iteration-budget design-handoff pattern (tick #135, RT-GAP-015) fixes a class of distributed-worker failure: an implementation worker hitting a mid-task design fork and silently burning its iteration budget re-deriving both paths. The pattern's contract:
DecisionRequest; the foreman decides; the worker resumes. The worker never re-iterates over a design question.iteration_budget_handoff/handoff.py)class IterationBudget:
def __init__(self, max_passes: int): self.max_passes = max_passes; self.used = 0
def spend(self) -> bool: # a pass = one deliver->verify cycle
if self.used >= self.max_passes: return False
self.used += 1; return True
@property
def exhausted(self): return self.used >= self.max_passes
@dataclass
class HandoffContract:
deliverable: str
acceptance_criteria: list[str]
decision_points: dict[str, list[str]] # {decision_name: [allowed options]}
max_escalations: int = 2
@dataclass
class DecisionRequest: # the escalation unit
name: str; options: list[str]; context: dict; resolved: str | None = None
def render_worker_prompt(budget, contract) -> str:
return "\n".join([
"ROLE: implementation worker (iteration-budget design-handoff)",
f"PASS BUDGET: {budget.remaining} verification pass(es) remaining (exhausted => no re-iteration)",
f"DELIVERABLE: {contract.deliverable}",
"HANDOFF RULE: on mid-task design ambiguity (Path A vs Path B), DO NOT re-iterate.",
" Escalate one DecisionRequest; the foreman decides; resume with the decision.",
f"DECISION POINTS: {json.dumps(contract.decision_points)}",
f"MAX ESCALATIONS: {contract.max_escalations}",
])
class Foreman:
def decide(self, req: DecisionRequest) -> str: # validates against contract
allowed = self.contract.decision_points.get(req.name)
if allowed is None: raise ContractViolation(...)
choice = self.decision_fn(req.name, req.context)
if choice not in allowed: raise ContractViolation(...)
req.resolved = choice
self.log.append({"event": "foreman_decision", "decision": req.name, "choice": choice})
return choice
def verify(self, artifact) -> list[str]: # [] == accepted
failures = self.verify_fn(artifact); self.log.append({"event": "verify", "failures": failures})
return failures
def repair_foreman_direct(self, artifact, failures): # fix injected, worker NOT re-iterated
...self.log.append({"event": "foreman_direct_repair", ...})
pipeline.py)def run(self, config, discovery) -> WorkerReport:
if len(discovery.source_types) > 1 or discovery.composite_pk_tables: # mid-task fork!
report.status = Status.ESCALATED
report.decision_requests.append(DecisionRequest(
name="pk_strategy",
options=[PATHS["PATH_A"], PATHS["PATH_B"]], # A: per-source adapters
context={"wire_formats": discovery.source_types, # B: normalize in pump core
"composite_pk_tables": discovery.composite_pk_tables}))
return report
def resume(self, report, foreman) -> WorkerReport:
while report.pending_decisions: # foreman decides, exactly once per fork
if self.escalations >= self.contract.max_escalations:
report.status = Status.BUDGET_EXHAUSTED; return report
req = report.pending_decisions[0]
foreman.decide(req) # NOT self.iterate_both_paths()
self.escalations += 1
... # build deliverable per decision; template extractor ships for verification
cdc_pump.py)Bug (found by foreman verification): extract_pk_naive assumes a single id column in payload.after — it crashes on composite PKs (order_id|line_id), Debezium-style payload.key JSON envelopes, and delete (before-only) rows.
Fix (foreman-direct, Path B): one resolver in the pump core, configured per table with dotted PK paths:
def extract_pk_fixed(event, pk_paths) -> str:
payload = event.get("payload", event)
raw_key = payload.get("key") # Debezium envelope: JSON string
if isinstance(raw_key, str) and raw_key:
key_obj = json.loads(raw_key) if ... else None
if isinstance(key_obj, dict) and all(p in key_obj for p in pk_paths):
return "|".join(str(key_obj[p]) for p in pk_paths)
row = payload.get("after") or payload.get("before")# before-only rows (deletes)
return "|".join(str(_dotted_get(row, p)) for p in pk_paths) # composite -> "7|2"
def make_resolver(pk_config): # single choke point, per-table paths
def resolver(event):
table = event.get("table")
if table not in pk_config: raise PkExtractionError(f"no pk config for '{table}'")
return extract_pk_fixed(event, pk_config[table])
return resolver
failures = foreman.verify(artifact) # delivery check (free)
while failures:
if not budget.spend(): return BUDGET_EXHAUSTED # each repair cycle costs one pass
artifact = foreman.repair_foreman_direct(artifact, failures) # canonical fix, no worker round-trip
failures = foreman.verify(artifact)
return DELIVERED
**Verification method:** 24 pytest tests (`~/iteration_budget_handoff/tests/test_handoff.py`), all passing (`24 passed in 0.01s`), run against Python 3.14 / pytest 9.0.2.
**End-to-end tick replay** (recon → design-handoff → implementation → foreman):
```
FINAL STATUS : delivered BUDGET: 1/3
FOREMAN LOG:
foreman_decision pk_strategy -> normalize_in_pump_core # Path B, decided once
verify failures: well-formed drops [naive extractor: DROP 'id', DROP 'id', DROP 'after']
foreman_direct_repair triggered_by: [pk-extraction drops] # fix injected, worker not re-iterated
verify failures: []
DELIVERED OUTBOX : [('c','1001'), ('u','7|2'), ('d','42')] | drops: 0
```
**Pattern invariants tested:**
- **Foreman decides, worker doesn't**: mid-task fork → exactly 1 `DecisionRequest` escalated, `budget.used == 0` while at the fork, foreman logs the decision (`test_foreman_decides_path_not_worker`).
- **Budget bounded**: `spend()` caps at `max_passes`; escalation cap stops the worker at `BUDGET_EXHAUSTED` without re-iterating (`test_escalation_cap_bounds_worker`, `test_budget_spend_bounded`).
- **Contract enforcement**: foreman choices outside declared options → `ContractViolation`; undeclared decisions rejected (`test_foreman_choice_must_satisfy_contract`).
- **Prompt carries budget + contract**: rendered prompt contains pass budget, handoff rule, decision points, escalation cap (`test_worker_prompt_carries_budget_and_contract`).
**Bug found + fix verified:**
- Naive extractor fails on composite PK, Debezium key envelope, and delete rows (3 tests, `KeyError` → clean DROP in pump).
- Fixed extractor passes: single PK, composite `7|2`, Debezium key `7|3`, before-only delete `42`, nested dotted path `tenant.org|id` `9|55` (5 tests).
- Edge cases: missing pk path → clean `PkExtractionError` (no crash); no before/after row → clean error; key JSON missing paths → falls back to row; empty key string ignored; unknown table rejected by resolver; unknown op and malformed rows drop with reasons without crashing (7 tests).
- Pipeline artifacts pass the full 7-event fixture with zero drops on well-formed rows (`test_pipeline_artifact_passes_full_fixture`), and deliver within a **1-pass budget** — one repair pass covers bug-find + fix + clean re-verify (`test_pipeline_budget_not_exceeded_under_repeated_verification`).
**Failure found and fixed during development:** the first run of the pipeline loop revealed batch-state leakage — `CdcPump.process()` accumulated `PumpResult` across verification batches, double-counting events on re-verify. Fixed by resetting the result at the start of each `process()` call (one call = one batch). The foreman's repair itself was proven working: the post-repair batch emitted correct composite keys `7|2, 7|3, 42, 9|55`.
---{"model": "deepseek-v4-flash", "problem_class": "distributed-systems-iteration-budget-design-handoff", "result": "passed", "tests": 24}