◐ Off-By-One · answer catalog

distributed-systems-iteration-budget-design-handoff

1 answer(s)godocker

distributed-systems-iteration-budget-design-handoff

📦 Source in repository (JSON)

Answer

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:

  1. Worker prompt carries an explicit pass budget + handoff contract. Passes are only consumed by foreman verification cycles. Design forks are not free-for-all iterations.
  2. On a mid-task fork (Path A vs Path B), the worker escalates one structured DecisionRequest; the foreman decides; the worker resumes. The worker never re-iterates over a design question.
  3. Foreman verification owns correctness. When it finds a bug (the pk-extraction bug), it repairs the artifact foreman-direct — the canonical fix is injected without sending the artifact back to the worker for another iteration loop.

Pattern core (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", ...})

Worker behavior — escalate, don't iterate (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

The bug and the foreman-direct fix (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

Foreman orchestration (bounded repair loop)

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

Evidence & signatures

**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}
Generated from the verified corpus · MIT licensedBack to the catalog