diff --git a/CHANGELOG.md b/CHANGELOG.md index 5b8f51c4..f936c0a5 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -9,6 +9,12 @@ later entries are regular releases. ### Added +- Staged trusted lineage producer primitives for supervised takeover and same-writer + delivery persistence, authenticated semantic publication, explicit builder-label + reconciliation, and transport-preserving attribution. Independent workflow and + runner assets are available for explicit materialization; normal init and + automatic adoption remain unchanged until #992. + - A pure typed builder-lineage contract with immutable exact targets, validated contribution chains, explicit comment history and authority accounts, and contributor-aware reviewer admission. Includes a standalone init support diff --git a/code-mower-package-manifest.json b/code-mower-package-manifest.json index a0f15c9c..6364b960 100644 --- a/code-mower-package-manifest.json +++ b/code-mower-package-manifest.json @@ -377,6 +377,11 @@ "source": "src/code_mower/builder_lineage.py", "target": "src/code_mower/builder_lineage.py" }, + { + "kind": "core", + "source": "src/code_mower/builder_lineage_producer.py", + "target": "src/code_mower/builder_lineage_producer.py" + }, { "kind": "core", "source": "src/code_mower/builder_runs.py", @@ -1617,6 +1622,11 @@ "source": "generated", "target": "src/code_mower/templates/lanes/generic.md" }, + { + "kind": "template", + "source": "src/code_mower/templates/lanes/lineage-producer.sh", + "target": "src/code_mower/templates/lanes/lineage-producer.sh" + }, { "kind": "lane-template", "source": "generated", @@ -1687,6 +1697,11 @@ "source": "generated", "target": "src/code_mower/templates/workflows/audit-label-cleanup.yml.j2" }, + { + "kind": "template", + "source": "src/code_mower/templates/workflows/builder-lineage-producer.yml.j2", + "target": "src/code_mower/templates/workflows/builder-lineage-producer.yml.j2" + }, { "kind": "workflow", "source": "generated", @@ -1862,6 +1877,11 @@ "source": "generated", "target": "templates/lanes/generic.md" }, + { + "kind": "template", + "source": "templates/lanes/lineage-producer.sh", + "target": "templates/lanes/lineage-producer.sh" + }, { "kind": "lane-template", "source": "generated", @@ -1992,6 +2012,11 @@ "source": "generated", "target": "templates/workflows/blind-review-artifacts-dry-run.yml.j2" }, + { + "kind": "template", + "source": "templates/workflows/builder-lineage-producer.yml.j2", + "target": "templates/workflows/builder-lineage-producer.yml.j2" + }, { "kind": "workflow", "source": "generated", diff --git a/src/code_mower/builder_lineage_producer.py b/src/code_mower/builder_lineage_producer.py new file mode 100644 index 00000000..f5d99105 --- /dev/null +++ b/src/code_mower/builder_lineage_producer.py @@ -0,0 +1,386 @@ +"""Explicit, staged lineage producers. No default CLI or runner calls this module. + +The broker owns policy, transport bindings, private stores and I/O adapters. +Provider text is never authority. Compatibility is refusal, not migration. +""" +from __future__ import annotations + +from dataclasses import dataclass +from itertools import chain +from pathlib import Path +import re +import subprocess + +from .builder_lineage import ( + Authorities, Chain, ContractError, Episode, History, Identity, LINEAGE_MARKER, + Target, parse_markers, render, resolve, +) +from .context_store import ContextStore, strict_json + + +class ProducerRefusal(ContractError): + """A bounded owner action, including any effects already attempted.""" + + def __init__(self, reason, *, comment_posted=False, labels_attempted=False): + super().__init__(reason) + self.comment_posted = comment_posted + self.labels_attempted = labels_attempted + + +def _supported(value, kind): + """The single compatibility boundary; unsupported history must be re-recorded.""" + legacy = False + if kind == "history": + legacy = any(LINEAGE_MARKER in c.body and not re.search( + LINEAGE_MARKER + r":", c.body) for c in value.comments) + elif kind == "episode": + legacy = isinstance(value, dict) and ( + "schema" in value or value.get("writer_state") == "self_quiescent") + elif kind == "record": + legacy = not isinstance(value, dict) or value.get("schema") != "code_mower.lineageProducer.v1" + if legacy: + raise ProducerRefusal("Unsupported lineage history; obtain verified re-recording.") + return value + + +def _arrivals(public, private, authorities): + _supported(public, "history") + for raw in chain(parse_markers(public, authorities), private): + yield _supported(raw, "episode") + + +@dataclass(frozen=True) +class Observation: + chain: Chain + decision: object + + +def observe(target, identity, authorities, history, private, *, author, labels): + """Consume all raw arrivals in one budget, then resolve exactly once.""" + if not isinstance(target, Target) or not isinstance(identity, Identity): + raise ProducerRefusal("Exact target and explicit identity policy required.") + if not isinstance(authorities, Authorities) or not authorities.accounts: + raise ProducerRefusal("Explicit publication authority required.") + if not isinstance(history, History): + raise ProducerRefusal("Readable authenticated history required.") + bound = Chain.from_arrivals(target, _arrivals(history, private, authorities)) + decision = resolve(bound, identity, author, labels) + if decision.status != "ready": + raise ProducerRefusal("Lineage is not ready: " + decision.reason) + return Observation(bound, decision) + + +def selected_history(payload, *, selected): + """An explicitly selected raw list; REST's numeric comments count is unrelated.""" + if not isinstance(payload, dict): + raise ProducerRefusal("PR metadata must be an object.") + return History(selected) + + +def decode_transport(raw): + """Strict JSON without normalizing null, duplicate keys or malformed ingress.""" + if not isinstance(raw, str): + raise ProducerRefusal("Transport must return JSON text.") + return strict_json('{"value":' + raw + '}')["value"] + + +def fetch_history(fetch_page, *, page_size=100, max_pages=8): + """Finite explicit requests, with one empty extra-page proof at a full cap.""" + if type(page_size) is not int or not 1 <= page_size <= 100: + raise ProducerRefusal("Invalid history page size.") + if type(max_pages) is not int or not 1 <= max_pages <= 8: + raise ProducerRefusal("Invalid history page cap.") + pages = [] + for page in range(1, max_pages + 2): + try: + raw = fetch_page(page, page_size) + except Exception: + raise ProducerRefusal("Authenticated history request failed.") from None + History(raw) # Validate *before* shape/terminal-page decisions. + if len(raw) > page_size or (page > max_pages and raw): + raise ProducerRefusal("Complete history exceeds the page cap.") + if page <= max_pages: + pages.append(raw) + if len(raw) < page_size: + return History.from_pages(pages) + raise AssertionError("finite page probe must terminate") + + +@dataclass(frozen=True) +class Snapshot: + target: Target + author: str + labels: tuple[str, ...] + + def __post_init__(self): + if not isinstance(self.target, Target) or not isinstance(self.author, str): + raise ProducerRefusal("Exact target and author read required.") + if not isinstance(self.labels, tuple) or any( + not isinstance(label, str) or not label.strip() for label in self.labels): + raise ProducerRefusal("Readable explicit labels required.") + + +def exact_snapshot(io, target): + snapshot = io.snapshot(target) + if not isinstance(snapshot, Snapshot) or snapshot.target != target: + raise ProducerRefusal("Fresh target, exact branch or head changed.") + return snapshot + + +def label_plan(observation, snapshot, identity): + if snapshot.target != observation.chain.target or observation.decision.status != "ready": + raise ProducerRefusal("Label plan requires the ready exact observation.") + writer = observation.decision.current_writer + desired = f"builder:{writer}" + if not writer or dict(identity.labels).get(desired) != writer: + raise ProducerRefusal("Current writer requires an explicit builder label mapping.") + active = [s for s in snapshot.labels if s.lower().startswith("builder:")] + remove = tuple(s for s in active if s != desired) + return desired, remove, desired not in active + + +@dataclass(frozen=True) +class Publication: + observation: Observation + comment_posted: bool + labels_attempted: bool + + +def publish(io, target, identity, authorities, private): + """Explicit POST -> public-only semantic readback -> fresh labels -> verification. + + A refusal after POST truthfully reports the partial effect. No retry occurs. + """ + posted = attempted = False + try: + initial = exact_snapshot(io, target) + public = io.history(target) + observation = observe(target, identity, authorities, public, private, + author=initial.author, labels=initial.labels) + expected = observation.chain.episodes + body = render(observation.chain) # Empty chains cannot be published. + # Validate the plan before POST, including configured destination label. + label_plan(observation, initial, identity) + public_only = Chain.from_arrivals(target, _arrivals(public, (), authorities)) + fresh = exact_snapshot(io, target) + if fresh != initial: + raise ProducerRefusal("Target metadata changed before publication.") + if public_only.episodes != expected: + io.post(target, body) + posted = True + readback = io.history(target) + # Private evidence must never repair a missing final public publication. + verified = Chain.from_arrivals(target, _arrivals(readback, (), authorities)) + if verified.episodes != expected: + raise ProducerRefusal("Public semantic lineage readback differs.") + fresh = exact_snapshot(io, target) + if fresh.author != initial.author: + raise ProducerRefusal("Author changed before reconciliation.") + desired, remove, add = label_plan(observation, fresh, identity) + if remove or add: + attempted = True + io.labels(target, desired, remove, add) + final = exact_snapshot(io, target) + active = [s for s in final.labels if s.lower().startswith("builder:")] + if active != [desired]: + raise ProducerRefusal("One active builder label was not verified.") + return Publication(observation, posted, attempted) + except Exception as exc: + reason = str(exc) if isinstance(exc, ContractError) else "Producer I/O unavailable; inspect partial outcome." + raise ProducerRefusal(reason, comment_posted=posted, labels_attempted=attempted) from None + + +@dataclass(frozen=True) +class Transport: + """Observed transport metadata supplied by the trusted launch/record adapter.""" + lane: str + provider: str + executor: str + integration: str + + def __post_init__(self): + allowed = { + ("codex", "codex", "codex_cli", "local_cli"), + ("claude", "claude", "claude_cli", "local_cli"), + ("devin", "devin_cli", "devin_cli", "local_cli"), + ("devin", "devin", "devin", "hosted_async_builder"), + } + if (self.lane, self.provider, self.executor, self.integration) not in allowed: + raise ProducerRefusal("Unsupported observed builder transport.") + + +def require_producer(transport, config, runtime_observation): + """Reuse role qualification; runtime is obtained from the broker's observer.""" + from .role_eligibility import decide_role, require_role + if not isinstance(transport, Transport): + raise ProducerRefusal("Observed transport required.") + runtime = runtime_observation() + decision = decide_role(transport.lane, "builder", config=config, + transport="devin_api_v3" if transport.integration == "hosted_async_builder" + else transport.executor, runtime=runtime, bounded=True) + require_role(decision, execution=True) + return decision + + +@dataclass(frozen=True, init=False) +class VerifiedDelivery: + """Only the trusted handoff/delivery converters mint this internal receipt.""" + episode: Episode + writer: str + round_id: str + transport: Transport + + +def _delivery(episode, writer, round_id, transport): + for value in (writer, round_id): + if not isinstance(value, str) or not re.fullmatch(r"[a-zA-Z0-9_-]{1,100}", value): + raise ProducerRefusal("Named writer and round identifiers required.") + if not isinstance(episode, Episode) or not isinstance(transport, Transport): + raise ProducerRefusal("Validated episode and observed transport required.") + if episode.destination_lane != transport.lane: + raise ProducerRefusal("Actual writer lane differs from delivery.") + receipt = object.__new__(VerifiedDelivery) + for key, value in dict(episode=episode, writer=writer, round_id=round_id, transport=transport).items(): + object.__setattr__(receipt, key, value) + return receipt + + +class ProducerStore: + """Bounded metadata in the existing private, atomic, descriptor-anchored store. + + Missing records refuse on read. Creation is explicit and only accepts a + verified first takeover; falsey/malformed existing payloads never mean empty. + """ + def __init__(self, root): + self.store = ContextStore(root) + + def _key(self, target): + from .lane_handoff import key + return key([target.repo, target.pr_number, target.branch]) + + def _read(self, raw, target): + _supported(raw, "record") + if set(raw) != {"schema", "repo", "pr_number", "branch", "episodes", "writer", "round_id", "transport"}: + raise ProducerRefusal("Private lineage record fields differ.") + if (raw["repo"], raw["pr_number"], raw["branch"]) != (target.repo, target.pr_number, target.branch): + raise ProducerRefusal("Private lineage target differs.") + if not isinstance(raw["episodes"], list) or not raw["episodes"] or len(raw["episodes"]) > 32: + raise ProducerRefusal("Private lineage episode list is invalid.") + transport = Transport(**raw["transport"]) + # Metadata validation without a separate Chain/budget or resolution. + _delivery(Episode.from_mapping(_supported(raw["episodes"][-1], "episode")), + raw["writer"], raw["round_id"], transport) + return raw + + def read(self, target): + with self.store.locked(self._key(target)) as locked: + return self._read(locked.read(), target) + + def record(self, delivery, target, identity, authorities, history, *, author, labels, + config, runtime_observation, create=False): + if not isinstance(delivery, VerifiedDelivery): + raise ProducerRefusal("Verified supervised delivery required.") + require_producer(delivery.transport, config, runtime_observation) + with self.store.locked(self._key(target)) as locked: + raw = locked.read() + if raw is None and create: + previous = [] + if delivery.episode.sequence != 1 or delivery.episode.kind != "handoff": + raise ProducerRefusal("Creation requires a first verified takeover.") + else: + raw = self._read(raw, target) + previous = raw["episodes"] + if delivery.episode.kind == "continuation" and ( + raw["writer"] != delivery.writer or raw["transport"] != delivery.transport.__dict__): + raise ProducerRefusal("Continuation must retain the actual writer and transport.") + if raw["round_id"] == delivery.round_id and previous[-1] != delivery.episode.to_mapping(): + raise ProducerRefusal("Conflicting replay of a supervised round.") + observation = observe(target, identity, authorities, history, + chain(previous, (delivery.episode,)), author=author, labels=labels) + episodes = [e.to_mapping() for e in observation.chain.episodes] + updated = dict(schema="code_mower.lineageProducer.v1", repo=target.repo, + pr_number=target.pr_number, branch=target.branch, episodes=episodes, + writer=delivery.writer, round_id=delivery.round_id, + transport=delivery.transport.__dict__) + if raw == updated: + return False + # An old round may not replace current writer metadata on an idempotent replay. + if previous and delivery.episode.sequence != len(episodes): + raise ProducerRefusal("Stale delivery replay.") + locked.write(updated) + return True + + +def projection(observation): + return {"status": observation.decision.status, "current_writer": observation.decision.current_writer, + "contributors": list(observation.decision.contributors), + "head_sha": observation.chain.target.head_sha, + "episode_count": len(observation.chain.episodes)} + + +class GitHub: + """Authenticated gh transport with finite requests; no checkout or config execution.""" + def _json(self, endpoint, *args): + result = subprocess.run(["gh", "api", endpoint, *args], check=True, text=True, + stdout=subprocess.PIPE, stderr=subprocess.DEVNULL, timeout=30) + return decode_transport(result.stdout) + + def snapshot(self, target): + raw = self._json(f"repos/{target.repo}/pulls/{target.pr_number}") + if not isinstance(raw, dict) or raw.get("state") != "open": + raise ProducerRefusal("Open PR target read required.") + observed = Target(raw["base"]["repo"]["full_name"], raw["number"], + raw["head"]["ref"], raw["head"]["sha"]) + labels = self._json(f"repos/{target.repo}/issues/{target.pr_number}/labels?per_page=100") + if not isinstance(labels, list) or len(labels) >= 100: + raise ProducerRefusal("Complete readable label list required.") + if any(not isinstance(item, dict) or not isinstance(item.get("name"), str) for item in labels): + raise ProducerRefusal("Malformed labels.") + return Snapshot(observed, raw["user"]["login"], tuple(item["name"] for item in labels)) + + def history(self, target): + return fetch_history(lambda page, size: self._json( + f"repos/{target.repo}/issues/{target.pr_number}/comments?per_page={size}&page={page}")) + + def post(self, target, body): + self._json(f"repos/{target.repo}/issues/{target.pr_number}/comments", + "--method", "POST", "-f", "body=" + body) + + def labels(self, target, desired, remove, add): + argv = ["gh", "pr", "edit", str(target.pr_number), "--repo", target.repo] + if add: + argv.extend(["--add-label", desired]) + for label in remove: + argv.extend(["--remove-label", label]) + subprocess.run(argv, check=True, stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL, timeout=30) + + +def staged_record(environ, *, io=None, clock=None): + """Real API for the independent staged artifacts; retrieval/attribution only. + + Policy/authority and observed transport are reviewed broker inputs. The + immutable policy revision is mandatory metadata, never a moving-ref read. + """ + from .builder_runs import record_lineage_builder + from datetime import datetime, timezone + io = io if io is not None else GitHub() + target = Target.from_mapping(decode_transport(environ["LINEAGE_TARGET_JSON"])) + policy = decode_transport(environ["LINEAGE_POLICY_JSON"]) + authority = Authorities(decode_transport(environ["LINEAGE_AUTHORITY_JSON"])) + identity = Identity.from_mapping(policy["identity"]) + # Validate a reviewed immutable base; neither base nor PR code is executed. + Target(target.repo, target.pr_number, target.branch, policy["base_sha"]) + transport = Transport(**decode_transport(environ["LINEAGE_TRANSPORT_JSON"])) + from .role_eligibility import decide_role, require_role + # Attribution does not launch a provider or claim a runtime/containment probe. + require_role(decide_role(transport.lane, "builder", config=policy["roles"], + transport="devin_api_v3" if transport.integration == "hosted_async_builder" + else transport.executor, bounded=True)) + initial = exact_snapshot(io, target) + history = io.history(target) + observation = observe(target, identity, authority, history, (), + author=initial.author, labels=initial.labels) + if exact_snapshot(io, target) != initial: + raise ProducerRefusal("Snapshot changed before attribution.") + now = clock() if clock else datetime.now(timezone.utc).replace(microsecond=0).isoformat() + return record_lineage_builder(observation, transport, Path(environ["LINEAGE_OUTPUT"]), created_at=now) diff --git a/src/code_mower/builder_runs.py b/src/code_mower/builder_runs.py index 0354a01b..7d38a82e 100644 --- a/src/code_mower/builder_runs.py +++ b/src/code_mower/builder_runs.py @@ -754,3 +754,21 @@ def main(argv: list[str] | None = None) -> int: if __name__ == "__main__": # pragma: no cover raise SystemExit(main()) + + +def record_lineage_builder(observation, transport, output, *, created_at): + """Explicit attribution after validated selected history; legacy auto-record is unchanged.""" + from .builder_lineage_producer import Observation, ProducerRefusal, Transport, projection + if not isinstance(observation, Observation) or not isinstance(transport, Transport): + raise ProducerRefusal("Validated complete observation and actual transport required.") + if (observation.decision.status != "ready" + or observation.decision.current_writer != transport.lane): + raise ProducerRefusal("Observed transport does not match the verified current writer.") + target = observation.chain.target + event = build_builder_run_event(provider=transport.provider, executor=transport.executor, + integration=transport.integration, repo=target.repo, pr=f"{target.repo}#{target.pr_number}", + branch=target.branch, builder_id=f"{transport.lane}-{target.pr_number}-{target.head_sha}", + created_at=created_at) + event["dimensions"]["lineage"] = projection(observation) + write_builder_run_event(event, output) + return event diff --git a/src/code_mower/lane_delivery.py b/src/code_mower/lane_delivery.py index b145ee38..5f8bd482 100644 --- a/src/code_mower/lane_delivery.py +++ b/src/code_mower/lane_delivery.py @@ -1474,3 +1474,111 @@ def _supervise_main(args: argparse.Namespace) -> int: if __name__ == "__main__": # pragma: no cover raise SystemExit(main()) + + +def lineage_target_state(repo, payload): + """Explicit exact producer snapshot; legacy classify normalization is unchanged.""" + from .builder_lineage import Target + from .builder_lineage_producer import ProducerRefusal, Snapshot + if (not isinstance(payload, dict) or payload.get("snapshot_complete") is not True + or payload.get("kind") != "pr" or payload.get("pr_state") != "OPEN" + or not isinstance(payload.get("labels"), list)): + raise ProducerRefusal("Complete exact PR snapshot required.") + number = payload.get("pr_number") + if not isinstance(number, str) or not number.isdigit() or payload.get("number") != number: + raise ProducerRefusal("Exact PR snapshot number required.") + target = Target(repo, int(number), payload.get("branch"), payload.get("head_sha")) + return Snapshot(target, payload.get("author"), tuple(payload["labels"])) + + +def _lineage_checkout(checkout, target): + from .builder_lineage_producer import ProducerRefusal + path = Path(checkout) + if path != path.resolve() or not (path / ".git").exists(): + raise ProducerRefusal("Known delivery checkout unavailable.") + def git(*args): + return subprocess.check_output(["git", "-C", str(path), *args], text=True, + timeout=10, stderr=subprocess.DEVNULL).strip() + if git("rev-parse", "HEAD") != target.head_sha or git("branch", "--show-current") != target.branch: + raise ProducerRefusal("Observed delivery checkout head or exact branch differs.") + + +class LineageRound: + """Explicit supervisor observer, inactive until a broker deliberately uses it. + + A round is registered before launch and passed as supervise_process(writer=). + Existing LocalWriter owns stop/reap observations. Stable writer identity is + distinct from the unique, non-reusable supervised round ID. + """ + def __init__(self, root, round_id, writer, before, transport, checkout, *, + config, runtime_observation): + from .builder_lineage import Target + from .builder_lineage_producer import ProducerRefusal, require_producer + from .lane_handoff import LocalWriter + require_producer(transport, config, runtime_observation) + if not isinstance(before, Target) or any( + not isinstance(v, str) or not re.fullmatch(r"[A-Za-z0-9_-]{1,100}", v) + for v in (round_id, writer)): + raise ProducerRefusal("Exact named supervised round required.") + _lineage_checkout(checkout, before) + self.control = LocalWriter(root, round_id) + self.control.register(repo=before.repo, lane=transport.lane, checkout=Path(checkout)) + self.round_id, self.writer, self.before, self.transport = round_id, writer, before, transport + with self.control.store.locked(self.control.key) as locked: + record = locked.read() + record["lineage_round"] = dict(round_id=round_id, writer=writer, + target={"repo": before.repo, "pr_number": before.pr_number, + "branch": before.branch, "head_sha": before.head_sha}, + transport=transport.__dict__) + locked.write(record) + + def stop_requested(self): + return self.control.stop_requested() + + def started(self, pid, pgid): + self.control.started(pid, pgid) + + def finish(self, *, quiescent): + self.control.finish(quiescent=quiescent) + + def observed(self, after): + from .builder_lineage_producer import ProducerRefusal + if (after.repo, after.pr_number, after.branch) != ( + self.before.repo, self.before.pr_number, self.before.branch): + raise ProducerRefusal("Delivery round target differs.") + with self.control.store.locked(self.control.key) as locked: + record = locked.read() + expected = dict(round_id=self.round_id, writer=self.writer, + target={"repo": self.before.repo, "pr_number": self.before.pr_number, + "branch": self.before.branch, "head_sha": self.before.head_sha}, + transport=self.transport.__dict__) + if (not isinstance(record, dict) or record.get("schema") != "code_mower.localWriter.v1" + or record.get("lineage_round") != expected + or record.get("repo") != after.repo or record.get("lane") != self.transport.lane + or record.get("finished") is not True or record.get("quiescent") is not True + or any(type(record.get(k)) is not int or record[k] <= 0 for k in ("pid", "pgid"))): + raise ProducerRefusal("Independent stopped/reaped named writer evidence required.") + _lineage_checkout(record["checkout"], after) + return self + + +def lineage_continuation(round_observer, after, previous): + """A stopped supervised round by the same actual writer, with chained heads.""" + from .builder_lineage import Episode + from .builder_lineage_producer import ProducerRefusal, _delivery + if not isinstance(round_observer, LineageRound): + raise ProducerRefusal("A supervised round observer is required.") + observed = round_observer.observed(after) + if (not isinstance(previous, dict) or previous.get("writer") != observed.writer + or previous.get("transport") != observed.transport.__dict__ + or previous.get("round_id") == observed.round_id): + raise ProducerRefusal("Continuation requires the same writer and a fresh supervised round.") + episode = Episode.from_mapping(previous["episodes"][-1]) + if (episode.repo, episode.pr_number, episode.branch, episode.resulting_head, episode.destination_lane) != ( + after.repo, after.pr_number, after.branch, observed.before.head_sha, observed.transport.lane): + raise ProducerRefusal("Continuation does not chain the exact prior delivery.") + return _delivery(Episode(sequence=episode.sequence + 1, repo=after.repo, + pr_number=after.pr_number, branch=after.branch, source_lane=observed.transport.lane, + destination_lane=observed.transport.lane, expected_head=observed.before.head_sha, + resulting_head=after.head_sha, writer_state="same_writer", kind="continuation"), + observed.writer, observed.round_id, observed.transport) diff --git a/src/code_mower/lane_handoff.py b/src/code_mower/lane_handoff.py index 6623adc2..941c91e8 100644 --- a/src/code_mower/lane_handoff.py +++ b/src/code_mower/lane_handoff.py @@ -206,3 +206,72 @@ def reserve_launch(handoff: Handoff, root: Path, *, head: Callable = observe_hea record["launch_reserved"] = True locked.write(record) return True + + +def observe_lineage_source(source, handoff): + """Read independent source exit evidence, without cancellation or discovery.""" + from .builder_lineage import Target + from .builder_lineage_producer import ProducerRefusal + from .lane_delivery import _lineage_checkout + repo, number = handoff.target_pr.split("#") + target = Target(repo, int(number), handoff.target_branch, handoff.expected_head) + if source.get("transport") == "local_process": + if set(source) != {"transport", "state_dir", "writer"}: + raise ProducerRefusal("Exact source writer binding required.") + writer = LocalWriter(Path(source["state_dir"]), source["writer"]) + with writer.store.locked(writer.key) as locked: + record = locked.read() + if (not isinstance(record, dict) or record.get("schema") != "code_mower.localWriter.v1" + or record.get("repo", "").lower() != target.repo + or record.get("lane") != handoff.source_lane + or record.get("finished") is not True or record.get("quiescent") is not True + or any(type(record.get(k)) is not int or record[k] <= 0 for k in ("pid", "pgid"))): + raise ProducerRefusal("Source supervisor has not proved named writer exit.") + _lineage_checkout(record["checkout"], target) + elif source.get("transport") == "remote_session": + from .remote_session import _key + if (set(source) != {"transport", "state_dir", "provider", "session"} + or source["provider"] not in {"fake", handoff.source_lane}): + raise ProducerRefusal("Exact remote writer binding required.") + engine = remote_engine(source) + with engine.store.locked(_key(source["session"])) as locked: + record = locked.read() + engine._writer_binding(record, target.repo) + if (record.get("writer_retired") is not True + or any(op["state"] == "pending" for op in record["operations"].values()) + or engine._writer_observation(record) != "terminated"): + raise ProducerRefusal("Remote source exit is unverified.") + else: + raise ProducerRefusal("Unsupported source writer transport.") + return "terminated" + + +def lineage_handoff(handoff, root, round_observer, after, *, source_branch_prefixes, + sequence=1): + """Convert accepted #962 launch binding plus independently observed exits.""" + from .builder_lineage import Episode + from .builder_lineage_producer import ProducerRefusal, _delivery + from .lane_delivery import LineageRound, validate_handoff + if not isinstance(round_observer, LineageRound): + raise ProducerRefusal("A supervised destination round is required.") + observed = round_observer.observed(after) + validated = validate_handoff(**handoff.as_dict(), running_lane=observed.transport.lane, + repo=after.repo, observed_head=observed.before.head_sha, + source_branch_prefixes=source_branch_prefixes) + if (validated != handoff or handoff.target_pr.lower() != f"{after.repo}#{after.pr_number}" + or handoff.target_branch != after.branch): + raise ProducerRefusal("Exact accepted handoff target differs.") + identity = key([handoff.target_pr.lower(), handoff.expected_head]) + with ContextStore(root).locked(identity) as locked: + record = locked.read() + if (not isinstance(record, dict) or record.get("accepted") is not True + or record.get("launch_reserved") is not True or record.get("handoff") != handoff.as_dict() + or record.get("fingerprint") != key([handoff.as_dict(), record.get("source")]) + or record.get("writer_state") != "terminated"): + raise ProducerRefusal("Verified accepted terminated-source launch binding required.") + state = observe_lineage_source(record["source"], handoff) + return _delivery(Episode(sequence=sequence, repo=after.repo, pr_number=after.pr_number, + branch=after.branch, source_lane=handoff.source_lane, + destination_lane=handoff.destination_lane, expected_head=handoff.expected_head, + resulting_head=after.head_sha, writer_state=state), observed.writer, + observed.round_id, observed.transport) diff --git a/src/code_mower/package_manifest.py b/src/code_mower/package_manifest.py index 0cdc6c52..c655aee8 100644 --- a/src/code_mower/package_manifest.py +++ b/src/code_mower/package_manifest.py @@ -9,6 +9,11 @@ DEFAULT_PACKAGE_CONFIG = "code-mower.example.yml" PACKAGE_FILES = ( + ("src/code_mower/builder_lineage_producer.py", "src/code_mower/builder_lineage_producer.py", "core"), + ("templates/workflows/builder-lineage-producer.yml.j2", "templates/workflows/builder-lineage-producer.yml.j2", "template"), + ("src/code_mower/templates/workflows/builder-lineage-producer.yml.j2", "src/code_mower/templates/workflows/builder-lineage-producer.yml.j2", "template"), + ("templates/lanes/lineage-producer.sh", "templates/lanes/lineage-producer.sh", "template"), + ("src/code_mower/templates/lanes/lineage-producer.sh", "src/code_mower/templates/lanes/lineage-producer.sh", "template"), ("src/code_mower/builder_lineage.py", "src/code_mower/builder_lineage.py", "core"), ("tools/builder_lineage.py", "tools/builder_lineage.py", "core"), ("tools/CODE_MOWER_APACHE_LICENSE.txt", "LICENSE", "package"), diff --git a/src/code_mower/templates/lanes/lineage-producer.sh b/src/code_mower/templates/lanes/lineage-producer.sh new file mode 100644 index 00000000..a913afc7 --- /dev/null +++ b/src/code_mower/templates/lanes/lineage-producer.sh @@ -0,0 +1,12 @@ +#!/usr/bin/env bash +# Independent staged primitive. Not called by a maintained runner or normal init. +# Requires a locally installed candidate and reviewed broker policy/authority/ +# actual transport inputs. Fresh exact target/branch/labels and bounded public +# history are read by the real API before an attribution artifact is written. +# Publication and label reconciliation remain separate explicit Python calls. +set -euo pipefail +python - <<'PY' +import os +from code_mower.builder_lineage_producer import staged_record +staged_record(os.environ) +PY diff --git a/src/code_mower/templates/workflows/builder-lineage-producer.yml.j2 b/src/code_mower/templates/workflows/builder-lineage-producer.yml.j2 new file mode 100644 index 00000000..1db32710 --- /dev/null +++ b/src/code_mower/templates/workflows/builder-lineage-producer.yml.j2 @@ -0,0 +1,44 @@ +# Staged only: explicitly materialize against an installed candidate package. +# Normal init does not emit this workflow. #992 owns maintained integration; +# #915 owns qualification of the eventual published package (existing pins stay). +# The trusted caller supplies reviewed immutable-base policy, authority accounts, +# and actual observed transport. No PR checkout, code or configuration execution. +name: Staged lineage attribution +on: + workflow_call: + inputs: + target_json: + required: true + type: string + policy_json: + required: true + type: string + authority_json: + required: true + type: string + transport_json: + required: true + type: string +permissions: + contents: read + pull-requests: read + issues: read +jobs: + attribute: + runs-on: ubuntu-latest + steps: + - name: Retrieve complete lineage and record attribution + env: + GH_TOKEN: ${{ github.token }} + LINEAGE_TARGET_JSON: ${{ inputs.target_json }} + LINEAGE_POLICY_JSON: ${{ inputs.policy_json }} + LINEAGE_AUTHORITY_JSON: ${{ inputs.authority_json }} + LINEAGE_TRANSPORT_JSON: ${{ inputs.transport_json }} + LINEAGE_OUTPUT: lineage-attribution.json + run: | + set -euo pipefail + python - <<'PY' + import os + from code_mower.builder_lineage_producer import staged_record + staged_record(os.environ) + PY diff --git a/templates/lanes/lineage-producer.sh b/templates/lanes/lineage-producer.sh new file mode 100644 index 00000000..a913afc7 --- /dev/null +++ b/templates/lanes/lineage-producer.sh @@ -0,0 +1,12 @@ +#!/usr/bin/env bash +# Independent staged primitive. Not called by a maintained runner or normal init. +# Requires a locally installed candidate and reviewed broker policy/authority/ +# actual transport inputs. Fresh exact target/branch/labels and bounded public +# history are read by the real API before an attribution artifact is written. +# Publication and label reconciliation remain separate explicit Python calls. +set -euo pipefail +python - <<'PY' +import os +from code_mower.builder_lineage_producer import staged_record +staged_record(os.environ) +PY diff --git a/templates/workflows/builder-lineage-producer.yml.j2 b/templates/workflows/builder-lineage-producer.yml.j2 new file mode 100644 index 00000000..1db32710 --- /dev/null +++ b/templates/workflows/builder-lineage-producer.yml.j2 @@ -0,0 +1,44 @@ +# Staged only: explicitly materialize against an installed candidate package. +# Normal init does not emit this workflow. #992 owns maintained integration; +# #915 owns qualification of the eventual published package (existing pins stay). +# The trusted caller supplies reviewed immutable-base policy, authority accounts, +# and actual observed transport. No PR checkout, code or configuration execution. +name: Staged lineage attribution +on: + workflow_call: + inputs: + target_json: + required: true + type: string + policy_json: + required: true + type: string + authority_json: + required: true + type: string + transport_json: + required: true + type: string +permissions: + contents: read + pull-requests: read + issues: read +jobs: + attribute: + runs-on: ubuntu-latest + steps: + - name: Retrieve complete lineage and record attribution + env: + GH_TOKEN: ${{ github.token }} + LINEAGE_TARGET_JSON: ${{ inputs.target_json }} + LINEAGE_POLICY_JSON: ${{ inputs.policy_json }} + LINEAGE_AUTHORITY_JSON: ${{ inputs.authority_json }} + LINEAGE_TRANSPORT_JSON: ${{ inputs.transport_json }} + LINEAGE_OUTPUT: lineage-attribution.json + run: | + set -euo pipefail + python - <<'PY' + import os + from code_mower.builder_lineage_producer import staged_record + staged_record(os.environ) + PY diff --git a/tests/lineage_producer_fixtures.py b/tests/lineage_producer_fixtures.py new file mode 100644 index 00000000..dbeb3e4e --- /dev/null +++ b/tests/lineage_producer_fixtures.py @@ -0,0 +1,291 @@ +"""Finite external-I/O fixtures for the staged producer's real contracts.""" +from contextlib import contextmanager +from copy import deepcopy +from pathlib import Path + +from code_mower.builder_lineage import Authorities, Chain, Episode, History, Identity, Target, render +from code_mower.builder_lineage_producer import Snapshot, Transport + +REPO = "owner/repo" +BRANCH = "codex/Topic" +AUTHORITY = Authorities(["lineage-publisher[bot]"]) +TRANSPORT = Transport("claude", "claude", "claude_cli", "local_cli") +POLICY = Identity.from_mapping({"enabled": True, + "authors": {"source-bot": "codex"}, + "labels": {"builder:codex": "codex", "builder:claude": "claude", "builder:devin": "devin"}, + "branch_prefixes": {"codex/": "codex"}, "require_verified_lineage": True}) + + +def sha(n): + return f"{n:040x}" + + +def target(n=1, **changes): + return Target(**(dict(repo=REPO, pr_number=42, branch=BRANCH, head_sha=sha(n)) | changes)) + + +def episode(n=1, **changes): + return Episode(**(dict(sequence=n, repo=REPO, pr_number=42, branch=BRANCH, + source_lane="codex" if n == 1 else "claude", destination_lane="claude", + expected_head=sha(n-1), resulting_head=sha(n), writer_state="terminated" if n == 1 else "same_writer", + kind="handoff" if n == 1 else "continuation") | changes)) + + +def comments(episodes, *, author="lineage-publisher[bot]"): + episodes = list(episodes) + bound = Chain.from_arrivals(target(episodes[-1].sequence), episodes) + return [{"user": {"login": author}, "body": render(bound)}] + + +def observation_args(n=1): + return dict(target=target(n), identity=POLICY, authorities=AUTHORITY, + history=History([]), private=[episode(i) for i in range(1, n+1)], + author="source-bot", labels=["builder:codex"]) + + +class MemoryStore: + """Stub only private storage I/O; all producer validation remains real.""" + records = None + effects = None + + def __init__(self, root): + self.root = str(root) + + @contextmanager + def locked(self, name): + store = self + key = (self.root, name) + class Locked: + def read(self): + return deepcopy(store.records.get(key)) + def write(self, value): + store.effects.append("write") + store.records[key] = deepcopy(value) + yield Locked() + + +class GitHubIO: + def __init__(self, n=1): + self.target = target(n) + self.current_labels = ("builder:codex", "keep") + self.public = [] + self.effects = [] + self.snapshots = 0 + self.reads = 0 + self.fail_snapshot = None + self.readback = None + self.fail_post = False + + def snapshot(self, requested): + self.snapshots += 1 + self.effects.append("snapshot") + if self.snapshots == self.fail_snapshot: + raise OSError("unreadable") + return Snapshot(self.target, "source-bot", self.current_labels) + + def history(self, requested): + self.reads += 1 + self.effects.append("history") + value = self.public if self.reads == 1 or self.readback is None else self.readback + return History(deepcopy(value)) + + def post(self, requested, body): + self.effects.append("post") + if self.fail_post: + raise OSError("unavailable") + self.public.append({"user": {"login": "lineage-publisher[bot]"}, "body": body}) + + def labels(self, requested, desired, remove, add): + self.effects.append("labels") + self.current_labels = tuple(s for s in self.current_labels if s not in remove) + if add: + self.current_labels += (desired,) + + +def round_fixture(root, n=1, writer="destination-writer", config=None, runtime="ready"): + from code_mower.lane_delivery import LineageRound + return LineageRound(Path(root), f"round-{n}", writer, target(n-1), TRANSPORT, + Path(root) / "checkout", config=config or {}, runtime_observation=lambda: runtime) + + +# Complete approved baseline from accepted commit e818a3b639dfe903bdc16aff3674af98a5a08233. +# Expectations are fixed independently of the candidate and available without Git history. +ACCEPTED_BASELINE = {'accepted_base': 'e818a3b639dfe903bdc16aff3674af98a5a08233', + 'modules': {'src/code_mower/builder_runs.py': {'definitions': {'BuilderInference': '4324ee0dd83a3091bcb788b89402e26337294acf9c735b55e46292ddfa6994c0', + 'PullRequestMetadata': '9026b80175ad20604b178df969a1b4afbc60fdb8f2d8f3a710d43d064870943f', + '_author_from_pr_payload': '26b8088bad8733d3822c8b27c979c061ba346f660cbf5984cfc68d3bdc345d64', + '_branch_from_pr_payload': 'fc07795121bfe64ac2b43a2b3cef32f9711575bf141c9ee6125dbfbd144f0e73', + '_builder_id_from_url': 'a2011e146f57eba697c32767b78bf7b14c8acd604dc27e465278fa23ed3c74c3', + '_builder_run_event_id': 'cf9409f40c62938b9c83c078f00f8a6373893825e02f00cf4d44746083822662', + '_cursor_agent_url': '2503edb4bee610da18b532fd7b4f358688e948d1fc2fabbffaabe01b7db2c73b', + '_default_output_path': '3942af5edc3c310369bbd998ac1f44ac17ddcb216d8083537a8d7ea73f69242d', + '_github_url': 'a67031f20d1ec25e56f6c9acf80ad637575067b804f2c7d05014aaf1e55a0f6d', + '_issue_url_from_ref': 'a2f14c33e838993294bc530c3bf7cc971d608efd71151b9e6af22298adb29f21', + '_load_json_object': 'de0bbf2032377f3f38cc1d460f4cd2315b47835fa20e1abde0e2bfa7206aab9b', + '_load_work_order_manifest': '1cf47e66e523ff7f52176b25a26c36d2c5e3e31b4c78fbc16640953345752142', + '_nested_text': 'fc40a445e4633ff5c37d12fb59afb7c9699560305ac7266f23d5237d89946ef4', + '_optional_nonnegative_float': '67db0dae03045ff109a2b7c7bf3170ad878958e4f41cc1734589086ecedb5d9e', + '_optional_nonnegative_int': '63f7b5c1332d17c1ea01982beba03a619b54bcea509f87568843627bab3b9239', + '_pr_url_from_ref': 'a2848bfc59ad175de26bd01db8dd34ee6aa605bcac2ed3f810f32410ee85055e', + '_record': '9e7ebe8763749508bbfa59b51dfe939fb4a810090ad93164a4bf9d08bd7cdeda', + '_repo_from_pr_payload': 'c9902153ecf27a5a7d6e5d7205c4c3f8f6cc279977ad0d11f82a8cabc13b3be3', + '_safe_slug': 'be862cd8926de03ae396ec2645470cac12fd91cf5ede25d492beb727e39d3c9d', + '_text': '9c73340bbc252d03137d6a83b18929b062f8e9304ef884007203a586437a511d', + '_utc_now': '043afb904b5bf3eb34f0c440f38ec6c16a15c650abcc43d5437ff54de04f3c70', + '_validate_anchor_repositories': 'b9907f68c7fcce6b1506b1c69b70b3da913a0c01484f19e89e57481cb3dd3bb3', + '_work_order_manifest_path': '11533340fb4efe2990fdfcb3515f14422f2af39b94f030415a1c23f02b8dced5', + 'build_auto_builder_run_event': '14182e08976610bb9a0ada12748b0bc865a09306fc2ac458f5f538449dece549', + 'build_builder_run_event': '1b1a4415bc0ad169ca12a572e795204a0dbb98cb3a54ddbd7867ced0adff7960', + 'infer_builder_from_pr': '4b0af12c30970e2dc376af44c4ec85e37c2a1e9c3fb486b6d2f3972de4290f8b', + 'load_pull_request_metadata': 'b53bc8d0408e09bbf8ef323a78b4a1276ad5c4a8a9190c51f62284a15f25b3f3', + 'main': '903e3edd24c6c19af1e29b8b7e9c501e501634fbaaed299d9cce57f68de08cdc', + 'write_builder_run_event': '9ea66af8043dbf80be2b506d9d9033521dccd50c6aef540ac1ca2e984e780456'}, + 'live_segments': ['10a56c918cb5194f3ae29ab8e3e31629b32c48984194a7a2793f5ea18333131a', + '5384bfdb2df380b6557cc7a71d16891415bccaa87699406e236f752c6415389f', + 'b53e5a5286172fb0a4772843ba3ed62331687bcb9939ec9c49e9066d4ee4272d', + '261a9112a2918e2aba03e19a2f221a027ab4b6f9e946a6c39814f54e18e07dfc', + '085ea91d7cfeb4f5695793f882c7463cbfdf7be9f98e6c5133847a245bbea366', + 'a830333312c4c9c2233bb02762bd498203b1c3bf0527614735693b4b79f6d167', + 'c517577851c489e45abae2591256c40404a05c6d19cd4d5ae7fd22b0084cec6c', + 'ab270e34de38d99f0364043ca1808983f498a9e3943133071931657340b5bc95', + '07ac27eea41ef7e0d14eeb4cc65ebf0beaa2041ad913826ebdf97d56d6febd5f', + 'bd3698cb9f306b035e45c51278a2dfff864360ff5b22ba137ee373b0821704bb', + '3cb9ffccba5146b2dd177a30c3057a4395ebd5c745d0789d2ed95d0cf4afe271', + '2f1f76a50065a201433318c3c80dbf0dda6cbd622c81283b7b13adf1508ea508', + 'ce2857e54587b6c88d98f7cd8f2a4692258523d77f499e9bef9909ff01190bb9', + '3a80dd6bb6886183e9072113b5be53d1ab4cc0e39ca9c4942c8f8f6c9d4c281e', + '93241e739b270886462ecc611e5dd2f20b260a5e9f0bd2a4b4e16bbd28c17d3a', + '858f7a09afbc43320b338e796579105689fc31b316f520574233d4af2fcb56f1', + 'fa869652896f7b59d67345f3e1aaaafd2152990364affe8a6abcfa211bb5f61b', + '53c376ef4ffb635cddeb6a148228838a0f4c7361e5885e9a98cf6366dd7170d4', + '423378d6bfd7535019d2730f8f091bf150fb75a6c05604a491cf48220c29b615', + 'af709bc9973f5e559c83abfc7ea6cf05c42ee4b871e938b8b5d0268bf03cf912', + '01617a4661f2cf67a5ad98e4beb386bd60f2233229331f51a490d770d8be7c5f', + '9d02f33e30a4947232fe0a447aabfd64ec9bde14b66d417c33226e0a6e5d3fb4', + 'e36aa3d0f4dbc062544aa00a56526b372eed03f467649dcfab0a46ba2ac1b05d']}, + 'src/code_mower/lane_delivery.py': {'definitions': {'DeliveryOutcome': '2c96fc216f2a371bdbdb15c917cd22542926268f8671d212a864031f87cb033c', + 'Handoff': '8dd7c68eca42cf0b33594059a1541070ae388e981290cd415808d57db249446d', + 'LaneDeliveryError': 'dde8efaa1162bda26f4c859f209e5eabd9f6f726ca0c10375bd539056faeeddb', + 'SupervisionResult': '8210ccfe373b85d4073c4790738eec62443d0487f35911ae1630ca11ac460c22', + 'TargetState': '1fda306e1efd7a44a87ed85da8d156eff5cad98f8a787c3f2e506579278f0198', + '_add_classify_parser': 'e1bcfdbd593639bc0a9dccd57300c874badd17d31e0934c748c2dee0b3539e80', + '_add_handoff_parser': '6c8d95a4a056651756cf9ceddf130f005b094ff6e38062c0c052d868bfa6ae03', + '_add_scan_prompt_parser': '570e8c30d178cf5ef5848310f4874ad3141d8f599e6ac98b9b552373fcc04264', + '_add_supervise_parser': '0ac827732ff11b217af15c15a19eec79e09f5efaaa76c3dcc4e7bb9166e824d5', + '_add_transition_parser': 'bbbab20d8d66eebb70a192b29e1f8f93f5245de1f61d5714cefbc78e4e39822d', + '_admit_builder_main': '2f51b690ce6d9c034c9c5ebf8fcbc52b35fc3f50b5464e6643afa29a9292c953', + '_assert_safe_metadata': '3fde171274a56d9ec5b1e4300c0b438eed1789bbcb121ef006738a956ce08594', + '_classify_main': 'fc679c5781e0c511ab2008e5fb55c5ecbf2634338a6547eccd87f2ab92fd767c', + '_default_is_group_alive': '9fefc020a14093b652085bcad0540ef8fcdb6edd3c5212b94a086de96b9745f0', + '_handoff_main': '9c149fabcb1009a3eee27b012409ee511eed0e23ee69b6acd85aaec51ed22f37', + '_load_state': '7e46824541822254b9d0f68121588f7326a8110ffcad8fbed7e8ad6a399edf7f', + '_scan_prompt_main': '8b9809c22d3743968b1711f014437249c3f0df3d35d7d23d21f0170144ddfe20', + '_supervise_main': '96e3729a167880c32e8bfd27d14ac6c15827c75f5012ff36719e5dc8d55cb192', + '_text': 'a8dee85fc038bdb336cd6dfd65f4656639b268af650315ec55a70c9cd2cbe86f', + '_transition_main': '28eb83382833cbd36c42819d5837958dcfbea0fe539395e039fa344de8a08c87', + '_utc_now': '043afb904b5bf3eb34f0c440f38ec6c16a15c650abcc43d5437ff54de04f3c70', + '_wait_for_group_exit': '6ee182644d242a62e9c685c15669b6734c2c9d98e11b65af056638b3266594c2', + 'assert_prompt_free_of_auth_material': 'a33ab3d2c9ff868aa5475f0a22257e216144a8622f3ae4b750f51041afa5fd5d', + 'authorize_branch_write': '24eb154c9648874df6e7edcc8199f15d732986e8082bf68559cc3a5b691694b4', + 'build_delivery_outcome_event': 'f8dd3431a9032d989d4fb495970639e490f092ec938a6482d783236f37e7db9d', + 'classify_delivery': '2c62d9e055b7dd21b5ecb3a7999964a3e996647f5c663d33042b478c1c399ee4', + 'main': '98705bde1e36af092c23a994e56819fc2e27bf92a683542943d5250fd3141473', + 'observed_transition': '2499b7f66ccfcd6105295a07f89d8cabd16ea5613b0471b10fc05dcc95dd729a', + 'scan_auth_material': '7e9b06358faa8ffc2234ece781396470d574212caadab614255c7197addf7fad', + 'supervise_process': 'fd369e456a5fa3654921b8c0e4069a35e1615e6a5cac35db1fa872adfcdec3b0', + 'terminate_process_group': '87b16c7b8b36ada7ff9da280152c265cfb5eba2ab79f693d5f47f89cc50c0984', + 'validate_handoff': '92dc1004123eccbec9ff05022c8d0150723e2922ebc9c60ab2d4398a93293a33', + 'write_delivery_outcome_event': 'ec2def82c067d57851d49bd55671e901e0561805a5582bf0e18169eced5ad0c5'}, + 'live_segments': ['0020ee2273c34f19d33ae5dd28ef6c69ddc71525aa5b1778b231255a431cd99a', + '5384bfdb2df380b6557cc7a71d16891415bccaa87699406e236f752c6415389f', + 'b53e5a5286172fb0a4772843ba3ed62331687bcb9939ec9c49e9066d4ee4272d', + 'b17c41ffb3939b0dc71fa9e725ac38888b83de0c14b844fb813deda3b5410663', + '261a9112a2918e2aba03e19a2f221a027ab4b6f9e946a6c39814f54e18e07dfc', + '3727adff524e0616022eadd8f4af21a0778b29fc4c77bdfefd1afce2cbf5e4b7', + 'a830333312c4c9c2233bb02762bd498203b1c3bf0527614735693b4b79f6d167', + '9a60319879ee88ae2b4a273d8f45962bc499df0b50ae75bfd829a71b0e51352f', + '5890cea262fd597e61cf1a114b65d2781c700c040e9bb53cd8696cf36830134d', + '1fd20c880d3cde1a258aa761090aabc445918b6cf2b91712928debfc04392cff', + '17ea74bb79ba96edc39a8382a97f22993b80cc4088b3e2717a98917c60c1dd5a', + 'c517577851c489e45abae2591256c40404a05c6d19cd4d5ae7fd22b0084cec6c', + '7b2ddb8a2e64c6166feedaf30ee50ac62165fe149fdd8b0d7b777ebee21ca34d', + 'ab270e34de38d99f0364043ca1808983f498a9e3943133071931657340b5bc95', + '07ac27eea41ef7e0d14eeb4cc65ebf0beaa2041ad913826ebdf97d56d6febd5f', + 'bd3698cb9f306b035e45c51278a2dfff864360ff5b22ba137ee373b0821704bb', + '3cb9ffccba5146b2dd177a30c3057a4395ebd5c745d0789d2ed95d0cf4afe271', + '7dd851d67f9f2afe2ccf20745c843c644449d1d0812f61a8c6d594c0aab33141', + '3a80dd6bb6886183e9072113b5be53d1ab4cc0e39ca9c4942c8f8f6c9d4c281e', + 'f0e1023d365a836dec3a903b3c22beee69b12f277423169358875d2245d13930', + '1a3daadc49a917e44aa47bdb61879c43943fc42db1871d06746bef0bd1e3ac2a', + '64067ade48b8909b78c82006819e87d5baab6c651372788617ecd533b1cd5ecb', + 'cd2ceb53d0bf8503f0791c84bee5693d4af6eb3426720739eff6dddfd9c7850f', + 'b3503e44b72d0f15293369506c0fd7b8c9b98cd25c6def90cd53ffc3e3977d69', + '5c4dd73916a9da18db978b468f8c515f768a82219d31605750280267519c2c2e', + '8253cdd8ea8b02ef2cba8c7beb280f21787bffc32028891d325ed0036c2a9640', + '04d095769a2a3affb49fa3a4c5b1c6c1879025a0ae6de6e12d04262d66bbabf8', + 'd3397f8ab3ff6d70d65ccf2d7f20a96e129849ac2269405bb3f6ab893f1891a3', + '65a1feda3e0ab3d45dcef88c348cec91ad2a89121892ea70287fb705d19b7314', + 'a9742388e70d4355d6e8455df8845d630e732fd36ad6446b9f51d2753bce3133', + '0f703b7341f2ed1b4e741a81c4240c975db9f316bd77f9b98cf923db5a214cca', + '734a29e70b5d95d80c1a58acab6df0b40d2f727129a2f32447bdac5db7424773', + '93968c5a741e0cab1ee23ec37a18f404ed1c83d979acfe9af0533466e32b5023', + '1c0c488917fd12d167958a913e4bee0d56004b897c7a69086cb8814e1853b78f', + 'b98c9b5b23af7449c14b4d85a45afb6584c1c9c2476128de773a8bd8f883fe97', + 'bfd622bd0ae48a943d072b9db94e114ef4c610783e998a0f5ab1efa9bd644814', + '2ae2071c9eaa29e6e8ba931ed8b9525600d21ab5d1c9b53824cad46a6a6e12fe', + '74b2efb213e46f61255c4707b931ed001c02fba7e5ba9177daa413cc449fd861', + '234ea0928540c34115f9a19a066aab87f22d15283ff94479a6c0ee953a799167', + 'f042f194945c8e5a11bb6ce2bbefd0020fa27e3463ac3f25db881ef74de0a557', + '530bbbf2290f92d40ff7446f2ea218edf767f7abc4c4a35a932669c880a1bea3', + 'c06d442a93bcc4f8ded5e06af83b5b235acf78bc1ef717f10d37f4b3ac18a62c', + '3e84c8a8e96c15e0defda059187db3ddb4275c4aecdb290ecc5d3088c7ef39fa', + '630f2eae48d515b5cf6ef5698b9f761efd92f2708a5309af84de4fa18b5748a1', + 'f41d0696020b1d1491d138e438cbf925261b30664c7ebe43e9b88e708aef6e95', + 'ab50daa88ec3d92f8cd09bbc1f1614637c4aa38f8439b6e131936e3b543fe399', + 'd07c84c155578873bff7f99a3fd5efb2a074b4c5fbceff6794c534d7690acc75', + '1a6a66f7b64c03a6e73a3e73cb4e9df7f2f85b8546a264832a48753f0f828882', + '58d540ae8fd412140d5e2489009e5dd7201b1d8b3b9d2a1c0c29af8db520e6b9', + '802868d81e7ad72d47aaa02360f865e734503918b4782ebabca2acf67610c252', + '5732ba765b9ea7d1b57dea54d73ba8bff19dc889b67f82b1f215e3955ad9664a', + 'd7ff743c92d70faacdba505ab49aa1e350798822d58ca3a5caa252c2713d77d6', + 'bf11c0d84460af33a13b0080f7d1aded574f5c3daa492d7df05bb9d3cd9a1239', + 'e36aa3d0f4dbc062544aa00a56526b372eed03f467649dcfab0a46ba2ac1b05d']}, + 'src/code_mower/lane_handoff.py': {'definitions': {'LocalWriter': 'ceaf03283017c1681a1e3cb2513de83587a73c737c202da142b1a1fdaa672ab6', + 'default_root': '94bcc017d44dc68fa7cc42a304352cefd884fe6d9c9b8d85f6c111aa3a498025', + 'key': '9ac80299cff2ae7814fe5f0f78100829b2110d190941d7b0c338df79d8a2c2a3', + 'observe_head': 'bdfb7a540f7913b6a99bf5ee6bb077c0185b04d7a408527c1e75bacf69848c06', + 'prepare': '631376ac3f53bc69c55cc4cc07eb805ecb778a35ccd3089378482ebd230288dd', + 'quiesce': '83a53943cde81c6745ce65463266df73f00402c9a794d172eef31dc10507ddd7', + 'read_source': '6b94b4c6193e4f65e04730fbc63f7079e8f81364dcfdc37b588d48eeaa32481e', + 'remote_engine': 'fe4c7d5fdcd24bca72e975e49609d9aba229ae20ea4d0d99db76695ab4823175', + 'reserve_launch': 'a209248ba74f3e8e916fc065a77a35bf35b8f2eee96938c2ede512d9209abf7d'}, + 'live_segments': ['f889e85620eec8de2c238812dd247410d2e75f8d5001f4eecdb79b3865165c21', + '5384bfdb2df380b6557cc7a71d16891415bccaa87699406e236f752c6415389f', + '3b3ea3cdc7068a9e97cfcddbcd9d8e0c54fe4d068347b80b158fdd8b14676e3b', + '261a9112a2918e2aba03e19a2f221a027ab4b6f9e946a6c39814f54e18e07dfc', + '3727adff524e0616022eadd8f4af21a0778b29fc4c77bdfefd1afce2cbf5e4b7', + '1fd20c880d3cde1a258aa761090aabc445918b6cf2b91712928debfc04392cff', + '17ea74bb79ba96edc39a8382a97f22993b80cc4088b3e2717a98917c60c1dd5a', + '7b2ddb8a2e64c6166feedaf30ee50ac62165fe149fdd8b0d7b777ebee21ca34d', + '3cb9ffccba5146b2dd177a30c3057a4395ebd5c745d0789d2ed95d0cf4afe271', + 'd30710bca9af95f1af88e54567bf9cb3571abb4430759486de0f61753166b6a4', + '14cef7b5ee10ef0b5dd80bed58c2b59a196568aea24768a9a105b5bbc0817de6', + '00ee41c0a680c0bafc712cebfd45a249372b727fb6ff6c25ae75d1d01f9c1656', + '7f0d608f78e3f586988f58d336cfabb575f9cd24462c6dbc9197c8c98213ba6f']}}, + 'unchanged_files': {'.github/workflows/ci.yml': 'df0083c724ae7007528951557363ede1adbe6bce9ddce48dc753bc3751f2e581', + '.github/workflows/claude-audit-labeler.yml': 'ac21c697ddf2a134f4e89cc6738cbaeb770f6ea76abf76823edd409ed0618e5e', + '.github/workflows/claude-clear-stale.yml': '7381379101486bc9ee97805b071301f01eb4b2655ca796cfde84110ea3784b58', + '.github/workflows/cloud-dogfood.yml': '9fe6b19a1214449ea8537efed12cfc6c6146c8e2f11a407ca782030e9b903b95', + '.github/workflows/code-mower-agent-pr-labeler.yml': '0555ccf6d8896cadc60b271725e0cdd8f3005ff2fa4eaf231b0fd23451e37545', + '.github/workflows/code-mower-fix-round-dispatch.yml': '304c952b3f5a6ed4373a091797a53877cc1e665c7b9548606d26150231190dd7', + '.github/workflows/code-mower-gate.yml': 'd6ed0b9a7a37ca329868c0f7f634c915d500233250ba6ffbf921a72ac4066780', + '.github/workflows/codex-audit-labeler.yml': '306b07779111d24cbe200c99fe5309978707299dbe2c2f56cbedbad640da984e', + '.github/workflows/codex-clear-stale.yml': 'fe832a2b795d70d4eab8bbcd89bfc069104d6f80854b86652e2aa397f2b489ca', + '.github/workflows/dispatch-lanes.yml': 'cf6d9997cfbe9468e0883497ba68f223196332bf6fc99f4b6b8e8b1068cabdcd', + '.github/workflows/lane-mac-runner.yml': 'd771366c6e2ce3afc30a4919abd168b70e4a83fd19906e8ba534fed58ad606a7', + '.github/workflows/local-cli-audit.yml': '5305326fd50cc183d1f3d960bac2cf96019ebad4bf4dd5663a73bad545e19ccc', + '.github/workflows/release.yml': '6a6b87b30c096ff46027e5b39b6e9f6cfb4c54aa6b42bbdedcc45c1f1cd8eb4b', + 'src/code_mower/init.py': '429d12d80f9c9b43dc65781ccd8ee65d7994a45c4652a0a2705c0fcf369e265b', + 'src/code_mower/templates/lanes/run_mac_lane.sh': '89139edb988d58104972c06ed7eeb5dc598804e001b8d9df2249ad4e4a101a5b', + 'templates/lanes/run_mac_lane.sh': '89139edb988d58104972c06ed7eeb5dc598804e001b8d9df2249ad4e4a101a5b', + 'tools/lanes/run_mac_lane.sh': '8b2475bde3c95b3fdb1d8f3c118a006364cc000a637c1a1f3ca80825c73b6506'}} diff --git a/tests/test_lane_delivery_contract.py b/tests/test_lane_delivery_contract.py index e1407b55..ced75282 100644 --- a/tests/test_lane_delivery_contract.py +++ b/tests/test_lane_delivery_contract.py @@ -2313,3 +2313,19 @@ def test_the_undelivered_note_names_an_unbrokered_declaration(self) -> None: if __name__ == "__main__": # pragma: no cover unittest.main() + + +class ExplicitLineageSnapshotTests(unittest.TestCase): + def test_exact_producer_snapshot_preserves_branch_case_and_legacy_defaults(self): + raw = {"kind": "pr", "number": "42", "pr_number": "42", "head_sha": "a" * 40, + "pr_state": "OPEN", "labels": [], "author": "source-bot", + "snapshot_complete": True, "branch": "codex/Topic"} + snapshot = lane_delivery.lineage_target_state("Owner/Repo", raw) + self.assertEqual(snapshot.target.branch, "codex/Topic") + self.assertEqual(snapshot.target.repo, "owner/repo") + legacy = lane_delivery.TargetState.from_mapping(raw) + self.assertNotIn("branch", legacy.as_dict()) + for change in ({"branch": None}, {"branch": " codex/Topic"}, {"labels": None}, + {"snapshot_complete": False}, {"head_sha": "bad"}, {"pr_number": True}): + with self.subTest(change=change), self.assertRaises(ValueError): + lane_delivery.lineage_target_state("owner/repo", raw | change) diff --git a/tests/test_lineage_producer_artifacts.py b/tests/test_lineage_producer_artifacts.py new file mode 100644 index 00000000..46a7bc4a --- /dev/null +++ b/tests/test_lineage_producer_artifacts.py @@ -0,0 +1,267 @@ +"""Execute staged artifacts against an isolated locally built/installed package.""" +import ast +import hashlib +import json +import os +from pathlib import Path +import shlex +import subprocess +import sys +import tempfile +import unittest +from urllib.parse import parse_qs, urlsplit + +import yaml +from code_mower import init, package +from code_mower.config import load_config +from lineage_producer_fixtures import ACCEPTED_BASELINE, AUTHORITY, POLICY, TRANSPORT, comments, episode, target + +ROOT = Path(__file__).resolve().parents[1] +BASE = "e818a3b639dfe903bdc16aff3674af98a5a08233" +CORE_HASH = "2842d106bc7c3ecb95007b7de6ff3447d8ea3a49e8995e422ff279463ae5653c" +ASSETS = ("workflows/builder-lineage-producer.yml.j2", "lanes/lineage-producer.sh") + + +class ArtifactTests(unittest.TestCase): + def accepted_baseline(self): + baseline = ACCEPTED_BASELINE + self.assertIsInstance(baseline, dict, 'Accepted baseline must be a complete mapping') + self.assertEqual(set(baseline), {'accepted_base', 'modules', 'unchanged_files'}, + 'Accepted baseline fields are missing or unsupported') + self.assertEqual(baseline['accepted_base'], BASE, 'Accepted baseline commit differs') + modules = baseline['modules'] + self.assertIsInstance(modules, dict, 'Accepted module baseline must be a mapping') + counts = {'src/code_mower/builder_runs.py': (29, 23), + 'src/code_mower/lane_delivery.py': (33, 54), + 'src/code_mower/lane_handoff.py': (9, 13)} + self.assertEqual(set(modules), set(counts), 'Accepted module inventory differs') + for path, (definition_count, live_count) in counts.items(): + module = modules[path] + self.assertIsInstance(module, dict, f'{path}: malformed accepted module baseline') + self.assertEqual(set(module), {'definitions', 'live_segments'}, + f'{path}: accepted source-segment fields differ') + self.assertIsInstance(module['definitions'], dict, f'{path}: definitions must be a mapping') + self.assertEqual(len(module['definitions']), definition_count, + f'{path}: accepted definitions are missing or extra') + self.assertIsInstance(module['live_segments'], list, f'{path}: live segments must be an ordered list') + self.assertEqual(len(module['live_segments']), live_count, + f'{path}: accepted live segments are missing or extra') + for name in module['definitions']: + self.assertTrue(isinstance(name, str) and name.isidentifier(), + f'{path}: malformed accepted definition name') + for digest in [*module['definitions'].values(), *module['live_segments']]: + self.assertIsInstance(digest, str, f'{path}: accepted source digest must be text') + self.assertRegex(digest, r'^[0-9a-f]{64}\Z', f'{path}: malformed accepted source digest') + files = baseline['unchanged_files'] + self.assertIsInstance(files, dict, 'Accepted unchanged-file baseline must be a mapping') + self.assertEqual(len(files), 17, 'Accepted unchanged-file inventory must contain all 17 files') + for path, digest in files.items(): + self.assertIsInstance(path, str, 'Accepted file path must be text') + self.assertIsInstance(digest, str, f'{path}: accepted file digest must be text') + self.assertRegex(digest, r'^[0-9a-f]{64}\Z', f'{path}: malformed accepted file digest') + serialized = (json.dumps(baseline, indent=2, sort_keys=True) + '\n').encode('utf-8') + self.assertEqual(hashlib.sha256(serialized).hexdigest(), + '6036e7daccf07b3ab5f458315e5e8a04feeb00190bae35c4fb1e243ef519e55e', + 'Complete accepted baseline differs from the independently approved value') + return baseline + + def module_source_hashes(self, path): + source = (ROOT/path).read_bytes().decode('utf-8') + lines = source.splitlines(keepends=True) + definitions, live_segments = {}, [] + for node in ast.parse(source, filename=path).body: + is_definition = isinstance(node, (ast.FunctionDef, ast.AsyncFunctionDef, ast.ClassDef)) + start = min([node.lineno, *(decorator.lineno for decorator in node.decorator_list)]) if is_definition else node.lineno + segment = ''.join(lines[start - 1:node.end_lineno]).encode('utf-8') + digest = hashlib.sha256(segment).hexdigest() + if is_definition: + self.assertNotIn(node.name, definitions, f'{path}: duplicate top-level definition {node.name}') + definitions[node.name] = digest + else: + live_segments.append(digest) + return definitions, live_segments + + @classmethod + def setUpClass(cls): + cls.tmp = tempfile.TemporaryDirectory() + cls.addClassCleanup(cls.tmp.cleanup) + cls.root = Path(cls.tmp.name) + cls.wheels = cls.root / "wheels" + cls.installed = cls.root / "installed" + result = subprocess.run([sys.executable, "-m", "pip", "wheel", "--no-deps", "--wheel-dir", str(cls.wheels), str(ROOT)], + cwd=ROOT, capture_output=True, text=True, timeout=120) + if result.returncode: + raise AssertionError(result.stdout + result.stderr) + wheel, = cls.wheels.glob("*.whl") + result = subprocess.run([sys.executable, "-m", "pip", "install", "--no-deps", "--no-compile", "--target", str(cls.installed), str(wheel)], + capture_output=True, text=True, timeout=60) + if result.returncode: + raise AssertionError(result.stdout + result.stderr) + cls.bin = cls.root / "bin" + cls.bin.mkdir() + bootstrap = f''' +import datetime +import sys +from pathlib import Path +from unittest.mock import patch +sys.path.insert(0, {str(cls.installed)!r}) +import code_mower.builder_lineage_producer as producer +assert Path(producer.__file__).resolve().is_relative_to(Path({str(cls.installed)!r})) +assert {str(ROOT / 'src')!r} not in sys.path +class Clock(datetime.datetime): + @classmethod + def now(cls, tz=None): + return cls(2026, 1, 1, tzinfo=tz) +with patch('datetime.datetime', Clock): + exec(compile(sys.stdin.read(), '', 'exec')) +''' + wrapper = cls.bin / "python" + wrapper.write_text("#!/bin/sh\nexec " + shlex.quote(sys.executable) + " -I -S -c " + shlex.quote(bootstrap) + "\n") + wrapper.chmod(0o755) + gh = cls.bin / "gh" + gh_code = cls.bin / "gh_fixture.py" + gh_code.write_text( ''' +import json, os, sys +from pathlib import Path +fixture = json.loads(Path(os.environ['FIXTURE']).read_text()) +endpoint = sys.argv[2] +with Path(os.environ['EFFECTS']).open('a') as sink: + sink.write(json.dumps(sys.argv[1:]) + '\\n') +assert sys.argv[1] == 'api' and len(sys.argv) == 3 +mode = fixture['mode'] +if '/pulls/' in endpoint: + raw = fixture['pr'] + if mode == 'target-failed': sys.exit(1) + if mode == 'target-null': raw = None + if mode == 'branch-missing': del raw['head']['ref'] + if mode == 'branch-case': raw['head']['ref'] = raw['head']['ref'].lower() + if mode == 'head-race': + reads = Path(os.environ['EFFECTS']).read_text().count('/pulls/') + if reads > 1: raw['head']['sha'] = 'f' * 40 +elif '/labels?' in endpoint: + raw = [{'name': 'builder:codex'}] + if mode == 'labels-failed': sys.exit(1) + if mode == 'labels-null': raw = None + if mode == 'labels-malformed': raw = [{'name': False}] +else: + assert '/comments?' in endpoint + page = int(endpoint.rsplit('page=', 1)[1]) + raw = fixture['comments'] if page == 1 else [] + if mode == 'history-failed': sys.exit(1) + if mode == 'history-null': raw = None + if mode == 'history-object': raw = {} + if mode == 'history-mixed': raw = [None, {'user': None}] + if mode == 'history-bad-body': raw = [{'body': None}] + if mode == 'page-cap': raw = [{}] * 100 +print(json.dumps(raw)) +''') + gh.write_text("#!/bin/sh\nexec " + shlex.quote(sys.executable) + " -I -S " + shlex.quote(str(gh_code)) + " \"$@\"\n") + gh.chmod(0o755) + + def materialized(self, relative): + source = self.installed / "code_mower/templates" / relative + entry = {"source": "explicit-staged-asset", "copy_from_path": str(source)} + result = init._materialize_generated_file(entry, relative, self.root / "rendered", + source_root=self.root) + self.assertFalse(result.placeholder) + return init._render_workflow_template(result.text, {}) + + def test_real_rendered_workflow_and_runner_failure_rows(self): + workflow = yaml.safe_load(self.materialized(ASSETS[0])) + self.assertEqual(set(workflow["permissions"].values()), {"read"}) + body = workflow["jobs"]["attribute"]["steps"][0]["run"] + scripts = [body, self.materialized(ASSETS[1])] + for index, script in enumerate(scripts): + for mode in ("success", "target-failed", "target-null", "branch-missing", "branch-case", + "labels-failed", "labels-null", "labels-malformed", "history-failed", "history-null", + "history-object", "history-mixed", "history-bad-body", "page-cap", "head-race", + "no-authority", "policy-denied"): + with self.subTest(artifact=index, mode=mode): + case = self.root / f"case-{index}-{mode}" + case.mkdir() + fixture = case / "fixture.json" + observed = target() + fixture.write_text(json.dumps({"mode": mode, "comments": comments([episode()]), + "pr": {"state": "open", "number": 42, "base": {"repo": {"full_name": observed.repo}}, + "head": {"ref": observed.branch, "sha": observed.head_sha}, "user": {"login": "source-bot"}}})) + roles = {"role_policy": {"claude": {"builder": {"enabled": False}}}} if mode == "policy-denied" else {} + env = os.environ | {"PATH": str(self.bin) + os.pathsep + os.environ["PATH"], + "FIXTURE": str(fixture), "EFFECTS": str(case / "effects.jsonl"), + "LINEAGE_TARGET_JSON": json.dumps(observed.__dict__) if hasattr(observed, '__dict__') else json.dumps( + {"repo": observed.repo, "pr_number": observed.pr_number, "branch": observed.branch, "head_sha": observed.head_sha}), + "LINEAGE_POLICY_JSON": json.dumps({"base_sha": BASE, "identity": POLICY.to_mapping(), "roles": roles}), + "LINEAGE_AUTHORITY_JSON": json.dumps([] if mode == "no-authority" else sorted(AUTHORITY.accounts)), + "LINEAGE_TRANSPORT_JSON": json.dumps(TRANSPORT.__dict__), "LINEAGE_OUTPUT": str(case/"event.json")} + result = subprocess.run(["bash", "-c", script], cwd=case, env=env, + capture_output=True, text=True, timeout=30) + self.assertEqual(result.returncode == 0, mode == "success", result.stderr) + self.assertEqual((case/"event.json").exists(), mode == "success") + if mode == "success": + event = json.loads((case/"event.json").read_text()) + self.assertEqual(event["dimensions"]["lineage"]["current_writer"], "claude") + self.assertEqual(event["tool"]["executor"], "claude_cli") + self.assertEqual(event["created_at"], "2026-01-01T00:00:00+00:00") + if mode == "page-cap": + effects = [json.loads(line) for line in (case/"effects.jsonl").read_text().splitlines()] + queries = [parse_qs(urlsplit(args[1]).query) for args in effects + if args[0] == 'api' and urlsplit(args[1]).path.endswith('/comments')] + requested_pages = [query['page'] for query in queries] + self.assertEqual(requested_pages, [[str(page)] for page in range(1, 10)]) + self.assertNotIn(['10'], requested_pages) + + def test_manifest_mirror_core_and_existing_definition_parity(self): + for asset in ASSETS: + self.assertEqual((ROOT/"templates"/asset).read_bytes(), (ROOT/"src/code_mower/templates"/asset).read_bytes()) + self.assertEqual((ROOT/"templates"/asset).read_bytes(), (self.installed/"code_mower/templates"/asset).read_bytes()) + for path in ("src/code_mower/builder_lineage.py", "tools/builder_lineage.py"): + self.assertEqual(hashlib.sha256((ROOT/path).read_bytes()).hexdigest(), CORE_HASH) + for path, baseline in self.accepted_baseline()['modules'].items(): + definitions, live_segments = self.module_source_hashes(path) + self.assertFalse(set(baseline['definitions']) - set(definitions), + f'{path}: missing accepted definitions: {sorted(set(baseline["definitions"]) - set(definitions))}') + for name, digest in baseline['definitions'].items(): + self.assertEqual(definitions[name], digest, f'{path}: accepted definition source changed: {name}') + self.assertEqual(live_segments, baseline['live_segments'], + f'{path}: ordered module-level live source changed') + expected = package.committed_package_manifest_text(package.generate_committed_package_manifest(ROOT)) + self.assertEqual(expected, (ROOT/"code-mower-package-manifest.json").read_text()) + + def test_default_init_emitted_helper_isolated_without_producer(self): + config = load_config(ROOT / "src/code_mower/templates/code-mower.example.yml") + plan = init.render_init_plan(config, package_mode=True, repo_root=ROOT) + output = self.root / "default-init" + init.apply_init_plan(plan, output, source_root=ROOT) + all_paths = [p.relative_to(output).as_posix() for p in output.rglob('*') if p.is_file()] + self.assertFalse(any('lineage-producer' in p or 'builder_lineage_producer' in p for p in all_paths)) + # Existing live init/runner/workflow inputs are byte-identical to accepted base. + paths = ['src/code_mower/init.py', 'tools/lanes/run_mac_lane.sh', + 'templates/lanes/run_mac_lane.sh', 'src/code_mower/templates/lanes/run_mac_lane.sh'] + paths.extend(p.relative_to(ROOT).as_posix() for p in (ROOT/'.github/workflows').glob('*')) + baseline = self.accepted_baseline()['unchanged_files'] + self.assertCountEqual(paths, baseline, 'Actual init/runner/workflow inventory differs from accepted baseline') + for path in paths: + self.assertEqual(hashlib.sha256((ROOT/path).read_bytes()).hexdigest(), baseline[path], + f'{path}: unchanged accepted file bytes differ') + # Normal init emits the pure tools helper, not the package delivery modules. + helper = output/'tools/builder_lineage.py' + self.assertTrue(helper.is_file()) + program = ''' +import importlib.util +import sys +from pathlib import Path +assert importlib.util.find_spec('code_mower') is None +emitted = Path.cwd()/'tools/builder_lineage.py' +sys.path.insert(0, str(emitted.parent)) +import builder_lineage +assert Path(builder_lineage.__file__).resolve() == emitted.resolve() +assert 'builder_lineage_producer' not in sys.modules +assert 'code_mower.builder_lineage_producer' not in sys.modules +assert 'code_mower' not in sys.modules +assert importlib.util.find_spec('code_mower') is None +assert importlib.util.find_spec('builder_lineage_producer') is None +assert not (emitted.parent/'builder_lineage_producer.py').exists() +assert not (Path.cwd()/'src/code_mower/builder_lineage_producer.py').exists() +''' + result = subprocess.run([sys.executable, '-I', '-S', '-c', program], cwd=output, + capture_output=True, text=True, timeout=30) + self.assertEqual(result.returncode, 0, result.stderr) diff --git a/tests/test_lineage_producer_publication.py b/tests/test_lineage_producer_publication.py new file mode 100644 index 00000000..76f4a911 --- /dev/null +++ b/tests/test_lineage_producer_publication.py @@ -0,0 +1,231 @@ +"""Raw ingress, shared bounds, public readback and exact effect ordering.""" +from pathlib import Path +import tempfile +import unittest + +from code_mower.builder_lineage import Authorities, ContractError, History, Identity +from code_mower.builder_lineage_producer import ( + ProducerRefusal, Snapshot, decode_transport, fetch_history, observe, publish, selected_history, +) +from code_mower.builder_runs import record_lineage_builder +from lineage_producer_fixtures import ( + AUTHORITY, GitHubIO, POLICY, TRANSPORT, comments, episode, observation_args, sha, target, +) + + +class PublicationTests(unittest.TestCase): + def publish(self, io, private=None, **changes): + args = dict(target=target(), identity=POLICY, authorities=AUTHORITY, + private=[episode()] if private is None else private) + args.update(changes) + return publish(io, **args) + + def test_success_semantic_duplicate_and_idempotent_replay(self): + io = GitHubIO() + io.readback = comments([episode()]) * 2 + outcome = self.publish(io) + self.assertTrue(outcome.comment_posted) + self.assertEqual(io.effects, ["snapshot", "history", "snapshot", "post", "history", "snapshot", "labels", "snapshot"]) + io.effects.clear() + outcome = self.publish(io) + self.assertFalse(outcome.comment_posted) + self.assertFalse(outcome.labels_attempted) + self.assertNotIn("post", io.effects) + self.assertNotIn("labels", io.effects) + + def test_initial_ingress_failure_means_no_post_labels_or_artifact(self): + malformed = [None, {}, [None], [False], [[{}]], [{"body": None}], [{"user": []}], + [{"body": 1}], [{"author": {"login": 4}}]] + for raw in malformed: + with self.subTest(raw=raw): + io = GitHubIO() + io.public = raw + with self.assertRaises(ProducerRefusal) as raised: + self.publish(io) + self.assertFalse(raised.exception.comment_posted) + self.assertEqual(io.effects, ["snapshot", "history"]) + for raw in ([], [{"user": None}], [{"author": None}], [{}]): + io = GitHubIO() + io.public = raw + self.assertTrue(self.publish(io).comment_posted) + + def test_selected_history_never_uses_numeric_rest_count(self): + for selected in (None, 3, {}, [1]): + with self.assertRaises(ContractError): + selected_history({"comments": 3}, selected=selected) + self.assertEqual(selected_history({"comments": 3}, selected=[]), History([])) + for raw in (None, {}, [[{}], None], [[{}], {}], [None]): + with self.assertRaises(ContractError): + History.from_pages(raw) + self.assertEqual(History.from_pages([[], [{"user": None}]]), History([{"user": None}])) + + def test_marker_trust_shapes_legacy_and_no_authority(self): + body = comments([episode()])[0]["body"] + invalid = [body.replace('"schema":', '"schema":"duplicate","schema":'), + body.replace('"episodes":[', '"episodes":[],"extra":['), + '', + '', body + body, + body.replace('CODE_MOWER_BUILDER_LINEAGE:', 'CODE_MOWER_BUILDER_LINEAGE'), + body.replace('"sequence":1', '"schema":"old","sequence":1')] + for bad in invalid: + io = GitHubIO() + io.public = [{"user": {"login": "lineage-publisher[bot]"}, "body": bad}] + with self.subTest(body=bad), self.assertRaises(ProducerRefusal): + self.publish(io) + self.assertNotIn("post", io.effects) + self.assertNotIn("labels", io.effects) + io = GitHubIO() + io.public = comments([episode()], author="outsider") + self.assertTrue(self.publish(io).comment_posted) + io = GitHubIO() + with self.assertRaises(ProducerRefusal): + self.publish(io, authorities=Authorities([])) + self.assertNotIn("post", io.effects) + + def test_post_then_missing_untrusted_conflicting_or_malformed_readback(self): + bad_readbacks = [[], comments([episode()], author="outsider"), + comments([episode(expected_head=sha(9))]), + [{"body": "", "user": {"login": "lineage-publisher[bot]"}}], + [{"body": False}]] + for raw in bad_readbacks: + io = GitHubIO() + io.readback = raw + with self.subTest(raw=raw), self.assertRaises(ProducerRefusal) as caught: + self.publish(io) + self.assertTrue(caught.exception.comment_posted) + self.assertFalse(caught.exception.labels_attempted) + self.assertEqual(io.effects, ["snapshot", "history", "snapshot", "post", "history"]) + + def test_public_only_missing_final_control(self): + io = GitHubIO(2) + io.public = comments([episode()]) + io.readback = comments([episode()]) + with self.assertRaises(ProducerRefusal) as caught: + self.publish(io, private=[episode(), episode(2)], target=target(2)) + self.assertTrue(caught.exception.comment_posted) + self.assertNotIn("labels", io.effects) + + def test_fresh_label_failure_and_head_race_ordering(self): + for stage, expected in ((1, []), (2, ["history"]), (3, ["post", "history"]), (4, ["labels"])): + io = GitHubIO() + io.fail_snapshot = stage + with self.subTest(stage=stage), self.assertRaises(ProducerRefusal) as caught: + self.publish(io) + for effect in expected: + self.assertIn(effect, io.effects) + self.assertEqual(caught.exception.comment_posted, stage >= 3) + self.assertEqual(caught.exception.labels_attempted, stage == 4) + if stage <= 3: + self.assertNotIn("labels", io.effects) + for stage in (1, 2, 3, 4): + io = GitHubIO() + original = io.snapshot + def snapshot(t, original=original, io=io, stage=stage): + value = original(t) + return Snapshot(target(9), value.author, value.labels) if io.snapshots == stage else value + io.snapshot = snapshot + with self.subTest(race=stage), self.assertRaises(ProducerRefusal): + self.publish(io) + if stage <= 2: + self.assertNotIn("post", io.effects) + if stage <= 3: + self.assertNotIn("labels", io.effects) + + def test_exact_target_and_empty_chain_branch_policy(self): + for changes in ({"repo": "wrong/repo"}, {"pr_number": 43}, {"branch": "codex/topic"}, {"head_sha": sha(9)}): + io = GitHubIO() + with self.assertRaises(ProducerRefusal): + self.publish(io, target=target(**changes)) + self.assertEqual(io.effects, ["snapshot"]) + args = observation_args() + args.update(private=[], author="unmapped", labels=["builder:claude"]) + with self.assertRaises(ProducerRefusal): + observe(**args) + for bad in (None, {}, {"repo": "owner/repo"}): + with self.assertRaises(ProducerRefusal): + observe(**(args | {"target": bad})) + + def test_identity_aliases_and_reviewer_floor_are_preserved(self): + policy = POLICY.with_reviewer_floor("claude", ["reviewer-bot"]) + self.assertEqual(policy.branch_prefixes, POLICY.branch_prefixes) + self.assertTrue(policy.require_verified_lineage) + self.assertNotIn("reviewer-bot", AUTHORITY.accounts) + with self.assertRaises(ContractError): + POLICY.with_reviewer_floor("claude", ["source-bot"]) + with self.assertRaises(ContractError): + Identity.from_mapping({"authors": {"BOT": "codex", " bot ": "claude"}}) + + +class BoundsTests(unittest.TestCase): + def test_single_560_budget_and_561st_probe_without_562nd_read(self): + episodes = [episode(n) for n in range(1, 33)] + public = History([row for n in range(1, 33) for row in comments(episodes[:n])]) + args = observation_args(32) | {"history": public, "private": episodes} + observation = observe(**args) + self.assertEqual(observation.chain.raw_arrival_count, 560) + self.assertEqual(len(observation.chain.episodes), 32) + read = [] + def arrivals(): + for n in range(34): + read.append(n) + yield episodes[n] if n < 32 else episodes[-1] + with self.assertRaisesRegex(ContractError, "arrival budget"): + observe(**(args | {"private": arrivals()})) + self.assertEqual(len(read), 33) + # Intermediate public head is not separately resolved before later private evidence. + args["history"] = History(comments(episodes[:16])) + self.assertEqual(observe(**args).decision.status, "ready") + + def test_finite_page_terminal_full_cap_extra_probe_and_failure(self): + for pages, success, expected in (([[]], True, [1]), + ([[{}, {}], [{}]], True, [1, 2]), + ([[{}, {}], [{}, {}], []], True, [1, 2, 3]), + ([[{}, {}], [{}, {}], [{}]], False, [1, 2, 3]), + ([[{}, {}], None], False, [1, 2]), + ([[{}, {}], {}], False, [1, 2])): + calls = [] + def fetch(page, size, calls=calls, pages=pages): + calls.append(page) + return pages[page-1] + if success: + fetch_history(fetch, page_size=2, max_pages=2) + else: + with self.assertRaises(ValueError): + fetch_history(fetch, page_size=2, max_pages=2) + self.assertEqual(calls, expected) + def failed(page, size): + raise OSError("request failed") + with self.assertRaises(ProducerRefusal): + fetch_history(failed) + for raw in ('null', '{"x":1,"x":2}', '[{"body":null}]'): + with self.assertRaises(ValueError): + History(decode_transport(raw)) + + +class AttributionTests(unittest.TestCase): + def test_verified_writer_and_actual_local_hosted_transport_remain_distinct(self): + from code_mower.builder_lineage_producer import Transport + for transport in (TRANSPORT, Transport("devin", "devin_cli", "devin_cli", "local_cli"), + Transport("devin", "devin", "devin", "hosted_async_builder")): + args = observation_args() + args["private"] = [episode(destination_lane=transport.lane)] + observation = observe(**args) + with tempfile.TemporaryDirectory() as tmp: + event = record_lineage_builder(observation, transport, Path(tmp)/"event.json", created_at="2026-01-01T00:00:00+00:00") + self.assertEqual(event["provider"], transport.provider) + self.assertEqual(event["tool"]["executor"], transport.executor) + self.assertEqual(event["tool"]["integration"], transport.integration) + + def test_no_lineage_control_and_malformed_selected_history_no_artifact(self): + from code_mower.builder_lineage_producer import Transport + args = observation_args() | {"private": []} + observation = observe(**args) + with tempfile.TemporaryDirectory() as tmp: + output = Path(tmp)/"event.json" + event = record_lineage_builder(observation, Transport("codex", "codex", "codex_cli", "local_cli"), + output, created_at="2026-01-01T00:00:00+00:00") + self.assertEqual(event["dimensions"]["lineage"]["episode_count"], 0) + output.unlink() + with self.assertRaises(ValueError): + selected_history({"comments": 100, "body": "takeover by claude"}, selected=None) + self.assertFalse(output.exists()) diff --git a/tests/test_lineage_producer_state.py b/tests/test_lineage_producer_state.py new file mode 100644 index 00000000..8dee34ad --- /dev/null +++ b/tests/test_lineage_producer_state.py @@ -0,0 +1,159 @@ +"""Private record, supervised exit, role admission and chained delivery rows.""" +from copy import deepcopy +from pathlib import Path +import unittest +from unittest.mock import patch + +from code_mower import lane_delivery, lane_handoff +from code_mower.builder_lineage import ContractError, History +from code_mower.builder_lineage_producer import ProducerRefusal, ProducerStore, _delivery +from lineage_producer_fixtures import ( + AUTHORITY, BRANCH, MemoryStore, POLICY, TRANSPORT, episode, round_fixture, sha, target, +) + + +class StateTests(unittest.TestCase): + def setUp(self): + MemoryStore.records, MemoryStore.effects = {}, [] + self.addCleanup(patch.stopall) + patch("code_mower.builder_lineage_producer.ContextStore", MemoryStore).start() + patch("code_mower.lane_handoff.ContextStore", MemoryStore).start() + self.checkout = patch("code_mower.lane_delivery._lineage_checkout").start() + self.store = ProducerStore(Path("/producer-state")) + + def record(self, receipt, *, n=1, create=False, **overrides): + kwargs = dict(author="source-bot", labels=["builder:codex"], config={}, + runtime_observation=lambda: "ready", create=create) + kwargs.update(overrides) + return self.store.record(receipt, target(n), POLICY, AUTHORITY, History([]), **kwargs) + + def takeover(self): + handoff = lane_delivery.validate_handoff(source_lane="codex", destination_lane="claude", + target_pr="owner/repo#42", expected_head=sha(0), running_lane="claude", + repo="owner/repo", observed_head=sha(0), target_branch=BRANCH, + source_branch_prefixes=["codex/"]) + source = {"transport": "local_process", "state_dir": "/source-state", "writer": "source-writer"} + old = lane_handoff.LocalWriter(Path(source["state_dir"]), source["writer"]) + old.register(repo="owner/repo", lane="codex", checkout=Path("/source-checkout")) + old.started(10, 10) + old.finish(quiescent=True) + self.assertTrue(lane_handoff.prepare(handoff, source, Path("/handoffs"), + stop=lambda *args: "terminated", head=lambda h: sha(0))["accepted"]) + self.assertTrue(lane_handoff.reserve_launch(handoff, Path("/handoffs"), + stop=lambda *args: "terminated", head=lambda h: sha(0))) + current = round_fixture("/rounds") + current.started(11, 11) + current.finish(quiescent=True) + return handoff, source, current + + def receipt(self): + handoff, source, current = self.takeover() + return lane_handoff.lineage_handoff(handoff, Path("/handoffs"), current, + target(), source_branch_prefixes=["codex/"]) + + def test_takeover_then_same_writer_and_idempotent_replay(self): + receipt = self.receipt() + self.assertTrue(self.record(receipt, create=True)) + writes = len(MemoryStore.effects) + self.assertFalse(self.record(receipt)) + self.assertEqual(writes, len(MemoryStore.effects)) + previous = self.store.read(target()) + current = round_fixture("/rounds", 2) + current.started(12, 12) + current.finish(quiescent=True) + continuation = lane_delivery.lineage_continuation(current, target(2), previous) + self.assertEqual(continuation.episode.writer_state, "same_writer") + self.assertTrue(self.record(continuation, n=2)) + self.assertEqual(len(self.store.read(target(2))["episodes"]), 2) + with self.assertRaises(ContractError): + self.record(receipt, n=2) + + def test_missing_falsey_malformed_legacy_store_refuses_without_write(self): + receipt = self.receipt() + key = ("/producer-state", self.store._key(target())) + for raw in (None, {}, [], False, 0, "", {"schema": "code_mower.builderLineage.v1"}): + with self.subTest(raw=raw): + MemoryStore.records[key] = raw + before = len(MemoryStore.effects) + with self.assertRaises((ValueError, TypeError)): + self.record(receipt) + self.assertEqual(before, len(MemoryStore.effects)) + for raw in ({}, False, []): + MemoryStore.records[key] = raw + with self.assertRaises(ValueError): + self.record(receipt, create=True) + + def test_wrong_target_exact_branch_and_conflicting_replay(self): + receipt = self.receipt() + self.record(receipt, create=True) + for changes in ({"branch": BRANCH.lower()}, {"repo": "other/repo"}, {"pr_number": 43}, {"head_sha": sha(9)}): + with self.subTest(changes=changes), self.assertRaises(ValueError): + self.store.record(receipt, target(**changes), POLICY, AUTHORITY, History([]), + author="source-bot", labels=[], config={}, runtime_observation=lambda: "ready") + bad = _delivery(episode(resulting_head=sha(4)), receipt.writer, receipt.round_id, TRANSPORT) + with self.assertRaises(ValueError): + self.record(bad, n=4) + + def test_state_strings_are_not_independent_exit_proof(self): + handoff, source, current = self.takeover() + intent_key = ("/handoffs", lane_handoff.key([handoff.target_pr.lower(), handoff.expected_head])) + original = deepcopy(MemoryStore.records[intent_key]) + for state in ("completed", "cancelled", "suspended", "running", None): + MemoryStore.records[intent_key] = original | {"writer_state": state} + with self.subTest(state=state), self.assertRaises(ProducerRefusal): + lane_handoff.lineage_handoff(handoff, Path("/handoffs"), current, + target(), source_branch_prefixes=["codex/"]) + MemoryStore.records[intent_key] = original + source_key = (source["state_dir"], lane_handoff.key(source["writer"])) + original_source = deepcopy(MemoryStore.records[source_key]) + for changes in ({"quiescent": False}, {"finished": False}, {"lane": "claude"}, + {"pid": None}, {"pgid": True}, {"repo": "wrong/repo"}): + MemoryStore.records[source_key] = original_source | changes + with self.subTest(changes=changes), self.assertRaises(ProducerRefusal): + lane_handoff.lineage_handoff(handoff, Path("/handoffs"), current, + target(), source_branch_prefixes=["codex/"]) + + def test_continuation_wrong_writer_unstopped_round_and_stale_binding(self): + self.record(self.receipt(), create=True) + previous = self.store.read(target()) + current = round_fixture("/rounds", 2, writer="wrong-writer") + current.started(12, 12) + current.finish(quiescent=True) + with self.assertRaises(ProducerRefusal): + lane_delivery.lineage_continuation(current, target(2), previous) + current = round_fixture("/other-rounds", 2) + for state in ({}, {"finished": True}, {"finished": True, "quiescent": True}): + key = ("/other-rounds", current.control.key) + MemoryStore.records[key].update(state) + with self.assertRaises(ProducerRefusal): + lane_delivery.lineage_continuation(current, target(2), previous) + current.started(12, 12) + current.finish(quiescent=True) + for changes in ({"writer": "other"}, {"round_id": "round-2"}, + {"episodes": [episode(resulting_head=sha(9)).to_mapping()]}): + with self.assertRaises(ValueError): + lane_delivery.lineage_continuation(current, target(2), previous | changes) + + def test_role_runtime_and_capability_fail_before_storage_or_launch(self): + for config, runtime in (({"role_policy": {"claude": {"builder": {"enabled": False}}}}, "ready"), + ({}, "unchecked"), ({}, "unavailable")): + with self.subTest(config=config, runtime=runtime), self.assertRaises(ValueError): + round_fixture("/rounds", config=config, runtime=runtime) + self.assertEqual(MemoryStore.effects, []) + receipt = _delivery(episode(), "writer", "round", TRANSPORT) + with self.assertRaises(ValueError): + self.record(receipt, create=True, config={"role_policy": {"claude": {"builder": {"enabled": False}}}}) + self.assertEqual(MemoryStore.effects, []) + + def test_real_supervisor_observer_path_with_external_process_io_stub(self): + # The existing supervisor's callback is the evidence source; no caller + # state string is accepted by the conversion. This finite completed + # process tests the real supervisor and its stopped/reaped callbacks. + import sys + import tempfile + current = round_fixture("/rounds") + with tempfile.TemporaryDirectory() as tmp: + result = lane_delivery.supervise_process([sys.executable, "-c", "print('finite')"], + log_path=Path(tmp)/"round.log", timeout_seconds=5, writer=current) + self.assertEqual(result.exit_code, 0) + self.assertIs(current.observed(target()), current)